airbyte_cdk.sources.streams.http

 1#
 2# Copyright (c) 2023 Airbyte, Inc., all rights reserved.
 3#
 4
 5# Initialize Streams Package
 6from .exceptions import UserDefinedBackoffException
 7from .http import HttpStream, HttpSubStream
 8from .http_client import HttpClient
 9
10__all__ = ["HttpClient", "HttpStream", "HttpSubStream", "UserDefinedBackoffException"]
class HttpClient:
139class HttpClient:
140    _DEFAULT_MAX_RETRY: int = 5
141    _DEFAULT_MAX_TIME: int = 60 * 10
142    # Backoff used in place of a rate-limit wait when another credential can serve the retry.
143    # Kept non-zero so a misreporting authenticator degrades to a slow retry, not a hot loop
144    # (the retry handler adds a second on top of whatever is returned).
145    TOKEN_ROTATION_BACKOFF: float = 0.1
146    _ACTIONS_TO_RETRY_ON = {
147        ResponseAction.RETRY,
148        ResponseAction.RATE_LIMITED,
149        ResponseAction.REFRESH_TOKEN_THEN_RETRY,
150    }
151
152    def __init__(
153        self,
154        name: str,
155        logger: logging.Logger,
156        error_handler: Optional[ErrorHandler] = None,
157        api_budget: Optional[APIBudget] = None,
158        session: Optional[Union[requests.Session, requests_cache.CachedSession]] = None,
159        authenticator: Optional[AuthBase] = None,
160        use_cache: bool = False,
161        backoff_strategy: Optional[Union[BackoffStrategy, List[BackoffStrategy]]] = None,
162        error_message_parser: Optional[ErrorMessageParser] = None,
163        disable_retries: bool = False,
164        message_repository: Optional[MessageRepository] = None,
165        request_timeout: Optional[Union[float, Tuple[Optional[float], Optional[float]]]] = None,
166        use_tcp_keepalive: bool = False,
167    ):
168        self._name = name
169        self._api_budget: APIBudget = api_budget or APIBudget(policies=[])
170        if session:
171            self._session = session
172        else:
173            self._use_cache = use_cache
174            self._session = self._request_session()
175            if use_tcp_keepalive:
176                tcp_keepalive_adapter = TcpKeepaliveHTTPAdapter(
177                    pool_connections=MAX_CONNECTION_POOL_SIZE,
178                    pool_maxsize=MAX_CONNECTION_POOL_SIZE,
179                )
180                self._session.mount("https://", tcp_keepalive_adapter)
181                self._session.mount("http://", tcp_keepalive_adapter)
182            else:
183                self._session.mount(
184                    "https://",
185                    requests.adapters.HTTPAdapter(
186                        pool_connections=MAX_CONNECTION_POOL_SIZE,
187                        pool_maxsize=MAX_CONNECTION_POOL_SIZE,
188                    ),
189                )
190        if isinstance(authenticator, AuthBase):
191            self._session.auth = authenticator
192        self._logger = logger
193        self._error_handler = error_handler or HttpStatusErrorHandler(self._logger)
194        if backoff_strategy is not None:
195            if isinstance(backoff_strategy, list):
196                self._backoff_strategies = backoff_strategy
197            else:
198                self._backoff_strategies = [backoff_strategy]
199        else:
200            self._backoff_strategies = [DefaultBackoffStrategy()]
201        self._error_message_parser = error_message_parser or JsonErrorMessageParser()
202        self._request_attempt_count: Dict[requests.PreparedRequest, int] = {}
203        self._token_refresh_outcomes: Dict[requests.PreparedRequest, bool] = {}
204        self._disable_retries = disable_retries
205        self._message_repository = message_repository
206        self._authenticator_update_failed = False
207        self._request_timeout = request_timeout
208
209    @property
210    def cache_filename(self) -> str:
211        """
212        Override if needed. Return the name of cache file
213        Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.
214        """
215        return f"{self._name}.sqlite"
216
217    def _request_session(self) -> requests.Session:
218        """
219        Session factory based on use_cache property and call rate limits (api_budget parameter)
220        :return: instance of request-based session
221        """
222        if self._use_cache:
223            cache_dir = os.getenv(ENV_REQUEST_CACHE_PATH)
224            # Use in-memory cache if cache_dir is not set
225            # This is a non-obvious interface, but it ensures we don't write sql files when running unit tests
226            # Use in-memory cache if cache_dir is not set
227            # This is a non-obvious interface, but it ensures we don't write sql files when running unit tests
228            sqlite_path = (
229                str(Path(cache_dir) / self.cache_filename)
230                if cache_dir
231                else "file::memory:?cache=shared"
232            )
233            # By using `PRAGMA synchronous=OFF` and `PRAGMA journal_mode=WAL`, we reduce the possible occurrences of `database table is locked` errors.
234            # Note that those were blindly added at the same time and one or the other might be sufficient to prevent the issues but we have seen good results with both. Feel free to revisit given more information.
235            # There are strong signals that `fast_save` might create problems but if the sync crashes, we start back from the beginning in terms of sqlite anyway so the impact should be minimal. Signals are:
236            # * https://github.com/requests-cache/requests-cache/commit/7fa89ffda300331c37d8fad7f773348a3b5b0236#diff-f43db4a5edf931647c32dec28ea7557aae4cae8444af4b26c8ecbe88d8c925aaR238
237            # * https://github.com/requests-cache/requests-cache/commit/7fa89ffda300331c37d8fad7f773348a3b5b0236#diff-2e7f95b7d7be270ff1a8118f817ea3e6663cdad273592e536a116c24e6d23c18R164-R168
238            # * `If the application running SQLite crashes, the data will be safe, but the database [might become corrupted](https://www.sqlite.org/howtocorrupt.html#cfgerr) if the operating system crashes or the computer loses power before that data has been written to the disk surface.` in [this description](https://www.sqlite.org/pragma.html#pragma_synchronous).
239            # The backend creates its tables when built, so build it under the lock too
240            with _SQLITE_CACHE_LOCK:
241                backend = requests_cache.SQLiteCache(sqlite_path, fast_save=True, wal=True)
242                backend.responses._lock = backend.redirects._lock = _SQLITE_CACHE_LOCK
243            return CachedLimiterSession(
244                cache_name=sqlite_path,
245                backend=backend,
246                api_budget=self._api_budget,
247                match_headers=True,
248            )
249        else:
250            return LimiterSession(api_budget=self._api_budget)
251
252    def clear_cache(self) -> None:
253        """
254        Clear cached requests for current session, can be called any time
255        """
256        if isinstance(self._session, requests_cache.CachedSession):
257            self._session.cache.clear()  # type: ignore # cache.clear is not typed
258
259    def _dedupe_query_params(
260        self, url: str, params: Optional[Mapping[str, str]]
261    ) -> Mapping[str, str]:
262        """
263        Remove query parameters from params mapping if they are already encoded in the URL.
264        :param url: URL with
265        :param params:
266        :return:
267        """
268        if params is None:
269            params = {}
270        query_string = urllib.parse.urlparse(url).query
271        query_dict = {k: v[0] for k, v in urllib.parse.parse_qs(query_string).items()}
272
273        duplicate_keys_with_same_value = {
274            k for k in query_dict.keys() if str(params.get(k)) == str(query_dict[k])
275        }
276        return {k: v for k, v in params.items() if k not in duplicate_keys_with_same_value}
277
278    def _create_prepared_request(
279        self,
280        http_method: str,
281        url: str,
282        dedupe_query_params: bool = False,
283        headers: Optional[Mapping[str, str]] = None,
284        params: Optional[Mapping[str, str]] = None,
285        json: Optional[Mapping[str, Any]] = None,
286        data: Optional[Union[str, Mapping[str, Any]]] = None,
287    ) -> requests.PreparedRequest:
288        if dedupe_query_params:
289            query_params = self._dedupe_query_params(url, params)
290        else:
291            query_params = params or {}
292        args = {"method": http_method, "url": url, "headers": headers, "params": query_params}
293        if http_method.upper() in BODY_REQUEST_METHODS:
294            if json and data:
295                raise RequestBodyException(
296                    "At the same time only one of the 'request_body_data' and 'request_body_json' functions can return data"
297                )
298            elif json:
299                args["json"] = json
300            elif data:
301                args["data"] = data
302        prepared_request: requests.PreparedRequest = self._session.prepare_request(
303            requests.Request(**args)
304        )
305
306        return prepared_request
307
308    @property
309    def _max_retries(self) -> int:
310        """
311        Determines the max retries based on the provided error handler.
312        """
313        max_retries = None
314        if self._disable_retries:
315            max_retries = 0
316        else:
317            max_retries = self._error_handler.max_retries
318        return max_retries if max_retries is not None else self._DEFAULT_MAX_RETRY
319
320    @property
321    def _max_time(self) -> int:
322        """
323        Determines the max time based on the provided error handler.
324        """
325        return (
326            self._error_handler.max_time
327            if self._error_handler.max_time is not None
328            else self._DEFAULT_MAX_TIME
329        )
330
331    def _send_with_retry(
332        self,
333        request: requests.PreparedRequest,
334        request_kwargs: Mapping[str, Any],
335        log_formatter: Optional[Callable[[requests.Response], Any]] = None,
336        exit_on_rate_limit: Optional[bool] = False,
337    ) -> requests.Response:
338        """
339        Sends a request with retry logic.
340
341        Args:
342            request (requests.PreparedRequest): The prepared HTTP request to send.
343            request_kwargs (Mapping[str, Any]): Additional keyword arguments for the request.
344
345        Returns:
346            requests.Response: The HTTP response received from the server after retries.
347        """
348
349        max_retries = self._max_retries
350        max_tries = max(0, max_retries) + 1
351        max_time = self._max_time
352
353        user_backoff_handler = user_defined_backoff_handler(max_tries=max_tries, max_time=max_time)(
354            self._send
355        )
356        rate_limit_backoff_handler = rate_limit_default_backoff_handler(max_tries=max_tries)
357        backoff_handler = http_client_default_backoff_handler(
358            max_tries=max_tries, max_time=max_time
359        )
360        # backoff handlers wrap _send, so it will always return a response -- except when all retries are exhausted
361        try:
362            response = backoff_handler(rate_limit_backoff_handler(user_backoff_handler))(
363                request,
364                request_kwargs,
365                log_formatter=log_formatter,
366                exit_on_rate_limit=exit_on_rate_limit,
367            )  # type: ignore # mypy can't infer that backoff_handler wraps _send
368
369            return response
370        except BaseBackoffException as e:
371            self._logger.error("Retries exhausted with backoff exception.", exc_info=True)
372
373            is_rate_limited = (
374                isinstance(e.response, requests.Response)
375                and e.response.status_code == requests.codes.too_many_requests
376            )
377
378            if is_rate_limited:
379                raise AirbyteTracedException(
380                    internal_message=f"Rate limit retry budget exhausted. Last exception: {e}",
381                    message="API rate limit exceeded.",
382                    failure_type=FailureType.transient_error,
383                    exception=e,
384                    stream_descriptor=StreamDescriptor(name=self._name),
385                )
386
387            raise AirbyteTracedException(
388                internal_message=f"Exhausted available request attempts. Exception: {e}",
389                message=f"Exhausted available request attempts. Please see logs for more details. Exception: {e}",
390                failure_type=e.failure_type or FailureType.system_error,
391                exception=e,
392                stream_descriptor=StreamDescriptor(name=self._name),
393            )
394
395    def _can_retry_on_another_token(self, request: requests.PreparedRequest) -> bool:
396        """Whether the authenticator can serve this request from a different credential now.
397
398        Opted into by implementing `TokenRotatingAuthenticator`.
399        """
400        authenticator = getattr(self._session, "auth", None)
401        if not isinstance(authenticator, TokenRotatingAuthenticator):
402            return False
403        try:
404            return bool(authenticator.has_alternative_token(request))
405        except Exception:
406            # Falling back to the computed wait is always safe, so never fail a retry over this.
407            self._logger.debug(
408                "Authenticator failed to report credential availability", exc_info=True
409            )
410            return False
411
412    def _update_authenticator_from_response(
413        self, request: requests.PreparedRequest, response: requests.Response
414    ) -> None:
415        """Let a quota-tracking authenticator reconcile its state against the server.
416
417        Authenticators only ever see requests, so an authenticator that tracks per-token quota
418        has no way to learn that the server disagrees with its local bookkeeping. This is the
419        feedback channel, mirroring what `LimiterMixin.send` does for the API budget.
420
421        Opted into by implementing `ResponseAwareAuthenticator`.
422        """
423        authenticator = getattr(self._session, "auth", None)
424        if not isinstance(authenticator, ResponseAwareAuthenticator):
425            return
426        if getattr(response, "from_cache", False):
427            # A replayed cached response carries the rate-limit headers from whenever it was
428            # first fetched and consumed no quota of its own.
429            return
430        try:
431            authenticator.update_from_response(request, response)
432        except Exception:
433            # Quota bookkeeping must never turn an otherwise fine response into a failure. Warn
434            # once so a persistently broken update -- which silently degrades the connector back
435            # to single-token behaviour -- is at least diagnosable from default-level logs.
436            if not self._authenticator_update_failed:
437                self._authenticator_update_failed = True
438                self._logger.warning(
439                    "Authenticator failed to update quota state from a response; token rotation "
440                    "may fall back to local counters only. Further occurrences log at debug.",
441                    exc_info=True,
442                )
443            else:
444                self._logger.debug(
445                    "Authenticator failed to update quota state from response", exc_info=True
446                )
447
448    def _send(
449        self,
450        request: requests.PreparedRequest,
451        request_kwargs: Mapping[str, Any],
452        log_formatter: Optional[Callable[[requests.Response], Any]] = None,
453        exit_on_rate_limit: Optional[bool] = False,
454    ) -> requests.Response:
455        if request not in self._request_attempt_count:
456            self._request_attempt_count[request] = 1
457        else:
458            self._request_attempt_count[request] += 1
459            if hasattr(self._session, "auth") and isinstance(self._session.auth, AuthBase):
460                self._session.auth(request)
461
462        self._logger.debug(
463            "Making outbound API request",
464            extra={"headers": request.headers, "url": request.url, "request_body": request.body},
465        )
466
467        response: Optional[requests.Response] = None
468        exc: Optional[requests.RequestException] = None
469
470        if self._request_timeout is not None and "timeout" not in request_kwargs:
471            request_kwargs = {**request_kwargs, "timeout": self._request_timeout}
472
473        try:
474            response = self._session.send(request, **request_kwargs)
475        except requests.RequestException as e:
476            exc = e
477
478        if response is not None:
479            self._update_authenticator_from_response(request, response)
480
481        error_resolution: ErrorResolution = self._error_handler.interpret_response(
482            response if response is not None else exc
483        )
484
485        # Evaluation of response.text can be heavy, for example, if streaming a large response
486        # Do it only in debug mode
487        if self._logger.isEnabledFor(logging.DEBUG) and response is not None:
488            if request_kwargs.get("stream"):
489                self._logger.debug(
490                    "Receiving response, but not logging it as the response is streamed",
491                    extra={"headers": response.headers, "status": response.status_code},
492                )
493            else:
494                self._logger.debug(
495                    "Receiving response",
496                    extra={
497                        "headers": response.headers,
498                        "status": response.status_code,
499                        "body": response.text,
500                    },
501                )
502
503        # Request/response logging for declarative cdk
504        if (
505            log_formatter is not None
506            and response is not None
507            and self._message_repository is not None
508        ):
509            formatter = log_formatter
510            # A response resolving to REDUCE_PAGE_SIZE or SPLIT_REQUEST_WINDOW is not a page of the stream:
511            # the retriever discards it and re-issues the request with a smaller page size or a narrower
512            # window. Logging it as an auxiliary request keeps it visible in the Connector Builder while
513            # keeping it out of the per-slice page count, which would otherwise report "limit reached" on a
514            # read that only retried.
515            log_as_auxiliary = error_resolution.response_action in (
516                ResponseAction.REDUCE_PAGE_SIZE,
517                ResponseAction.SPLIT_REQUEST_WINDOW,
518            )
519            self._message_repository.log_message(
520                Level.DEBUG,
521                lambda: _as_auxiliary_request_log(
522                    formatter(response),
523                    title=(
524                        f"Stream '{self._name}' page rejected, retrying with a smaller page size"
525                        if error_resolution.response_action == ResponseAction.REDUCE_PAGE_SIZE
526                        else f"Stream '{self._name}' request window rejected, retrying with a smaller window"
527                    ),
528                    description=(
529                        (
530                            f"Request for stream '{self._name}' whose response asked for a smaller page. The "
531                            f"same page is requested again with a reduced page size, so this request produced "
532                            f"no records."
533                        )
534                        if error_resolution.response_action == ResponseAction.REDUCE_PAGE_SIZE
535                        else (
536                            f"Request for stream '{self._name}' whose response asked for a smaller request "
537                            f"window. The window is split and re-read as smaller children, so this request "
538                            f"produced no records."
539                        )
540                    ),
541                )
542                if log_as_auxiliary
543                else formatter(response),
544            )
545
546        self._handle_error_resolution(
547            response=response,
548            exc=exc,
549            request=request,
550            error_resolution=error_resolution,
551            exit_on_rate_limit=exit_on_rate_limit,
552        )
553
554        return response  # type: ignore # will either return a valid response of type requests.Response or raise an exception
555
556    def _get_response_body(self, response: requests.Response) -> Optional[JsonType]:
557        """
558        Extracts and returns the body of an HTTP response.
559
560        This method attempts to parse the response body as JSON. If the response
561        body is not valid JSON, it falls back to decoding the response content
562        as a UTF-8 string. If both attempts fail, it returns None.
563
564        Args:
565            response (requests.Response): The HTTP response object.
566
567        Returns:
568            Optional[JsonType]: The parsed JSON object as a string, the decoded
569            response content as a string, or None if both parsing attempts fail.
570        """
571        try:
572            return str(response.json())
573        except requests.exceptions.JSONDecodeError:
574            try:
575                return response.content.decode("utf-8")
576            except Exception:
577                return "The Content of the Response couldn't be decoded."
578
579    def _evict_key(self, prepared_request: requests.PreparedRequest) -> None:
580        """
581        Addresses high memory consumption when enabling concurrency in https://github.com/airbytehq/oncall/issues/6821.
582
583        The `_request_attempt_count` attribute keeps growing as multiple requests are made using the same `http_client`.
584        To mitigate this issue, we evict keys for completed requests once we confirm that no further retries are needed.
585        This helps manage memory usage more efficiently while maintaining the necessary logic for retry attempts.
586        """
587        if prepared_request in self._request_attempt_count:
588            del self._request_attempt_count[prepared_request]
589        self._token_refresh_outcomes.pop(prepared_request, None)
590
591    def _auth_header_changed_since(self, request: requests.PreparedRequest) -> bool:
592        """Whether the authenticator's current Authorization header differs from the one the request was sent with."""
593        # request.headers is None on an unprepared request
594        sent = request.headers.get("Authorization") if request.headers else None
595        if not sent or not hasattr(self._session.auth, "get_auth_header"):
596            return False
597        current = self._session.auth.get_auth_header().get("Authorization")  # type: ignore[union-attr]
598        return current is not None and current != sent
599
600    @staticmethod
601    def _is_token_endpoint_rejection(error: requests.exceptions.RequestException) -> bool:
602        """Whether the token endpoint answered with a 4xx other than 429, i.e. it rejected the credentials rather than failing transiently."""
603        response = error.response
604        return (
605            response is not None
606            and 400 <= response.status_code < 500
607            and response.status_code != 429
608        )
609
610    def _handle_error_resolution(
611        self,
612        response: Optional[requests.Response],
613        exc: Optional[requests.RequestException],
614        request: requests.PreparedRequest,
615        error_resolution: ErrorResolution,
616        exit_on_rate_limit: Optional[bool] = False,
617    ) -> None:
618        if error_resolution.response_action not in self._ACTIONS_TO_RETRY_ON:
619            self._evict_key(request)
620
621        if error_resolution.response_action == ResponseAction.RESET_PAGINATION:
622            raise PaginationResetRequiredException()
623
624        if error_resolution.response_action == ResponseAction.REDUCE_PAGE_SIZE:
625            raise PageSizeReductionRequiredException(
626                stream_name=self._name, error_message=error_resolution.error_message
627            )
628
629        if error_resolution.response_action == ResponseAction.SPLIT_REQUEST_WINDOW:
630            raise RequestWindowSplitRequiredException(
631                stream_name=self._name,
632                error_message=error_resolution.error_message,
633                failure_type=error_resolution.failure_type,
634            )
635
636        # Emit stream status RUNNING with the reason RATE_LIMITED to log that the rate limit has been reached
637        if error_resolution.response_action == ResponseAction.RATE_LIMITED:
638            # TODO: Update to handle with message repository when concurrent message repository is ready
639            reasons = [AirbyteStreamStatusReason(type=AirbyteStreamStatusReasonType.RATE_LIMITED)]
640            message = orjson.dumps(
641                AirbyteMessageSerializer.dump(
642                    stream_status_as_airbyte_message(
643                        StreamDescriptor(name=self._name), AirbyteStreamStatus.RUNNING, reasons
644                    )
645                )
646            ).decode()
647
648            # Simply printing the stream status is a temporary solution and can cause future issues. Currently, the _send method is
649            # wrapped with backoff decorators, and we can only emit messages by iterating record_iterator in the abstract source at the
650            # end of the retry decorator behavior. This approach does not allow us to emit messages in the queue before exiting the
651            # backoff retry loop. Adding `\n` to the message and ignore 'end' ensure that few messages are printed at the same time.
652            print(f"{message}\n", end="", flush=True)
653
654        # Handle REFRESH_TOKEN_THEN_RETRY: force refresh the OAuth token before retrying,
655        # at most once per request. A config error from the refresh (e.g. bad credentials,
656        # including a 4xx rejection from the token endpoint) or a request rejected again
657        # after a refresh attempt fails fast instead of refreshing in a loop; the failure
658        # type reflects whether the refresh succeeded. Transient refresh failures (network
659        # errors, 5xx, 429) keep the retry transient. Non-OAuth auth types (e.g.,
660        # BearerAuthenticator) fall through to normal retry.
661        if error_resolution.response_action == ResponseAction.REFRESH_TOKEN_THEN_RETRY:
662            status = (
663                f"status code '{response.status_code}'"
664                if response is not None
665                else f"exception '{exc}'"
666            )
667            if request in self._token_refresh_outcomes:
668                refreshed = self._token_refresh_outcomes[request]
669                self._evict_key(request)
670                if refreshed:
671                    internal_message = f"'{request.method}' request to '{request.url}' was rejected with {status} again after the OAuth token was refreshed; not refreshing again."
672                    failure_type = FailureType.config_error
673                    message = "Refreshed OAuth access token is rejected by the API."
674                else:
675                    internal_message = f"'{request.method}' request to '{request.url}' was rejected with {status} and the OAuth token could not be refreshed; not refreshing again."
676                    failure_type = FailureType.transient_error
677                    message = (
678                        "API rejects the current OAuth access token and the token refresh failed."
679                    )
680                if error_resolution.error_message:
681                    internal_message += (
682                        f" Error handler message: '{error_resolution.error_message}'"
683                    )
684                self._logger.error(internal_message)
685                raise AirbyteTracedException(
686                    internal_message=internal_message,
687                    message=message,
688                    failure_type=failure_type,
689                )
690            if (
691                hasattr(self._session, "auth")
692                and self._session.auth is not None
693                and hasattr(self._session.auth, "refresh_and_set_access_token")
694            ):
695                self._token_refresh_outcomes[request] = False
696                try:
697                    if self._auth_header_changed_since(request):
698                        self._logger.info(
699                            "OAuth token was already replaced since this request was sent; retrying with the current token without refreshing again."
700                        )
701                    else:
702                        self._session.auth.refresh_and_set_access_token()  # type: ignore[union-attr]
703                        self._logger.info(
704                            "Refreshed OAuth token due to REFRESH_TOKEN_THEN_RETRY response action"
705                        )
706                    self._token_refresh_outcomes[request] = True
707                except AirbyteTracedException as refresh_error:
708                    if refresh_error.failure_type == FailureType.config_error:
709                        self._evict_key(request)
710                        raise
711                    self._logger.warning(
712                        f"Failed to refresh OAuth token: {refresh_error}. Proceeding with retry using existing token."
713                    )
714                except requests.exceptions.RequestException as refresh_error:
715                    if self._is_token_endpoint_rejection(refresh_error):
716                        self._evict_key(request)
717                        internal_message = (
718                            f"'{request.method}' request to '{request.url}' was rejected with {status} and the OAuth token endpoint "
719                            f"rejected the refresh request: {refresh_error}"
720                        )
721                        self._logger.error(internal_message)
722                        raise AirbyteTracedException(
723                            internal_message=internal_message,
724                            message="OAuth token refresh request is rejected by the token endpoint.",
725                            failure_type=FailureType.config_error,
726                        ) from refresh_error
727                    self._logger.warning(
728                        f"Failed to refresh OAuth token: {refresh_error}. Proceeding with retry using existing token."
729                    )
730                except Exception as refresh_error:
731                    self._logger.warning(
732                        f"Failed to refresh OAuth token: {refresh_error}. Proceeding with retry using existing token."
733                    )
734            else:
735                self._logger.warning(
736                    "REFRESH_TOKEN_THEN_RETRY action received but authenticator does not support token refresh. "
737                    "Proceeding with normal retry."
738                )
739
740        if error_resolution.response_action == ResponseAction.FAIL:
741            if response is not None:
742                filtered_response_message = filter_secrets(
743                    f"Request (body): '{str(request.body)}'. Response (body): '{self._get_response_body(response)}'. Response (headers): '{response.headers}'."
744                )
745                error_message = f"'{request.method}' request to '{request.url}' failed with status code '{response.status_code}' and error message: '{self._error_message_parser.parse_response_error_message(response)}'. {filtered_response_message}"
746            else:
747                error_message = (
748                    f"'{request.method}' request to '{request.url}' failed with exception: '{exc}'"
749                )
750
751            # ensure the exception message is emitted before raised
752            self._logger.error(error_message)
753
754            raise AirbyteTracedException(
755                internal_message=error_message,
756                message=error_resolution.error_message or error_message,
757                failure_type=error_resolution.failure_type,
758            )
759
760        elif error_resolution.response_action == ResponseAction.IGNORE:
761            if response is not None:
762                log_message = f"Ignoring response for '{request.method}' request to '{request.url}' with response code '{response.status_code}'"
763            else:
764                log_message = f"Ignoring response for '{request.method}' request to '{request.url}' with error '{exc}'"
765
766            self._logger.info(error_resolution.error_message or log_message)
767
768        # TODO: Consider dynamic retry count depending on subsequent error codes
769        elif error_resolution.response_action in (
770            ResponseAction.RETRY,
771            ResponseAction.RATE_LIMITED,
772            ResponseAction.REFRESH_TOKEN_THEN_RETRY,
773        ):
774            user_defined_backoff_time = None
775            # Asked before the strategies, not after. The backoff they compute describes the
776            # credential the server just rejected, so when another credential can serve the
777            # retry that wait is irrelevant -- and a strategy is allowed to refuse a wait by
778            # raising (`max_waiting_time_in_seconds`), which would otherwise end the stream
779            # before rotation was ever considered. Rotating is strictly the better outcome
780            # there: it is the same retry, seconds from now, on a credential with quota.
781            #
782            # Two consequences of not calling the strategies, both deliberate. A rate limit
783            # that yields no backoff at all now rotates too, rather than falling through to
784            # the default exponential retry -- on a rotating credential that is the better
785            # behaviour, and `has_alternative_token` only answers True when the retry will
786            # rotate -- the CDK's own authenticator narrows that further, to a sending credential
787            # that is tracked and spent. And a `max_waiting_time_in_seconds` the manifest got
788            # wrong -- one that cannot be evaluated -- is not reported from here, since that
789            # error is raised from inside the strategy. Both capped strategies therefore resolve
790            # the field once in `__post_init__` too, so a manifest mistake fails at startup
791            # rather than waiting for a rate limit that finds no spare credential.
792            rotate_instead_of_waiting = (
793                error_resolution.response_action == ResponseAction.RATE_LIMITED
794                and self._can_retry_on_another_token(request)
795            )
796
797            if rotate_instead_of_waiting:
798                # Says that a wait was skipped without the number, which is no longer computed,
799                # and names the cap explicitly: a connector that configured one gets no other
800                # signal that the retry went ahead without consulting it.
801                self._logger.info(
802                    "Rate limited on the current credential; retrying in "
803                    f"{self.TOKEN_ROTATION_BACKOFF}s with another one instead of waiting for the "
804                    "rate limit to reset. Any configured backoff, including a wait cap, is not "
805                    "evaluated for this retry."
806                )
807                user_defined_backoff_time = self.TOKEN_ROTATION_BACKOFF
808            else:
809                for backoff_strategy in self._backoff_strategies:
810                    backoff_time = backoff_strategy.backoff_time(
811                        response_or_exception=response if response is not None else exc,
812                        attempt_count=self._request_attempt_count[request],
813                    )
814                    if backoff_time:
815                        user_defined_backoff_time = backoff_time
816                        break
817
818            error_message = (
819                error_resolution.error_message
820                or f"Request to {request.url} failed with failure type {error_resolution.failure_type}, response action {error_resolution.response_action}."
821            )
822
823            retry_endlessly = (
824                error_resolution.response_action == ResponseAction.RATE_LIMITED
825                and not exit_on_rate_limit
826            )
827
828            if user_defined_backoff_time:
829                raise UserDefinedBackoffException(
830                    backoff=user_defined_backoff_time,
831                    request=request,
832                    response=(response if response is not None else exc),
833                    error_message=error_message,
834                    failure_type=error_resolution.failure_type,
835                )
836
837            elif retry_endlessly:
838                raise RateLimitBackoffException(
839                    request=request,
840                    response=(response if response is not None else exc),
841                    error_message=error_message,
842                    failure_type=error_resolution.failure_type,
843                )
844
845            raise DefaultBackoffException(
846                request=request,
847                response=(response if response is not None else exc),
848                error_message=error_message,
849                failure_type=error_resolution.failure_type,
850            )
851
852        elif response:
853            try:
854                response.raise_for_status()
855            except requests.HTTPError as e:
856                self._logger.error(response.text)
857                raise e
858
859    @property
860    def name(self) -> str:
861        return self._name
862
863    def send_request(
864        self,
865        http_method: str,
866        url: str,
867        request_kwargs: Mapping[str, Any],
868        headers: Optional[Mapping[str, str]] = None,
869        params: Optional[Mapping[str, str]] = None,
870        json: Optional[Mapping[str, Any]] = None,
871        data: Optional[Union[str, Mapping[str, Any]]] = None,
872        dedupe_query_params: bool = False,
873        log_formatter: Optional[Callable[[requests.Response], Any]] = None,
874        exit_on_rate_limit: Optional[bool] = False,
875    ) -> Tuple[requests.PreparedRequest, requests.Response]:
876        """
877        Prepares and sends request and return request and response objects.
878        """
879
880        request: requests.PreparedRequest = self._create_prepared_request(
881            http_method=http_method,
882            url=url,
883            dedupe_query_params=dedupe_query_params,
884            headers=headers,
885            params=params,
886            json=json,
887            data=data,
888        )
889
890        env_settings = self._session.merge_environment_settings(
891            url=request.url,
892            proxies=request_kwargs.get("proxies", {}),
893            stream=request_kwargs.get("stream"),
894            verify=request_kwargs.get("verify"),
895            cert=request_kwargs.get("cert"),
896        )
897        request_kwargs = {**request_kwargs, **env_settings}
898
899        response: requests.Response = self._send_with_retry(
900            request=request,
901            request_kwargs=request_kwargs,
902            log_formatter=log_formatter,
903            exit_on_rate_limit=exit_on_rate_limit,
904        )
905
906        return request, response
HttpClient( name: str, logger: logging.Logger, error_handler: Optional[airbyte_cdk.sources.streams.http.error_handlers.ErrorHandler] = None, api_budget: Optional[airbyte_cdk.sources.streams.call_rate.APIBudget] = None, session: Union[requests.sessions.Session, requests_cache.session.CachedSession, NoneType] = None, authenticator: Optional[requests.auth.AuthBase] = None, use_cache: bool = False, backoff_strategy: Union[airbyte_cdk.BackoffStrategy, List[airbyte_cdk.BackoffStrategy], NoneType] = None, error_message_parser: Optional[airbyte_cdk.sources.streams.http.error_handlers.ErrorMessageParser] = None, disable_retries: bool = False, message_repository: Optional[airbyte_cdk.MessageRepository] = None, request_timeout: Union[float, Tuple[Optional[float], Optional[float]], NoneType] = None, use_tcp_keepalive: bool = False)
152    def __init__(
153        self,
154        name: str,
155        logger: logging.Logger,
156        error_handler: Optional[ErrorHandler] = None,
157        api_budget: Optional[APIBudget] = None,
158        session: Optional[Union[requests.Session, requests_cache.CachedSession]] = None,
159        authenticator: Optional[AuthBase] = None,
160        use_cache: bool = False,
161        backoff_strategy: Optional[Union[BackoffStrategy, List[BackoffStrategy]]] = None,
162        error_message_parser: Optional[ErrorMessageParser] = None,
163        disable_retries: bool = False,
164        message_repository: Optional[MessageRepository] = None,
165        request_timeout: Optional[Union[float, Tuple[Optional[float], Optional[float]]]] = None,
166        use_tcp_keepalive: bool = False,
167    ):
168        self._name = name
169        self._api_budget: APIBudget = api_budget or APIBudget(policies=[])
170        if session:
171            self._session = session
172        else:
173            self._use_cache = use_cache
174            self._session = self._request_session()
175            if use_tcp_keepalive:
176                tcp_keepalive_adapter = TcpKeepaliveHTTPAdapter(
177                    pool_connections=MAX_CONNECTION_POOL_SIZE,
178                    pool_maxsize=MAX_CONNECTION_POOL_SIZE,
179                )
180                self._session.mount("https://", tcp_keepalive_adapter)
181                self._session.mount("http://", tcp_keepalive_adapter)
182            else:
183                self._session.mount(
184                    "https://",
185                    requests.adapters.HTTPAdapter(
186                        pool_connections=MAX_CONNECTION_POOL_SIZE,
187                        pool_maxsize=MAX_CONNECTION_POOL_SIZE,
188                    ),
189                )
190        if isinstance(authenticator, AuthBase):
191            self._session.auth = authenticator
192        self._logger = logger
193        self._error_handler = error_handler or HttpStatusErrorHandler(self._logger)
194        if backoff_strategy is not None:
195            if isinstance(backoff_strategy, list):
196                self._backoff_strategies = backoff_strategy
197            else:
198                self._backoff_strategies = [backoff_strategy]
199        else:
200            self._backoff_strategies = [DefaultBackoffStrategy()]
201        self._error_message_parser = error_message_parser or JsonErrorMessageParser()
202        self._request_attempt_count: Dict[requests.PreparedRequest, int] = {}
203        self._token_refresh_outcomes: Dict[requests.PreparedRequest, bool] = {}
204        self._disable_retries = disable_retries
205        self._message_repository = message_repository
206        self._authenticator_update_failed = False
207        self._request_timeout = request_timeout
TOKEN_ROTATION_BACKOFF: float = 0.1
cache_filename: str
209    @property
210    def cache_filename(self) -> str:
211        """
212        Override if needed. Return the name of cache file
213        Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.
214        """
215        return f"{self._name}.sqlite"

