airbyte_cdk.sources.declarative.retrievers.simple_retriever

  1#
  2# Copyright (c) 2025 Airbyte, Inc., all rights reserved.
  3#
  4
  5import datetime
  6import json
  7import logging
  8from collections import defaultdict
  9from dataclasses import InitVar, dataclass, field
 10from functools import partial
 11from typing import (
 12    Any,
 13    Callable,
 14    Iterable,
 15    List,
 16    Mapping,
 17    MutableMapping,
 18    Optional,
 19    Set,
 20    Tuple,
 21    Union,
 22)
 23
 24import requests
 25from typing_extensions import deprecated
 26
 27from airbyte_cdk.models import FailureType
 28from airbyte_cdk.sources.declarative.extractors.http_selector import HttpSelector
 29from airbyte_cdk.sources.declarative.extractors.record_filter import (
 30    ClientSideIncrementalRecordFilterDecorator,
 31)
 32from airbyte_cdk.sources.declarative.interpolation import InterpolatedString
 33from airbyte_cdk.sources.declarative.partition_routers.single_partition_router import (
 34    SinglePartitionRouter,
 35)
 36from airbyte_cdk.sources.declarative.requesters.paginators.no_pagination import NoPagination
 37from airbyte_cdk.sources.declarative.requesters.paginators.paginator import (
 38    Paginator,
 39    page_size_override_kwargs,
 40    stream_slice_kwargs,
 41)
 42from airbyte_cdk.sources.declarative.requesters.query_properties import QueryProperties
 43from airbyte_cdk.sources.declarative.requesters.request_options import (
 44    DefaultRequestOptionsProvider,
 45    RequestOptionsProvider,
 46)
 47from airbyte_cdk.sources.declarative.requesters.requester import Requester
 48from airbyte_cdk.sources.declarative.retrievers.page_size_reducer import (
 49    PageSizeReducer,
 50    PageSizeReduction,
 51)
 52from airbyte_cdk.sources.declarative.retrievers.pagination_tracker import PaginationTracker
 53from airbyte_cdk.sources.declarative.retrievers.request_window_splitting import (
 54    RequestWindowSplitting,
 55)
 56from airbyte_cdk.sources.declarative.retrievers.retriever import Retriever
 57from airbyte_cdk.sources.declarative.stream_slicers.stream_slicer import StreamSlicer
 58from airbyte_cdk.sources.source import ExperimentalClassWarning
 59from airbyte_cdk.sources.streams.core import StreamData
 60from airbyte_cdk.sources.streams.http.page_size_reduction_exception import (
 61    PageSizeReductionNotSupportedException,
 62    PageSizeReductionRequiredException,
 63)
 64from airbyte_cdk.sources.streams.http.pagination_reset_exception import (
 65    PaginationResetRequiredException,
 66)
 67from airbyte_cdk.sources.streams.http.request_window_split_exception import (
 68    RequestWindowSplitNotSupportedException,
 69    RequestWindowSplitRequiredException,
 70)
 71from airbyte_cdk.sources.types import Config, Record, StreamSlice
 72from airbyte_cdk.utils.mapping_helpers import combine_mappings
 73from airbyte_cdk.utils.traced_exception import AirbyteTracedException
 74
 75FULL_REFRESH_SYNC_COMPLETE_KEY = "__ab_full_refresh_sync_complete"
 76LOGGER = logging.getLogger("airbyte")
 77
 78# A defense-in-depth bound independent of any `request_window_splitter`'s own no-progress guard: protects
 79# against a misbehaving custom cursor whose children don't actually shrink the window. Not user-configurable.
 80_MAX_REQUEST_WINDOW_SPLIT_DEPTH = 10
 81
 82
 83@dataclass
 84class SimpleRetriever(Retriever):
 85    """
 86    Retrieves records by synchronously sending requests to fetch records.
 87
 88    The retriever acts as an orchestrator between the requester, the record selector, the paginator, and the stream slicer.
 89
 90    For each stream slice, submit requests until there are no more pages of records to fetch.
 91
 92    This retriever currently inherits from HttpStream to reuse the request submission and pagination machinery.
 93    As a result, some of the parameters passed to some methods are unused.
 94    The two will be decoupled in a future release.
 95
 96    Attributes:
 97        stream_name (str): The stream's name
 98        stream_primary_key (Optional[Union[str, List[str], List[List[str]]]]): The stream's primary key
 99        requester (Requester): The HTTP requester
100        record_selector (HttpSelector): The record selector
101        paginator (Optional[Paginator]): The paginator
102        stream_slicer (Optional[StreamSlicer]): The stream slicer
103        parameters (Mapping[str, Any]): Additional runtime parameters to be used for string interpolation
104        post_pagination_filter (Optional[ClientSideIncrementalRecordFilterDecorator]): Set for data feed streams only.
105            Records the cursor considers already synced are dropped once pagination has observed them
106        page_size_reduction (Optional[PageSizeReduction]): How much to shrink the page size when an error handler
107            resolves to `ResponseAction.REDUCE_PAGE_SIZE`. `None` disables page size reduction entirely.
108            It is immutable configuration; the page size in effect lives in a `PageSizeReducer` that
109            `_read_pages` creates per call, so the retriever and its paginator - both shared by every
110            partition of the stream, read concurrently - stay stateless. When `page_size_reduction` is
111            `None` no reducer is created and `_read_pages` keeps its previous behaviour
112        request_window_splitting (Optional[RequestWindowSplitting]): Policy applied when an error handler
113            resolves to `ResponseAction.SPLIT_REQUEST_WINDOW`, or when `RequestWindowSplitRequiredException`
114            is raised directly by custom code. `None` disables request-window splitting entirely: `read_records`
115            converts the exception into `RequestWindowSplitNotSupportedException` instead of attempting a
116            split it cannot perform.
117        request_window_splitter: Bound method - normally the stream's own cursor's `split_request_window` -
118            asked to replace a failing `StreamSlice` with smaller children. `None` has the same effect as
119            `request_window_splitting` being `None`.
120    """
121
122    requester: Requester
123    record_selector: HttpSelector
124    config: Config
125    parameters: InitVar[Mapping[str, Any]]
126    name: str
127    _name: Union[InterpolatedString, str] = field(init=False, repr=False, default="")
128    primary_key: Optional[Union[str, List[str], List[List[str]]]]
129    _primary_key: str = field(init=False, repr=False, default="")
130    paginator: Optional[Paginator] = None
131    stream_slicer: StreamSlicer = field(
132        default_factory=lambda: SinglePartitionRouter(parameters={})
133    )
134    request_option_provider: RequestOptionsProvider = field(
135        default_factory=lambda: DefaultRequestOptionsProvider(parameters={})
136    )
137    ignore_stream_slicer_parameters_on_paginated_requests: bool = False
138    additional_query_properties: Optional[QueryProperties] = None
139    log_formatter: Optional[Callable[[requests.Response], Any]] = None
140    pagination_tracker_factory: Callable[[], PaginationTracker] = field(
141        default_factory=lambda: lambda: PaginationTracker()
142    )
143    post_pagination_filter: Optional[ClientSideIncrementalRecordFilterDecorator] = None
144    page_size_reduction: Optional[PageSizeReduction] = None
145    request_window_splitting: Optional[RequestWindowSplitting] = None
146    request_window_splitter: Optional[
147        Callable[[StreamSlice, Optional[datetime.timedelta]], Optional[List[StreamSlice]]]
148    ] = None
149
150    def __post_init__(self, parameters: Mapping[str, Any]) -> None:
151        self._paginator = self.paginator or NoPagination(parameters=parameters)
152        self._parameters = parameters
153        self._name = (
154            InterpolatedString(self._name, parameters=parameters)
155            if isinstance(self._name, str)
156            else self._name
157        )
158
159    @property  # type: ignore
160    def name(self) -> str:
161        """
162        :return: Stream name
163        """
164        return (
165            str(self._name.eval(self.config))
166            if isinstance(self._name, InterpolatedString)
167            else self._name
168        )
169
170    @name.setter
171    def name(self, value: str) -> None:
172        if not isinstance(value, property):
173            self._name = value
174
175    def _get_mapping(
176        self, method: Callable[..., Optional[Union[Mapping[str, Any], str]]], **kwargs: Any
177    ) -> Tuple[Union[Mapping[str, Any], str], Set[str]]:
178        """
179        Get mapping from the provided method, and get the keys of the mapping.
180        If the method returns a string, it will return the string and an empty set.
181        If the method returns a dict, it will return the dict and its keys.
182        """
183        mapping = method(**kwargs) or {}
184        keys = set(mapping.keys()) if not isinstance(mapping, str) else set()
185        return mapping, keys
186
187    def _get_request_options(
188        self,
189        stream_slice: Optional[StreamSlice],
190        next_page_token: Optional[Mapping[str, Any]],
191        paginator_method: Callable[..., Optional[Union[Mapping[str, Any], str]]],
192        stream_slicer_method: Callable[..., Optional[Union[Mapping[str, Any], str]]],
193        page_size_override: Optional[int] = None,
194    ) -> Union[Mapping[str, Any], str]:
195        """
196        Get the request_option from the paginator and the stream slicer.
197        Raise a ValueError if there's a key collision
198        Returned merged mapping otherwise
199        """
200        is_body_json = paginator_method.__name__ == "get_request_body_json"
201
202        mappings = [
203            paginator_method(
204                stream_slice=stream_slice,
205                next_page_token=next_page_token,
206                **page_size_override_kwargs(page_size_override),
207            ),
208        ]
209        if not next_page_token or not self.ignore_stream_slicer_parameters_on_paginated_requests:
210            mappings.append(
211                stream_slicer_method(
212                    stream_slice=stream_slice,
213                    next_page_token=next_page_token,
214                )
215            )
216        return combine_mappings(mappings, allow_same_value_merge=is_body_json)
217
218    def _request_headers(
219        self,
220        stream_slice: Optional[StreamSlice] = None,
221        next_page_token: Optional[Mapping[str, Any]] = None,
222        page_size_override: Optional[int] = None,
223    ) -> Mapping[str, Any]:
224        """
225        Specifies request headers.
226        Authentication headers will overwrite any overlapping headers returned from this method.
227        """
228        headers = self._get_request_options(
229            stream_slice,
230            next_page_token,
231            self._paginator.get_request_headers,
232            self.request_option_provider.get_request_headers,
233            **page_size_override_kwargs(page_size_override),
234        )
235        if isinstance(headers, str):
236            raise ValueError("Request headers cannot be a string")
237        return {str(k): str(v) for k, v in headers.items()}
238
239    def _request_params(
240        self,
241        stream_slice: Optional[StreamSlice] = None,
242        next_page_token: Optional[Mapping[str, Any]] = None,
243        page_size_override: Optional[int] = None,
244    ) -> Mapping[str, Any]:
245        """
246        Specifies the query parameters that should be set on an outgoing HTTP request given the inputs.
247
248        E.g: you might want to define query parameters for paging if next_page_token is not None.
249        """
250        params = self._get_request_options(
251            stream_slice,
252            next_page_token,
253            self._paginator.get_request_params,
254            self.request_option_provider.get_request_params,
255            **page_size_override_kwargs(page_size_override),
256        )
257        if isinstance(params, str):
258            raise ValueError("Request params cannot be a string")
259        return params
260
261    def _request_body_data(
262        self,
263        stream_slice: Optional[StreamSlice] = None,
264        next_page_token: Optional[Mapping[str, Any]] = None,
265        page_size_override: Optional[int] = None,
266    ) -> Union[Mapping[str, Any], str]:
267        """
268        Specifies how to populate the body of the request with a non-JSON payload.
269
270        If returns a ready text that it will be sent as is.
271        If returns a dict that it will be converted to a urlencoded form.
272        E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2"
273
274        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
275        """
276        return self._get_request_options(
277            stream_slice,
278            next_page_token,
279            self._paginator.get_request_body_data,
280            self.request_option_provider.get_request_body_data,
281            **page_size_override_kwargs(page_size_override),
282        )
283
284    def _request_body_json(
285        self,
286        stream_slice: Optional[StreamSlice] = None,
287        next_page_token: Optional[Mapping[str, Any]] = None,
288        page_size_override: Optional[int] = None,
289    ) -> Optional[Mapping[str, Any]]:
290        """
291        Specifies how to populate the body of the request with a JSON payload.
292
293        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
294        """
295        body_json = self._get_request_options(
296            stream_slice,
297            next_page_token,
298            self._paginator.get_request_body_json,
299            self.request_option_provider.get_request_body_json,
300            **page_size_override_kwargs(page_size_override),
301        )
302        if isinstance(body_json, str):
303            raise ValueError("Request body json cannot be a string")
304        return body_json
305
306    def _paginator_path(
307        self,
308        next_page_token: Optional[Mapping[str, Any]] = None,
309        stream_slice: Optional[StreamSlice] = None,
310    ) -> Optional[str]:
311        """
312        If the paginator points to a path, follow it, else return nothing so the requester is used.
313        :param next_page_token:
314        :return:
315        """
316        return self._paginator.path(
317            next_page_token=next_page_token,
318            stream_state={},  # stream_state as an interpolation context is deprecated
319            stream_slice=stream_slice,
320        )
321
322    def _parse_response(
323        self,
324        response: Optional[requests.Response],
325        records_schema: Mapping[str, Any],
326        stream_slice: Optional[StreamSlice] = None,
327        next_page_token: Optional[Mapping[str, Any]] = None,
328    ) -> Iterable[Record]:
329        if not response:
330            yield from []
331        else:
332            yield from self.record_selector.select_records(
333                response=response,
334                stream_state={},  # stream_state as an interpolation context is deprecated
335                records_schema=records_schema,
336                stream_slice=stream_slice,
337                next_page_token=next_page_token,
338            )
339
340    @property  # type: ignore
341    def primary_key(self) -> Optional[Union[str, List[str], List[List[str]]]]:
342        """The stream's primary key"""
343        return self._primary_key
344
345    @primary_key.setter
346    def primary_key(self, value: str) -> None:
347        if not isinstance(value, property):
348            self._primary_key = value
349
350    def _next_page_token(
351        self,
352        response: requests.Response,
353        last_page_size: int,
354        last_record: Optional[Record],
355        last_page_token_value: Optional[Any],
356        page_size_override: Optional[int] = None,
357        stream_slice: Optional[StreamSlice] = None,
358    ) -> Optional[Mapping[str, Any]]:
359        """
360        Specifies a pagination strategy.
361
362        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.
363
364        :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.
365        """
366        return self._paginator.next_page_token(
367            response=response,
368            last_page_size=last_page_size,
369            last_record=last_record,
370            last_page_token_value=last_page_token_value,
371            **page_size_override_kwargs(page_size_override),
372            **stream_slice_kwargs(self._paginator.next_page_token, stream_slice),
373        )
374
375    def _fetch_next_page(
376        self,
377        stream_slice: StreamSlice,
378        next_page_token: Optional[Mapping[str, Any]] = None,
379        page_size_override: Optional[int] = None,
380    ) -> Optional[requests.Response]:
381        return self.requester.send_request(
382            path=self._paginator_path(
383                next_page_token=next_page_token,
384                stream_slice=stream_slice,
385            ),
386            stream_state={},  # stream_state as an interpolation context is deprecated
387            stream_slice=stream_slice,
388            next_page_token=next_page_token,
389            request_headers=self._request_headers(
390                stream_slice=stream_slice,
391                next_page_token=next_page_token,
392                **page_size_override_kwargs(page_size_override),
393            ),
394            request_params=self._request_params(
395                stream_slice=stream_slice,
396                next_page_token=next_page_token,
397                **page_size_override_kwargs(page_size_override),
398            ),
399            request_body_data=self._request_body_data(
400                stream_slice=stream_slice,
401                next_page_token=next_page_token,
402                **page_size_override_kwargs(page_size_override),
403            ),
404            request_body_json=self._request_body_json(
405                stream_slice=stream_slice,
406                next_page_token=next_page_token,
407                **page_size_override_kwargs(page_size_override),
408            ),
409            log_formatter=self.log_formatter,
410        )
411
412    # This logic is similar to _read_pages in the HttpStream class. When making changes here, consider making changes there as well.
413    def _read_pages(
414        self,
415        records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]],
416        stream_slice: StreamSlice,
417    ) -> Iterable[Record]:
418        original_stream_slice = stream_slice
419        pagination_tracker = self.pagination_tracker_factory()
420        page_size_reducer = (
421            PageSizeReducer(
422                self.page_size_reduction,
423                self._paginator.get_page_size(),
424                stream_name=self.name,
425            )
426            if self.page_size_reduction
427            else None
428        )
429        reset_pagination = False
430        reduce_page_size = False
431        next_page_token = self._get_initial_next_page_token()
432        while True:
433            page_size_override = page_size_reducer.page_size_override if page_size_reducer else None
434            merged_records: MutableMapping[str, Any] = defaultdict(dict)
435            last_page_size = 0
436            last_record: Optional[Record] = None
437
438            response = None
439            try:
440                if self.additional_query_properties:
441                    for (
442                        properties
443                    ) in self.additional_query_properties.get_request_property_chunks():
444                        stream_slice = StreamSlice(
445                            partition=stream_slice.partition or {},
446                            cursor_slice=stream_slice.cursor_slice or {},
447                            extra_fields={"query_properties": properties},
448                        )
449                        response = self._fetch_next_page(
450                            stream_slice,
451                            next_page_token,
452                            **page_size_override_kwargs(page_size_override),
453                        )
454
455                        for current_record in records_generator_fn(response):
456                            if self.additional_query_properties.property_chunking:
457                                merge_key = self.additional_query_properties.property_chunking.get_merge_key(
458                                    current_record
459                                )
460                                if merge_key:
461                                    _deep_merge(merged_records[merge_key], current_record)
462                                else:
463                                    # We should still emit records even if the record did not have a merge key
464                                    pagination_tracker.observe(current_record)
465                                    last_page_size += 1
466                                    last_record = current_record
467                                    yield current_record
468                            else:
469                                pagination_tracker.observe(current_record)
470                                last_page_size += 1
471                                last_record = current_record
472                                yield current_record
473
474                    for merged_record in merged_records.values():
475                        record = Record(
476                            data=merged_record, stream_name=self.name, associated_slice=stream_slice
477                        )
478                        pagination_tracker.observe(record)
479                        last_page_size += 1
480                        last_record = record
481                        yield record
482                else:
483                    response = self._fetch_next_page(
484                        stream_slice,
485                        next_page_token,
486                        **page_size_override_kwargs(page_size_override),
487                    )
488                    for current_record in records_generator_fn(response):
489                        pagination_tracker.observe(current_record)
490                        last_page_size += 1
491                        last_record = current_record
492                        yield current_record
493            except PaginationResetRequiredException:
494                reset_pagination = True
495            except PageSizeReductionRequiredException:
496                if page_size_reducer is None:
497                    # The action can be attached to a requester we cannot validate at config time, such as one
498                    # built by a custom error handler or a CustomRequester. Re-raise as the misconfiguration it
499                    # is: the exception being handled is the neutral "the API asked for a smaller page" signal.
500                    raise PageSizeReductionNotSupportedException(stream_name=self.name)
501                if last_page_size:
502                    # The reduction is safe only because it re-issues a page whose records were not emitted.
503                    # Every in-CDK way of reaching this raises from `_fetch_next_page`, before the record loop,
504                    # and the factory rejects the manifest constructs that would not, but a custom extractor,
505                    # filter or transformation can issue its own request from inside the record generator.
506                    # Re-issuing the page then duplicates the records already yielded, so fail instead.
507                    raise AirbyteTracedException(
508                        internal_message=f"Stream {self.name} requested a page size reduction after {last_page_size} records of the page had already been emitted",
509                        message=f"Stream {self.name} asked for a smaller page size in the middle of a page. The page cannot be requested again without duplicating the records already read from it. Move the REDUCE_PAGE_SIZE action to the error handler of the stream's main requester.",
510                        failure_type=FailureType.config_error,
511                    )
512                # Raises once the page size cannot be reduced any further and the retries allowed at that
513                # floor are spent, which is what stops the loop when the API keeps failing.
514                page_size_reducer.reduce()
515                reduce_page_size = True
516            else:
517                if page_size_reducer:
518                    page_size_reducer.on_successful_page()
519                if not response:
520                    break
521
522            if reduce_page_size:
523                # Retry the very same page: neither the token nor the slice change, only the page size does -
524                # and not even that once the reducer is at its floor and only waiting is left.
525                reduce_page_size = False
526                continue
527
528            if reset_pagination or pagination_tracker.has_reached_limit():
529                next_page_token = self._get_initial_next_page_token()
530                previous_slice = stream_slice
531                stream_slice = pagination_tracker.reduce_slice_range_if_possible(
532                    stream_slice, original_stream_slice
533                )
534                LOGGER.info(
535                    f"Hitting PaginationReset event. StreamSlice used will go from {previous_slice} to {stream_slice}"
536                )
537                reset_pagination = False
538            else:
539                last_page_token_value = (
540                    next_page_token.get("next_page_token") if next_page_token else None
541                )
542                next_page_token = self._next_page_token(
543                    response=response,  # type:ignore # we are breaking from the loop on the try/else if there are no response so this should be fine
544                    last_page_size=last_page_size,
545                    last_record=last_record,
546                    last_page_token_value=last_page_token_value,
547                    **page_size_override_kwargs(page_size_override),
548                    **stream_slice_kwargs(self._next_page_token, stream_slice),
549                )
550                if not next_page_token:
551                    break
552
553        # Always return an empty generator just in case no records were ever yielded
554        yield from []
555
556    def _get_initial_next_page_token(self) -> Optional[Mapping[str, Any]]:
557        initial_token = self._paginator.get_initial_token()
558        next_page_token = {"next_page_token": initial_token} if initial_token is not None else None
559        return next_page_token
560
561    def read_records(
562        self,
563        records_schema: Mapping[str, Any],
564        stream_slice: Optional[StreamSlice] = None,
565    ) -> Iterable[StreamData]:
566        """
567        Fetch a stream's records from an HTTP API source
568
569        :param records_schema: json schema to describe record
570        :param stream_slice: The stream slice to read data for
571        :return: The records read from the API source
572        """
573        _slice = stream_slice or StreamSlice(partition={}, cursor_slice={})  # None-check
574        yield from self._read_records_or_split_request_window(
575            records_schema, _slice, _slice, depth=0
576        )
577
578    def _read_records_or_split_request_window(
579        self,
580        records_schema: Mapping[str, Any],
581        stream_slice: StreamSlice,
582        original_slice: StreamSlice,
583        depth: int,
584    ) -> Iterable[StreamData]:
585        """
586        Read `stream_slice` to completion, replacing it with smaller children and recursing into each of them in
587        turn when a `RequestWindowSplitRequiredException` is raised while reading it.
588
589        Each recursive call re-enters this method - and, through it, `_read_pages` - from scratch for the child
590        slice it is given, so every child gets its own paginator token, `PaginationTracker`, and
591        `PageSizeReducer` for free: nothing here needs to reset that state explicitly. Because the whole
592        recursion lives inside the single generator `DeclarativePartition.read()` consumes,
593        `PartitionReader.process_partition()` only calls `cursor.close_partition()` once, after this generator is
594        fully exhausted - so a failure anywhere in the recursion (a child that cannot be read, or a window that
595        cannot be split any further) propagates out without ever checkpointing the original partition.
596
597        `original_slice` is threaded through the recursion unchanged so every record yielded, however deep the
598        recursion went to produce it, is re-stamped with the partition's own slice rather than the child slice it
599        was actually read against - see `_reassociate_with_original_slice`. `depth` bounds the recursion
600        independently of `request_window_splitter`'s own no-progress guard: it is enforced here so a
601        misbehaving custom cursor cannot recurse indefinitely regardless of what that guard does or does not
602        catch.
603        """
604        record_generator = partial(
605            self._parse_records,
606            stream_slice=stream_slice,
607            records_schema=records_schema,
608        )
609
610        emitted_count = 0
611        try:
612            records: Iterable[Mapping[str, Any]] = self._read_pages(record_generator, stream_slice)
613            if self.post_pagination_filter:
614                # A data feed paginates until it reaches a record older than the cursor, so the page that triggers the stop
615                # condition still holds already-synced records. Those are filtered here rather than in the record selector
616                # so that the paginator keeps seeing the whole page: the stop condition is evaluated on the last record of
617                # the page, which is precisely one of the records being dropped. Two consequences of filtering this late:
618                # the pagination tracker observes the dropped records, and a `file_uploader` on the record selector has
619                # already uploaded their files by the time they are dropped.
620                records = self.post_pagination_filter.filter_records(
621                    records,
622                    # the filter is only used for its cursor comparison, which does not read the stream state
623                    stream_state={},
624                    stream_slice=stream_slice,
625                )
626            for record in records:
627                emitted_count += 1
628                yield self._reassociate_with_original_slice(record, original_slice)
629        except RequestWindowSplitRequiredException as exception:
630            if self.request_window_splitting is None or self.request_window_splitter is None:
631                raise RequestWindowSplitNotSupportedException(stream_name=self.name) from exception
632
633            if emitted_count:
634                # Splitting and re-reading the window re-emits these records: a failed partition is never
635                # checkpointed, so the next attempt would re-emit them anyway. Always replaying makes forward
636                # progress instead of retrying the same oversized window forever; primary-key dedup downstream
637                # handles the duplicates either way.
638                LOGGER.warning(
639                    f"Stream {self.name} already emitted {emitted_count} record(s) from {stream_slice} before "
640                    f"it was rejected; splitting and re-reading the window may re-emit them."
641                )
642
643            if depth >= _MAX_REQUEST_WINDOW_SPLIT_DEPTH:
644                # Logged separately from the exception below so it stays visible even if the trace message's
645                # internal_message isn't surfaced by whatever catches it.
646                LOGGER.warning(
647                    f"Stream {self.name} hit the maximum request window split depth "
648                    f"({_MAX_REQUEST_WINDOW_SPLIT_DEPTH}) while splitting {original_slice}."
649                )
650                raise AirbyteTracedException(
651                    internal_message=f"Stream {self.name} exceeded the maximum request window split depth of {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} while splitting {original_slice}",
652                    # `transient_error`, so the only remediation is the connector's own, if it defined one.
653                    message=self._with_request_window_failure_message(
654                        f"Stream {self.name} could not split its request window to a size the API accepts "
655                        f"within {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} splits. This usually means the stream's "
656                        f"cursor is not actually shrinking the window on each split; if it uses a custom "
657                        f"cursor, check its `split_request_window` implementation."
658                    ),
659                    failure_type=FailureType.transient_error,
660                ) from exception
661
662            children = self.request_window_splitter(
663                stream_slice, self.request_window_splitting.min_split_window
664            )
665            if children is None:
666                min_split_window_note = (
667                    f", its configured `min_split_window` ({self.request_window_splitting.min_split_window})"
668                    if self.request_window_splitting.min_split_window
669                    else ""
670                )
671                raise AirbyteTracedException(
672                    internal_message=f"Stream {self.name} could not split its request window {stream_slice} any further",
673                    # `transient_error` regardless of how the triggering response was classified: exhaustion is
674                    # never the user's fault, matching `PageSizeReducer`'s equivalent exhaustion branch.
675                    message=self._with_request_window_failure_message(
676                        f"The API kept rejecting stream {self.name}'s request window even at the smallest window "
677                        f"its cursor allows{min_split_window_note}, or splitting the window further would not "
678                        f"make progress."
679                    ),
680                    failure_type=FailureType.transient_error,
681                ) from exception
682
683            LOGGER.info(
684                f"Stream {self.name}: the API rejected request window {stream_slice} (split depth {depth}); "
685                f"reducing it to {children} and reading each in turn."
686            )
687            for child in children:
688                yield from self._read_records_or_split_request_window(
689                    records_schema, child, original_slice, depth + 1
690                )
691
692    def _with_request_window_failure_message(self, message: str) -> str:
693        """
694        :return: the message followed by `request_window_splitting`'s configured `failure_message`, when it
695            defined one - mirroring `PageSizeReducer._with_failure_message`
696        """
697        assert self.request_window_splitting is not None
698        failure_message = (self.request_window_splitting.failure_message or "").strip()
699        if not failure_message:
700            return message
701        return f"{message} {failure_message}"
702
703    @staticmethod
704    def _reassociate_with_original_slice(
705        record: StreamData, original_slice: StreamSlice
706    ) -> StreamData:
707        """
708        Records read while recursing into a split child window are parsed against that child's `StreamSlice`,
709        so `RecordSelector` stamps them with the child as `associated_slice`. `DeclarativePartition.read()`
710        passes an already-built `Record` through unchanged rather than re-wrapping it, so left uncorrected the
711        child slice - not the partition's own slice - is what `ConcurrentCursor.observe()` would key its
712        per-partition bookkeeping by. `close_partition()` always looks that bookkeeping up by the partition's own
713        slice, so it would never find it, silently losing the precise most-recently-observed cursor value for a
714        split partition (recoverable in practice today only because callers of that lookup fall back to the
715        slice's end boundary when it is missing). Re-stamping every record with `original_slice` here restores
716        the same tracking a non-split read already gets, rather than relying on that fallback.
717        """
718        if isinstance(record, Record) and record.associated_slice is not original_slice:
719            return Record(
720                data=record.data,
721                stream_name=record.stream_name,
722                associated_slice=original_slice,
723                file_reference=record.file_reference,
724            )
725        return record
726
727    def _parse_records(
728        self,
729        response: Optional[requests.Response],
730        records_schema: Mapping[str, Any],
731        stream_slice: Optional[StreamSlice],
732    ) -> Iterable[Record]:
733        yield from self._parse_response(
734            response,
735            stream_slice=stream_slice,
736            records_schema=records_schema,
737        )
738
739    def must_deduplicate_query_params(self) -> bool:
740        return True
741
742    @staticmethod
743    def _to_partition_key(to_serialize: Any) -> str:
744        # separators have changed in Python 3.4. To avoid being impacted by further change, we explicitly specify our own value
745        return json.dumps(to_serialize, indent=None, separators=(",", ":"), sort_keys=True)
746
747
748def _deep_merge(
749    target: MutableMapping[str, Any], source: Union[Record, MutableMapping[str, Any]]
750) -> None:
751    """
752    Recursively merge two dictionaries, combining nested dictionaries instead of overwriting them.
753
754    :param target: The dictionary to merge into (modified in place)
755    :param source: The dictionary to merge from
756    """
757    for key, value in source.items():
758        if (
759            key in target
760            and isinstance(target[key], MutableMapping)
761            and isinstance(value, MutableMapping)
762        ):
763            _deep_merge(target[key], value)
764        else:
765            target[key] = value
766
767
768@deprecated(
769    "This class is experimental. Use at your own risk.",
770    category=ExperimentalClassWarning,
771)
772@dataclass
773class LazySimpleRetriever(SimpleRetriever):
774    """
775    A retriever that supports lazy loading from parent streams.
776    """
777
778    def _read_pages(
779        self,
780        records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]],
781        stream_slice: StreamSlice,
782    ) -> Iterable[Record]:
783        response = stream_slice.extra_fields["child_response"]
784        if response:
785            last_page_size, last_record = 0, None
786            for record in records_generator_fn(response):  # type: ignore[call-arg] # only _parse_records expected as a func
787                last_page_size += 1
788                last_record = record
789                yield record
790
791            next_page_token = self._next_page_token(
792                response,
793                last_page_size,
794                last_record,
795                None,
796                **stream_slice_kwargs(self._next_page_token, stream_slice),
797            )
798            if next_page_token:
799                yield from self._paginate(
800                    next_page_token,
801                    records_generator_fn,
802                    stream_slice,
803                )
804
805            yield from []
806        else:
807            # coderabbit detected an interesting bug/gap where if we were to not get a child_response, we
808            # might recurse forever. This might not be the case, but it is worth noting that this code path
809            # isn't comprehensively tested.
810            yield from self._read_pages(records_generator_fn, stream_slice)
811
812    def _paginate(
813        self,
814        next_page_token: Any,
815        records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]],
816        stream_slice: StreamSlice,
817    ) -> Iterable[Record]:
818        """Handle pagination by fetching subsequent pages."""
819        pagination_complete = False
820
821        while not pagination_complete:
822            response = self._fetch_next_page(stream_slice, next_page_token)
823            last_page_size, last_record = 0, None
824
825            for record in records_generator_fn(response):  # type: ignore[call-arg] # only _parse_records expected as a func
826                last_page_size += 1
827                last_record = record
828                yield record
829
830            if not response:
831                pagination_complete = True
832            else:
833                last_page_token_value = (
834                    next_page_token.get("next_page_token") if next_page_token else None
835                )
836                next_page_token = self._next_page_token(
837                    response,
838                    last_page_size,
839                    last_record,
840                    last_page_token_value,
841                    **stream_slice_kwargs(self._next_page_token, stream_slice),
842                )
843
844                if not next_page_token:
845                    pagination_complete = True
FULL_REFRESH_SYNC_COMPLETE_KEY = '__ab_full_refresh_sync_complete'
LOGGER = <Logger airbyte (INFO)>
@dataclass
class SimpleRetriever(airbyte_cdk.sources.declarative.retrievers.retriever.Retriever):
 84@dataclass
 85class SimpleRetriever(Retriever):
 86    """
 87    Retrieves records by synchronously sending requests to fetch records.
 88
 89    The retriever acts as an orchestrator between the requester, the record selector, the paginator, and the stream slicer.
 90
 91    For each stream slice, submit requests until there are no more pages of records to fetch.
 92
 93    This retriever currently inherits from HttpStream to reuse the request submission and pagination machinery.
 94    As a result, some of the parameters passed to some methods are unused.
 95    The two will be decoupled in a future release.
 96
 97    Attributes:
 98        stream_name (str): The stream's name
 99        stream_primary_key (Optional[Union[str, List[str], List[List[str]]]]): The stream's primary key
100        requester (Requester): The HTTP requester
101        record_selector (HttpSelector): The record selector
102        paginator (Optional[Paginator]): The paginator
103        stream_slicer (Optional[StreamSlicer]): The stream slicer
104        parameters (Mapping[str, Any]): Additional runtime parameters to be used for string interpolation
105        post_pagination_filter (Optional[ClientSideIncrementalRecordFilterDecorator]): Set for data feed streams only.
106            Records the cursor considers already synced are dropped once pagination has observed them
107        page_size_reduction (Optional[PageSizeReduction]): How much to shrink the page size when an error handler
108            resolves to `ResponseAction.REDUCE_PAGE_SIZE`. `None` disables page size reduction entirely.
109            It is immutable configuration; the page size in effect lives in a `PageSizeReducer` that
110            `_read_pages` creates per call, so the retriever and its paginator - both shared by every
111            partition of the stream, read concurrently - stay stateless. When `page_size_reduction` is
112            `None` no reducer is created and `_read_pages` keeps its previous behaviour
113        request_window_splitting (Optional[RequestWindowSplitting]): Policy applied when an error handler
114            resolves to `ResponseAction.SPLIT_REQUEST_WINDOW`, or when `RequestWindowSplitRequiredException`
115            is raised directly by custom code. `None` disables request-window splitting entirely: `read_records`
116            converts the exception into `RequestWindowSplitNotSupportedException` instead of attempting a
117            split it cannot perform.
118        request_window_splitter: Bound method - normally the stream's own cursor's `split_request_window` -
119            asked to replace a failing `StreamSlice` with smaller children. `None` has the same effect as
120            `request_window_splitting` being `None`.
121    """
122
123    requester: Requester
124    record_selector: HttpSelector
125    config: Config
126    parameters: InitVar[Mapping[str, Any]]
127    name: str
128    _name: Union[InterpolatedString, str] = field(init=False, repr=False, default="")
129    primary_key: Optional[Union[str, List[str], List[List[str]]]]
130    _primary_key: str = field(init=False, repr=False, default="")
131    paginator: Optional[Paginator] = None
132    stream_slicer: StreamSlicer = field(
133        default_factory=lambda: SinglePartitionRouter(parameters={})
134    )
135    request_option_provider: RequestOptionsProvider = field(
136        default_factory=lambda: DefaultRequestOptionsProvider(parameters={})
137    )
138    ignore_stream_slicer_parameters_on_paginated_requests: bool = False
139    additional_query_properties: Optional[QueryProperties] = None
140    log_formatter: Optional[Callable[[requests.Response], Any]] = None
141    pagination_tracker_factory: Callable[[], PaginationTracker] = field(
142        default_factory=lambda: lambda: PaginationTracker()
143    )
144    post_pagination_filter: Optional[ClientSideIncrementalRecordFilterDecorator] = None
145    page_size_reduction: Optional[PageSizeReduction] = None
146    request_window_splitting: Optional[RequestWindowSplitting] = None
147    request_window_splitter: Optional[
148        Callable[[StreamSlice, Optional[datetime.timedelta]], Optional[List[StreamSlice]]]
149    ] = None
150
151    def __post_init__(self, parameters: Mapping[str, Any]) -> None:
152        self._paginator = self.paginator or NoPagination(parameters=parameters)
153        self._parameters = parameters
154        self._name = (
155            InterpolatedString(self._name, parameters=parameters)
156            if isinstance(self._name, str)
157            else self._name
158        )
159
160    @property  # type: ignore
161    def name(self) -> str:
162        """
163        :return: Stream name
164        """
165        return (
166            str(self._name.eval(self.config))
167            if isinstance(self._name, InterpolatedString)
168            else self._name
169        )
170
171    @name.setter
172    def name(self, value: str) -> None:
173        if not isinstance(value, property):
174            self._name = value
175
176    def _get_mapping(
177        self, method: Callable[..., Optional[Union[Mapping[str, Any], str]]], **kwargs: Any
178    ) -> Tuple[Union[Mapping[str, Any], str], Set[str]]:
179        """
180        Get mapping from the provided method, and get the keys of the mapping.
181        If the method returns a string, it will return the string and an empty set.
182        If the method returns a dict, it will return the dict and its keys.
183        """
184        mapping = method(**kwargs) or {}
185        keys = set(mapping.keys()) if not isinstance(mapping, str) else set()
186        return mapping, keys
187
188    def _get_request_options(
189        self,
190        stream_slice: Optional[StreamSlice],
191        next_page_token: Optional[Mapping[str, Any]],
192        paginator_method: Callable[..., Optional[Union[Mapping[str, Any], str]]],
193        stream_slicer_method: Callable[..., Optional[Union[Mapping[str, Any], str]]],
194        page_size_override: Optional[int] = None,
195    ) -> Union[Mapping[str, Any], str]:
196        """
197        Get the request_option from the paginator and the stream slicer.
198        Raise a ValueError if there's a key collision
199        Returned merged mapping otherwise
200        """
201        is_body_json = paginator_method.__name__ == "get_request_body_json"
202
203        mappings = [
204            paginator_method(
205                stream_slice=stream_slice,
206                next_page_token=next_page_token,
207                **page_size_override_kwargs(page_size_override),
208            ),
209        ]
210        if not next_page_token or not self.ignore_stream_slicer_parameters_on_paginated_requests:
211            mappings.append(
212                stream_slicer_method(
213                    stream_slice=stream_slice,
214                    next_page_token=next_page_token,
215                )
216            )
217        return combine_mappings(mappings, allow_same_value_merge=is_body_json)
218
219    def _request_headers(
220        self,
221        stream_slice: Optional[StreamSlice] = None,
222        next_page_token: Optional[Mapping[str, Any]] = None,
223        page_size_override: Optional[int] = None,
224    ) -> Mapping[str, Any]:
225        """
226        Specifies request headers.
227        Authentication headers will overwrite any overlapping headers returned from this method.
228        """
229        headers = self._get_request_options(
230            stream_slice,
231            next_page_token,
232            self._paginator.get_request_headers,
233            self.request_option_provider.get_request_headers,
234            **page_size_override_kwargs(page_size_override),
235        )
236        if isinstance(headers, str):
237            raise ValueError("Request headers cannot be a string")
238        return {str(k): str(v) for k, v in headers.items()}
239
240    def _request_params(
241        self,
242        stream_slice: Optional[StreamSlice] = None,
243        next_page_token: Optional[Mapping[str, Any]] = None,
244        page_size_override: Optional[int] = None,
245    ) -> Mapping[str, Any]:
246        """
247        Specifies the query parameters that should be set on an outgoing HTTP request given the inputs.
248
249        E.g: you might want to define query parameters for paging if next_page_token is not None.
250        """
251        params = self._get_request_options(
252            stream_slice,
253            next_page_token,
254            self._paginator.get_request_params,
255            self.request_option_provider.get_request_params,
256            **page_size_override_kwargs(page_size_override),
257        )
258        if isinstance(params, str):
259            raise ValueError("Request params cannot be a string")
260        return params
261
262    def _request_body_data(
263        self,
264        stream_slice: Optional[StreamSlice] = None,
265        next_page_token: Optional[Mapping[str, Any]] = None,
266        page_size_override: Optional[int] = None,
267    ) -> Union[Mapping[str, Any], str]:
268        """
269        Specifies how to populate the body of the request with a non-JSON payload.
270
271        If returns a ready text that it will be sent as is.
272        If returns a dict that it will be converted to a urlencoded form.
273        E.g. {"key1": "value1", "key2": "value2"} => "key1=value1&key2=value2"
274
275        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
276        """
277        return self._get_request_options(
278            stream_slice,
279            next_page_token,
280            self._paginator.get_request_body_data,
281            self.request_option_provider.get_request_body_data,
282            **page_size_override_kwargs(page_size_override),
283        )
284
285    def _request_body_json(
286        self,
287        stream_slice: Optional[StreamSlice] = None,
288        next_page_token: Optional[Mapping[str, Any]] = None,
289        page_size_override: Optional[int] = None,
290    ) -> Optional[Mapping[str, Any]]:
291        """
292        Specifies how to populate the body of the request with a JSON payload.
293
294        At the same time only one of the 'request_body_data' and 'request_body_json' functions can be overridden.
295        """
296        body_json = self._get_request_options(
297            stream_slice,
298            next_page_token,
299            self._paginator.get_request_body_json,
300            self.request_option_provider.get_request_body_json,
301            **page_size_override_kwargs(page_size_override),
302        )
303        if isinstance(body_json, str):
304            raise ValueError("Request body json cannot be a string")
305        return body_json
306
307    def _paginator_path(
308        self,
309        next_page_token: Optional[Mapping[str, Any]] = None,
310        stream_slice: Optional[StreamSlice] = None,
311    ) -> Optional[str]:
312        """
313        If the paginator points to a path, follow it, else return nothing so the requester is used.
314        :param next_page_token:
315        :return:
316        """
317        return self._paginator.path(
318            next_page_token=next_page_token,
319            stream_state={},  # stream_state as an interpolation context is deprecated
320            stream_slice=stream_slice,
321        )
322
323    def _parse_response(
324        self,
325        response: Optional[requests.Response],
326        records_schema: Mapping[str, Any],
327        stream_slice: Optional[StreamSlice] = None,
328        next_page_token: Optional[Mapping[str, Any]] = None,
329    ) -> Iterable[Record]:
330        if not response:
331            yield from []
332        else:
333            yield from self.record_selector.select_records(
334                response=response,
335                stream_state={},  # stream_state as an interpolation context is deprecated
336                records_schema=records_schema,
337                stream_slice=stream_slice,
338                next_page_token=next_page_token,
339            )
340
341    @property  # type: ignore
342    def primary_key(self) -> Optional[Union[str, List[str], List[List[str]]]]:
343        """The stream's primary key"""
344        return self._primary_key
345
346    @primary_key.setter
347    def primary_key(self, value: str) -> None:
348        if not isinstance(value, property):
349            self._primary_key = value
350
351    def _next_page_token(
352        self,
353        response: requests.Response,
354        last_page_size: int,
355        last_record: Optional[Record],
356        last_page_token_value: Optional[Any],
357        page_size_override: Optional[int] = None,
358        stream_slice: Optional[StreamSlice] = None,
359    ) -> Optional[Mapping[str, Any]]:
360        """
361        Specifies a pagination strategy.
362
363        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.
364
365        :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.
366        """
367        return self._paginator.next_page_token(
368            response=response,
369            last_page_size=last_page_size,
370            last_record=last_record,
371            last_page_token_value=last_page_token_value,
372            **page_size_override_kwargs(page_size_override),
373            **stream_slice_kwargs(self._paginator.next_page_token, stream_slice),
374        )
375
376    def _fetch_next_page(
377        self,
378        stream_slice: StreamSlice,
379        next_page_token: Optional[Mapping[str, Any]] = None,
380        page_size_override: Optional[int] = None,
381    ) -> Optional[requests.Response]:
382        return self.requester.send_request(
383            path=self._paginator_path(
384                next_page_token=next_page_token,
385                stream_slice=stream_slice,
386            ),
387            stream_state={},  # stream_state as an interpolation context is deprecated
388            stream_slice=stream_slice,
389            next_page_token=next_page_token,
390            request_headers=self._request_headers(
391                stream_slice=stream_slice,
392                next_page_token=next_page_token,
393                **page_size_override_kwargs(page_size_override),
394            ),
395            request_params=self._request_params(
396                stream_slice=stream_slice,
397                next_page_token=next_page_token,
398                **page_size_override_kwargs(page_size_override),
399            ),
400            request_body_data=self._request_body_data(
401                stream_slice=stream_slice,
402                next_page_token=next_page_token,
403                **page_size_override_kwargs(page_size_override),
404            ),
405            request_body_json=self._request_body_json(
406                stream_slice=stream_slice,
407                next_page_token=next_page_token,
408                **page_size_override_kwargs(page_size_override),
409            ),
410            log_formatter=self.log_formatter,
411        )
412
413    # This logic is similar to _read_pages in the HttpStream class. When making changes here, consider making changes there as well.
414    def _read_pages(
415        self,
416        records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]],
417        stream_slice: StreamSlice,
418    ) -> Iterable[Record]:
419        original_stream_slice = stream_slice
420        pagination_tracker = self.pagination_tracker_factory()
421        page_size_reducer = (
422            PageSizeReducer(
423                self.page_size_reduction,
424                self._paginator.get_page_size(),
425                stream_name=self.name,
426            )
427            if self.page_size_reduction
428            else None
429        )
430        reset_pagination = False
431        reduce_page_size = False
432        next_page_token = self._get_initial_next_page_token()
433        while True:
434            page_size_override = page_size_reducer.page_size_override if page_size_reducer else None
435            merged_records: MutableMapping[str, Any] = defaultdict(dict)
436            last_page_size = 0
437            last_record: Optional[Record] = None
438
439            response = None
440            try:
441                if self.additional_query_properties:
442                    for (
443                        properties
444                    ) in self.additional_query_properties.get_request_property_chunks():
445                        stream_slice = StreamSlice(
446                            partition=stream_slice.partition or {},
447                            cursor_slice=stream_slice.cursor_slice or {},
448                            extra_fields={"query_properties": properties},
449                        )
450                        response = self._fetch_next_page(
451                            stream_slice,
452                            next_page_token,
453                            **page_size_override_kwargs(page_size_override),
454                        )
455
456                        for current_record in records_generator_fn(response):
457                            if self.additional_query_properties.property_chunking:
458                                merge_key = self.additional_query_properties.property_chunking.get_merge_key(
459                                    current_record
460                                )
461                                if merge_key:
462                                    _deep_merge(merged_records[merge_key], current_record)
463                                else:
464                                    # We should still emit records even if the record did not have a merge key
465                                    pagination_tracker.observe(current_record)
466                                    last_page_size += 1
467                                    last_record = current_record
468                                    yield current_record
469                            else:
470                                pagination_tracker.observe(current_record)
471                                last_page_size += 1
472                                last_record = current_record
473                                yield current_record
474
475                    for merged_record in merged_records.values():
476                        record = Record(
477                            data=merged_record, stream_name=self.name, associated_slice=stream_slice
478                        )
479                        pagination_tracker.observe(record)
480                        last_page_size += 1
481                        last_record = record
482                        yield record
483                else:
484                    response = self._fetch_next_page(
485                        stream_slice,
486                        next_page_token,
487                        **page_size_override_kwargs(page_size_override),
488                    )
489                    for current_record in records_generator_fn(response):
490                        pagination_tracker.observe(current_record)
491                        last_page_size += 1
492                        last_record = current_record
493                        yield current_record
494            except PaginationResetRequiredException:
495                reset_pagination = True
496            except PageSizeReductionRequiredException:
497                if page_size_reducer is None:
498                    # The action can be attached to a requester we cannot validate at config time, such as one
499                    # built by a custom error handler or a CustomRequester. Re-raise as the misconfiguration it
500                    # is: the exception being handled is the neutral "the API asked for a smaller page" signal.
501                    raise PageSizeReductionNotSupportedException(stream_name=self.name)
502                if last_page_size:
503                    # The reduction is safe only because it re-issues a page whose records were not emitted.
504                    # Every in-CDK way of reaching this raises from `_fetch_next_page`, before the record loop,
505                    # and the factory rejects the manifest constructs that would not, but a custom extractor,
506                    # filter or transformation can issue its own request from inside the record generator.
507                    # Re-issuing the page then duplicates the records already yielded, so fail instead.
508                    raise AirbyteTracedException(
509                        internal_message=f"Stream {self.name} requested a page size reduction after {last_page_size} records of the page had already been emitted",
510                        message=f"Stream {self.name} asked for a smaller page size in the middle of a page. The page cannot be requested again without duplicating the records already read from it. Move the REDUCE_PAGE_SIZE action to the error handler of the stream's main requester.",
511                        failure_type=FailureType.config_error,
512                    )
513                # Raises once the page size cannot be reduced any further and the retries allowed at that
514                # floor are spent, which is what stops the loop when the API keeps failing.
515                page_size_reducer.reduce()
516                reduce_page_size = True
517            else:
518                if page_size_reducer:
519                    page_size_reducer.on_successful_page()
520                if not response:
521                    break
522
523            if reduce_page_size:
524                # Retry the very same page: neither the token nor the slice change, only the page size does -
525                # and not even that once the reducer is at its floor and only waiting is left.
526                reduce_page_size = False
527                continue
528
529            if reset_pagination or pagination_tracker.has_reached_limit():
530                next_page_token = self._get_initial_next_page_token()
531                previous_slice = stream_slice
532                stream_slice = pagination_tracker.reduce_slice_range_if_possible(
533                    stream_slice, original_stream_slice
534                )
535                LOGGER.info(
536                    f"Hitting PaginationReset event. StreamSlice used will go from {previous_slice} to {stream_slice}"
537                )
538                reset_pagination = False
539            else:
540                last_page_token_value = (
541                    next_page_token.get("next_page_token") if next_page_token else None
542                )
543                next_page_token = self._next_page_token(
544                    response=response,  # type:ignore # we are breaking from the loop on the try/else if there are no response so this should be fine
545                    last_page_size=last_page_size,
546                    last_record=last_record,
547                    last_page_token_value=last_page_token_value,
548                    **page_size_override_kwargs(page_size_override),
549                    **stream_slice_kwargs(self._next_page_token, stream_slice),
550                )
551                if not next_page_token:
552                    break
553
554        # Always return an empty generator just in case no records were ever yielded
555        yield from []
556
557    def _get_initial_next_page_token(self) -> Optional[Mapping[str, Any]]:
558        initial_token = self._paginator.get_initial_token()
559        next_page_token = {"next_page_token": initial_token} if initial_token is not None else None
560        return next_page_token
561
562    def read_records(
563        self,
564        records_schema: Mapping[str, Any],
565        stream_slice: Optional[StreamSlice] = None,
566    ) -> Iterable[StreamData]:
567        """
568        Fetch a stream's records from an HTTP API source
569
570        :param records_schema: json schema to describe record
571        :param stream_slice: The stream slice to read data for
572        :return: The records read from the API source
573        """
574        _slice = stream_slice or StreamSlice(partition={}, cursor_slice={})  # None-check
575        yield from self._read_records_or_split_request_window(
576            records_schema, _slice, _slice, depth=0
577        )
578
579    def _read_records_or_split_request_window(
580        self,
581        records_schema: Mapping[str, Any],
582        stream_slice: StreamSlice,
583        original_slice: StreamSlice,
584        depth: int,
585    ) -> Iterable[StreamData]:
586        """
587        Read `stream_slice` to completion, replacing it with smaller children and recursing into each of them in
588        turn when a `RequestWindowSplitRequiredException` is raised while reading it.
589
590        Each recursive call re-enters this method - and, through it, `_read_pages` - from scratch for the child
591        slice it is given, so every child gets its own paginator token, `PaginationTracker`, and
592        `PageSizeReducer` for free: nothing here needs to reset that state explicitly. Because the whole
593        recursion lives inside the single generator `DeclarativePartition.read()` consumes,
594        `PartitionReader.process_partition()` only calls `cursor.close_partition()` once, after this generator is
595        fully exhausted - so a failure anywhere in the recursion (a child that cannot be read, or a window that
596        cannot be split any further) propagates out without ever checkpointing the original partition.
597
598        `original_slice` is threaded through the recursion unchanged so every record yielded, however deep the
599        recursion went to produce it, is re-stamped with the partition's own slice rather than the child slice it
600        was actually read against - see `_reassociate_with_original_slice`. `depth` bounds the recursion
601        independently of `request_window_splitter`'s own no-progress guard: it is enforced here so a
602        misbehaving custom cursor cannot recurse indefinitely regardless of what that guard does or does not
603        catch.
604        """
605        record_generator = partial(
606            self._parse_records,
607            stream_slice=stream_slice,
608            records_schema=records_schema,
609        )
610
611        emitted_count = 0
612        try:
613            records: Iterable[Mapping[str, Any]] = self._read_pages(record_generator, stream_slice)
614            if self.post_pagination_filter:
615                # A data feed paginates until it reaches a record older than the cursor, so the page that triggers the stop
616                # condition still holds already-synced records. Those are filtered here rather than in the record selector
617                # so that the paginator keeps seeing the whole page: the stop condition is evaluated on the last record of
618                # the page, which is precisely one of the records being dropped. Two consequences of filtering this late:
619                # the pagination tracker observes the dropped records, and a `file_uploader` on the record selector has
620                # already uploaded their files by the time they are dropped.
621                records = self.post_pagination_filter.filter_records(
622                    records,
623                    # the filter is only used for its cursor comparison, which does not read the stream state
624                    stream_state={},
625                    stream_slice=stream_slice,
626                )
627            for record in records:
628                emitted_count += 1
629                yield self._reassociate_with_original_slice(record, original_slice)
630        except RequestWindowSplitRequiredException as exception:
631            if self.request_window_splitting is None or self.request_window_splitter is None:
632                raise RequestWindowSplitNotSupportedException(stream_name=self.name) from exception
633
634            if emitted_count:
635                # Splitting and re-reading the window re-emits these records: a failed partition is never
636                # checkpointed, so the next attempt would re-emit them anyway. Always replaying makes forward
637                # progress instead of retrying the same oversized window forever; primary-key dedup downstream
638                # handles the duplicates either way.
639                LOGGER.warning(
640                    f"Stream {self.name} already emitted {emitted_count} record(s) from {stream_slice} before "
641                    f"it was rejected; splitting and re-reading the window may re-emit them."
642                )
643
644            if depth >= _MAX_REQUEST_WINDOW_SPLIT_DEPTH:
645                # Logged separately from the exception below so it stays visible even if the trace message's
646                # internal_message isn't surfaced by whatever catches it.
647                LOGGER.warning(
648                    f"Stream {self.name} hit the maximum request window split depth "
649                    f"({_MAX_REQUEST_WINDOW_SPLIT_DEPTH}) while splitting {original_slice}."
650                )
651                raise AirbyteTracedException(
652                    internal_message=f"Stream {self.name} exceeded the maximum request window split depth of {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} while splitting {original_slice}",
653                    # `transient_error`, so the only remediation is the connector's own, if it defined one.
654                    message=self._with_request_window_failure_message(
655                        f"Stream {self.name} could not split its request window to a size the API accepts "
656                        f"within {_MAX_REQUEST_WINDOW_SPLIT_DEPTH} splits. This usually means the stream's "
657                        f"cursor is not actually shrinking the window on each split; if it uses a custom "
658                        f"cursor, check its `split_request_window` implementation."
659                    ),
660                    failure_type=FailureType.transient_error,
661                ) from exception
662
663            children = self.request_window_splitter(
664                stream_slice, self.request_window_splitting.min_split_window
665            )
666            if children is None:
667                min_split_window_note = (
668                    f", its configured `min_split_window` ({self.request_window_splitting.min_split_window})"
669                    if self.request_window_splitting.min_split_window
670                    else ""
671                )
672                raise AirbyteTracedException(
673                    internal_message=f"Stream {self.name} could not split its request window {stream_slice} any further",
674                    # `transient_error` regardless of how the triggering response was classified: exhaustion is
675                    # never the user's fault, matching `PageSizeReducer`'s equivalent exhaustion branch.
676                    message=self._with_request_window_failure_message(
677                        f"The API kept rejecting stream {self.name}'s request window even at the smallest window "
678                        f"its cursor allows{min_split_window_note}, or splitting the window further would not "
679                        f"make progress."
680                    ),
681                    failure_type=FailureType.transient_error,
682                ) from exception
683
684            LOGGER.info(
685                f"Stream {self.name}: the API rejected request window {stream_slice} (split depth {depth}); "
686                f"reducing it to {children} and reading each in turn."
687            )
688            for child in children:
689                yield from self._read_records_or_split_request_window(
690                    records_schema, child, original_slice, depth + 1
691                )
692
693    def _with_request_window_failure_message(self, message: str) -> str:
694        """
695        :return: the message followed by `request_window_splitting`'s configured `failure_message`, when it
696            defined one - mirroring `PageSizeReducer._with_failure_message`
697        """
698        assert self.request_window_splitting is not None
699        failure_message = (self.request_window_splitting.failure_message or "").strip()
700        if not failure_message:
701            return message
702        return f"{message} {failure_message}"
703
704    @staticmethod
705    def _reassociate_with_original_slice(
706        record: StreamData, original_slice: StreamSlice
707    ) -> StreamData:
708        """
709        Records read while recursing into a split child window are parsed against that child's `StreamSlice`,
710        so `RecordSelector` stamps them with the child as `associated_slice`. `DeclarativePartition.read()`
711        passes an already-built `Record` through unchanged rather than re-wrapping it, so left uncorrected the
712        child slice - not the partition's own slice - is what `ConcurrentCursor.observe()` would key its
713        per-partition bookkeeping by. `close_partition()` always looks that bookkeeping up by the partition's own
714        slice, so it would never find it, silently losing the precise most-recently-observed cursor value for a
715        split partition (recoverable in practice today only because callers of that lookup fall back to the
716        slice's end boundary when it is missing). Re-stamping every record with `original_slice` here restores
717        the same tracking a non-split read already gets, rather than relying on that fallback.
718        """
719        if isinstance(record, Record) and record.associated_slice is not original_slice:
720            return Record(
721                data=record.data,
722                stream_name=record.stream_name,
723                associated_slice=original_slice,
724                file_reference=record.file_reference,
725            )
726        return record
727
728    def _parse_records(
729        self,
730        response: Optional[requests.Response],
731        records_schema: Mapping[str, Any],
732        stream_slice: Optional[StreamSlice],
733    ) -> Iterable[Record]:
734        yield from self._parse_response(
735            response,
736            stream_slice=stream_slice,
737            records_schema=records_schema,
738        )
739
740    def must_deduplicate_query_params(self) -> bool:
741        return True
742
743    @staticmethod
744    def _to_partition_key(to_serialize: Any) -> str:
745        # separators have changed in Python 3.4. To avoid being impacted by further change, we explicitly specify our own value
746        return json.dumps(to_serialize, indent=None, separators=(",", ":"), sort_keys=True)

