Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 9 additions & 6 deletions docs/CATALOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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-<timestamp>.parquet` (the selection, which
registered step's latest run settled at its current version, then writes two
immutable files: `manifests/clean-<timestamp>.parquet` (the selection, which
`hflow export snapshot --manifest` reads) and `manifests/clean-<timestamp>.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
Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion src/hflow/curation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
"""
Expand Down
29 changes: 13 additions & 16 deletions src/hflow/dataset.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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.
"""
Expand All @@ -228,7 +225,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"
Expand Down
34 changes: 34 additions & 0 deletions tests/test_catalog_curation.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,13 +31,15 @@
from hflow.checks import camera_frame_stats
from hflow.cli import main as cli_main
from hflow.curation import (
CheckCoverage,
CurationReport,
NonSingleSelectQueryError,
curate,
open_catalog_connection,
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

Expand Down Expand Up @@ -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:
Expand Down
33 changes: 33 additions & 0 deletions tests/test_dataset.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
Loading