From 7d7777c7b3cb8060757be911c871263943972615 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Thu, 20 Aug 2026 02:40:59 +0000 Subject: [PATCH 1/4] fix: support gzip decoder pass-through Co-Authored-By: syed.khadeer@airbyte.io --- .../decoders/composite_raw_decoder.py | 47 +++++++++++++++---- .../parsers/model_to_component_factory.py | 7 +-- .../decoders/test_composite_decoder.py | 36 ++++++++++++++ .../test_model_to_component_factory.py | 45 ++++++++++++++++++ 4 files changed, 122 insertions(+), 13 deletions(-) diff --git a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py index 0561369a7..d8000f873 100644 --- a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py +++ b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py @@ -15,6 +15,7 @@ import ijson import orjson import requests +from typing_extensions import Buffer from airbyte_cdk.models import FailureType from airbyte_cdk.sources.declarative.decoders.decoder import DECODER_OUTPUT_TYPE, Decoder @@ -29,23 +30,53 @@ logger = logging.getLogger("airbyte") +class _PrefixedStream(io.RawIOBase): + def __init__(self, prefix: bytes, stream: BufferedIOBase) -> None: + self._prefix = prefix + self._stream = stream + + def readable(self) -> bool: + return True + + def readinto(self, buffer: Buffer) -> int: + buffer_view = memoryview(buffer) + prefix_size = min(len(self._prefix), len(buffer_view)) + if prefix_size: + buffer_view[:prefix_size] = self._prefix[:prefix_size] + self._prefix = self._prefix[prefix_size:] + + if prefix_size == len(buffer_view): + return prefix_size + + data = self._stream.read(len(buffer_view) - prefix_size) + if not data: + return prefix_size + + buffer_view[prefix_size : prefix_size + len(data)] = data + return prefix_size + len(data) + + @dataclass class GzipParser(Parser): inner_parser: Parser def parse(self, data: BufferedIOBase) -> PARSER_OUTPUT_TYPE: - """ - Decompress gzipped bytes and pass decompressed data to the inner parser. + """Decompress gzipped data or pass uncompressed data through unchanged. - IMPORTANT: - - If the data is not gzipped, reset the pointer and pass the data to the inner parser as is. + Args: + data: A byte stream containing compressed or uncompressed data. - Note: - - The data is not decoded by default. + Yields: + Records parsed by the inner parser. """ + prefix = data.read(2) + prefixed_data = io.BufferedReader(_PrefixedStream(prefix, data)) - with gzip.GzipFile(fileobj=data, mode="rb") as gzipobj: - yield from self.inner_parser.parse(gzipobj) + if prefix == b"\x1f\x8b": + with gzip.GzipFile(fileobj=prefixed_data, mode="rb") as gzipobj: + yield from self.inner_parser.parse(gzipobj) + else: + yield from self.inner_parser.parse(prefixed_data) @dataclass 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 fbecdacb3..8a2f4bdc2 100644 --- a/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py +++ b/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py @@ -2783,15 +2783,12 @@ def create_gzip_decoder( gzip_parser: GzipParser = ModelToComponentFactory._get_parser(model, config) # type: ignore # based on the model, we know this will be a GzipParser if self._emit_connector_builder_messages: - # This is very surprising but if the response is not streamed, - # CompositeRawDecoder calls response.content and the requests library actually uncompress the data as opposed to response.raw, - # which uses urllib3 directly and does not uncompress the data. - return CompositeRawDecoder(gzip_parser.inner_parser, False) + return CompositeRawDecoder(gzip_parser, False) return CompositeRawDecoder.by_headers( [({"Content-Encoding", "Content-Type"}, _compressed_response_types, gzip_parser)], stream_response=True, - fallback_parser=gzip_parser.inner_parser, + fallback_parser=gzip_parser, ) @staticmethod diff --git a/unit_tests/sources/declarative/decoders/test_composite_decoder.py b/unit_tests/sources/declarative/decoders/test_composite_decoder.py index 8af9a4b4c..9d3fff603 100644 --- a/unit_tests/sources/declarative/decoders/test_composite_decoder.py +++ b/unit_tests/sources/declarative/decoders/test_composite_decoder.py @@ -73,6 +73,42 @@ def generate_csv( return csv_data.encode(encoding) +class NonSeekableBytesIO(BytesIO): + def seekable(self) -> bool: + return False + + +def test_gzip_parser_decompresses_gzip_payload(): + parser = GzipParser(inner_parser=CsvParser()) + + assert list(parser.parse(BytesIO(compress_with_gzip("date,units\n2026-08-01,42\n")))) == [ + {"date": "2026-08-01", "units": "42"} + ] + + +@pytest.mark.parametrize("stream_class", [BytesIO, NonSeekableBytesIO]) +def test_gzip_parser_passes_through_non_gzip_payload(stream_class): + parser = GzipParser(inner_parser=CsvParser()) + + assert list(parser.parse(stream_class(b"date,units\n2026-08-01,42\n"))) == [ + {"date": "2026-08-01", "units": "42"} + ] + + +def test_gzip_parser_handles_empty_payload(): + parser = GzipParser(inner_parser=CsvParser()) + + assert list(parser.parse(BytesIO())) == [] + + +def test_nested_gzip_parser_decompresses_single_gzip_payload(): + parser = GzipParser(inner_parser=GzipParser(inner_parser=CsvParser())) + + assert list(parser.parse(BytesIO(compress_with_gzip("date,units\n2026-08-01,42\n")))) == [ + {"date": "2026-08-01", "units": "42"} + ] + + @pytest.mark.parametrize("encoding", ["utf-8", "utf", "iso-8859-1"]) def test_composite_raw_decoder_gzip_csv_parser(requests_mock, encoding: str): requests_mock.register_uri( diff --git a/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py b/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py index 2d9694db2..e7d9575c6 100644 --- a/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py +++ b/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py @@ -1,6 +1,8 @@ # # Copyright (c) 2023 Airbyte, Inc., all rights reserved. # +import gzip +import io import json import logging from copy import deepcopy @@ -22,6 +24,7 @@ ) from freezegun.api import FakeDatetime from pydantic.v1 import ValidationError +from urllib3 import HTTPResponse from airbyte_cdk.legacy.sources.declarative.declarative_stream import DeclarativeStream from airbyte_cdk.legacy.sources.declarative.incremental import DatetimeBasedCursor @@ -96,12 +99,18 @@ from airbyte_cdk.sources.declarative.models.declarative_component_schema import ( ConstantBackoffStrategy as ConstantBackoffStrategyModel, ) +from airbyte_cdk.sources.declarative.models.declarative_component_schema import ( + CsvDecoder as CsvDecoderModel, +) from airbyte_cdk.sources.declarative.models.declarative_component_schema import ( CustomRequester as CustomRequesterModel, ) from airbyte_cdk.sources.declarative.models.declarative_component_schema import ( ExponentialBackoffStrategy as ExponentialBackoffStrategyModel, ) +from airbyte_cdk.sources.declarative.models.declarative_component_schema import ( + GzipDecoder as GzipDecoderModel, +) from airbyte_cdk.sources.declarative.models.declarative_component_schema import ( OffsetIncrement as OffsetIncrementModel, ) @@ -261,6 +270,42 @@ def test_create_check_stream(): assert check.stream_names == ["list_stream"] +@pytest.mark.parametrize( + "headers", + [ + {"Content-Type": "application/gzip"}, + {"Content-Type": "application/x-gzip"}, + {"Content-Type": "text/csv"}, + {"Content-Type": "binary/octet-stream"}, + {"Content-Encoding": "gzip"}, + ], +) +@pytest.mark.parametrize("emit_connector_builder_messages", [False, True]) +def test_create_gzip_decoder_handles_compressed_response( + headers: Mapping[str, str], emit_connector_builder_messages: bool +): + csv_data = b"date,units\n2026-08-01,42\n" + response = requests.Response() + response.status_code = 200 + response.headers.update(headers) + response.raw = HTTPResponse( + body=io.BytesIO(gzip.compress(csv_data)), + headers=headers, + status=200, + preload_content=False, + ) + + model = GzipDecoderModel( + type="GzipDecoder", + decoder=CsvDecoderModel(type="CsvDecoder"), + ) + decoder = ModelToComponentFactory( + emit_connector_builder_messages=emit_connector_builder_messages + ).create_gzip_decoder(model, {}) + + assert list(decoder.decode(response)) == [{"date": "2026-08-01", "units": "42"}] + + def test_create_component_type_mismatch(): manifest = {"check": {"type": "MismatchType", "stream_names": ["list_stream"]}} From 0dbb3759bdd66884d7db77d09480517b0b5f9afb Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Thu, 20 Aug 2026 02:43:30 +0000 Subject: [PATCH 2/4] fix: handle short gzip header reads Co-Authored-By: syed.khadeer@airbyte.io --- .../declarative/decoders/composite_raw_decoder.py | 9 ++++++++- .../decoders/test_composite_decoder.py | 15 +++++++++++++++ 2 files changed, 23 insertions(+), 1 deletion(-) diff --git a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py index d8000f873..930949e29 100644 --- a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py +++ b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py @@ -31,6 +31,8 @@ class _PrefixedStream(io.RawIOBase): + """Restore consumed header bytes ahead of the remaining stream.""" + def __init__(self, prefix: bytes, stream: BufferedIOBase) -> None: self._prefix = prefix self._stream = stream @@ -69,7 +71,12 @@ def parse(self, data: BufferedIOBase) -> PARSER_OUTPUT_TYPE: Yields: Records parsed by the inner parser. """ - prefix = data.read(2) + prefix = b"" + while len(prefix) < 2: + chunk = data.read(2 - len(prefix)) + if not chunk: + break + prefix += chunk prefixed_data = io.BufferedReader(_PrefixedStream(prefix, data)) if prefix == b"\x1f\x8b": diff --git a/unit_tests/sources/declarative/decoders/test_composite_decoder.py b/unit_tests/sources/declarative/decoders/test_composite_decoder.py index 9d3fff603..5d08ff1d2 100644 --- a/unit_tests/sources/declarative/decoders/test_composite_decoder.py +++ b/unit_tests/sources/declarative/decoders/test_composite_decoder.py @@ -78,6 +78,13 @@ def seekable(self) -> bool: return False +class OneByteAtATimeBytesIO(BytesIO): + def read(self, size: int = -1) -> bytes: + if size > 0: + size = 1 + return super().read(size) + + def test_gzip_parser_decompresses_gzip_payload(): parser = GzipParser(inner_parser=CsvParser()) @@ -86,6 +93,14 @@ def test_gzip_parser_decompresses_gzip_payload(): ] +def test_gzip_parser_handles_short_reads(): + parser = GzipParser(inner_parser=CsvParser()) + + assert list( + parser.parse(OneByteAtATimeBytesIO(compress_with_gzip("date,units\n2026-08-01,42\n"))) + ) == [{"date": "2026-08-01", "units": "42"}] + + @pytest.mark.parametrize("stream_class", [BytesIO, NonSeekableBytesIO]) def test_gzip_parser_passes_through_non_gzip_payload(stream_class): parser = GzipParser(inner_parser=CsvParser()) From 36b6431983290db7669b2308e884b0100ee06663 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Thu, 20 Aug 2026 02:50:26 +0000 Subject: [PATCH 3/4] fix: initialize prefixed gzip stream Co-Authored-By: syed.khadeer@airbyte.io --- .../sources/declarative/decoders/composite_raw_decoder.py | 1 + 1 file changed, 1 insertion(+) diff --git a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py index 930949e29..e29f863c8 100644 --- a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py +++ b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py @@ -34,6 +34,7 @@ class _PrefixedStream(io.RawIOBase): """Restore consumed header bytes ahead of the remaining stream.""" def __init__(self, prefix: bytes, stream: BufferedIOBase) -> None: + super().__init__() self._prefix = prefix self._stream = stream From 511654b9c62c7d63815115050649eeec1c1ea4a4 Mon Sep 17 00:00:00 2001 From: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Date: Fri, 21 Aug 2026 19:36:21 +0000 Subject: [PATCH 4/4] fix: strip transport and payload gzip layers independently Co-Authored-By: bot_apk --- .../parsers/model_to_component_factory.py | 6 +++- .../decoders/test_composite_decoder.py | 12 ++++++-- .../test_model_to_component_factory.py | 29 +++++++++++++++++++ 3 files changed, 44 insertions(+), 3 deletions(-) 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 8a2f4bdc2..e575a4c79 100644 --- a/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py +++ b/airbyte_cdk/sources/declarative/parsers/model_to_component_factory.py @@ -2785,8 +2785,12 @@ def create_gzip_decoder( if self._emit_connector_builder_messages: return CompositeRawDecoder(gzip_parser, False) + transport_gzip_parser = GzipParser(inner_parser=gzip_parser) return CompositeRawDecoder.by_headers( - [({"Content-Encoding", "Content-Type"}, _compressed_response_types, gzip_parser)], + [ + ({"Content-Encoding"}, {"gzip"}, transport_gzip_parser), + ({"Content-Type"}, _compressed_response_types, gzip_parser), + ], stream_response=True, fallback_parser=gzip_parser, ) diff --git a/unit_tests/sources/declarative/decoders/test_composite_decoder.py b/unit_tests/sources/declarative/decoders/test_composite_decoder.py index 5d08ff1d2..848ec1e9b 100644 --- a/unit_tests/sources/declarative/decoders/test_composite_decoder.py +++ b/unit_tests/sources/declarative/decoders/test_composite_decoder.py @@ -3,6 +3,7 @@ # import csv import gzip +import io import json import socket from http.server import BaseHTTPRequestHandler, HTTPServer @@ -77,6 +78,12 @@ class NonSeekableBytesIO(BytesIO): def seekable(self) -> bool: return False + def seek(self, *args, **kwargs) -> int: + raise io.UnsupportedOperation("seek") + + def tell(self) -> int: + raise io.UnsupportedOperation("tell") + class OneByteAtATimeBytesIO(BytesIO): def read(self, size: int = -1) -> bytes: @@ -85,10 +92,11 @@ def read(self, size: int = -1) -> bytes: return super().read(size) -def test_gzip_parser_decompresses_gzip_payload(): +@pytest.mark.parametrize("stream_class", [BytesIO, NonSeekableBytesIO]) +def test_gzip_parser_decompresses_gzip_payload(stream_class): parser = GzipParser(inner_parser=CsvParser()) - assert list(parser.parse(BytesIO(compress_with_gzip("date,units\n2026-08-01,42\n")))) == [ + assert list(parser.parse(stream_class(compress_with_gzip("date,units\n2026-08-01,42\n")))) == [ {"date": "2026-08-01", "units": "42"} ] diff --git a/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py b/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py index e7d9575c6..21c99adc7 100644 --- a/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py +++ b/unit_tests/sources/declarative/parsers/test_model_to_component_factory.py @@ -293,6 +293,35 @@ def test_create_gzip_decoder_handles_compressed_response( headers=headers, status=200, preload_content=False, + decode_content=False, + ) + + model = GzipDecoderModel( + type="GzipDecoder", + decoder=CsvDecoderModel(type="CsvDecoder"), + ) + decoder = ModelToComponentFactory( + emit_connector_builder_messages=emit_connector_builder_messages + ).create_gzip_decoder(model, {}) + + assert list(decoder.decode(response)) == [{"date": "2026-08-01", "units": "42"}] + + +@pytest.mark.parametrize("emit_connector_builder_messages", [False, True]) +def test_create_gzip_decoder_handles_transport_and_content_gzip( + emit_connector_builder_messages: bool, +): + csv_data = b"date,units\n2026-08-01,42\n" + headers = {"Content-Encoding": "gzip", "Content-Type": "application/gzip"} + response = requests.Response() + response.status_code = 200 + response.headers.update(headers) + response.raw = HTTPResponse( + body=io.BytesIO(gzip.compress(gzip.compress(csv_data))), + headers=headers, + status=200, + preload_content=False, + decode_content=False, ) model = GzipDecoderModel(