Retrieves records by synchronously sending requests to fetch records.

The retriever acts as an orchestrator between the requester, the record selector, the paginator, and the stream slicer.

For each stream slice, submit requests until there are no more pages of records to fetch.

This retriever currently inherits from HttpStream to reuse the request submission and pagination machinery. As a result, some of the parameters passed to some methods are unused. The two will be decoupled in a future release.

Attributes:
  • stream_name (str): The stream's name
  • stream_primary_key (Optional[Union[str, List[str], List[List[str]]]]): The stream's primary key
  • requester (Requester): The HTTP requester
  • record_selector (HttpSelector): The record selector
  • paginator (Optional[Paginator]): The paginator
  • stream_slicer (Optional[StreamSlicer]): The stream slicer
  • parameters (Mapping[str, Any]): Additional runtime parameters to be used for string interpolation
  • post_pagination_filter (Optional[ClientSideIncrementalRecordFilterDecorator]): Set for data feed streams only. Records the cursor considers already synced are dropped once pagination has observed them
  • page_size_reduction (Optional[PageSizeReduction]): How much to shrink the page size when an error handler resolves to ResponseAction.REDUCE_PAGE_SIZE. None disables page size reduction entirely. It is immutable configuration; the page size in effect lives in a PageSizeReducer that _read_pages creates per call, so the retriever and its paginator - both shared by every partition of the stream, read concurrently - stay stateless. When page_size_reduction is None no reducer is created and _read_pages keeps its previous behaviour
  • request_window_splitting (Optional[RequestWindowSplitting]): Policy applied when an error handler resolves to ResponseAction.SPLIT_REQUEST_WINDOW, or when RequestWindowSplitRequiredException is raised directly by custom code. None disables request-window splitting entirely: read_records converts the exception into RequestWindowSplitNotSupportedException instead of attempting a split it cannot perform.
  • request_window_splitter: Bound method - normally the stream's own cursor's split_request_window - asked to replace a failing StreamSlice with smaller children. None has the same effect as request_window_splitting being None.
