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
24 changes: 11 additions & 13 deletions packages/hflow-server/src/hflow_server/_catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -607,24 +607,22 @@ 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,
connection.execute(
"""
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
""",
Expand Down Expand Up @@ -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],
),
Expand Down
100 changes: 100 additions & 0 deletions packages/hflow-server/tests/test_server_episode_dossier.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"] == []
Loading