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
7 changes: 5 additions & 2 deletions docs/CATALOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -301,8 +301,11 @@ these views registered (it is what `curate()` uses; take it and explore):
column per measurement key** (latest value; booleans as 0/1).
- `episodes_raw`, `check_runs`, `measurements`, `observations`, `tags`, `intervals`,
`ingest_failures`: the long tables, exactly as stored.
- `episodes_latest`, `measurements_latest`: one row per episode / per
(episode, key), most recent append wins.
- `episodes_latest`: one row per episode, most recent append wins.
- `measurements_latest`: one row per `(episode_id, key)`, taken only from the
latest run of each `(episode_id, check_name)`. A key that the latest run of
its check did not record (a newer version omitted it, or the run errored) is
withdrawn: the wide `episodes` view keeps the column and reads `NULL`.
- `observations_latest`: all observation fields from the latest run of each
`(episode_id, check_name)`; it switches the repeated result as a unit, so
fields omitted by a newer check version do not leak in from an older one.
Expand Down
59 changes: 36 additions & 23 deletions src/hflow/curation.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,12 @@
own source. This view is the only place that expression lives -- every
consumer, ``stale_episodes`` included, reads it rather than re-deriving
the window.
- ``measurements_latest`` -- one row per (episode_id, key), most recent by
the OWNING episode's recorded_at (joined in), not the measurement row's
own -- the latter can go stale independently of the episode it belongs to.
- ``measurements_latest`` -- one row per (episode_id, key), drawn only from
the latest run of each (episode, check) in ``check_runs_latest``, so a key a
newer check version omits is withdrawn rather than left at its old value.
Ranked by the OWNING episode's recorded_at (joined in), not the measurement
row's own -- the latter can go stale independently of the episode it
belongs to.
- ``observations_latest`` -- every timestamped observation field from the
latest run of each (episode, check), switched as one coherent result.
- ``episodes`` -- the wide view for everyday queries: latest episode rows,
Expand Down Expand Up @@ -364,16 +367,16 @@ def _register_catalog_relations(
)
connection.execute(
"""
CREATE VIEW measurements_latest AS
CREATE VIEW check_runs_latest AS
SELECT * EXCLUDE (row_rank) FROM (
SELECT
m.* EXCLUDE (recorded_at),
c.* EXCLUDE (recorded_at),
-- Rank (and report) by the EPISODE's recorded_at, not this
-- table's own: episodes/<file_stem>.parquet is the only file
-- create-if-absent guarantees a single writer for, so it is
-- the sole trustworthy recorded_at once episodes_latest has
-- already picked a winner. A repair pass that wins the
-- episodes race can still crash before reaching measurements
-- episodes race can still crash before reaching dependents
-- (see #51); replayed appends now reconcile stale dependents
-- (catalog._reconcile_replayed_append), but until a replay
-- happens -- and on mirrors synced before it -- this table's
Expand All @@ -387,32 +390,33 @@ def _register_catalog_relations(
-- proves an append complete" idiom.
e.recorded_at,
row_number() OVER (
PARTITION BY m.episode_id, m.key
ORDER BY e.recorded_at DESC, m.run_fingerprint DESC
PARTITION BY c.episode_id, c.check_name
ORDER BY e.recorded_at DESC, c.run_fingerprint DESC
) AS row_rank
FROM measurements m
FROM check_runs c
JOIN episodes_raw e USING (episode_id, run_fingerprint)
) WHERE row_rank = 1
"""
)
connection.execute(
"""
CREATE VIEW check_runs_latest AS
CREATE VIEW measurements_latest AS
SELECT * EXCLUDE (row_rank) FROM (
SELECT
c.* EXCLUDE (recorded_at),
-- Ranked by the EPISODE's recorded_at for the same reason
-- measurements_latest is: episodes is the only table whose
-- single writer is guaranteed, so ranking off this table's own
-- recorded_at could pick a different run_fingerprint than
-- episodes_latest did and stitch two runs together.
e.recorded_at,
m.* EXCLUDE (recorded_at),
latest_check_run.recorded_at,
row_number() OVER (
PARTITION BY c.episode_id, c.check_name
ORDER BY e.recorded_at DESC, c.run_fingerprint DESC
PARTITION BY m.episode_id, m.key
ORDER BY latest_check_run.recorded_at DESC, m.run_fingerprint DESC
) AS row_rank
FROM check_runs c
JOIN episodes_raw e USING (episode_id, run_fingerprint)
FROM measurements m
-- Only the latest run of each check speaks for it, as in
-- observations_latest: ranking every row per key alone would let
-- a key a newer version stopped emitting keep its old value.
-- The per-key rank then still yields one row per key when a key
-- moved between checks.
JOIN check_runs_latest latest_check_run
USING (episode_id, run_fingerprint, check_name, check_version)
) WHERE row_rank = 1
"""
)
Expand All @@ -425,10 +429,19 @@ def _register_catalog_relations(
USING (episode_id, run_fingerprint, check_name, check_version)
"""
)
# Every key any completed run recorded, not only the current ones: a key
# withdrawn by a newer check version stays a (NULL) column so queries
# naming it still bind, and the case-collision guard keeps covering every
# append, as docs/CATALOG.md promises.
measurement_keys = [
str(row[0])
for row in connection.execute(
"SELECT DISTINCT key FROM measurements_latest ORDER BY key"
"""
SELECT DISTINCT m.key
FROM measurements m
JOIN episodes_raw e USING (episode_id, run_fingerprint)
ORDER BY m.key
"""
).fetchall()
]
_raise_if_measurement_keys_case_collide(measurement_keys)
Expand Down Expand Up @@ -544,8 +557,8 @@ def _refresh_local_catalog_connection(
derived_view_names = (
"episodes",
"observations_latest",
"check_runs_latest",
"measurements_latest",
"check_runs_latest",
"episodes_latest",
)
try:
Expand Down
86 changes: 86 additions & 0 deletions tests/test_catalog_curation.py
Original file line number Diff line number Diff line change
Expand Up @@ -547,6 +547,92 @@ def test_rerunning_a_changed_check_appends_new_version_rows(tmp_path: Path) -> N
assert wide_row == (2.0,)


def _measured_row(
check_name: str, version: str, measurements: dict[str, hflow.MeasurementValue]
) -> CheckRunRow:
return CheckRunRow(
check_name=check_name,
check_version=version,
critical=False,
status=hflow.CheckStatus.MEASURED,
duration_s=0.01,
measurements=measurements,
)


def _append_runs(tmp_path: Path, *runs: list[CheckRunRow]) -> Path:
catalog = Catalog(tmp_path / "catalog")
canonical = write_fake_canonical(tmp_path)
for check_rows in runs:
assert catalog.append_episode(
canonical_path=canonical,
stamps=FAKE_STAMPS,
episode_metadata={},
check_rows=check_rows,
).written
return catalog.root


def test_a_key_a_newer_check_version_omits_is_withdrawn(tmp_path: Path) -> None:
catalog_root = _append_runs(
tmp_path,
[
_measured_row("blur", "v1", {"blur_score": 0.9, "frames": 100.0}),
_measured_row("steady", "v1", {"steady/score": 0.5}),
],
[_measured_row("blur", "v2", {"frames": 100.0})],
)

with open_catalog_connection(catalog_root) as connection:
assert connection.execute(
"SELECT check_name, check_version, key, value_double "
"FROM measurements_latest ORDER BY key"
).fetchall() == [
("blur", "v2", "frames", 100.0),
("steady", "v1", "steady/score", 0.5),
]
assert connection.execute(
'SELECT blur_score, frames, "steady/score" FROM episodes'
).fetchall() == [(None, 100.0, 0.5)]
report = curate(catalog_root, "SELECT episode_id FROM episodes WHERE blur_score > 0.8")
assert report.row_count == 0


def test_an_errored_latest_run_withdraws_that_checks_measurements(tmp_path: Path) -> None:
catalog_root = _append_runs(
tmp_path,
[_measured_row("remote_check", "v1", {"score": 1.0})],
[
CheckRunRow(
check_name="remote_check",
check_version="v1",
critical=False,
status=hflow.CheckStatus.ERROR,
duration_s=0.1,
error="temporary timeout",
)
],
)

with open_catalog_connection(catalog_root) as connection:
assert connection.execute("SELECT key FROM measurements_latest").fetchall() == []
assert connection.execute("SELECT score FROM episodes").fetchall() == [(None,)]


def test_a_key_moved_to_another_check_keeps_one_latest_row(tmp_path: Path) -> None:
catalog_root = _append_runs(
tmp_path,
[_measured_row("old_check", "v1", {"score": 1.0})],
[_measured_row("new_check", "v1", {"score": 2.0})],
)

with open_catalog_connection(catalog_root) as connection:
assert connection.execute(
"SELECT check_name, value_double FROM measurements_latest"
).fetchall() == [("new_check", 2.0)]
assert connection.execute("SELECT score FROM episodes").fetchall() == [(2.0,)]


def test_successful_retry_after_error_appends_repaired_outcome(tmp_path: Path) -> None:
catalog = Catalog(tmp_path / "catalog")
canonical = write_fake_canonical(tmp_path)
Expand Down
Loading