SimpleRetriever( requester: airbyte_cdk.Requester, record_selector: airbyte_cdk.sources.declarative.extractors.HttpSelector, config: Mapping[str, Any], parameters: dataclasses.InitVar[typing.Mapping[str, typing.Any]], name: str = <property object>, primary_key: Union[str, List[str], List[List[str]], NoneType] = <property object>, paginator: Optional[airbyte_cdk.sources.declarative.requesters.paginators.Paginator] = None, stream_slicer: airbyte_cdk.sources.declarative.stream_slicers.StreamSlicer = <factory>, request_option_provider: airbyte_cdk.sources.declarative.requesters.request_options.RequestOptionsProvider = <factory>, ignore_stream_slicer_parameters_on_paginated_requests: bool = False, additional_query_properties: Optional[airbyte_cdk.sources.declarative.requesters.query_properties.QueryProperties] = None, log_formatter: Optional[Callable[[requests.models.Response], Any]] = None, pagination_tracker_factory: Callable[[], airbyte_cdk.sources.declarative.retrievers.pagination_tracker.PaginationTracker] = <factory>, post_pagination_filter: Optional[airbyte_cdk.sources.declarative.extractors.record_filter.ClientSideIncrementalRecordFilterDecorator] = None, page_size_reduction: Optional[airbyte_cdk.sources.declarative.retrievers.page_size_reducer.PageSizeReduction] = None, request_window_splitting: Optional[airbyte_cdk.sources.declarative.retrievers.request_window_splitting.RequestWindowSplitting] = None, request_window_splitter: Optional[Callable[[airbyte_cdk.StreamSlice, Optional[datetime.timedelta]], Optional[List[airbyte_cdk.StreamSlice]]]] = None)
config: Mapping[str, Any]
parameters: dataclasses.InitVar[typing.Mapping[str, typing.Any]]
name: str
160    @property  # type: ignore
161    def name(self) -> str:
162        """
163        :return: Stream name
164        """
165        return (
166            str(self._name.eval(self.config))
167            if isinstance(self._name, InterpolatedString)
168            else self._name
169        )
Returns

