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
2 changes: 1 addition & 1 deletion docs/design/cot_vintage.md
Original file line number Diff line number Diff line change
Expand Up @@ -581,7 +581,7 @@ Four further findings, all fixed:
Scheduler while the store gained nothing, with `failed` being terminal. Now non-zero, in
a single combined message: an earlier draft raised on the failure first and silently
swallowed a restatement suspect in the same run, which is the more serious of the two.
**`ingest --retry-failed`** is the way back, and it is the only one: nothing else ever
**`ingest --retry`** is the way back, and it is the only one: nothing else ever
wrote `parse_status` back to `pending`, and two separate defects had stranded whole
backlogs there.
- **Disaggregated and TFF were fetched from 2006**, but cftc.gov serves 404 for
Expand Down
47 changes: 33 additions & 14 deletions src/cotdata/vintage_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
cotdata-schedule backfill
"""
import argparse
from collections import Counter

from . import config

Expand Down Expand Up @@ -88,18 +89,32 @@ def _cmd_ingest(args) -> int:
snaps = vintage.read_snapshots()
if args.snapshot:
snaps = [s for s in snaps if s.get("snapshot_id") == args.snapshot]
elif args.retry_failed:
# The ONLY route back from parse_status="failed", which is otherwise terminal:
# --pending does not select it and nothing else ever writes the field back. Two
# separate defects have reached that state en masse (a stale lock, and a path that
# would not resolve), and in both cases the entire backlog was stranded with no
# command able to recover it. Raw bytes are retained, so a retry is always safe.
for s_ in [x for x in snaps if x.get("parse_status") == "failed"]:
elif args.retry:
# The ONLY route back from the two terminal states, since --pending selects
# neither and nothing else ever writes the field back:
#
# failed the parse raised. Retry after fixing the cause. Two separate defects
# have driven whole backlogs here (a stale lock, and a path that would
# not resolve), each time with no command able to recover them.
# skipped no canonicaliser existed for that report type at ingest time. This is
# the one that matters going forward: the design notes promise that
# adding a canonicaliser later just means re-marking these pending and
# re-running, and until now there was no way to do the re-marking. The
# weekly static is the live case, because every retained copy of it is
# skipped and would be stranded the day it gets a canonicaliser.
#
# Raw bytes are retained in both states, so a retry is always safe. Re-skipping is
# harmless and quiet: a source that still has no canonicaliser simply drains again.
stuck = [x for x in snaps if x.get("parse_status") in ("failed", "skipped")]
by_state = Counter(x.get("parse_status") for x in stuck)
for s_ in stuck:
vintage.update_snapshot(s_["snapshot_id"], parse_status="pending",
parse_error=None)
snaps = vintage.read_snapshots()
snaps = [x for x in snaps if x.get("parse_status") == "pending"]
print(f"vintage ingest: reset {len(snaps)} failed snapshot(s) to pending.")
snaps = [x for x in vintage.read_snapshots()
if x.get("parse_status") == "pending"]
detail = ", ".join(f"{n} {state}" for state, n in sorted(by_state.items()))
print(f"vintage ingest: reset {len(stuck)} snapshot(s) to pending"
+ (f" ({detail})." if detail else "."))
elif not args.all_snapshots:
# Default to pending even with no flag. A bare `ingest` used to select EVERY
# snapshot ever recorded and re-ingest each one, which is not a no-op: replaying
Expand Down Expand Up @@ -184,7 +199,7 @@ def _cmd_ingest(args) -> int:
if n_failed:
parts.append(f"{n_failed} snapshot(s) FAILED to parse (raw bytes retained, so "
f"nothing is lost: fix the cause and re-run with "
f"'cotdata-vintage ingest --retry-failed', which is the only way back "
f"'cotdata-vintage ingest --retry', which is the only way back "
f"since 'failed' is not 'pending')")
raise SystemExit(
"cotdata-vintage: " + "; ".join(parts)
Expand Down Expand Up @@ -313,9 +328,13 @@ def main(argv=None) -> int:
g.add_argument("--snapshot", default=None, help="Ingest one snapshot id.")
g.add_argument("--pending", action="store_true",
help="Ingest all parse_status=pending. This is also the default.")
g.add_argument("--retry-failed", action="store_true", dest="retry_failed",
help="Reset parse_status=failed snapshots to pending and re-ingest "
"them. The only way back from a failed parse.")
# --retry-failed kept as an alias: it shipped earlier the same day under that name,
# and the widened behaviour is a superset, so nothing anyone already wrote breaks.
g.add_argument("--retry", "--retry-failed", action="store_true", dest="retry",
help="Reset snapshots stuck in parse_status=failed OR =skipped back to "
"pending and re-ingest them. The only way back from either: "
"--pending selects neither. Use after fixing a parse failure, or "
"after adding a canonicaliser for a report type that was skipped.")
g.add_argument("--all-snapshots", action="store_true", dest="all_snapshots",
help="Replay EVERY recorded snapshot, including already-ingested ones. "
"Only for rebuilding a store from retained raw bytes.")
Expand Down
43 changes: 43 additions & 0 deletions tests/test_vintage_ingest.py
Original file line number Diff line number Diff line change
Expand Up @@ -408,3 +408,46 @@ def test_a_run_where_every_snapshot_failed_does_not_report_success(store_env):
"retrieved_at": "2026-07-31T00:00:00Z"} for i in range(3)]})
with _pytest.raises(SystemExit, match="3 snapshot"):
vintage_cli.main(["ingest", "--pending"])


def test_retry_also_recovers_skipped_snapshots(store_env):
"""`skipped` means no canonicaliser existed for that report type at ingest time, and
it was as terminal as `failed`: --pending selects neither and nothing else wrote the
field back. The design notes promise that adding a canonicaliser later just means
re-marking these pending and re-running, and there was no way to do the re-marking.

