diff --git a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py index 0561369a7..440e4b147 100644 --- a/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py +++ b/airbyte_cdk/sources/declarative/decoders/composite_raw_decoder.py @@ -28,6 +28,9 @@ logger = logging.getLogger("airbyte") +# Leading bytes that identify a gzip stream (RFC 1952). +GZIP_MAGIC_NUMBER = b"\x1f\x8b" + @dataclass class GzipParser(Parser): @@ -38,14 +41,26 @@ def parse(self, data: BufferedIOBase) -> PARSER_OUTPUT_TYPE: Decompress gzipped bytes and pass decompressed data to the inner parser. IMPORTANT: - - If the data is not gzipped, reset the pointer and pass the data to the inner parser as is. + - If the data is not gzipped, pass the data to the inner parser as is. Note: - The data is not decoded by default. """ - with gzip.GzipFile(fileobj=data, mode="rb") as gzipobj: - yield from self.inner_parser.parse(gzipobj) + # Wrap in a BufferedReader so we can peek at the leading bytes without + # consuming them. This lets nested/fallback gzip parsers gracefully handle + # payloads that are not actually gzip-compressed (e.g. a response that + # intermittently arrives with fewer gzip layers than the decoder expects) + # instead of raising BadGzipFile. + buffered_data = data if isinstance(data, io.BufferedReader) else io.BufferedReader(data) # type: ignore[arg-type] + if ( + buffered_data.peek(len(GZIP_MAGIC_NUMBER))[: len(GZIP_MAGIC_NUMBER)] + == GZIP_MAGIC_NUMBER + ): + with gzip.GzipFile(fileobj=buffered_data, mode="rb") as gzipobj: + yield from self.inner_parser.parse(gzipobj) + else: + yield from self.inner_parser.parse(buffered_data) @dataclass diff --git a/unit_tests/sources/declarative/decoders/test_composite_decoder.py b/unit_tests/sources/declarative/decoders/test_composite_decoder.py index 8af9a4b4c..ad77c7568 100644 --- a/unit_tests/sources/declarative/decoders/test_composite_decoder.py +++ b/unit_tests/sources/declarative/decoders/test_composite_decoder.py @@ -93,6 +93,74 @@ def test_composite_raw_decoder_gzip_csv_parser(requests_mock, encoding: str): assert counter == 3 +@pytest.mark.parametrize("encoding", ["utf-8", "iso-8859-1"]) +def test_gzip_parser_passes_through_uncompressed_data(requests_mock, encoding: str): + """A GzipParser should gracefully parse data that is not gzip-compressed. + + This mirrors the documented behavior: if the payload is not gzipped, the data + is passed to the inner parser as-is instead of raising BadGzipFile. + """ + requests_mock.register_uri( + "GET", + "https://airbyte.io/", + content=generate_csv(encoding=encoding, delimiter="\t", should_compress=False), + ) + response = requests.get("https://airbyte.io/", stream=True) + + parser = GzipParser(inner_parser=CsvParser(encoding=encoding, delimiter="\\t")) + composite_raw_decoder = CompositeRawDecoder(parser=parser) + + assert len(list(composite_raw_decoder.decode(response))) == 3 + + +@pytest.mark.parametrize("should_compress", [True, False]) +def test_nested_gzip_parser_tolerates_variable_gzip_layers(requests_mock, should_compress: bool): + """Nested GzipParser(GzipParser(CsvParser)) should handle a payload that has + either one or two gzip layers. + + App Store Connect style responses can intermittently arrive with fewer gzip + layers than the decoder was configured for. The outer/inner parsers should + fall through to the inner parser when a layer is not actually gzipped rather + than raising BadGzipFile. + """ + csv_bytes = generate_csv(should_compress=True) # inner gzip layer (the report file) + if should_compress: + buf = BytesIO() + with gzip.GzipFile(fileobj=buf, mode="wb") as f: + f.write(csv_bytes) + content = buf.getvalue() # second (outer) gzip layer + else: + content = csv_bytes + + requests_mock.register_uri("GET", "https://airbyte.io/", content=content) + response = requests.get("https://airbyte.io/", stream=True) + + parser = GzipParser(inner_parser=GzipParser(inner_parser=CsvParser())) + composite_raw_decoder = CompositeRawDecoder(parser=parser) + + assert len(list(composite_raw_decoder.decode(response))) == 3 + + +def test_gzip_fallback_parser_handles_uncompressed_response(requests_mock): + """When header-based selection falls back to a GzipParser but the response is + not gzipped, records should still be parsed without raising BadGzipFile.""" + requests_mock.register_uri( + "GET", + "https://airbyte.io/", + content="".join(generate_jsonlines()).encode("utf-8"), + headers={"Content-Encoding": "identity"}, + ) + response = requests.get("https://airbyte.io/", stream=True) + + composite_raw_decoder = CompositeRawDecoder.by_headers( + [({"Content-Encoding"}, {"gzip"}, GzipParser(inner_parser=JsonLineParser()))], + stream_response=True, + fallback_parser=GzipParser(inner_parser=JsonLineParser()), + ) + + assert len(list(composite_raw_decoder.decode(response))) == 3 + + def generate_jsonlines() -> Iterable[str]: """ Generator function to yield data in JSON Lines format.