Override if needed. Return the name of cache file Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.

def clear_cache(self) -> None:
252    def clear_cache(self) -> None:
253        """
254        Clear cached requests for current session, can be called any time
255        """
256        if isinstance(self._session, requests_cache.CachedSession):
257            self._session.cache.clear()  # type: ignore # cache.clear is not typed

Clear cached requests for current session, can be called any time

name: str
859    @property
860    def name(self) -> str:
861        return self._name
def send_request( self, http_method: str, url: str, request_kwargs: Mapping[str, Any], headers: Optional[Mapping[str, str]] = None, params: Optional[Mapping[str, str]] = None, json: Optional[Mapping[str, Any]] = None, data: Union[str, Mapping[str, Any], NoneType] = None, dedupe_query_params: bool = False, log_formatter: Optional[Callable[[requests.models.Response], Any]] = None, exit_on_rate_limit: Optional[bool] = False) -> Tuple[requests.models.PreparedRequest, requests.models.Response]:
863    def send_request(
864        self,
865        http_method: str,
866        url: str,
867        request_kwargs: Mapping[str, Any],
868        headers: Optional[Mapping[str, str]] = None,
869        params: Optional[Mapping[str, str]] = None,
870        json: Optional[Mapping[str, Any]] = None,
871        data: Optional[Union[str, Mapping[str, Any]]] = None,
872        dedupe_query_params: bool = False,
873        log_formatter: Optional[Callable[[requests.Response], Any]] = None,
874        exit_on_rate_limit: Optional[bool] = False,
875    ) -> Tuple[requests.PreparedRequest, requests.Response]:
876        """
877        Prepares and sends request and return request and response objects.
878        """
879
880        request: requests.PreparedRequest = self._create_prepared_request(
881            http_method=http_method,
882            url=url,
883            dedupe_query_params=dedupe_query_params,
884            headers=headers,
885            params=params,
886            json=json,
887            data=data,
888        )
889
890        env_settings = self._session.merge_environment_settings(
891            url=request.url,
892            proxies=request_kwargs.get("proxies", {}),
893            stream=request_kwargs.get("stream"),
894            verify=request_kwargs.get("verify"),
895            cert=request_kwargs.get("cert"),
896        )
897        request_kwargs = {**request_kwargs, **env_settings}
898
899        response: requests.Response = self._send_with_retry(
900            request=request,
901            request_kwargs=request_kwargs,
902            log_formatter=log_formatter,
903            exit_on_rate_limit=exit_on_rate_limit,
904        )
905
906        return request, response

Prepares and sends request and return request and response objects.

 45class HttpStream(Stream, CheckpointMixin, ABC):
 46    """
 47    Base abstract class for an Airbyte Stream using the HTTP protocol. Basic building block for users building an Airbyte source for a HTTP API.
 48    """
 49
 50    source_defined_cursor = True  # Most HTTP streams use a source defined cursor (i.e: the user can't configure it like on a SQL table)
 51    page_size: Optional[int] = (
 52        None  # Use this variable to define page size for API http requests with pagination support
 53    )
 54
 55    def __init__(
 56        self, authenticator: Optional[AuthBase] = None, api_budget: Optional[APIBudget] = None
 57    ):
 58        self._exit_on_rate_limit: bool = False
 59        self._http_client = HttpClient(
 60            name=self.name,
 61            logger=self.logger,
 62            error_handler=self.get_error_handler(),
 63            api_budget=api_budget or APIBudget(policies=[]),
 64            authenticator=authenticator,
 65            use_cache=self.use_cache,
 66            backoff_strategy=self.get_backoff_strategy(),
 67            message_repository=InMemoryMessageRepository(),
 68        )
 69
 70        # There are three conditions that dictate if RFR should automatically be applied to a stream
 71        # 1. Streams that explicitly initialize their own cursor should defer to it and not automatically apply RFR
 72        # 2. Streams with at least one cursor_field are incremental and thus a superior sync to RFR.
 73        # 3. Streams overriding read_records() do not guarantee that they will call the parent implementation which can perform
 74        #    per-page checkpointing so RFR is only supported if a stream use the default `HttpStream.read_records()` method
 75        if (
 76            not self.cursor
 77            and len(self.cursor_field) == 0
 78            and type(self).read_records is HttpStream.read_records
 79        ):
 80            self.cursor = ResumableFullRefreshCursor()
 81
 82    @property
 83    def exit_on_rate_limit(self) -> bool:
 84        """
 85        :return: False if the stream will retry endlessly when rate limited
 86        """
 87        return self._exit_on_rate_limit
 88
 89    @exit_on_rate_limit.setter
 90    def exit_on_rate_limit(self, value: bool) -> None:
 91        self._exit_on_rate_limit = value
 92
 93    @property
 94    def cache_filename(self) -> str:
 95        """
 96        Override if needed. Return the name of cache file
 97        Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.
 98        """
 99        return f"{self.name}.sqlite"
