diff --git a/src/hflow/reader.py b/src/hflow/reader.py index 32f88c2d..f12aa2ad 100644 --- a/src/hflow/reader.py +++ b/src/hflow/reader.py @@ -29,16 +29,18 @@ from mcap.reader import McapReader, make_reader from mcap.records import Attachment from mcap.stream_reader import CRCValidationError +from zstandard import ZstdError logger = logging.getLogger(__name__) DEFAULT_BATCH_MAX_MESSAGES = 1024 DEFAULT_BATCH_MAX_BYTES = 32 * 1024 * 1024 -# The named reason a file fails its own integrity stamp, returned by +# Named reasons a file fails its own integrity stamp, returned by # :func:`verify_canonical_integrity` and recorded on the check lane's refusal # row, so downstream tooling can filter for damaged canonicals by exact value. CANONICAL_CRC_MISMATCH_REASON = "canonical-crc-mismatch" +CANONICAL_DECOMPRESSION_FAILED_REASON = "canonical-decompression-failed" @dataclass(frozen=True) @@ -339,7 +341,7 @@ def open_reader(path: Path | str, *, validate_crcs: bool = False) -> EpisodeRead def verify_canonical_integrity(path: Path | str) -> tuple[bool, str | None]: - """Validate one episode file's chunk CRCs with a strict full read. + """Validate one episode file's decompression and chunk CRCs with a strict full read. The check lane's front door. ``Episode`` reads run with CRC validation off (the reader docstring's trust argument covers bytes identified by @@ -350,11 +352,11 @@ def verify_canonical_integrity(path: Path | str) -> tuple[bool, str | None]: certify. Returns ``(is_valid, reason)``: ``(True, None)`` when every chunk - matches its stored CRC, and ``(False, CANONICAL_CRC_MISMATCH_REASON)`` - when the file refuses its own integrity stamp. ``CRCValidationError`` - is caught by its precise type -- it subclasses ``ValueError``, and the - broader type would also swallow unrelated boundary errors this function - must not answer for. + decompresses and matches its stored CRC, or ``(False, reason)`` for a + CRC mismatch or zstd decompression failure. Both exceptions are caught + by precise type: MCAP propagates ``ZstdError`` directly from the chunk + decompressor, before it can validate the CRC. Filesystem failures and + unrelated reader errors still propagate to the caller. """ with Path(path).open("rb") as stream: try: @@ -363,4 +365,6 @@ def verify_canonical_integrity(path: Path | str) -> tuple[bool, str | None]: pass except CRCValidationError: return (False, CANONICAL_CRC_MISMATCH_REASON) + except ZstdError: + return (False, CANONICAL_DECOMPRESSION_FAILED_REASON) return (True, None) diff --git a/tests/reuse_test_helpers.py b/tests/reuse_test_helpers.py index a8c69fa9..12e38ff4 100644 --- a/tests/reuse_test_helpers.py +++ b/tests/reuse_test_helpers.py @@ -47,3 +47,27 @@ def flip_chunk_payload_bytes(episode_path: Path, *, count: int = 4) -> None: def content_id_differs_from_delivery_receipt(episode_path: Path, receipt_content_id: str) -> bool: """True when the file on disk no longer matches the recorded content id.""" return content_episode_id(episode_path) != receipt_content_id + + +def corrupt_zstd_chunk_payload(episode_path: Path) -> None: + """Invalidate the first chunk's zstd frame magic, inside compressed records. + + Use the summary chunk index and the same MCAP field layout as + ``flip_chunk_payload_bytes``. Only the first byte of the zstd frame magic + changes (0x28 -> 0x29), forcing the actual decompressor to reject it. + The stored CRC, compression name, record lengths, and all MCAP headers, + indexes, metadata, and footer remain byte-for-byte intact. + """ + from mcap.reader import make_reader + + data = bytearray(episode_path.read_bytes()) + summary = make_reader(io.BytesIO(bytes(data))).get_summary() + assert summary is not None and summary.chunk_indexes + chunk_start = summary.chunk_indexes[0].chunk_start_offset + compression_length = struct.unpack_from(" None: + data_root = tmp_path / "data" + app, probe_runs, caption_runs = _app_with_probe_check(data_root) + source_uri = "episodes-in/episode.mcap" + source = synthesize_episode(data_root / source_uri, SPEC) + synced = app.process(source, stages={hflow.Stage.SYNC}, record=False) + assert verify_canonical_integrity(synced.canonical_path) == (True, None) + corrupt_zstd_chunk_payload(synced.canonical_path) + + reason = CANONICAL_DECOMPRESSION_FAILED_REASON + assert reason == "canonical-decompression-failed" + assert verify_canonical_integrity(synced.canonical_path) == (False, reason) + refused = app.process(source, stages="metadata_backfill") + assert refused.refusal_reason == reason + assert refused.has_errors + assert refused.checks == [] + assert f"REFUSED: {reason}" in refused.summary() + + relabel_refused = app.process(source, stages="relabel") + assert relabel_refused.refusal_reason == reason + assert relabel_refused.enrichments == [] + assert process_stage_batch(app, [source_uri], "meta") == { + "processed": 0, + "quarantined": 0, + "errors": 1, + } + assert probe_runs == [] + assert caption_runs == [] + + connection = open_catalog_connection(data_root / "catalog") + try: + rows = connection.execute( + "SELECT check_name, status, critical, error FROM check_runs" + ).fetchall() + episode_status = connection.execute("SELECT status FROM episodes").fetchall() + failures = connection.execute("SELECT failure_kind FROM ingest_failures").fetchall() + finally: + connection.close() + assert rows == [(CANONICAL_INTEGRITY_STEP_NAME, "error", True, reason)] + assert episode_status == [("unverified",)] + assert failures == [] + + +def test_one_error_filter_finds_both_canonical_corruption_species(tmp_path: Path) -> None: + data_root = tmp_path / "data" + app, probe_runs, _caption_runs = _app_with_probe_check(data_root) + for name, corrupt in ( + ("crc", _corrupt_first_chunk_crc), + ("zstd", corrupt_zstd_chunk_payload), + ): + source = synthesize_episode(data_root / "episodes-in" / f"{name}.mcap", SPEC) + synced = app.process(source, stages={hflow.Stage.SYNC}, record=False) + corrupt(synced.canonical_path) + app.process(source, stages="metadata_backfill") + + connection = open_catalog_connection(data_root / "catalog") + try: + rows = connection.execute( + "SELECT check_name, status, error FROM check_runs WHERE error IN (?, ?) ORDER BY error", + [CANONICAL_CRC_MISMATCH_REASON, CANONICAL_DECOMPRESSION_FAILED_REASON], + ).fetchall() + finally: + connection.close() + assert rows == [ + (CANONICAL_INTEGRITY_STEP_NAME, "error", CANONICAL_CRC_MISMATCH_REASON), + (CANONICAL_INTEGRITY_STEP_NAME, "error", CANONICAL_DECOMPRESSION_FAILED_REASON), + ] + assert probe_runs == [] + + +def test_missing_canonical_remains_an_infrastructure_failure(tmp_path: Path) -> None: + data_root = tmp_path / "data" + app, probe_runs, _caption_runs = _app_with_probe_check(data_root) + source_uri = "episodes-in/episode.mcap" + source = synthesize_episode(data_root / source_uri, SPEC) + synced = app.process(source, stages={hflow.Stage.SYNC}, record=False) + synced.canonical_path.unlink() + + with pytest.raises(FileNotFoundError): + verify_canonical_integrity(synced.canonical_path) + with pytest.raises(FileNotFoundError, match="no canonical episode exists"): + app.process(source, stages="metadata_backfill") + assert process_stage_batch(app, [source_uri], "meta") == { + "processed": 0, + "quarantined": 0, + "errors": 1, + } + assert probe_runs == [] + + connection = open_catalog_connection(data_root / "catalog") + try: + failures = connection.execute( + "SELECT source_uri, stage, failure_kind, error_type FROM ingest_failures" + ).fetchall() + refusal_rows = connection.execute("SELECT error FROM check_runs").fetchall() + finally: + connection.close() + assert failures == [(source_uri, "meta", "infrastructure", "FileNotFoundError")] + assert refusal_rows == []