Stream name

primary_key: Union[str, List[str], List[List[str]], NoneType]
341    @property  # type: ignore
342    def primary_key(self) -> Optional[Union[str, List[str], List[List[str]]]]:
343        """The stream's primary key"""
344        return self._primary_key

The stream's primary key

ignore_stream_slicer_parameters_on_paginated_requests: bool = False
log_formatter: Optional[Callable[[requests.models.Response], Any]] = None
request_window_splitter: Optional[Callable[[airbyte_cdk.StreamSlice, Optional[datetime.timedelta]], Optional[List[airbyte_cdk.StreamSlice]]]] = None
def read_records( self, records_schema: Mapping[str, Any], stream_slice: Optional[airbyte_cdk.StreamSlice] = None) -> Iterable[Union[Mapping[str, Any], airbyte_cdk.AirbyteMessage]]:
562    def read_records(
563        self,
564        records_schema: Mapping[str, Any],
565        stream_slice: Optional[StreamSlice] = None,
566    ) -> Iterable[StreamData]:
567        """
568        Fetch a stream's records from an HTTP API source
569
570        :param records_schema: json schema to describe record
571        :param stream_slice: The stream slice to read data for
572        :return: The records read from the API source
573        """
574        _slice = stream_slice or StreamSlice(partition={}, cursor_slice={})  # None-check
575        yield from self._read_records_or_split_request_window(
576            records_schema, _slice, _slice, depth=0
577        )