100
101    @property
102    def use_cache(self) -> bool:
103        """
104        Override if needed. If True, all records will be cached.
105        Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.
106        """
107        return False
108
109    @property
110    @abstractmethod
111    def url_base(self) -> str:
112        """
113        :return: URL base for the  API endpoint e.g: if you wanted to hit https://myapi.com/v1/some_entity then this should return "https://myapi.com/v1/"
114        """
115
116    @property
117    def http_method(self) -> str:
118        """
119        Override if needed. See get_request_data/get_request_json if using POST/PUT/PATCH.
120        """
121        return "GET"
122
123    @property
124    @deprecated(
125        "Deprecated as of CDK version 3.0.0. "
126        "You should set error_handler explicitly in HttpStream.get_error_handler() instead."
127    )
128    def raise_on_http_errors(self) -> bool:
129        """
130        Override if needed. If set to False, allows opting-out of raising HTTP code exception.
131        """
132        return True
133
134    @property
135    @deprecated(
136        "Deprecated as of CDK version 3.0.0. "
137        "You should set backoff_strategies explicitly in HttpStream.get_backoff_strategy() instead."
138    )
139    def max_retries(self) -> Union[int, None]:
140        """
141        Override if needed. Specifies maximum amount of retries for backoff policy. Return None for no limit.
142        """
143        return 5
144
145    @property
146    @deprecated(
147        "Deprecated as of CDK version 3.0.0. "
148        "You should set backoff_strategies explicitly in HttpStream.get_backoff_strategy() instead."
149    )
150    def max_time(self) -> Union[int, None]:
151        """
152        Override if needed. Specifies maximum total waiting time (in seconds) for backoff policy. Return None for no limit.
153        """
154        return 60 * 10
155
156    @property
157    @deprecated(
158        "Deprecated as of CDK version 3.0.0. "
159        "You should set backoff_strategies explicitly in HttpStream.get_backoff_strategy() instead."
160    )
161    def retry_factor(self) -> float:
162        """
163        Override if needed. Specifies factor for backoff policy.
164        """
165        return 5
166
167    @abstractmethod
168    def next_page_token(self, response: requests.Response) -> Optional[Mapping[str, Any]]:
169        """
170        Override this method to define a pagination strategy.
171
172        The value returned from this method is passed to most other methods in this class. Use it to form a request e.g: set headers or query params.
173
174        :return: The token for the next page from the input response object. Returning None means there are no more pages to read in this response.
175        """
176
177    @abstractmethod
178    def path(
179        self,
180        *,
181        stream_state: Optional[Mapping[str, Any]] = None,
182        stream_slice: Optional[Mapping[str, Any]] = None,
183        next_page_token: Optional[Mapping[str, Any]] = None,
184    ) -> str:
185        """
186        Returns the URL path for the API endpoint e.g: if you wanted to hit https://myapi.com/v1/some_entity then this should return "some_entity"
187        """
188
189    def request_params(
190        self,
191        stream_state: Optional[Mapping[str, Any]],
192        stream_slice: Optional[Mapping[str, Any]] = None,
193        next_page_token: Optional[Mapping[str, Any]] = None,
194    ) -> MutableMapping[str, Any]:
195        """
196        Override this method to define the query parameters that should be set on an outgoing HTTP request given the inputs.
197
198        E.g: you might want to define query parameters for paging if next_page_token is not None.
199        """
200        return {}
201
202    def request_headers(
203        self,
204        stream_state: Optional[Mapping[str, Any]],
205        stream_slice: Optional[Mapping[str, Any]] = None,
206        next_page_token: Optional[Mapping[str, Any]] = None,
207    ) -> Mapping[str, Any]:
208        """
209        Override to return any non-auth headers. Authentication headers will overwrite any overlapping headers returned from this method.
210        """
211        return {}
212
213    def request_body_data(
214        self,
215        stream_state: Optional[Mapping[str, Any]],
216        stream_slice: Optional[Mapping[str, Any]] = None,
217        next_page_token: Optional[Mapping[str, Any]] = None,
218    ) -> Optional[Union[Mapping[str, Any], str]]:
219        """
220        Override when creating POST/PUT/PATCH requests to populate the body of the request with a non-JSON payload.
221
222        If returns a ready text that it will be sent as is.
223        If returns a dict that it will be converted to a urlencoded form.
224        E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2"
225
226        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
227        """
228        return None
229
230    def request_body_json(
231        self,
232        stream_state: Optional[Mapping[str, Any]],
233        stream_slice: Optional[Mapping[str, Any]] = None,
234        next_page_token: Optional[Mapping[str, Any]] = None,
235    ) -> Optional[Mapping[str, Any]]:
236        """
237        Override when creating POST/PUT/PATCH requests to populate the body of the request with a JSON payload.
238
239        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
240        """
241        return None
242
243    def request_kwargs(
244        self,
245        stream_state: Optional[Mapping[str, Any]],
246        stream_slice: Optional[Mapping[str, Any]] = None,
247        next_page_token: Optional[Mapping[str, Any]] = None,
248    ) -> Mapping[str, Any]:
249        """
250        Override to return a mapping of keyword arguments to be used when creating the HTTP request.
251        Any option listed in https://docs.python-requests.org/en/latest/api/#requests.adapters.BaseAdapter.send for can be returned from
252        this method. Note that these options do not conflict with request-level options such as headers, request params, etc..
253        """
254        return {}
255
256    @abstractmethod
257    def parse_response(
258        self,
259        response: requests.Response,
260        *,
261        stream_state: Mapping[str, Any],
262        stream_slice: Optional[Mapping[str, Any]] = None,
263        next_page_token: Optional[Mapping[str, Any]] = None,
264    ) -> Iterable[Mapping[str, Any]]:
265        """
266        Parses the raw response object into a list of records.
267        By default, this returns an iterable containing the input. Override to parse differently.
268        :param response:
269        :param stream_state:
270        :param stream_slice:
271        :param next_page_token:
272        :return: An iterable containing the parsed response
273        """
274
275    def get_backoff_strategy(self) -> Optional[Union[BackoffStrategy, List[BackoffStrategy]]]:
276        """
277        Used to initialize Adapter to avoid breaking changes.
278        If Stream has a `backoff_time` method implementation, we know this stream uses old (pre-HTTPClient) backoff handlers and thus an adapter is needed.
279
280        Override to provide custom BackoffStrategy
281        :return Optional[BackoffStrategy]:
282        """
283        if hasattr(self, "backoff_time"):
284            return HttpStreamAdapterBackoffStrategy(self)
285        else:
286            return None
287
288    def get_error_handler(self) -> Optional[ErrorHandler]:
289        """
290        Used to initialize Adapter to avoid breaking changes.
291        If Stream has a `should_retry` method implementation, we know this stream uses old (pre-HTTPClient) error handlers and thus an adapter is needed.
292
293        Override to provide custom ErrorHandler
294        :return Optional[ErrorHandler]:
295        """
296        if hasattr(self, "should_retry"):
297            error_handler = HttpStreamAdapterHttpStatusErrorHandler(
298                stream=self,
299                logger=logging.getLogger(),
300                max_retries=self.max_retries,
301                max_time=timedelta(seconds=self.max_time or 0),
302            )
303            return error_handler
304        else:
305            return None
306
307    @classmethod
308    def _join_url(cls, url_base: str, path: str) -> str:
309        return urljoin(url_base, path)
310
311    @classmethod
312    def parse_response_error_message(cls, response: requests.Response) -> Optional[str]:
313        """
314        Parses the raw response object from a failed request into a user-friendly error message.
315        By default, this method tries to grab the error message from JSON responses by following common API patterns. Override to parse differently.
316
317        :param response:
318        :return: A user-friendly message that indicates the cause of the error
319        """
320
321        # default logic to grab error from common fields
322        def _try_get_error(value: Optional[JsonType]) -> Optional[str]:
323            if isinstance(value, str):
324                return value
325            elif isinstance(value, list):
326                errors_in_value = [_try_get_error(v) for v in value]
327                return ", ".join(v for v in errors_in_value if v is not None)
328            elif isinstance(value, dict):
329                new_value = (
330                    value.get("message")
331                    or value.get("messages")
332                    or value.get("error")
333                    or value.get("errors")
334                    or value.get("failures")
335                    or value.get("failure")
336                    or value.get("detail")
337                )
338                return _try_get_error(new_value)
339            return None
340
341        try:
342            body = response.json()
343            return _try_get_error(body)
344        except requests.exceptions.JSONDecodeError:
345            return None
346
347    def get_error_display_message(self, exception: BaseException) -> Optional[str]:
348        """
349        Retrieves the user-friendly display message that corresponds to an exception.
350        This will be called when encountering an exception while reading records from the stream, and used to build the AirbyteTraceMessage.
351
352        The default implementation of this method only handles HTTPErrors by passing the response to self.parse_response_error_message().
353        The method should be overriden as needed to handle any additional exception types.
354
355        :param exception: The exception that was raised
356        :return: A user-friendly message that indicates the cause of the error
357        """
358        if isinstance(exception, requests.HTTPError) and exception.response is not None:
359            return self.parse_response_error_message(exception.response)
360        return None
361
362    def read_records(
363        self,
364        sync_mode: SyncMode,
365        cursor_field: Optional[List[str]] = None,
366        stream_slice: Optional[Mapping[str, Any]] = None,
367        stream_state: Optional[Mapping[str, Any]] = None,
368    ) -> Iterable[StreamData]:
369        # A cursor_field indicates this is an incremental stream which offers better checkpointing than RFR enabled via the cursor
370        if self.cursor_field or not isinstance(self.get_cursor(), ResumableFullRefreshCursor):
371            yield from self._read_pages(
372                lambda req, res, state, _slice: self.parse_response(
373                    res, stream_slice=_slice, stream_state=state
374                ),
375                stream_slice,
376                stream_state,
377            )
378        else:
379            yield from self._read_single_page(
380                lambda req, res, state, _slice: self.parse_response(
381                    res, stream_slice=_slice, stream_state=state
382                ),
383                stream_slice,
384                stream_state,
385            )
386
387    @property
388    def state(self) -> MutableMapping[str, Any]:
389        cursor = self.get_cursor()
390        if cursor:
391            return cursor.get_stream_state()  # type: ignore
392        return self._state
393
394    @state.setter
395    def state(self, value: MutableMapping[str, Any]) -> None:
396        cursor = self.get_cursor()
397        if cursor:
398            cursor.set_initial_state(value)
399        self._state = value
400
401    def get_cursor(self) -> Optional[Cursor]:
402        # I don't love that this is semi-stateful but not sure what else to do. We don't know exactly what type of cursor to
403        # instantiate when creating the class. We can make a few assumptions like if there is a cursor_field which implies
404        # incremental, but we don't know until runtime if this is a substream. Ideally, a stream should explicitly define
405        # its cursor, but because we're trying to automatically apply RFR we're stuck with this logic where we replace the
406        # cursor at runtime once we detect this is a substream based on self.has_multiple_slices being reassigned
407        if self.has_multiple_slices and isinstance(self.cursor, ResumableFullRefreshCursor):
408            self.cursor = SubstreamResumableFullRefreshCursor()
409            return self.cursor
410        else:
411            return self.cursor
412
413    def _read_pages(
414        self,
415        records_generator_fn: Callable[
416            [
417                requests.PreparedRequest,
418                requests.Response,
419                Mapping[str, Any],
420                Optional[Mapping[str, Any]],
421            ],
422            Iterable[StreamData],
423        ],
424        stream_slice: Optional[Mapping[str, Any]] = None,
425        stream_state: Optional[Mapping[str, Any]] = None,
426    ) -> Iterable[StreamData]:
427        stream_state = stream_state or {}
428        pagination_complete = False
429        next_page_token = None
430        while not pagination_complete:
431            request, response = self._fetch_next_page(stream_slice, stream_state, next_page_token)
432            yield from records_generator_fn(request, response, stream_state, stream_slice)
433
434            next_page_token = self.next_page_token(response)
435            if not next_page_token:
436                pagination_complete = True
437
438        cursor = self.get_cursor()
439        if cursor and isinstance(cursor, SubstreamResumableFullRefreshCursor):
440            partition, _, _ = self._extract_slice_fields(stream_slice=stream_slice)
441            # Substreams checkpoint state by marking an entire parent partition as completed so that on the subsequent attempt
442            # after a failure, completed parents are skipped and the sync can make progress
443            cursor.close_slice(StreamSlice(cursor_slice={}, partition=partition))
444
445        # Always return an empty generator just in case no records were ever yielded
446        yield from []
447
448    def _read_single_page(
449        self,
450        records_generator_fn: Callable[
451            [
452                requests.PreparedRequest,
453                requests.Response,
454                Mapping[str, Any],
455                Optional[Mapping[str, Any]],
456            ],
457            Iterable[StreamData],
458        ],
459        stream_slice: Optional[Mapping[str, Any]] = None,
460        stream_state: Optional[Mapping[str, Any]] = None,
461    ) -> Iterable[StreamData]:
462        partition, cursor_slice, remaining_slice = self._extract_slice_fields(
463            stream_slice=stream_slice
464        )
465        stream_state = stream_state or {}
466        next_page_token = cursor_slice or None
467
468        request, response = self._fetch_next_page(remaining_slice, stream_state, next_page_token)
469        yield from records_generator_fn(request, response, stream_state, remaining_slice)
470
471        next_page_token = self.next_page_token(response) or {
472            "__ab_full_refresh_sync_complete": True
473        }
474
475        cursor = self.get_cursor()
476        if cursor:
477            cursor.close_slice(StreamSlice(cursor_slice=next_page_token, partition=partition))
478
479        # Always return an empty generator just in case no records were ever yielded
480        yield from []
481
482    @staticmethod
483    def _extract_slice_fields(
484        stream_slice: Optional[Mapping[str, Any]],
485    ) -> tuple[Mapping[str, Any], Mapping[str, Any], Mapping[str, Any]]:
486        if not stream_slice:
487            return {}, {}, {}
488
489        if isinstance(stream_slice, StreamSlice):
490            partition = stream_slice.partition
491            cursor_slice = stream_slice.cursor_slice
492            remaining = {k: v for k, v in stream_slice.items()}
493        else:
494            # RFR streams that implement stream_slices() to generate stream slices in the legacy mapping format are converted into a
495            # structured stream slice mapping by the LegacyCursorBasedCheckpointReader. The structured mapping object has separate
496            # fields for the partition and cursor_slice value
497            partition = stream_slice.get("partition", {})
498            cursor_slice = stream_slice.get("cursor_slice", {})
499            remaining = {
500                key: val
501                for key, val in stream_slice.items()
502                if key != "partition" and key != "cursor_slice"
503            }
504        return partition, cursor_slice, remaining
505
506    def _fetch_next_page(
507        self,
508        stream_slice: Optional[Mapping[str, Any]] = None,
509        stream_state: Optional[Mapping[str, Any]] = None,
510        next_page_token: Optional[Mapping[str, Any]] = None,
511    ) -> Tuple[requests.PreparedRequest, requests.Response]:
512        request, response = self._http_client.send_request(
513            http_method=self.http_method,
514            url=self._join_url(
515                self.url_base,
516                self.path(
517                    stream_state=stream_state,
518                    stream_slice=stream_slice,
519                    next_page_token=next_page_token,
520                ),
521            ),
522            request_kwargs=self.request_kwargs(
523                stream_state=stream_state,
524                stream_slice=stream_slice,
525                next_page_token=next_page_token,
526            ),
527            headers=self.request_headers(
528                stream_state=stream_state,
529                stream_slice=stream_slice,
530                next_page_token=next_page_token,
531            ),
532            params=self.request_params(
533                stream_state=stream_state,
534                stream_slice=stream_slice,
535                next_page_token=next_page_token,
536            ),
537            json=self.request_body_json(
538                stream_state=stream_state,
539                stream_slice=stream_slice,
540                next_page_token=next_page_token,
541            ),
542            data=self.request_body_data(
543                stream_state=stream_state,
544                stream_slice=stream_slice,
545                next_page_token=next_page_token,
546            ),
547            dedupe_query_params=True,
548            log_formatter=self.get_log_formatter(),
549            exit_on_rate_limit=self.exit_on_rate_limit,
550        )
551
552        return request, response
553
554    def get_log_formatter(self) -> Optional[Callable[[requests.Response], Any]]:
555        """
556
557        :return Optional[Callable[[requests.Response], Any]]: Function that will be used in logging inside HttpClient
558        """
559        return None

