From 73ca7af0cc21b19bbdfcf5b2d497d839afacb9ac Mon Sep 17 00:00:00 2001 From: k0shir0 Date: Sat, 3 Oct 2026 18:06:19 -0500 Subject: [PATCH 1/2] fix(curation,dataset): query check_runs_latest in coverage and dataset SQL Fixes #676. --- src/hflow/curation.py | 2 +- src/hflow/dataset.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/hflow/curation.py b/src/hflow/curation.py index 7dc82558..6336f6f5 100644 --- a/src/hflow/curation.py +++ b/src/hflow/curation.py @@ -749,7 +749,7 @@ def _collect_coverage(connection: duckdb.DuckDBPyConnection) -> tuple[int, list[ coverage_rows = connection.execute( f""" SELECT check_name, count(DISTINCT episode_id) AS episodes_ran - FROM check_runs WHERE status IN ({status_list}) + FROM check_runs_latest WHERE status IN ({status_list}) AND episode_id IN (SELECT episode_id FROM episodes_latest) GROUP BY check_name ORDER BY check_name """ diff --git a/src/hflow/dataset.py b/src/hflow/dataset.py index 94a42379..39a42683 100644 --- a/src/hflow/dataset.py +++ b/src/hflow/dataset.py @@ -228,7 +228,7 @@ def default_dataset_sql(application: "App") -> str: ) predicates.append( "episode_id IN (\n" - " SELECT episode_id FROM check_runs\n" + " SELECT episode_id FROM check_runs_latest\n" f" WHERE status IN ({settled_statuses})\n" f" AND ({identity_clauses})\n" " GROUP BY episode_id\n" From 06bf4bc738f271313b05de52a455ef27a253e504 Mon Sep 17 00:00:00 2001 From: k0shir0 Date: Sat, 3 Oct 2026 21:54:19 -0500 Subject: [PATCH 2/2] add replay error regression tests, clarify latest check results --- docs/CATALOG.md | 15 +++++++++------ src/hflow/dataset.py | 27 ++++++++++++--------------- tests/test_catalog_curation.py | 34 ++++++++++++++++++++++++++++++++++ tests/test_dataset.py | 33 +++++++++++++++++++++++++++++++++ 4 files changed, 88 insertions(+), 21 deletions(-) diff --git a/docs/CATALOG.md b/docs/CATALOG.md index 68c921f5..5018256e 100644 --- a/docs/CATALOG.md +++ b/docs/CATALOG.md @@ -449,13 +449,15 @@ hflow dataset create clean --print-sql # see the policy, change nothing It selects the current generation of every source recording that is not quarantined, was produced by this pipeline's current transform, and has every -registered step recorded at its current version, then writes two immutable -files: `manifests/clean-.parquet` (the selection, which +registered step's latest run settled at its current version, then writes two +immutable files: `manifests/clean-.parquet` (the selection, which `hflow export snapshot --manifest` reads) and `manifests/clean-.json` (the effective SQL, the version stamps it required, the row count, and the coverage). Nothing is hidden -- `--print-sql` gives you the query to edit into a sharper one, and `--sql` runs yours instead while keeping the artifact and -the provenance record. +the provenance record. If a check or enrichment errors on replay, its earlier +settled result no longer qualifies the episode, even when the step is noncritical +and the episode's status remains `ok`. Two subtleties worth knowing if you write the equivalent yourself, because both of them silently return an empty dataset that reads like a policy @@ -632,9 +634,10 @@ coverage (episodes each check ran on): A statistic over half a delivery must not look like a statistic over all of it: steps skip when an episode quarantines upstream, so partial coverage is -normal, and it must be visible, not inferred. "Ran" means a `check_runs` -status of `passed`, `failed`, or `measured` (a failed verdict still ran; -`skipped` and `error` did not produce evidence). +normal, and it must be visible, not inferred. "Ran" means the step's latest +`check_runs_latest` status is `passed`, `failed`, or `measured` (a failed verdict +still ran; `skipped` and `error` did not produce evidence). An error on replay +withdraws the earlier run's coverage; the episode stays in the denominator. Registration order therefore decides how much evidence a quarantine costs. A critical check that quarantines skips every check registered **after** it, and diff --git a/src/hflow/dataset.py b/src/hflow/dataset.py index 39a42683..a0f6fd2e 100644 --- a/src/hflow/dataset.py +++ b/src/hflow/dataset.py @@ -179,23 +179,20 @@ def default_dataset_sql(application: "App") -> str: transform change does not mix two canonical behaviors in one dataset. 2. **Status is ``ok``.** Excludes two different things. ``quarantined`` is the pipeline's own critical checks rejecting the episode. ``unverified`` - is a critical check that crashed, so nobody actually checked it, and - that is the half rule 3 cannot see: a crash leaves no settled row, but - an EARLIER settled run of the same check satisfies rule 3 on its own, - and the episode would otherwise land in the dataset on the strength of - a result that a later run withdrew. - 3. **Every registered step settled, at its current version.** A check added - last week that has not been backfilled leaves its episodes out rather - than silently reporting a dataset with a hole in it. + is a critical check that crashed, so nobody actually checked it. + 3. **Every registered step's latest run settled, at its current version.** + A check added last week that has not been backfilled leaves its episodes + out rather than silently reporting a dataset with a hole in it. A later + error withdraws an earlier settled result, including for noncritical + checks and enrichments whose errors leave the episode's status ``ok``. 4. **One row per source recording**, which the ``episodes`` view already guarantees, so a reprocessed recording contributes its current generation and not both. - Rules 2 and 3 overlap without either being redundant. An episode whose - critical check ONLY ever crashed is dropped by rule 3, which needs a - settled row and never gets one, and rule 2 agrees. An episode that settled - once and crashed later is dropped by rule 2 alone. Neither rule subsumes - the other, so both stay. + Rules 2 and 3 serve different purposes. A critical check that rejects + an episode still settled with ``failed``, so rule 2 excludes what rule 3 + would allow. A noncritical check that errors leaves the episode's status + ``ok``, so rule 3 excludes what rule 2 would allow. Rule 3 has two traps in it, and both of them yield an EMPTY dataset that looks like a policy decision: @@ -204,9 +201,9 @@ def default_dataset_sql(application: "App") -> str: built-in library is evidence-only and records ``measured``, so a ``status = 'passed'`` reading selects nothing at all. - "Settled" is wider than "ran", because a default check the pipeline - supersedes records ``skipped`` on every episode forever -- and wrapping + supersedes records ``superseded`` on every episode forever -- and wrapping a built-in to configure it is the documented way to configure one, so - reading ``skipped`` as an unfilled hole empties the dataset of the + reading ``superseded`` as an unfilled hole empties the dataset of the pipelines most likely to want it. See :data:`hflow.steps.SETTLED_STATUSES` for why that is safe here and where the two differ. """ diff --git a/tests/test_catalog_curation.py b/tests/test_catalog_curation.py index 0f49f372..23173fc8 100644 --- a/tests/test_catalog_curation.py +++ b/tests/test_catalog_curation.py @@ -31,6 +31,7 @@ from hflow.checks import camera_frame_stats from hflow.cli import main as cli_main from hflow.curation import ( + CheckCoverage, CurationReport, NonSingleSelectQueryError, curate, @@ -38,6 +39,7 @@ reject_non_single_select, ) from hflow.format import CATALOG_FORMAT_VERSION +from hflow.stage_execution import run_stages_directly from hflow.testing import SyntheticEpisodeSpec, synthesize_episode from hflow.transform import EpisodeStamps @@ -718,6 +720,38 @@ def test_coverage_denominators(recorded_data_root: Path) -> None: assert "late_check: 1/2 (50%)" in report.summary() +def test_a_check_that_errors_on_replay_loses_coverage( + tmp_path: Path, one_second_camera_less_episode: Path +) -> None: + data_root = tmp_path / "data" + episode_uri = "episodes-in/episode_0001.mcap" + episode_path = data_root / episode_uri + episode_path.parent.mkdir(parents=True) + episode_path.write_bytes(one_second_camera_less_episode.read_bytes()) + app = hflow.App("errored-check-coverage", data_root=data_root, default_checks=()) + should_fail = False + + @app.check(version="1") + async def flaky_check(ep: hflow.Episode) -> hflow.CheckResult: + if should_fail: + raise RuntimeError("simulated check failure") + return hflow.CheckResult(verdict=True) + + asyncio.run(run_stages_directly(app, [episode_uri], hflow.RUN_PROFILES["full"])) + before = curate(app.workspace.catalog_root, "SELECT episode_id FROM episodes") + assert before.total_episodes == 1 + assert before.coverage == [CheckCoverage("flaky_check", 1, 1)] + + should_fail = True + with pytest.raises(RuntimeError, match="1 of 1 episodes had processing errors"): + asyncio.run(run_stages_directly(app, [episode_uri], {hflow.Stage.META})) + + after = curate(app.workspace.catalog_root, "SELECT episode_id FROM episodes") + assert after.total_episodes == 1 + assert after.row_count == 1 + assert after.coverage == [] + + def test_cli_curate( recorded_data_root: Path, tmp_path: Path, capsys: pytest.CaptureFixture[str] ) -> None: diff --git a/tests/test_dataset.py b/tests/test_dataset.py index 7693c354..2c8bcd4f 100644 --- a/tests/test_dataset.py +++ b/tests/test_dataset.py @@ -16,6 +16,7 @@ dataset_slug, default_dataset_sql, ) +from hflow.stage_execution import run_stages_directly from hflow.storage import BucketStorageRoot from hflow.testing import SyntheticEpisodeSpec, synthesize_episode from hflow.workspace import Workspace @@ -58,6 +59,38 @@ def _ingest(project: Path, monkeypatch: pytest.MonkeyPatch) -> None: class TestDefaultPolicy: + def test_a_check_that_errors_on_replay_excludes_its_episode( + self, tmp_path: Path, one_second_camera_less_episode: Path + ) -> None: + data_root = tmp_path / "data" + episode_uri = "episodes-in/episode_0001.mcap" + episode_path = data_root / episode_uri + episode_path.parent.mkdir(parents=True) + episode_path.write_bytes(one_second_camera_less_episode.read_bytes()) + app = hflow.App("errored-check-dataset", data_root=data_root, default_checks=()) + should_fail = False + + @app.check(version="1") + async def flaky_check(ep: hflow.Episode) -> hflow.CheckResult: + if should_fail: + raise RuntimeError("simulated check failure") + return hflow.CheckResult(verdict=True) + + asyncio.run(run_stages_directly(app, [episode_uri], hflow.RUN_PROFILES["full"])) + sql = f"SELECT episode_id FROM ({default_dataset_sql(app)})" + with hflow.open_catalog_connection(app.workspace.catalog_root) as connection: + assert len(connection.execute(sql).fetchall()) == 1 + + should_fail = True + with pytest.raises(RuntimeError, match="1 of 1 episodes had processing errors"): + asyncio.run(run_stages_directly(app, [episode_uri], {hflow.Stage.META})) + + with hflow.open_catalog_connection(app.workspace.catalog_root) as connection: + # A noncritical error leaves the episode ok; the latest-step rule + # must withdraw it rather than relying on the critical-check gate. + assert connection.execute("SELECT status FROM episodes").fetchall() == [("ok",)] + assert connection.execute(sql).fetchall() == [] + def test_a_step_that_never_ran_excludes_its_episodes( self, ingested_project: Path, monkeypatch: pytest.MonkeyPatch ) -> None: