From c2e49dba5801599d4e76eb3640da0d2a0dfeccf1 Mon Sep 17 00:00:00 2001 From: Baz <22987674+bazarnov@users.noreply.github.com> Date: Fri, 2 Oct 2026 12:08:52 +0300 Subject: [PATCH] fix(low-code): send HttpComponentsResolver partition router request options and support AsyncRetriever --- .../manifest_declarative_source.py | 3 +- .../concurrent_declarative_source.py | 3 +- .../parsers/model_to_component_factory.py | 9 +- .../resolvers/http_components_resolver.py | 4 + .../test_http_components_resolver.py | 89 ++++++++++++++++++- 5 files changed, 103 insertions(+), 5 deletions(-) diff --git a/airbyte_cdk/legacy/sources/declarative/manifest_declarative_source.py b/airbyte_cdk/legacy/sources/declarative/manifest_declarative_source.py index dea33879ac..5a20bdf1f3 100644 --- a/airbyte_cdk/legacy/sources/declarative/manifest_declarative_source.py +++ b/airbyte_cdk/legacy/sources/declarative/manifest_declarative_source.py @@ -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 diff --git a/airbyte_cdk/sources/declarative/concurrent_declarative_source.py b/airbyte_cdk/sources/declarative/concurrent_declarative_source.py index 52f83a644d..0c5e9da08c 100644 --- a/airbyte_cdk/sources/declarative/concurrent_declarative_source.py +++ b/airbyte_cdk/sources/declarative/concurrent_declarative_source.py @@ -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 diff --git a/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py b/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py index 6c8b630676..4ea2df6967 100644 --- a/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py +++ b/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py @@ -5237,12 +5237,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, transformations=[], ) @@ -5262,7 +5264,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 {}, diff --git a/airbyte_cdk/sources/declarative/resolvers/http_components_resolver.py b/airbyte_cdk/sources/declarative/resolvers/http_components_resolver.py index 11952b963b..09ffa60785 100644 --- a/airbyte_cdk/sources/declarative/resolvers/http_components_resolver.py +++ b/airbyte_cdk/sources/declarative/resolvers/http_components_resolver.py @@ -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, @@ -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] diff --git a/unit_tests/sources/declarative/resolvers/test_http_components_resolver.py b/unit_tests/sources/declarative/resolvers/test_http_components_resolver.py index 991539e1ea..24f4b9f0bd 100644 --- a/unit_tests/sources/declarative/resolvers/test_http_components_resolver.py +++ b/unit_tests/sources/declarative/resolvers/test_http_components_resolver.py @@ -4,7 +4,7 @@ import json from copy import deepcopy -from unittest.mock import MagicMock +from unittest.mock import MagicMock, patch import pytest @@ -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, ) @@ -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"]