Base abstract class for an Airbyte Stream using the HTTP protocol. Basic building block for users building an Airbyte source for a HTTP API.

source_defined_cursor = True

Return False if the cursor can be configured by the user.

page_size: Optional[int] = None
exit_on_rate_limit: bool
82    @property
83    def exit_on_rate_limit(self) -> bool:
84        """
85        :return: False if the stream will retry endlessly when rate limited
86        """
87        return self._exit_on_rate_limit
Returns

False if the stream will retry endlessly when rate limited

cache_filename: str
93    @property
94    def cache_filename(self) -> str:
95        """
96        Override if needed. Return the name of cache file
97        Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.
98        """
99        return f"{self.name}.sqlite"

Override if needed. Return the name of cache file Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.

use_cache: bool
101    @property
102    def use_cache(self) -> bool:
103        """
104        Override if needed. If True, all records will be cached.
105        Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.
106        """
107        return False

Override if needed. If True, all records will be cached. Note that if the environment variable REQUEST_CACHE_PATH is not set, the cache will be in-memory only.

url_base: str
109    @property
110    @abstractmethod
111    def url_base(self) -> str:
112        """
113        :return: URL base for the  API endpoint e.g: if you wanted to hit https://myapi.com/v1/some_entity then this should return "https://myapi.com/v1/"
114        """
Returns

