airbyte_cdk.sources.declarative.requesters.paginators.paginator

  1#
  2# Copyright (c) 2023 Airbyte, Inc., all rights reserved.
  3#
  4
  5import inspect
  6from abc import ABC, abstractmethod
  7from dataclasses import dataclass
  8from functools import lru_cache
  9from typing import Any, Callable, Dict, Mapping, Optional
 10
 11import requests
 12
 13from airbyte_cdk.sources.declarative.requesters.request_options.request_options_provider import (
 14    RequestOptionsProvider,
 15)
 16from airbyte_cdk.sources.types import Record, StreamSlice
 17
 18
 19def page_size_override_kwargs(page_size_override: Optional[int]) -> Dict[str, Any]:
 20    """
 21    Build the `page_size_override` keyword argument only when there is an override to pass.
 22
 23    Paginators and pagination strategies defined outside of the CDK may not accept the argument, and they only
 24    need to when the stream actually reduces its page size (see `ResponseAction.REDUCE_PAGE_SIZE`).
 25    """
 26    return {"page_size_override": page_size_override} if page_size_override is not None else {}
 27
 28
 29def stream_slice_kwargs(
 30    next_page_token: Callable[..., Any], stream_slice: Optional[StreamSlice]
 31) -> Dict[str, Any]:
 32    """
 33    Build the `stream_slice` keyword argument only for a `next_page_token` that accepts it.
 34
 35    Unlike `page_size_override`, the slice is almost always present, so passing it only when it is set would still
 36    break a paginator or pagination strategy defined outside of the CDK whose signature predates the argument. The
 37    signature is checked instead, and a callee that does not declare `stream_slice` (or `**kwargs`) is called as
 38    before.
 39    """
 40    if stream_slice is None:
 41        return {}
 42    function = getattr(next_page_token, "__func__", next_page_token)
 43    try:
 44        accepts_stream_slice = _accepts_stream_slice(function)
 45    except TypeError:
 46        # Unhashable callables cannot be cached; they are rare enough to inspect on every call.
 47        accepts_stream_slice = _accepts_stream_slice.__wrapped__(function)
 48    return {"stream_slice": stream_slice} if accepts_stream_slice else {}
 49
 50
 51@lru_cache(maxsize=None)
 52def _accepts_stream_slice(function: Callable[..., Any]) -> bool:
 53    try:
 54        parameters = inspect.signature(function).parameters
 55    except (TypeError, ValueError):
 56        return False
 57    return "stream_slice" in parameters or any(
 58        parameter.kind == inspect.Parameter.VAR_KEYWORD for parameter in parameters.values()
 59    )
 60
 61
 62@dataclass
 63class Paginator(ABC, RequestOptionsProvider):
 64    """
 65    Defines the token to use to fetch the next page of records from the API.
 66
 67    If needed, the Paginator will set request options to be set on the HTTP request to fetch the next page of records.
 68    If the next_page_token is the path to the next page of records, then it should be accessed through the `path` method
 69    """
 70
 71    @abstractmethod
 72    def get_initial_token(self) -> Optional[Any]:
 73        """
 74        Get the page token that should be included in the request to get the first page of records
 75        """
 76
 77    @abstractmethod
 78    def next_page_token(
 79        self,
 80        response: requests.Response,
 81        last_page_size: int,
 82        last_record: Optional[Record],
 83        last_page_token_value: Optional[Any],
 84        page_size_override: Optional[int] = None,
 85        stream_slice: Optional[StreamSlice] = None,
 86    ) -> Optional[Mapping[str, Any]]:
 87        """
 88        Returns the next_page_token to use to fetch the next page of records.
 89
 90        :param response: the response to process
 91        :param last_page_size: the number of records read from the response
 92        :param last_record: the last record extracted from the response
 93        :param last_page_token_value: The current value of the page token made on the last request
 94        :param page_size_override: the page size that was actually requested, when it differs from the configured
 95            one because of a `REDUCE_PAGE_SIZE` response action
 96        :param stream_slice: the slice the page was read for, so that a stop condition can compare the page against
 97            the slice's own window or partition
 98        :return: A mapping {"next_page_token": <token>} for the next page from the input response object. Returning None means there are no more pages to read in this response.
 99        """