Fetch a stream's records from an HTTP API source

Parameters
  • records_schema: json schema to describe record
  • stream_slice: The stream slice to read data for
Returns

The records read from the API source

def must_deduplicate_query_params(self) -> bool:
740    def must_deduplicate_query_params(self) -> bool:
741        return True
@deprecated('This class is experimental. Use at your own risk.', category=ExperimentalClassWarning)
@dataclass
class LazySimpleRetriever(SimpleRetriever):
769@deprecated(
770    "This class is experimental. Use at your own risk.",
771    category=ExperimentalClassWarning,
772)
773@dataclass
774class LazySimpleRetriever(SimpleRetriever):
775    """
776    A retriever that supports lazy loading from parent streams.
777    """
778
779    def _read_pages(
780        self,
781        records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]],
782        stream_slice: StreamSlice,
783    ) -> Iterable[Record]:
784        response = stream_slice.extra_fields["child_response"]
785        if response:
786            last_page_size, last_record = 0, None
787            for record in records_generator_fn(response):  # type: ignore[call-arg] # only _parse_records expected as a func
788                last_page_size += 1
789                last_record = record
790                yield record
791
792            next_page_token = self._next_page_token(
793                response,
794                last_page_size,
795                last_record,
796                None,
797                **stream_slice_kwargs(self._next_page_token, stream_slice),
798            )
799            if next_page_token:
800                yield from self._paginate(
801                    next_page_token,
802                    records_generator_fn,
803                    stream_slice,
804                )
805
806            yield from []
807        else:
808            # coderabbit detected an interesting bug/gap where if we were to not get a child_response, we
809            # might recurse forever. This might not be the case, but it is worth noting that this code path
810            # isn't comprehensively tested.
811            yield from self._read_pages(records_generator_fn, stream_slice)
812
813    def _paginate(
814        self,
815        next_page_token: Any,
816        records_generator_fn: Callable[[Optional[requests.Response]], Iterable[Record]],
817        stream_slice: StreamSlice,
818    ) -> Iterable[Record]:
819        """Handle pagination by fetching subsequent pages."""
820        pagination_complete = False
821
822        while not pagination_complete:
823            response = self._fetch_next_page(stream_slice, next_page_token)
824            last_page_size, last_record = 0, None
825
826            for record in records_generator_fn(response):  # type: ignore[call-arg] # only _parse_records expected as a func
827                last_page_size += 1
828                last_record = record
829                yield record
830
831            if not response:
832                pagination_complete = True
833            else:
834                last_page_token_value = (
835                    next_page_token.get("next_page_token") if next_page_token else None
836                )
837                next_page_token = self._next_page_token(
838                    response,
839                    last_page_size,
840                    last_record,
841                    last_page_token_value,
842                    **stream_slice_kwargs(self._next_page_token, stream_slice),
843                )
844
845                if not next_page_token:
846                    pagination_complete = True

