1# Copyright 2016 Google LLC
2#
3# Licensed under the Apache License, Version 2.0 (the "License");
4# you may not use this file except in compliance with the License.
5# You may obtain a copy of the License at
6#
7# http://www.apache.org/licenses/LICENSE-2.0
8#
9# Unless required by applicable law or agreed to in writing, software
10# distributed under the License is distributed on an "AS IS" BASIS,
11# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12# See the License for the specific language governing permissions and
13# limitations under the License.
14
15"""Transport adapter for Requests."""
16
17from __future__ import absolute_import
18
19import functools
20import http.client as http_client
21import logging
22import numbers
23import time
24from typing import Optional
25
26try:
27 import requests
28except ImportError as caught_exc: # pragma: NO COVER
29 raise ImportError(
30 "The requests library is not installed from please install the requests package to use the requests transport."
31 ) from caught_exc
32import requests.adapters # pylint: disable=ungrouped-imports
33import requests.exceptions # pylint: disable=ungrouped-imports
34from requests.packages.urllib3.util.ssl_ import ( # type: ignore
35 create_urllib3_context,
36) # pylint: disable=ungrouped-imports
37
38import google.auth.transport._mtls_helper
39from google.auth import _helpers, exceptions, transport
40from google.auth.transport import _mtls_helper
41from google.oauth2 import service_account
42
43_LOGGER = logging.getLogger(__name__)
44
45_DEFAULT_TIMEOUT = 120 # in seconds
46
47
48class _Response(transport.Response):
49 """Requests transport response adapter.
50
51 Args:
52 response (requests.Response): The raw Requests response.
53 """
54
55 def __init__(self, response):
56 self._response = response
57
58 @property
59 def status(self):
60 return self._response.status_code
61
62 @property
63 def headers(self):
64 return self._response.headers
65
66 @property
67 def data(self):
68 return self._response.content
69
70
71class TimeoutGuard(object):
72 """A context manager raising an error if the suite execution took too long.
73
74 Args:
75 timeout (Union[None, Union[float, Tuple[float, float]]]):
76 The maximum number of seconds a suite can run without the context
77 manager raising a timeout exception on exit. If passed as a tuple,
78 the smaller of the values is taken as a timeout. If ``None``, a
79 timeout error is never raised.
80 timeout_error_type (Optional[Exception]):
81 The type of the error to raise on timeout. Defaults to
82 :class:`requests.exceptions.Timeout`.
83 """
84
85 def __init__(self, timeout, timeout_error_type=requests.exceptions.Timeout):
86 self._timeout = timeout
87 self.remaining_timeout = timeout
88 self._timeout_error_type = timeout_error_type
89
90 def __enter__(self):
91 self._start = time.time()
92 return self
93
94 def __exit__(self, exc_type, exc_value, traceback):
95 if exc_value:
96 return # let the error bubble up automatically
97
98 if self._timeout is None:
99 return # nothing to do, the timeout was not specified
100
101 elapsed = time.time() - self._start
102 deadline_hit = False
103
104 if isinstance(self._timeout, numbers.Number):
105 self.remaining_timeout = self._timeout - elapsed
106 deadline_hit = self.remaining_timeout <= 0
107 else:
108 self.remaining_timeout = tuple(x - elapsed for x in self._timeout)
109 deadline_hit = min(self.remaining_timeout) <= 0
110
111 if deadline_hit:
112 raise self._timeout_error_type()
113
114
115class Request(transport.Request):
116 """Requests request adapter.
117
118 This class is used internally for making requests using various transports
119 in a consistent way. If you use :class:`AuthorizedSession` you do not need
120 to construct or use this class directly.
121
122 This class can be useful if you want to manually refresh a
123 :class:`~google.auth.credentials.Credentials` instance::
124
125 import google.auth.transport.requests
126 import requests
127
128 request = google.auth.transport.requests.Request()
129
130 credentials.refresh(request)
131
132 Args:
133 session (requests.Session): An instance :class:`requests.Session` used
134 to make HTTP requests. If not specified, a session will be created.
135
136 .. automethod:: __call__
137 """
138
139 def __init__(self, session: Optional[requests.Session] = None) -> None:
140 if not session:
141 session = requests.Session()
142
143 self.session = session
144
145 def __del__(self):
146 try:
147 if hasattr(self, "session") and self.session is not None:
148 self.session.close()
149 except TypeError:
150 # NOTE: For certain Python binary built, the queue.Empty exception
151 # might not be considered a normal Python exception causing
152 # TypeError.
153 pass
154
155 def __call__(
156 self,
157 url,
158 method="GET",
159 body=None,
160 headers=None,
161 timeout=_DEFAULT_TIMEOUT,
162 **kwargs,
163 ):
164 """Make an HTTP request using requests.
165
166 Args:
167 url (str): The URI to be requested.
168 method (str): The HTTP method to use for the request. Defaults
169 to 'GET'.
170 body (bytes): The payload or body in HTTP request.
171 headers (Mapping[str, str]): Request headers.
172 timeout (Optional[int]): The number of seconds to wait for a
173 response from the server. If not specified or if None, the
174 requests default timeout will be used.
175 kwargs: Additional arguments passed through to the underlying
176 requests :meth:`~requests.Session.request` method.
177
178 Returns:
179 google.auth.transport.Response: The HTTP response.
180
181 Raises:
182 google.auth.exceptions.TransportError: If any exception occurred.
183 """
184 try:
185 _helpers.request_log(_LOGGER, method, url, body, headers)
186 response = self.session.request(
187 method, url, data=body, headers=headers, timeout=timeout, **kwargs
188 )
189 _helpers.response_log(_LOGGER, response)
190 return _Response(response)
191 except requests.exceptions.RequestException as caught_exc:
192 new_exc = exceptions.TransportError(caught_exc)
193 raise new_exc from caught_exc
194
195
196class _MutualTlsAdapter(requests.adapters.HTTPAdapter):
197 """
198 A TransportAdapter that enables mutual TLS.
199
200 Args:
201 cert (bytes): client certificate in PEM format
202 key (bytes): client private key in PEM format
203
204 Raises:
205 ImportError: if certifi is not installed
206 google.auth.exceptions.MutualTLSChannelError: If the cert or key is invalid.
207 """
208
209 def __init__(self, cert, key, **kwargs):
210 import ssl
211
212 import certifi
213
214 ctx_poolmanager = create_urllib3_context()
215 ctx_poolmanager.load_verify_locations(cafile=certifi.where())
216
217 ctx_proxymanager = create_urllib3_context()
218 ctx_proxymanager.load_verify_locations(cafile=certifi.where())
219
220 try:
221 with _mtls_helper.secure_cert_key_paths(cert, key) as (
222 cert_path,
223 key_path,
224 passphrase,
225 ):
226 password = passphrase
227 ctx_poolmanager.load_cert_chain(
228 certfile=cert_path,
229 keyfile=key_path,
230 password=password,
231 )
232 ctx_proxymanager.load_cert_chain(
233 certfile=cert_path,
234 keyfile=key_path,
235 password=password,
236 )
237 except (
238 ssl.SSLError,
239 OSError,
240 IOError,
241 ValueError,
242 RuntimeError,
243 TypeError,
244 ) as exc:
245 raise exceptions.MutualTLSChannelError(
246 "Failed to configure client certificate and key for mTLS."
247 ) from exc
248
249 self._ctx_poolmanager = ctx_poolmanager
250 self._ctx_proxymanager = ctx_proxymanager
251
252 super(_MutualTlsAdapter, self).__init__(**kwargs)
253
254 def init_poolmanager(self, *args, **kwargs):
255 kwargs["ssl_context"] = self._ctx_poolmanager
256 super(_MutualTlsAdapter, self).init_poolmanager(*args, **kwargs)
257
258 def proxy_manager_for(self, *args, **kwargs):
259 kwargs["ssl_context"] = self._ctx_proxymanager
260 return super(_MutualTlsAdapter, self).proxy_manager_for(*args, **kwargs)
261
262
263class _MutualTlsOffloadAdapter(requests.adapters.HTTPAdapter):
264 """
265 A TransportAdapter that enables mutual TLS and offloads the client side
266 signing operation to the signing library.
267
268 Args:
269 enterprise_cert_file_path (str): the path to a enterprise cert JSON
270 file. The file should contain the following field:
271
272 {
273 "libs": {
274 "signer_library": "...",
275 "offload_library": "..."
276 }
277 }
278
279 Raises:
280 ImportError: if certifi is not installed
281 google.auth.exceptions.MutualTLSChannelError: If mutual TLS channel
282 creation failed for any reason.
283 """
284
285 def __init__(self, enterprise_cert_file_path):
286 import certifi
287
288 from google.auth.transport import _custom_tls_signer
289
290 self.signer = _custom_tls_signer.CustomTlsSigner(enterprise_cert_file_path)
291 self.signer.load_libraries()
292
293 poolmanager = create_urllib3_context()
294 poolmanager.load_verify_locations(cafile=certifi.where())
295 self.signer.attach_to_ssl_context(poolmanager)
296 self._ctx_poolmanager = poolmanager
297
298 proxymanager = create_urllib3_context()
299 proxymanager.load_verify_locations(cafile=certifi.where())
300 self.signer.attach_to_ssl_context(proxymanager)
301 self._ctx_proxymanager = proxymanager
302
303 super(_MutualTlsOffloadAdapter, self).__init__()
304
305 def init_poolmanager(self, *args, **kwargs):
306 kwargs["ssl_context"] = self._ctx_poolmanager
307 super(_MutualTlsOffloadAdapter, self).init_poolmanager(*args, **kwargs)
308
309 def proxy_manager_for(self, *args, **kwargs):
310 kwargs["ssl_context"] = self._ctx_proxymanager
311 return super(_MutualTlsOffloadAdapter, self).proxy_manager_for(*args, **kwargs)
312
313
314class AuthorizedSession(requests.Session):
315 """A Requests Session class with credentials.
316
317 This class is used to perform requests to API endpoints that require
318 authorization::
319
320 from google.auth.transport.requests import AuthorizedSession
321
322 authed_session = AuthorizedSession(credentials)
323
324 response = authed_session.request(
325 'GET', 'https://www.googleapis.com/storage/v1/b')
326
327
328 The underlying :meth:`request` implementation handles adding the
329 credentials' headers to the request and refreshing credentials as needed.
330
331 This class also supports mutual TLS via :meth:`configure_mtls_channel`
332 method. In order to use this method, the `GOOGLE_API_USE_CLIENT_CERTIFICATE`
333 environment variable must be explicitly set to ``true``, otherwise it does
334 nothing. Assume the environment is set to ``true``, the method behaves in the
335 following manner:
336
337 If client_cert_callback is provided, client certificate and private
338 key are loaded using the callback; if client_cert_callback is None,
339 application default SSL credentials will be used. Exceptions are raised if
340 there are problems with the certificate, private key, or the loading process,
341 so it should be called within a try/except block.
342
343 First we set the environment variable to ``true``, then create an :class:`AuthorizedSession`
344 instance and specify the endpoints::
345
346 regular_endpoint = 'https://pubsub.googleapis.com/v1/projects/{my_project_id}/topics'
347 mtls_endpoint = 'https://pubsub.mtls.googleapis.com/v1/projects/{my_project_id}/topics'
348
349 authed_session = AuthorizedSession(credentials)
350
351 Now we can pass a callback to :meth:`configure_mtls_channel`::
352
353 def my_cert_callback():
354 # some code to load client cert bytes and private key bytes, both in
355 # PEM format.
356 some_code_to_load_client_cert_and_key()
357 if loaded:
358 return cert, key
359 raise MyClientCertFailureException()
360
361 # Always call configure_mtls_channel within a try/except block.
362 try:
363 authed_session.configure_mtls_channel(my_cert_callback)
364 except:
365 # handle exceptions.
366
367 if authed_session.is_mtls:
368 response = authed_session.request('GET', mtls_endpoint)
369 else:
370 response = authed_session.request('GET', regular_endpoint)
371
372
373 You can alternatively use application default SSL credentials like this::
374
375 try:
376 authed_session.configure_mtls_channel()
377 except:
378 # handle exceptions.
379
380 Args:
381 credentials (google.auth.credentials.Credentials): The credentials to
382 add to the request.
383 refresh_status_codes (Sequence[int]): Which HTTP status codes indicate
384 that credentials should be refreshed and the request should be
385 retried.
386 max_refresh_attempts (int): The maximum number of times to attempt to
387 refresh the credentials and retry the request.
388 refresh_timeout (Optional[int]): The timeout value in seconds for
389 credential refresh HTTP requests.
390 auth_request (google.auth.transport.requests.Request):
391 (Optional) An instance of
392 :class:`~google.auth.transport.requests.Request` used when
393 refreshing credentials. If not passed,
394 an instance of :class:`~google.auth.transport.requests.Request`
395 is created.
396 default_host (Optional[str]): A host like "pubsub.googleapis.com".
397 This is used when a self-signed JWT is created from service
398 account credentials.
399 """
400
401 def __init__(
402 self,
403 credentials,
404 refresh_status_codes=transport.DEFAULT_REFRESH_STATUS_CODES,
405 max_refresh_attempts=transport.DEFAULT_MAX_REFRESH_ATTEMPTS,
406 refresh_timeout=None,
407 auth_request=None,
408 default_host=None,
409 ):
410 super(AuthorizedSession, self).__init__()
411 self.credentials = credentials
412 self._refresh_status_codes = refresh_status_codes
413 self._max_refresh_attempts = max_refresh_attempts
414 self._refresh_timeout = refresh_timeout
415 self._is_mtls = False
416 self._default_host = default_host
417
418 if auth_request is None:
419 self._auth_request_session = requests.Session()
420
421 # Using an adapter to make HTTP requests robust to network errors.
422 # This adapter retrys HTTP requests when network errors occur
423 # and the requests seems safely retryable.
424 retry_adapter = requests.adapters.HTTPAdapter(max_retries=3)
425 self._auth_request_session.mount("https://", retry_adapter)
426
427 # Do not pass `self` as the session here, as it can lead to
428 # infinite recursion.
429 auth_request = Request(self._auth_request_session)
430 else:
431 self._auth_request_session = None
432
433 # Request instance used by internal methods (for example,
434 # credentials.refresh).
435 self._auth_request = auth_request
436
437 # https://google.aip.dev/auth/4111
438 # Attempt to use self-signed JWTs when a service account is used.
439 if isinstance(self.credentials, service_account.Credentials):
440 self.credentials._create_self_signed_jwt(
441 "https://{}/".format(self._default_host) if self._default_host else None
442 )
443
444 def configure_mtls_channel(self, client_cert_callback=None):
445 """Configure the client certificate and key for SSL connection.
446
447 This method configures mTLS if client certificates are explicitly enabled
448 (via GOOGLE_API_USE_CLIENT_CERTIFICATE=true) or auto-enabled (when the env
449 variable is unset and workload certificates are discovered). In these cases,
450 if the client certificate and key are successfully obtained, a
451 :class:`_MutualTlsAdapter` instance will be mounted to the "https://" prefix.
452
453 Args:
454 client_cert_callback (Optional[Callable[[], (bytes, bytes)]]):
455 The optional callback returns the client certificate and private
456 key bytes both in PEM format.
457 If the callback is None, application default SSL credentials
458 will be used.
459
460 .. warning::
461 Calling this method mutates the underlying `requests.Session` adapter
462 dictionary. It is not thread-safe to call this explicitly while other
463 threads are making requests.
464
465 Raises:
466 google.auth.exceptions.MutualTLSChannelError: If mutual TLS channel
467 creation failed for any reason. The existing session state (such
468 as adapter mounts) remains unmodified if this error is raised.
469 """
470 use_client_cert = google.auth.transport._mtls_helper.check_use_client_cert()
471 if not use_client_cert:
472 return
473
474 try:
475 (
476 is_mtls,
477 cert,
478 key,
479 ) = google.auth.transport._mtls_helper.get_client_cert_and_key(
480 client_cert_callback
481 )
482
483 old_adapter = self.adapters.get("https://")
484
485 kwargs = {}
486 if old_adapter is not None:
487 kwargs["max_retries"] = getattr(old_adapter, "max_retries", 0)
488 kwargs["pool_connections"] = getattr(
489 old_adapter, "_pool_connections", requests.adapters.DEFAULT_POOLSIZE
490 )
491 kwargs["pool_maxsize"] = getattr(
492 old_adapter, "_pool_maxsize", requests.adapters.DEFAULT_POOLSIZE
493 )
494 kwargs["pool_block"] = getattr(
495 old_adapter, "_pool_block", requests.adapters.DEFAULT_POOLBLOCK
496 )
497
498 old_auth_adapter = None
499 auth_kwargs = {}
500 if self._auth_request_session is not None:
501 old_auth_adapter = self._auth_request_session.adapters.get("https://")
502
503 if old_auth_adapter is not None:
504 auth_kwargs["max_retries"] = getattr(
505 old_auth_adapter, "max_retries", 0
506 )
507 auth_kwargs["pool_connections"] = getattr(
508 old_auth_adapter,
509 "_pool_connections",
510 requests.adapters.DEFAULT_POOLSIZE,
511 )
512 auth_kwargs["pool_maxsize"] = getattr(
513 old_auth_adapter,
514 "_pool_maxsize",
515 requests.adapters.DEFAULT_POOLSIZE,
516 )
517 auth_kwargs["pool_block"] = getattr(
518 old_auth_adapter,
519 "_pool_block",
520 requests.adapters.DEFAULT_POOLBLOCK,
521 )
522
523 if is_mtls:
524 new_adapter = _MutualTlsAdapter(cert, key, **kwargs)
525 if self._auth_request_session is not None:
526 new_auth_adapter = _MutualTlsAdapter(cert, key, **auth_kwargs)
527 else:
528 new_auth_adapter = None
529 else:
530 new_adapter = requests.adapters.HTTPAdapter(**kwargs)
531 if self._auth_request_session is not None:
532 new_auth_adapter = requests.adapters.HTTPAdapter(**auth_kwargs)
533 else:
534 new_auth_adapter = None
535 except (
536 exceptions.ClientCertError,
537 ImportError,
538 OSError,
539 ValueError,
540 ) as caught_exc:
541 new_exc = exceptions.MutualTLSChannelError(caught_exc)
542 raise new_exc from caught_exc
543
544 self.mount("https://", new_adapter)
545
546 if old_adapter is not None and old_adapter is not new_adapter:
547 old_adapter.close()
548
549 if self._auth_request_session is not None and new_auth_adapter is not None:
550 self._auth_request_session.mount("https://", new_auth_adapter)
551
552 if (
553 old_auth_adapter is not None
554 and old_auth_adapter is not new_auth_adapter
555 ):
556 old_auth_adapter.close()
557
558 self._is_mtls = is_mtls
559 if is_mtls:
560 self._cached_cert = cert
561 else:
562 if hasattr(self, "_cached_cert"):
563 del self._cached_cert
564
565 def request(
566 self,
567 method,
568 url,
569 data=None,
570 headers=None,
571 max_allowed_time=None,
572 timeout=_DEFAULT_TIMEOUT,
573 **kwargs,
574 ):
575 """Implementation of Requests' request.
576
577 Args:
578 timeout (Optional[Union[float, Tuple[float, float]]]):
579 The amount of time in seconds to wait for the server response
580 with each individual request. Can also be passed as a tuple
581 ``(connect_timeout, read_timeout)``. See :meth:`requests.Session.request`
582 documentation for details.
583 max_allowed_time (Optional[float]):
584 If the method runs longer than this, a ``Timeout`` exception is
585 automatically raised. Unlike the ``timeout`` parameter, this
586 value applies to the total method execution time, even if
587 multiple requests are made under the hood.
588
589 Mind that it is not guaranteed that the timeout error is raised
590 at ``max_allowed_time``. It might take longer, for example, if
591 an underlying request takes a lot of time, but the request
592 itself does not timeout, e.g. if a large file is being
593 transmitted. The timeout error will be raised after such
594 request completes.
595 Raises:
596 google.auth.exceptions.MutualTLSChannelError: If mutual TLS
597 channel creation fails for any reason.
598 ValueError: If the client certificate is invalid.
599 """
600 # pylint: disable=arguments-differ
601 # Requests has a ton of arguments to request, but only two
602 # (method, url) are required. We pass through all of the other
603 # arguments to super, so no need to exhaustively list them here.
604
605 # Use a kwarg for this instead of an attribute to maintain
606 # thread-safety.
607 _credential_refresh_attempt = kwargs.pop("_credential_refresh_attempt", 0)
608
609 # Make a copy of the headers. They will be modified by the credentials
610 # and we want to pass the original headers if we recurse.
611 request_headers = headers.copy() if headers is not None else {}
612
613 # Do not apply the timeout unconditionally in order to not override the
614 # _auth_request's default timeout.
615 auth_request = (
616 self._auth_request
617 if timeout is None
618 else functools.partial(self._auth_request, timeout=timeout)
619 )
620
621 remaining_time = max_allowed_time
622
623 with TimeoutGuard(remaining_time) as guard:
624 self.credentials.before_request(auth_request, method, url, request_headers)
625 remaining_time = guard.remaining_timeout
626
627 with TimeoutGuard(remaining_time) as guard:
628 _helpers.request_log(_LOGGER, method, url, data, headers)
629 response = super(AuthorizedSession, self).request(
630 method,
631 url,
632 data=data,
633 headers=request_headers,
634 timeout=timeout,
635 **kwargs,
636 )
637 remaining_time = guard.remaining_timeout
638
639 # If the response indicated that the credentials needed to be
640 # refreshed, then refresh the credentials and re-attempt the
641 # request.
642 # A stored token may expire between the time it is retrieved and
643 # the time the request is made, so we may need to try twice.
644 if (
645 response.status_code in self._refresh_status_codes
646 and _credential_refresh_attempt < self._max_refresh_attempts
647 ):
648 # Handle unauthorized permission error(401 status code)
649 if response.status_code == http_client.UNAUTHORIZED:
650 use_mtls = self.is_mtls and _mtls_helper.is_mtls_endpoint(url)
651 if use_mtls:
652 (
653 call_cert_bytes,
654 call_key_bytes,
655 cached_fingerprint,
656 current_cert_fingerprint,
657 ) = _mtls_helper.check_parameters_for_unauthorized_response(
658 self._cached_cert
659 )
660 if cached_fingerprint != current_cert_fingerprint:
661 try:
662 _LOGGER.info(
663 "Client certificate has changed, reconfiguring mTLS "
664 "channel."
665 )
666 self.configure_mtls_channel(
667 lambda: (call_cert_bytes, call_key_bytes)
668 )
669 except Exception as e:
670 _LOGGER.error("Failed to reconfigure mTLS channel: %s", e)
671 raise exceptions.MutualTLSChannelError(
672 "Failed to reconfigure mTLS channel"
673 ) from e
674 else:
675 _LOGGER.info(
676 "Skipping reconfiguration of mTLS channel because the client"
677 " certificate has not changed."
678 )
679 _LOGGER.info(
680 "Refreshing credentials due to a %s response. Attempt %s/%s.",
681 response.status_code,
682 _credential_refresh_attempt + 1,
683 self._max_refresh_attempts,
684 )
685
686 # Do not apply the timeout unconditionally in order to not override the
687 # _auth_request's default timeout.
688 auth_request = (
689 self._auth_request
690 if timeout is None
691 else functools.partial(self._auth_request, timeout=timeout)
692 )
693
694 with TimeoutGuard(remaining_time) as guard:
695 self.credentials.refresh(auth_request)
696 remaining_time = guard.remaining_timeout
697
698 # Recurse. Pass in the original headers, not our modified set, but
699 # do pass the adjusted max allowed time (i.e. the remaining total time).
700 return self.request(
701 method,
702 url,
703 data=data,
704 headers=headers,
705 max_allowed_time=remaining_time,
706 timeout=timeout,
707 _credential_refresh_attempt=_credential_refresh_attempt + 1,
708 **kwargs,
709 )
710
711 return response
712
713 @property
714 def is_mtls(self):
715 """Indicates if the created SSL channel is mutual TLS."""
716 return self._is_mtls
717
718 def close(self):
719 if self._auth_request_session is not None:
720 self._auth_request_session.close()
721 super(AuthorizedSession, self).close()