URL base for the API endpoint e.g: if you wanted to hit https://myapi.com/v1/some_entity then this should return "https://myapi.com/v1/"

http_method: str
116    @property
117    def http_method(self) -> str:
118        """
119        Override if needed. See get_request_data/get_request_json if using POST/PUT/PATCH.
120        """
121        return "GET"

Override if needed. See get_request_data/get_request_json if using POST/PUT/PATCH.

raise_on_http_errors: bool
123    @property
124    @deprecated(
125        "Deprecated as of CDK version 3.0.0. "
126        "You should set error_handler explicitly in HttpStream.get_error_handler() instead."
127    )
128    def raise_on_http_errors(self) -> bool:
129        """
130        Override if needed. If set to False, allows opting-out of raising HTTP code exception.
131        """
132        return True

Override if needed. If set to False, allows opting-out of raising HTTP code exception.

max_retries: Optional[int]
134    @property
135    @deprecated(
136        "Deprecated as of CDK version 3.0.0. "
137        "You should set backoff_strategies explicitly in HttpStream.get_backoff_strategy() instead."
138    )
139    def max_retries(self) -> Union[int, None]:
140        """
141        Override if needed. Specifies maximum amount of retries for backoff policy. Return None for no limit.
142        """
143        return 5

Override if needed. Specifies maximum amount of retries for backoff policy. Return None for no limit.

max_time: Optional[int]
145    @property
146    @deprecated(
147        "Deprecated as of CDK version 3.0.0. "
148        "You should set backoff_strategies explicitly in HttpStream.get_backoff_strategy() instead."
149    )
150    def max_time(self) -> Union[int, None]:
151        """
152        Override if needed. Specifies maximum total waiting time (in seconds) for backoff policy. Return None for no limit.
153        """
154        return 60 * 10

Override if needed. Specifies maximum total waiting time (in seconds) for backoff policy. Return None for no limit.

retry_factor: float
156    @property
157    @deprecated(
158        "Deprecated as of CDK version 3.0.0. "
159        "You should set backoff_strategies explicitly in HttpStream.get_backoff_strategy() instead."
160    )
161    def retry_factor(self) -> float:
162        """
163        Override if needed. Specifies factor for backoff policy.
164        """
165        return 5

Override if needed. Specifies factor for backoff policy.

