diff --git a/packages/hflow-server/src/hflow_server/_catalog.py b/packages/hflow-server/src/hflow_server/_catalog.py index c615f74f..5a655e53 100644 --- a/packages/hflow-server/src/hflow_server/_catalog.py +++ b/packages/hflow-server/src/hflow_server/_catalog.py @@ -607,12 +607,13 @@ def find_canonical_uri(connection: duckdb.DuckDBPyConnection, episode_id: str) - def query_latest_run_intervals( connection: duckdb.DuckDBPyConnection, episode_id: str ) -> list[EpisodeIntervalRecord]: - """One episode's intervals from its LATEST run -- the current evidence. + """One episode's intervals from each check's LATEST run -- the current evidence. - ``check_version`` rides in from that run's ``check_runs`` row because the + ``check_version`` rides in from that run's ``check_runs_latest`` row because the intervals table does not carry one itself. One owner for this join: the dossier and the timeline must never disagree about which run's intervals - an episode "has". + an episode "has", and an earlier stage's intervals remain visible after + subsequent stages append. """ return _validated_records( EpisodeIntervalRecord, @@ -620,11 +621,8 @@ def query_latest_run_intervals( """ SELECT i.label, i.start_ns, i.end_ns, i.check_name, r.check_version FROM intervals AS i - JOIN episodes_latest AS e - ON i.episode_id = e.episode_id AND i.run_fingerprint = e.run_fingerprint - LEFT JOIN check_runs AS r - ON r.episode_id = i.episode_id AND r.run_fingerprint = i.run_fingerprint - AND r.check_name = i.check_name + JOIN check_runs_latest AS r + USING (episode_id, run_fingerprint, check_name) WHERE i.episode_id = ? ORDER BY i.start_ns, i.label """, @@ -680,16 +678,16 @@ def query_episode_dossier( [episode_id], ), ) - # Intervals and tags are the episode's LATEST run only -- the current - # evidence. + # Intervals and tags are each check's LATEST run only -- the current + # evidence, matching check_runs_latest. intervals = query_latest_run_intervals(connection, episode_id) tags = _validated_records( EpisodeTagRecord, connection.execute( - f"SELECT t.tag, t.check_name, {_recorded_at_as_iso_text('t.recorded_at')} " + f"SELECT t.tag, t.check_name, {_recorded_at_as_iso_text('r.recorded_at')} " "FROM tags AS t " - "JOIN episodes_latest AS e " - " ON t.episode_id = e.episode_id AND t.run_fingerprint = e.run_fingerprint " + "JOIN check_runs_latest AS r " + " USING (episode_id, run_fingerprint, check_name) " "WHERE t.episode_id = ? ORDER BY t.tag", [episode_id], ), diff --git a/packages/hflow-server/tests/test_server_episode_dossier.py b/packages/hflow-server/tests/test_server_episode_dossier.py index 7d4bdaee..7f4e3fcf 100644 --- a/packages/hflow-server/tests/test_server_episode_dossier.py +++ b/packages/hflow-server/tests/test_server_episode_dossier.py @@ -189,3 +189,103 @@ def test_dossier_reports_unverified_when_a_critical_check_crashed( episode = _dossier(api, append.episode_id)["episode"] assert episode["status"] == "unverified" assert episode["quarantine_tags"] == [] + + +def test_dossier_intervals_and_tags_survive_subsequent_stage_append( + tmp_path: Path, unbuilt_assets_dir: Path +) -> None: + """An earlier stage's intervals and tags remain visible after a subsequent stage appends.""" + data_root = tmp_path / "data" + catalog_root = data_root / "catalog" + episodes_directory = data_root / "episodes" + episodes_directory.mkdir(parents=True) + canonical = episodes_directory / "multi_stage.canonical.mcap" + canonical.write_bytes(b"canonical for multi-stage episode") + + catalog = Catalog(catalog_root) + # Stage 1: META records intervals and tags + meta_append = catalog.append_episode( + canonical_path=canonical, + stamps=STAMPS, + episode_metadata={"task": "fold_napkin"}, + check_rows=[ + CheckRunRow( + check_name="camera_freeze", + check_version="v1", + critical=False, + status=hflow.CheckStatus.MEASURED, + duration_s=0.01, + intervals=[hflow.Interval(start_ns=100, end_ns=500, label="frozen:front_camera")], + tags=["frozen:front_camera"], + ) + ], + ) + assert meta_append.written + + # Stage 2: LABELS appends another check to the same episode + labels_append = catalog.append_episode( + canonical_path=canonical, + stamps=STAMPS, + episode_metadata={"task": "fold_napkin"}, + check_rows=[ + CheckRunRow( + check_name="action_chunking", + check_version="v1", + critical=False, + status=hflow.CheckStatus.MEASURED, + duration_s=0.02, + ) + ], + ) + assert labels_append.written + assert labels_append.episode_id == meta_append.episode_id + assert labels_append.run_fingerprint != meta_append.run_fingerprint + + api = TestClient( + create_app(ServerSettings(data_root=str(data_root), assets_dir=unbuilt_assets_dir)) + ) + dossier = _dossier(api, meta_append.episode_id) + assert dossier["intervals"] == [ + { + "label": "frozen:front_camera", + "start_ns": 100, + "end_ns": 500, + "check_name": "camera_freeze", + "check_version": "v1", + } + ] + assert len(dossier["tags"]) == 1 + assert dossier["tags"][0]["tag"] == "frozen:front_camera" + assert dossier["tags"][0]["check_name"] == "camera_freeze" + + timeline_response = api.get(f"/api/v1/episodes/{meta_append.episode_id}/timeline") + assert timeline_response.status_code == 200 + timeline = timeline_response.json() + assert len(timeline["intervals"]) == 1 + assert timeline["intervals"][0]["label"] == "frozen:front_camera" + assert timeline["intervals"][0]["check_name"] == "camera_freeze" + assert timeline["intervals"][0]["kind"] == "frozen" + + # Stage 3: camera_freeze re-runs with a new version that omits intervals and tags + rerun_append = catalog.append_episode( + canonical_path=canonical, + stamps=STAMPS, + episode_metadata={"task": "fold_napkin"}, + check_rows=[ + CheckRunRow( + check_name="camera_freeze", + check_version="v2", + critical=False, + status=hflow.CheckStatus.MEASURED, + duration_s=0.01, + intervals=[], + tags=[], + ) + ], + ) + assert rerun_append.written + dossier_after_rerun = _dossier(api, meta_append.episode_id) + assert dossier_after_rerun["intervals"] == [] + assert dossier_after_rerun["tags"] == [] + timeline_after_rerun = api.get(f"/api/v1/episodes/{meta_append.episode_id}/timeline").json() + assert timeline_after_rerun["intervals"] == []