100        pass
101
102    def get_page_size(self) -> Optional[int]:
103        """
104        Evaluated against the config alone. A pagination strategy whose `page_size` template references the
105        response evaluates it with the response in context inside `next_page_token`, so the two can disagree -
106        pre-existing, and only observable for a `page_size` that is not a constant.
107
108        :return: the number of records this paginator asks for per page, or None if it does not define one
109        """
110        return None
111
112    @abstractmethod
113    def path(
114        self,
115        next_page_token: Optional[Mapping[str, Any]],
116        stream_state: Optional[Mapping[str, Any]] = None,
117        stream_slice: Optional[StreamSlice] = None,
118    ) -> Optional[str]:
119        """
120        Returns the URL path to hit to fetch the next page of records
121
122        e.g: if you wanted to hit https://myapi.com/v1/some_entity then this will return "some_entity"
123
124        :return: path to hit to fetch the next request. Returning None means the path is not defined by the next_page_token
125        """
126        pass
def page_size_override_kwargs(page_size_override: Optional[int]) -> Dict[str, Any]:
20def page_size_override_kwargs(page_size_override: Optional[int]) -> Dict[str, Any]:
21    """
22    Build the `page_size_override` keyword argument only when there is an override to pass.
23
24    Paginators and pagination strategies defined outside of the CDK may not accept the argument, and they only
25    need to when the stream actually reduces its page size (see `ResponseAction.REDUCE_PAGE_SIZE`).
26    """
27    return {"page_size_override": page_size_override} if page_size_override is not None else {}

Build the page_size_override keyword argument only when there is an override to pass.

Paginators and pagination strategies defined outside of the CDK may not accept the argument, and they only need to when the stream actually reduces its page size (see ResponseAction.REDUCE_PAGE_SIZE).

def stream_slice_kwargs( next_page_token: Callable[..., Any], stream_slice: Optional[airbyte_cdk.StreamSlice]) -> Dict[str, Any]:
30def stream_slice_kwargs(
31    next_page_token: Callable[..., Any], stream_slice: Optional[StreamSlice]
32) -> Dict[str, Any]:
33    """
34    Build the `stream_slice` keyword argument only for a `next_page_token` that accepts it.
35
36    Unlike `page_size_override`, the slice is almost always present, so passing it only when it is set would still
37    break a paginator or pagination strategy defined outside of the CDK whose signature predates the argument. The
38    signature is checked instead, and a callee that does not declare `stream_slice` (or `**kwargs`) is called as
39    before.
40    """
41    if stream_slice is None:
42        return {}
43    function = getattr(next_page_token, "__func__", next_page_token)
44    try:
45        accepts_stream_slice = _accepts_stream_slice(function)
46    except TypeError:
47        # Unhashable callables cannot be cached; they are rare enough to inspect on every call.
48        accepts_stream_slice = _accepts_stream_slice.__wrapped__(function)
49    return {"stream_slice": stream_slice} if accepts_stream_slice else {}

Build the stream_slice keyword argument only for a next_page_token that accepts it.

Unlike page_size_override, the slice is almost always present, so passing it only when it is set would still break a paginator or pagination strategy defined outside of the CDK whose signature predates the argument. The signature is checked instead, and a callee that does not declare stream_slice (or **kwargs) is called as before.

 63@dataclass
 64class Paginator(ABC, RequestOptionsProvider):
 65    """
 66    Defines the token to use to fetch the next page of records from the API.
 67
 68    If needed, the Paginator will set request options to be set on the HTTP request to fetch the next page of records.
 69    If the next_page_token is the path to the next page of records, then it should be accessed through the `path` method
 70    """
 71
 72    @abstractmethod
 73    def get_initial_token(self) -> Optional[Any]:
 74        """
 75        Get the page token that should be included in the request to get the first page of records
 76        """
 77
 78    @abstractmethod
 79    def next_page_token(
 80        self,
 81        response: requests.Response,
 82        last_page_size: int,
 83        last_record: Optional[Record],
 84        last_page_token_value: Optional[Any],
 85        page_size_override: Optional[int] = None,
 86        stream_slice: Optional[StreamSlice] = None,
 87    ) -> Optional[Mapping[str, Any]]:
 88        """
 89        Returns the next_page_token to use to fetch the next page of records.
 90
 91        :param response: the response to process
 92        :param last_page_size: the number of records read from the response
 93        :param last_record: the last record extracted from the response
 94        :param last_page_token_value: The current value of the page token made on the last request
 95        :param page_size_override: the page size that was actually requested, when it differs from the configured
 96            one because of a `REDUCE_PAGE_SIZE` response action
 97        :param stream_slice: the slice the page was read for, so that a stop condition can compare the page against
 98            the slice's own window or partition
 99        :return: A mapping {"next_page_token": <token>} for the next page from the input response object. Returning None means there are no more pages to read in this response.
100        """
101        pass
102
103    def get_page_size(self) -> Optional[int]:
104        """
105        Evaluated against the config alone. A pagination strategy whose `page_size` template references the
106        response evaluates it with the response in context inside `next_page_token`, so the two can disagree -
107        pre-existing, and only observable for a `page_size` that is not a constant.
108
109        :return: the number of records this paginator asks for per page, or None if it does not define one
110        """
111        return None
112
113    @abstractmethod
114    def path(
115        self,
116        next_page_token: Optional[Mapping[str, Any]],
117        stream_state: Optional[Mapping[str, Any]] = None,
118        stream_slice: Optional[StreamSlice] = None,
119    ) -> Optional[str]:
120        """
121        Returns the URL path to hit to fetch the next page of records
122
123        e.g: if you wanted to hit https://myapi.com/v1/some_entity then this will return "some_entity"
124
125        :return: path to hit to fetch the next request. Returning None means the path is not defined by the next_page_token
126        """
127        pass