A retriever that supports lazy loading from parent streams.

LazySimpleRetriever( requester: airbyte_cdk.Requester, record_selector: airbyte_cdk.sources.declarative.extractors.HttpSelector, config: Mapping[str, Any], parameters: dataclasses.InitVar[typing.Mapping[str, typing.Any]], name: str = <property object>, primary_key: Union[str, List[str], List[List[str]], NoneType] = <property object>, paginator: Optional[airbyte_cdk.sources.declarative.requesters.paginators.Paginator] = None, stream_slicer: airbyte_cdk.sources.declarative.stream_slicers.StreamSlicer = <factory>, request_option_provider: airbyte_cdk.sources.declarative.requesters.request_options.RequestOptionsProvider = <factory>, ignore_stream_slicer_parameters_on_paginated_requests: bool = False, additional_query_properties: Optional[airbyte_cdk.sources.declarative.requesters.query_properties.QueryProperties] = None, log_formatter: Optional[Callable[[requests.models.Response], Any]] = None, pagination_tracker_factory: Callable[[], airbyte_cdk.sources.declarative.retrievers.pagination_tracker.PaginationTracker] = <factory>, post_pagination_filter: Optional[airbyte_cdk.sources.declarative.extractors.record_filter.ClientSideIncrementalRecordFilterDecorator] = None, page_size_reduction: Optional[airbyte_cdk.sources.declarative.retrievers.page_size_reducer.PageSizeReduction] = None, request_window_splitting: Optional[airbyte_cdk.sources.declarative.retrievers.request_window_splitting.RequestWindowSplitting] = None, request_window_splitter: Optional[Callable[[airbyte_cdk.StreamSlice, Optional[datetime.timedelta]], Optional[List[airbyte_cdk.StreamSlice]]]] = None)