@abstractmethod
def next_page_token(self, response: requests.models.Response) -> Optional[Mapping[str, Any]]:
167    @abstractmethod
168    def next_page_token(self, response: requests.Response) -> Optional[Mapping[str, Any]]:
169        """
170        Override this method to define a pagination strategy.
171
172        The value returned from this method is passed to most other methods in this class. Use it to form a request e.g: set headers or query params.
173
174        :return: The token for the next page from the input response object. Returning None means there are no more pages to read in this response.
175        """

Override this method to define a pagination strategy.

The value returned from this method is passed to most other methods in this class. Use it to form a request e.g: set headers or query params.

Returns

The token for the next page from the input response object. Returning None means there are no more pages to read in this response.

@abstractmethod
def path( self, *, stream_state: Optional[Mapping[str, Any]] = None, stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> str:
177    @abstractmethod
178    def path(
179        self,
180        *,
181        stream_state: Optional[Mapping[str, Any]] = None,
182        stream_slice: Optional[Mapping[str, Any]] = None,
183        next_page_token: Optional[Mapping[str, Any]] = None,
184    ) -> str:
185        """
186        Returns the URL path for the API endpoint e.g: if you wanted to hit https://myapi.com/v1/some_entity then this should return "some_entity"
187        """

Returns the URL path for the API endpoint e.g: if you wanted to hit https://myapi.com/v1/some_entity then this should return "some_entity"

def request_params( self, stream_state: Optional[Mapping[str, Any]], stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> MutableMapping[str, Any]:
189    def request_params(
190        self,
191        stream_state: Optional[Mapping[str, Any]],
192        stream_slice: Optional[Mapping[str, Any]] = None,
193        next_page_token: Optional[Mapping[str, Any]] = None,
194    ) -> MutableMapping[str, Any]:
195        """
196        Override this method to define the query parameters that should be set on an outgoing HTTP request given the inputs.
197
198        E.g: you might want to define query parameters for paging if next_page_token is not None.
199        """
200        return {}

Override this method to define the query parameters that should be set on an outgoing HTTP request given the inputs.

E.g: you might want to define query parameters for paging if next_page_token is not None.

def request_headers( self, stream_state: Optional[Mapping[str, Any]], stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> Mapping[str, Any]:
202    def request_headers(
203        self,
204        stream_state: Optional[Mapping[str, Any]],
205        stream_slice: Optional[Mapping[str, Any]] = None,
206        next_page_token: Optional[Mapping[str, Any]] = None,
207    ) -> Mapping[str, Any]:
208        """
209        Override to return any non-auth headers. Authentication headers will overwrite any overlapping headers returned from this method.
210        """
211        return {}

Override to return any non-auth headers. Authentication headers will overwrite any overlapping headers returned from this method.

def request_body_data( self, stream_state: Optional[Mapping[str, Any]], stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> Union[str, Mapping[str, Any], NoneType]:
213    def request_body_data(
214        self,
215        stream_state: Optional[Mapping[str, Any]],
216        stream_slice: Optional[Mapping[str, Any]] = None,
217        next_page_token: Optional[Mapping[str, Any]] = None,
218    ) -> Optional[Union[Mapping[str, Any], str]]:
219        """
220        Override when creating POST/PUT/PATCH requests to populate the body of the request with a non-JSON payload.
221
222        If returns a ready text that it will be sent as is.
223        If returns a dict that it will be converted to a urlencoded form.
224        E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2"
225
226        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
227        """
228        return None

Override when creating POST/PUT/PATCH requests to populate the body of the request with a non-JSON payload.

If returns a ready text that it will be sent as is. If returns a dict that it will be converted to a urlencoded form. E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2"

At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.

def request_body_json( self, stream_state: Optional[Mapping[str, Any]], stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> Optional[Mapping[str, Any]]:
230    def request_body_json(
231        self,
232        stream_state: Optional[Mapping[str, Any]],
233        stream_slice: Optional[Mapping[str, Any]] = None,
234        next_page_token: Optional[Mapping[str, Any]] = None,
235    ) -> Optional[Mapping[str, Any]]:
236        """
237        Override when creating POST/PUT/PATCH requests to populate the body of the request with a JSON payload.
238
239        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
240        """
241        return None

Override when creating POST/PUT/PATCH requests to populate the body of the request with a JSON payload.

At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.

def request_kwargs( self, stream_state: Optional[Mapping[str, Any]], stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> Mapping[str, Any]:
243    def request_kwargs(
244        self,
245        stream_state: Optional[Mapping[str, Any]],
246        stream_slice: Optional[Mapping[str, Any]] = None,
247        next_page_token: Optional[Mapping[str, Any]] = None,
248    ) -> Mapping[str, Any]:
249        """
250        Override to return a mapping of keyword arguments to be used when creating the HTTP request.
251        Any option listed in https://docs.python-requests.org/en/latest/api/#requests.adapters.BaseAdapter.send for can be returned from
252        this method. Note that these options do not conflict with request-level options such as headers, request params, etc..
253        """
254        return {}

Override to return a mapping of keyword arguments to be used when creating the HTTP request. Any option listed in https://docs.python-requests.org/en/latest/api/#requests.adapters.BaseAdapter.send for can be returned from this method. Note that these options do not conflict with request-level options such as headers, request params, etc..

@abstractmethod
def parse_response( self, response: requests.models.Response, *, stream_state: Mapping[str, Any], stream_slice: Optional[Mapping[str, Any]] = None, next_page_token: Optional[Mapping[str, Any]] = None) -> Iterable[Mapping[str, Any]]:
256    @abstractmethod
257    def parse_response(
258        self,
259        response: requests.Response,
260        *,
261        stream_state: Mapping[str, Any],
262        stream_slice: Optional[Mapping[str, Any]] = None,
263        next_page_token: Optional[Mapping[str, Any]] = None,
264    ) -> Iterable[Mapping[str, Any]]:
265        """
266        Parses the raw response object into a list of records.
267        By default, this returns an iterable containing the input. Override to parse differently.
268        :param response:
269        :param stream_state:
270        :param stream_slice:
271        :param next_page_token:
272        :return: An iterable containing the parsed response
273        """

Parses the raw response object into a list of records. By default, this returns an iterable containing the input. Override to parse differently.

Parameters
  • response:
  • stream_state:
  • stream_slice:
  • next_page_token:
Returns

An iterable containing the parsed response

def get_backoff_strategy( self) -> Union[airbyte_cdk.BackoffStrategy, List[airbyte_cdk.BackoffStrategy], NoneType]:
275    def get_backoff_strategy(self) -> Optional[Union[BackoffStrategy, List[BackoffStrategy]]]:
276        """
277        Used to initialize Adapter to avoid breaking changes.
278        If Stream has a `backoff_time` method implementation, we know this stream uses old (pre-HTTPClient) backoff handlers and thus an adapter is needed.
279
280        Override to provide custom BackoffStrategy
281        :return Optional[BackoffStrategy]:
282        """
283        if hasattr(self, "backoff_time"):
284            return HttpStreamAdapterBackoffStrategy(self)
285        else:
286            return None

Used to initialize Adapter to avoid breaking changes. If Stream has a backoff_time method implementation, we know this stream uses old (pre-HTTPClient) backoff handlers and thus an adapter is needed.

Override to provide custom BackoffStrategy

Returns
def get_error_handler( self) -> Optional[airbyte_cdk.sources.streams.http.error_handlers.ErrorHandler]:
288    def get_error_handler(self) -> Optional[ErrorHandler]:
289        """
290        Used to initialize Adapter to avoid breaking changes.
291        If Stream has a `should_retry` method implementation, we know this stream uses old (pre-HTTPClient) error handlers and thus an adapter is needed.
292
293        Override to provide custom ErrorHandler
294        :return Optional[ErrorHandler]:
295        """
296        if hasattr(self, "should_retry"):
297            error_handler = HttpStreamAdapterHttpStatusErrorHandler(
298                stream=self,
299                logger=logging.getLogger(),
300                max_retries=self.max_retries,
301                max_time=timedelta(seconds=self.max_time or 0),
302            )
303            return error_handler
304        else:
305            return None

Used to initialize Adapter to avoid breaking changes. If Stream has a should_retry method implementation, we know this stream uses old (pre-HTTPClient) error handlers and thus an adapter is needed.

Override to provide custom ErrorHandler

Returns
@classmethod
def parse_response_error_message(cls, response: requests.models.Response) -> Optional[str]:
311    @classmethod
312    def parse_response_error_message(cls, response: requests.Response) -> Optional[str]:
313        """
314        Parses the raw response object from a failed request into a user-friendly error message.
315        By default, this method tries to grab the error message from JSON responses by following common API patterns. Override to parse differently.
316
317        :param response:
318        :return: A user-friendly message that indicates the cause of the error
319        """
320
321        # default logic to grab error from common fields
322        def _try_get_error(value: Optional[JsonType]) -> Optional[str]:
323            if isinstance(value, str):
324                return value
325            elif isinstance(value, list):
326                errors_in_value = [_try_get_error(v) for v in value]
327                return ", ".join(v for v in errors_in_value if v is not None)
328            elif isinstance(value, dict):
329                new_value = (
330                    value.get("message")
331                    or value.get("messages")
332                    or value.get("error")
333                    or value.get("errors")
334                    or value.get("failures")
335                    or value.get("failure")
336                    or value.get("detail")
337                )
338                return _try_get_error(new_value)
339            return None
340
341        try:
342            body = response.json()
343            return _try_get_error(body)
344        except requests.exceptions.JSONDecodeError:
345            return None

Parses the raw response object from a failed request into a user-friendly error message. By default, this method tries to grab the error message from JSON responses by following common API patterns. Override to parse differently.

Parameters
  • response:
Returns

A user-friendly message that indicates the cause of the error

def get_error_display_message(self, exception: BaseException) -> Optional[str]:
347    def get_error_display_message(self, exception: BaseException) -> Optional[str]:
348        """
349        Retrieves the user-friendly display message that corresponds to an exception.
350        This will be called when encountering an exception while reading records from the stream, and used to build the AirbyteTraceMessage.
351
352        The default implementation of this method only handles HTTPErrors by passing the response to self.parse_response_error_message().
353        The method should be overriden as needed to handle any additional exception types.
354
355        :param exception: The exception that was raised
356        :return: A user-friendly message that indicates the cause of the error
357        """
358        if isinstance(exception, requests.HTTPError) and exception.response is not None:
359            return self.parse_response_error_message(exception.response)
360        return None

Retrieves the user-friendly display message that corresponds to an exception. This will be called when encountering an exception while reading records from the stream, and used to build the AirbyteTraceMessage.

The default implementation of this method only handles HTTPErrors by passing the response to self.parse_response_error_message(). The method should be overriden as needed to handle any additional exception types.

Parameters
  • exception: The exception that was raised
Returns

A user-friendly message that indicates the cause of the error

def read_records( self, sync_mode: airbyte_protocol_dataclasses.models.airbyte_protocol.SyncMode, cursor_field: Optional[List[str]] = None, stream_slice: Optional[Mapping[str, Any]] = None, stream_state: Optional[Mapping[str, Any]] = None) -> Iterable[Union[Mapping[str, Any], airbyte_cdk.AirbyteMessage]]:
362    def read_records(
363        self,
364        sync_mode: SyncMode,
365        cursor_field: Optional[List[str]] = None,
366        stream_slice: Optional[Mapping[str, Any]] = None,
367        stream_state: Optional[Mapping[str, Any]] = None,
368    ) -> Iterable[StreamData]:
369        # A cursor_field indicates this is an incremental stream which offers better checkpointing than RFR enabled via the cursor
370        if self.cursor_field or not isinstance(self.get_cursor(), ResumableFullRefreshCursor):
371            yield from self._read_pages(
372                lambda req, res, state, _slice: self.parse_response(
373                    res, stream_slice=_slice, stream_state=state
374                ),
375                stream_slice,
376                stream_state,
377            )
378        else:
379            yield from self._read_single_page(
380                lambda req, res, state, _slice: self.parse_response(
381                    res, stream_slice=_slice, stream_state=state
382                ),
383                stream_slice,
384                stream_state,
385            )

This method should be overridden by subclasses to read records based on the inputs

state: MutableMapping[str, Any]
387    @property
388    def state(self) -> MutableMapping[str, Any]:
389        cursor = self.get_cursor()
390        if cursor:
391            return cursor.get_stream_state()  # type: ignore
392        return self._state

State getter, should return state in form that can serialized to a string and send to the output as a STATE AirbyteMessage.

A good example of a state is a cursor_value: { self.cursor_field: "cursor_value" }

State should try to be as small as possible but at the same time descriptive enough to restore syncing process from the point where it stopped.

def get_cursor(self) -> Optional[airbyte_cdk.sources.streams.checkpoint.Cursor]:
401    def get_cursor(self) -> Optional[Cursor]:
402        # I don't love that this is semi-stateful but not sure what else to do. We don't know exactly what type of cursor to
403        # instantiate when creating the class. We can make a few assumptions like if there is a cursor_field which implies
404        # incremental, but we don't know until runtime if this is a substream. Ideally, a stream should explicitly define
405        # its cursor, but because we're trying to automatically apply RFR we're stuck with this logic where we replace the
406        # cursor at runtime once we detect this is a substream based on self.has_multiple_slices being reassigned
407        if self.has_multiple_slices and isinstance(self.cursor, ResumableFullRefreshCursor):
408            self.cursor = SubstreamResumableFullRefreshCursor()
409            return self.cursor
410        else:
411            return self.cursor

A Cursor is an interface that a stream can implement to manage how its internal state is read and updated while reading records. Historically, Python connectors had no concept of a cursor to manage state. Python streams need to define a cursor implementation and override this method to manage state through a Cursor.

def get_log_formatter(self) -> Optional[Callable[[requests.models.Response], Any]]:
554    def get_log_formatter(self) -> Optional[Callable[[requests.Response], Any]]:
555        """
556
557        :return Optional[Callable[[requests.Response], Any]]: Function that will be used in logging inside HttpClient
558        """
559        return None
Returns

Function that will be used in logging inside HttpClient

class HttpSubStream(airbyte_cdk.sources.streams.http.HttpStream, abc.ABC):
562class HttpSubStream(HttpStream, ABC):
563    def __init__(self, parent: HttpStream, **kwargs: Any):
564        """
565        :param parent: should be the instance of HttpStream class
566        """
567        super().__init__(**kwargs)
568        self.parent = parent
569        self.has_multiple_slices = (
570            True  # Substreams are based on parent records which implies there are multiple slices
571        )
572
573        # There are three conditions that dictate if RFR should automatically be applied to a stream
574        # 1. Streams that explicitly initialize their own cursor should defer to it and not automatically apply RFR
575        # 2. Streams with at least one cursor_field are incremental and thus a superior sync to RFR.
576        # 3. Streams overriding read_records() do not guarantee that they will call the parent implementation which can perform
577        #    per-page checkpointing so RFR is only supported if a stream use the default `HttpStream.read_records()` method
578        if (
579            not self.cursor
580            and len(self.cursor_field) == 0
581            and type(self).read_records is HttpStream.read_records
582        ):
583            self.cursor = SubstreamResumableFullRefreshCursor()
584
585    def stream_slices(
586        self,
587        sync_mode: SyncMode,
588        cursor_field: Optional[List[str]] = None,
589        stream_state: Optional[Mapping[str, Any]] = None,
590    ) -> Iterable[Optional[Mapping[str, Any]]]:
591        # read_stateless() assumes the parent is not concurrent. This is currently okay since the concurrent CDK does
592        # not support either substreams or RFR, but something that needs to be considered once we do
593        for parent_record in self.parent.read_only_records(stream_state):
594            # Skip non-records (eg AirbyteLogMessage)
595            if isinstance(parent_record, AirbyteMessage):
596                if parent_record.type == MessageType.RECORD:
597                    parent_record = parent_record.record.data  # type: ignore [assignment, union-attr]  # Incorrect type for assignment
598                else:
599                    continue
600            elif isinstance(parent_record, Record):
601                parent_record = parent_record.data
602            yield {"parent": parent_record}

Base abstract class for an Airbyte Stream using the HTTP protocol. Basic building block for users building an Airbyte source for a HTTP API.

HttpSubStream( parent: HttpStream, **kwargs: Any)
563    def __init__(self, parent: HttpStream, **kwargs: Any):
564        """
565        :param parent: should be the instance of HttpStream class
566        """
567        super().__init__(**kwargs)
568        self.parent = parent
569        self.has_multiple_slices = (
570            True  # Substreams are based on parent records which implies there are multiple slices
571        )
572
573        # There are three conditions that dictate if RFR should automatically be applied to a stream
574        # 1. Streams that explicitly initialize their own cursor should defer to it and not automatically apply RFR
575        # 2. Streams with at least one cursor_field are incremental and thus a superior sync to RFR.
576        # 3. Streams overriding read_records() do not guarantee that they will call the parent implementation which can perform
577        #    per-page checkpointing so RFR is only supported if a stream use the default `HttpStream.read_records()` method
578        if (
579            not self.cursor
580            and len(self.cursor_field) == 0
581            and type(self).read_records is HttpStream.read_records
582        ):
583            self.cursor = SubstreamResumableFullRefreshCursor()
Parameters
  • parent: should be the instance of HttpStream class
parent
has_multiple_slices = False
def stream_slices( self, sync_mode: airbyte_protocol_dataclasses.models.airbyte_protocol.SyncMode, cursor_field: Optional[List[str]] = None, stream_state: Optional[Mapping[str, Any]] = None) -> Iterable[Optional[Mapping[str, Any]]]:
585    def stream_slices(
586        self,
587        sync_mode: SyncMode,
588        cursor_field: Optional[List[str]] = None,
589        stream_state: Optional[Mapping[str, Any]] = None,
590    ) -> Iterable[Optional[Mapping[str, Any]]]:
591        # read_stateless() assumes the parent is not concurrent. This is currently okay since the concurrent CDK does
592        # not support either substreams or RFR, but something that needs to be considered once we do
593        for parent_record in self.parent.read_only_records(stream_state):
594            # Skip non-records (eg AirbyteLogMessage)
595            if isinstance(parent_record, AirbyteMessage):
596                if parent_record.type == MessageType.RECORD:
597                    parent_record = parent_record.record.data  # type: ignore [assignment, union-attr]  # Incorrect type for assignment
598                else:
599                    continue
600            elif isinstance(parent_record, Record):
601                parent_record = parent_record.data
602            yield {"parent": parent_record}

Override to define the slices for this stream. See the stream slicing section of the docs for more information.

Parameters
  • sync_mode:
  • cursor_field:
  • stream_state:
Returns
class UserDefinedBackoffException(airbyte_cdk.sources.streams.http.exceptions.BaseBackoffException):
40class UserDefinedBackoffException(BaseBackoffException):
41    """
42    An exception that exposes how long it attempted to backoff
43    """
44
45    def __init__(
46        self,
47        backoff: Union[int, float],
48        request: requests.PreparedRequest,
49        response: Optional[Union[requests.Response, Exception]],
50        error_message: str = "",
51        failure_type: Optional[FailureType] = None,
52    ):
53        """
54        :param backoff: how long to backoff in seconds
55        :param request: the request that triggered this backoff exception
56        :param response: the response that triggered the backoff exception
57        """
58        self.backoff = backoff
59        super().__init__(
60            request=request,
61            response=response,
62            error_message=error_message,
63            failure_type=failure_type,
64        )

An exception that exposes how long it attempted to backoff

UserDefinedBackoffException( backoff: Union[int, float], request: requests.models.PreparedRequest, response: Union[requests.models.Response, Exception, NoneType], error_message: str = '', failure_type: Optional[airbyte_protocol_dataclasses.models.airbyte_protocol.FailureType] = None)
45    def __init__(
46        self,
47        backoff: Union[int, float],
48        request: requests.PreparedRequest,
49        response: Optional[Union[requests.Response, Exception]],
50        error_message: str = "",
51        failure_type: Optional[FailureType] = None,
52    ):
53        """
54        :param backoff: how long to backoff in seconds
55        :param request: the request that triggered this backoff exception
56        :param response: the response that triggered the backoff exception
57        """
58        self.backoff = backoff
59        super().__init__(
60            request=request,
61            response=response,
62            error_message=error_message,
63            failure_type=failure_type,
64        )
Parameters
  • backoff: how long to backoff in seconds
  • request: the request that triggered this backoff exception
  • response: the response that triggered the backoff exception
backoff