Defines the token to use to fetch the next page of records from the API.

If needed, the Paginator will set request options to be set on the HTTP request to fetch the next page of records. If the next_page_token is the path to the next page of records, then it should be accessed through the path method

@abstractmethod
def get_initial_token(self) -> Optional[Any]:
72    @abstractmethod
73    def get_initial_token(self) -> Optional[Any]:
74        """
75        Get the page token that should be included in the request to get the first page of records
76        """

Get the page token that should be included in the request to get the first page of records

@abstractmethod
def next_page_token( self, response: requests.models.Response, last_page_size: int, last_record: Optional[airbyte_cdk.Record], last_page_token_value: Optional[Any], page_size_override: Optional[int] = None, stream_slice: Optional[airbyte_cdk.StreamSlice] = None) -> Optional[Mapping[str, Any]]:
 78    @abstractmethod
 79    def next_page_token(
 80        self,
 81        response: requests.Response,
 82        last_page_size: int,
 83        last_record: Optional[Record],
 84        last_page_token_value: Optional[Any],
 85        page_size_override: Optional[int] = None,
 86        stream_slice: Optional[StreamSlice] = None,
 87    ) -> Optional[Mapping[str, Any]]:
 88        """
 89        Returns the next_page_token to use to fetch the next page of records.
 90
 91        :param response: the response to process
 92        :param last_page_size: the number of records read from the response
 93        :param last_record: the last record extracted from the response
 94        :param last_page_token_value: The current value of the page token made on the last request
 95        :param page_size_override: the page size that was actually requested, when it differs from the configured
 96            one because of a `REDUCE_PAGE_SIZE` response action
 97        :param stream_slice: the slice the page was read for, so that a stop condition can compare the page against
 98            the slice's own window or partition
 99        :return: A mapping {"next_page_token": <token>} for the next page from the input response object. Returning None means there are no more pages to read in this response.
100        """
101        pass

Returns the next_page_token to use to fetch the next page of records.

Parameters
  • response: the response to process
  • last_page_size: the number of records read from the response
  • last_record: the last record extracted from the response
  • last_page_token_value: The current value of the page token made on the last request
  • page_size_override: the page size that was actually requested, when it differs from the configured one because of a REDUCE_PAGE_SIZE response action
  • stream_slice: the slice the page was read for, so that a stop condition can compare the page against the slice's own window or partition
Returns

A mapping {"next_page_token": } for the next page from the input response object. Returning None means there are no more pages to read in this response.

def get_page_size(self) -> Optional[int]:
103    def get_page_size(self) -> Optional[int]:
104        """
105        Evaluated against the config alone. A pagination strategy whose `page_size` template references the
106        response evaluates it with the response in context inside `next_page_token`, so the two can disagree -
107        pre-existing, and only observable for a `page_size` that is not a constant.
108
109        :return: the number of records this paginator asks for per page, or None if it does not define one
110        """
111        return None

Evaluated against the config alone. A pagination strategy whose page_size template references the response evaluates it with the response in context inside next_page_token, so the two can disagree - pre-existing, and only observable for a page_size that is not a constant.

Returns

the number of records this paginator asks for per page, or None if it does not define one

@abstractmethod
def path( self, next_page_token: Optional[Mapping[str, Any]], stream_state: Optional[Mapping[str, Any]] = None, stream_slice: Optional[airbyte_cdk.StreamSlice] = None) -> Optional[str]:
113    @abstractmethod
114    def path(
115        self,
116        next_page_token: Optional[Mapping[str, Any]],
117        stream_state: Optional[Mapping[str, Any]] = None,
118        stream_slice: Optional[StreamSlice] = None,
119    ) -> Optional[str]:
120        """
121        Returns the URL path to hit to fetch the next page of records
122
123        e.g: if you wanted to hit https://myapi.com/v1/some_entity then this will return "some_entity"
124
125        :return: path to hit to fetch the next request. Returning None means the path is not defined by the next_page_token
126        """
127        pass

Returns the URL path to hit to fetch the next page of records

e.g: if you wanted to hit https://myapi.com/v1/some_entity then this will return "some_entity"

Returns

path to hit to fetch the next request. Returning None means the path is not defined by the next_page_token