The live case is the weekly static: every retained copy of it is `skipped`, so all of
them would have been stranded the day it gets a canonicaliser."""
import pytest as _pytest

from cotdata import vintage, vintage_cli
vintage._write_manifest({"schema_version": 1, "snapshots": [
{"snapshot_id": "w1", "report_type": "legacy", "source_kind": "weekly_static",
"local_path": "y", "parse_status": "pending",
"retrieved_at": "2026-07-31T00:00:00Z"},
{"snapshot_id": "d1", "report_type": "legacy", "source_kind": "annual_zip",
"local_path": "gone.zip", "parse_status": "pending",
"retrieved_at": "2026-07-31T00:00:00Z"},
]})
with _pytest.raises(SystemExit): # w1 skips, d1 fails
vintage_cli.main(["ingest", "--pending"])
after = {s["snapshot_id"]: s["parse_status"] for s in vintage.read_snapshots()}
assert after == {"w1": "skipped", "d1": "failed"}

assert vintage_cli.main(["ingest", "--pending"]) == 0 # both terminal

with _pytest.raises(SystemExit): # --retry reaches both
vintage_cli.main(["ingest", "--retry"])
again = {s["snapshot_id"]: s["parse_status"] for s in vintage.read_snapshots()}
assert again == {"w1": "skipped", "d1": "failed"} # re-skipping is quiet+safe


def test_retry_failed_still_works_as_an_alias(store_env):
"""It shipped earlier the same day under the narrower name; the widened behaviour is a
superset, so nothing already written breaks."""
from cotdata import vintage, vintage_cli
vintage._write_manifest({"schema_version": 1, "snapshots": [
{"snapshot_id": "w1", "report_type": "legacy", "source_kind": "weekly_static",
"local_path": "y", "parse_status": "skipped",
"retrieved_at": "2026-07-31T00:00:00Z"}]})
assert vintage_cli.main(["ingest", "--retry-failed"]) == 0
Loading