Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -556,7 +556,8 @@ def _dynamic_stream_configs(
f"Expected one of {list(COMPONENTS_RESOLVER_TYPE_MAPPING.keys())}."
)

if "retriever" in components_resolver_config:
# AsyncRetriever has no `requester`, and a cached polling request replays a stale status
if "requester" in components_resolver_config.get("retriever", {}):
components_resolver_config["retriever"]["requester"]["use_cache"] = True

# Create a resolver for dynamic components based on type
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -890,7 +890,8 @@ def _dynamic_stream_configs(
f"Expected one of {list(COMPONENTS_RESOLVER_TYPE_MAPPING.keys())}."
)

if "retriever" in components_resolver_config:
# AsyncRetriever has no `requester`, and a cached polling request replays a stale status
if "requester" in components_resolver_config.get("retriever", {}):
components_resolver_config["retriever"]["requester"]["use_cache"] = True

# Create a resolver for dynamic components based on type
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5268,12 +5268,14 @@ def create_components_mapping_definition(
def create_http_components_resolver(
self, model: HttpComponentsResolverModel, config: Config, stream_name: Optional[str] = None
) -> Any:
partition_router = self._build_stream_slicer_from_partition_router(model.retriever, config)
retriever = self._create_component_from_model(
model=model.retriever,
config=config,
name=f"{stream_name if stream_name else '__http_components_resolver'}",
primary_key=None,
stream_slicer=self._build_stream_slicer_from_partition_router(model.retriever, config),
stream_slicer=partition_router,
partition_router=partition_router,
Comment thread
bazarnov marked this conversation as resolved.
transformations=[],
)

Expand All @@ -5293,7 +5295,10 @@ def create_http_components_resolver(

return HttpComponentsResolver(
retriever=retriever,
stream_slicer=self._build_stream_slicer_from_partition_router(model.retriever, config),
# AsyncRetriever reads records only from the job slices its own slicer yields
stream_slicer=retriever.stream_slicer
if isinstance(retriever, AsyncRetriever)
else partition_router,
config=config,
components_mapping=components_mapping,
parameters=model.parameters or {},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import dpath
from typing_extensions import deprecated

from airbyte_cdk.models import AirbyteMessage
from airbyte_cdk.sources.declarative.interpolation import InterpolatedString
from airbyte_cdk.sources.declarative.resolvers.components_resolver import (
ComponentMappingDefinition,
Expand Down Expand Up @@ -95,6 +96,9 @@ def resolve_components(
for components_values in self.retriever.read_records(
records_schema={}, stream_slice=stream_slice
):
# AsyncRetriever emits a slice log message before its records
if isinstance(components_values, AirbyteMessage):
continue
updated_config = deepcopy(stream_template_config)
kwargs["components_values"] = components_values # type: ignore[assignment] # component_values will always be of type Mapping[str, Any]
kwargs["stream_slice"] = stream_slice # type: ignore[assignment] # stream_slice will always be of type Mapping[str, Any]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

import json
from copy import deepcopy
from unittest.mock import MagicMock
from unittest.mock import MagicMock, patch

import pytest

Expand All @@ -14,6 +14,7 @@
DestinationSyncMode,
Type,
)
from airbyte_cdk.sources.declarative.async_job.job_orchestrator import AsyncJobOrchestrator
from airbyte_cdk.sources.declarative.concurrent_declarative_source import (
ConcurrentDeclarativeSource,
)
Expand Down Expand Up @@ -629,3 +630,89 @@ def test_dynamic_streams_with_http_components_resolver_retriever_with_parent_str
actual_record_stream_names.sort()

assert actual_record_stream_names == expected_stream_names


def test_dynamic_streams_with_http_components_resolver_partition_router_request_option():
manifest = deepcopy(_MANIFEST)
manifest["dynamic_streams"][0]["components_resolver"]["retriever"]["partition_router"] = {
"type": "ListPartitionRouter",
"cursor_field": "p",
"values": ["p1"],
"request_option": {
"type": "RequestOption",
"inject_into": "request_parameter",
"field_name": "p",
},
}
with HttpMocker() as http_mocker:
http_mocker.get(
HttpRequest(url="https://api.test.com/items?p=p1"),
HttpResponse(body=json.dumps([{"id": 1, "name": "item_1"}])),
)

source = ConcurrentDeclarativeSource(
source_config=manifest, config=_CONFIG, catalog=None, state=None
)
actual_catalog = source.discover(logger=source.logger, config=_CONFIG)

assert [stream.name for stream in actual_catalog.streams] == ["item_1"]


@patch.object(AsyncJobOrchestrator, "_WAIT_TIME_BETWEEN_STATUS_UPDATE_IN_SECONDS", 0)
def test_dynamic_streams_with_http_components_resolver_async_retriever():
manifest = deepcopy(_MANIFEST)
manifest["dynamic_streams"][0]["components_resolver"]["retriever"] = {
"type": "AsyncRetriever",
"status_mapping": {
"failed": ["failed"],
"running": ["pending"],
"timeout": ["timeout"],
"completed": ["ready"],
},
"status_extractor": {"type": "DpathExtractor", "field_path": ["status"]},
"download_target_extractor": {"type": "DpathExtractor", "field_path": ["urls"]},
"record_selector": {
"type": "RecordSelector",
"extractor": {"type": "DpathExtractor", "field_path": []},
},
"creation_requester": {
"type": "HttpRequester",
"url": "https://api.test.com/items_job",
"http_method": "POST",
},
"polling_requester": {
"type": "HttpRequester",
"url": "https://api.test.com/items_job/{{ creation_response['id'] }}",
"http_method": "GET",
},
"download_requester": {
"type": "HttpRequester",
"url": "{{ download_target }}",
"http_method": "GET",
},
}
with HttpMocker() as http_mocker:
http_mocker.post(
HttpRequest(url="https://api.test.com/items_job"),
HttpResponse(body=json.dumps({"id": "job_1"})),
)
http_mocker.get(
HttpRequest(url="https://api.test.com/items_job/job_1"),
HttpResponse(
body=json.dumps({"status": "ready", "urls": ["https://api.test.com/download/1"]})
),
)
http_mocker.get(
HttpRequest(url="https://api.test.com/download/1"),
HttpResponse(
body=json.dumps([{"id": 1, "name": "item_1"}, {"id": 2, "name": "item_2"}])
),
)

source = ConcurrentDeclarativeSource(
source_config=manifest, config=_CONFIG, catalog=None, state=None
)
actual_catalog = source.discover(logger=source.logger, config=_CONFIG)

# the slice log message AsyncRetriever emits before its records must not resolve a stream
assert [stream.name for stream in actual_catalog.streams] == ["item_1", "item_2"]
Loading