From 0fee807f54404f9109f2555ff1cafb2510e844c1 Mon Sep 17 00:00:00 2001 From: Matt Spinola Date: Sun, 26 Jul 2026 17:05:27 -0400 Subject: [PATCH] feat: stop writing the legacy manifest, add a one-shot migration The legacy manifest.json held both producer halves in ONE file. That was unsafe two ways: the update is a read-modify-write, so two producers lose each other's entries, and a file-level sync between two stores resolves the file last-writer-wins and silently discards one side. The second one stopped being theoretical when the Windows producer moved from a redirected drive to a local store, making the two stores separate copies merged by a sync. Nothing writes manifest.json now. `cotdata-update --migrate-manifests` splits an existing one into the per-half files. It is idempotent, entries already in a half file win so a re-run cannot resurrect stale bookkeeping, and it never touches data. Until a store is migrated, a domain absent from the per-half files is still read from the aggregate, with a warning naming the domains and the command to run. THE FALLBACK IS PER DOMAIN, NOT PER HALF. Found by dry-running the migration against a copy of the real store: manifests/prices.json held `prices` but not `metadata`, because the price producer had run on the new code while the metadata producer had not. A per-half rule treated the whole prices half as migrated and hid `metadata` entirely. Verified before and after migration now show the same 241 entries. reconcile_manifest is rewritten to prune each manifest FILE in place rather than the merged view. It previously wrote its result to the legacy aggregate, which would now be a write to a file nothing reads. It prunes the aggregate too when present, so an un-migrated store can still be cleaned. Suite 134 passed. Co-Authored-By: Claude Opus 5 --- README.md | 17 +++- src/cotdata/store.py | 160 ++++++++++++++++++++++++++---------- src/cotdata/update.py | 19 +++++ tests/test_manifest_seam.py | 113 +++++++++++++++++++++---- 4 files changed, 247 insertions(+), 62 deletions(-) diff --git a/README.md b/README.md index 973b14c..adbf707 100644 --- a/README.md +++ b/README.md @@ -185,9 +185,20 @@ become a second COT producer racing the first. `--check` and `--reconcile` are r and work from either. Each half also owns its own manifest (`manifests/cot.json`, `manifests/prices.json`). -The manifest update is a read-modify-write, so two producers sharing one file eventually -lose an entry. The legacy top-level `manifest.json` is still written for consumers pinned -to an older cotdata, but current readers prefer the per-half files. See ADR-0007. +The legacy top-level `manifest.json` held both halves in ONE file, which was unsafe two +ways: the update is a read-modify-write, so two producers lose each other's entries, and +a file-level sync between two stores resolves it last-writer-wins and silently discards +one side. The per-half files are disjoint, so both problems go away. + +**Nothing writes `manifest.json` any more.** Migrate a store once: + +```bash +cotdata-update --migrate-manifests +``` + +Idempotent, and it never touches data. Until you run it, a domain missing from the +per-half files is still read from the aggregate with a warning. Delete `manifest.json` +once every consumer of that store is on this version. See ADR-0007. ### Scheduling on Windows (Task Scheduler) diff --git a/src/cotdata/store.py b/src/cotdata/store.py index d8994f8..42d6932 100644 --- a/src/cotdata/store.py +++ b/src/cotdata/store.py @@ -2,6 +2,7 @@ contract between producers (write) and consumers (read).""" import datetime as dt import json +import logging import os import tempfile from pathlib import Path @@ -10,6 +11,8 @@ from . import config +logger = logging.getLogger(__name__) + def _atomic_write_parquet(df: pd.DataFrame, path: Path) -> None: """Write to a temp file in the same dir, then os.replace — so a consumer @@ -141,33 +144,59 @@ def _read_json(p) -> dict: return {} +def _overlay(dest: dict, src: dict, only_domains=None) -> None: + for key, value in src.items(): + if key == "schema_version": + dest["schema_version"] = max(int(dest.get("schema_version", 0) or 0), + int(value or 0)) + continue + if only_domains is not None and key not in only_domains: + continue + if isinstance(value, dict): + dest.setdefault(key, {}).update(value) + else: + dest[key] = value + + def load_manifest() -> dict: - """The merged manifest. + """The manifest, assembled from the per-half files. + + Nothing writes the legacy aggregate any more. It is still read, but ONLY for a + half whose own file does not exist yet, so an un-migrated store keeps working + while a migrated one never touches it. The fallback is per-half rather than + all-or-nothing because a store can be half migrated: the price producer may have + run on the new code while the COT producer has not, leaving prices in + manifests/prices.json and everything else still in the aggregate. - Reads the legacy aggregate first, then overlays each per-half manifest, so fresher - per-half data wins and a half whose producer has not run yet still resolves from the - legacy file. Consumers see the same shape as before and need no change. + Run ``cotdata-update --migrate-manifests`` once to convert a store and silence + the warning. The fallback goes away in a later release. """ - merged = _read_json(config.manifest_path()) - found_half = False + merged: dict = {} for half in HALVES: - part = _read_json(config.manifest_path_for(half)) - if not part: - continue - found_half = True - for key, value in part.items(): - if key == "schema_version": - continue - if isinstance(value, dict): - merged.setdefault(key, {}).update(value) - else: - merged[key] = value + _overlay(merged, _read_json(config.manifest_path_for(half))) + + legacy = _read_json(config.manifest_path()) + if legacy: + # Fall back per DOMAIN, not per half. A half file existing does not mean the + # half is complete: a real store was found with manifests/prices.json holding + # `prices` but not `metadata`, because the price producer had run on the new + # code while the metadata producer had not. A per-half rule would have hidden + # `metadata` entirely until the migration ran. + missing = {k for k, v in legacy.items() + if k != "schema_version" and isinstance(v, dict) and k not in merged} + if missing: + _overlay(merged, legacy, only_domains=missing) + logger.warning( + "reading the legacy manifest.json for %s. Nothing writes that file any " + "more, and a file-level sync between two stores resolves it " + "last-writer-wins. Run 'cotdata-update --migrate-manifests' once on " + "this store.", ", ".join(sorted(missing))) merged["schema_version"] = max(int(merged.get("schema_version", 0) or 0), - int(part.get("schema_version", 0) or 0)) + int(legacy.get("schema_version", 0) or 0)) + if not merged: return _empty_manifest() - if not found_half and "schema_version" not in merged: - merged["schema_version"] = config.SCHEMA_VERSION + merged.setdefault("schema_version", config.SCHEMA_VERSION) return merged @@ -191,12 +220,47 @@ def require_schema(min_version: int) -> None: ) +def migrate_manifests() -> dict: + """Split a legacy ``manifest.json`` into per-half files. Idempotent. + + Run once per store when upgrading. Entries already present in a half file win, so + re-running cannot resurrect stale bookkeeping and an already-migrated store is + left alone. Returns ``{half: n_entries_added}``. + """ + legacy = _read_json(config.manifest_path()) + added = {h: 0 for h in HALVES} + if not legacy: + return added + for half in HALVES: + part = _read_json(config.manifest_path_for(half)) + for kind, entries in legacy.items(): + if kind == "schema_version" or not isinstance(entries, dict): + continue + try: + if half_for(kind) != half: + continue + except ValueError: + continue # a retired domain the current code no longer declares + dest = part.setdefault(kind, {}) + for name, entry in entries.items(): + if name not in dest: + dest[name] = entry + added[half] += 1 + if part: + part["schema_version"] = max(int(part.get("schema_version", 0) or 0), + int(legacy.get("schema_version", 0) or 0)) + _atomic_write_json(part, config.manifest_path_for(half)) + return added + + def _touch_manifest(kind: str, name: str, df: pd.DataFrame, source: str) -> None: - """Record one entry, into this domain's half AND the legacy aggregate. + """Record one entry into its half's manifest. - Dual-write on purpose. The per-half file is the one current readers prefer, so it is - safe from the moment this ships. The legacy aggregate keeps consumers pinned to an - older cotdata working, and is dropped once they have moved. + Only the half file is written. The legacy aggregate held both halves in ONE file, + which made it unsafe in two ways: two producers doing a read-modify-write on it + lose each other's entries, and a file-level sync between two stores resolves it + last-writer-wins and silently discards one side. The per-half files are disjoint, + so both problems go away by construction. """ half = half_for(kind) last = None @@ -215,11 +279,6 @@ def _touch_manifest(kind: str, name: str, df: pd.DataFrame, source: str) -> None part["schema_version"] = config.SCHEMA_VERSION _atomic_write_json(part, config.manifest_path_for(half)) - legacy = _read_json(config.manifest_path()) - legacy.setdefault(kind, {})[name] = entry - legacy["schema_version"] = config.SCHEMA_VERSION - _atomic_write_json(legacy, config.manifest_path()) - def _atomic_write_json(m: dict, path) -> None: tmp = path.with_suffix(".json.tmp") @@ -254,20 +313,35 @@ def reconcile_manifest() -> dict: retired ``cot`` domain, …) — and drop domains left empty. Returns ``{domain: [pruned names]}``. - Provably safe: only removes bookkeeping for files that do not exist on disk; + Provably safe: only removes bookkeeping for files that do not exist on disk, never deletes or renames data. + + Operates on each manifest FILE in place rather than on the merged view, so a + prune is written back to the file the entry actually lives in. The legacy + aggregate is pruned too when it is still present, so an un-migrated store can be + cleaned without migrating first. """ - m = load_manifest() pruned: dict = {} - for domain in [k for k, v in m.items() if isinstance(v, dict)]: - d = _domain_dir(domain) - gone = [name for name in m[domain] if not (d / f"{name}.parquet").exists()] - if gone: - for name in gone: - del m[domain][name] - pruned[domain] = sorted(gone) - if not m[domain]: - del m[domain] - if pruned: - _write_manifest(m) - return pruned + + def _prune(doc: dict) -> bool: + changed = False + for domain in [k for k, v in doc.items() if isinstance(v, dict)]: + d = _domain_dir(domain) + gone = [n for n in doc[domain] if not (d / f"{n}.parquet").exists()] + if gone: + for n in gone: + del doc[domain][n] + pruned.setdefault(domain, []).extend(gone) + changed = True + if not doc[domain]: + del doc[domain] + changed = True + return changed + + targets = [config.manifest_path_for(h) for h in HALVES] + [config.manifest_path()] + for path in targets: + doc = _read_json(path) + if doc and _prune(doc): + _atomic_write_json(doc, path) + + return {k: sorted(set(v)) for k, v in pruned.items()} diff --git a/src/cotdata/update.py b/src/cotdata/update.py index 8dd7902..b3a6ff1 100644 --- a/src/cotdata/update.py +++ b/src/cotdata/update.py @@ -65,6 +65,10 @@ def main(argv=None, half=None) -> None: p.add_argument("--check", action="store_true", help="Print store status (row counts, newest data, staleness) from " "the manifest and exit. Read-only, cross-platform, no network.") + p.add_argument("--migrate-manifests", action="store_true", + help="One-shot: split a legacy manifest.json into the per-half " + "manifests/ files, then exit. Idempotent, read-only on data. " + "Run once per store after upgrading.") p.add_argument("--reconcile", action="store_true", help="Prune manifest entries whose parquet file is missing (ghosts " "from old naming), refresh status.json, and exit. Never touches data.") @@ -79,6 +83,21 @@ def main(argv=None, half=None) -> None: status.print_check() return + if args.migrate_manifests: + from . import store + added = store.migrate_manifests() + total = sum(added.values()) + if not total: + print("manifest migration: nothing to do (already migrated, or no legacy " + "manifest.json).") + else: + for half, n in sorted(added.items()): + print(f" {half:<8} +{n} entries -> manifests/{half}.json") + print(f"manifest migration: moved {total} entries out of the legacy " + f"aggregate. manifest.json is no longer written and can be deleted " + f"once every consumer is on this version.") + return + if args.reconcile: from . import status, store pruned = store.reconcile_manifest() diff --git a/tests/test_manifest_seam.py b/tests/test_manifest_seam.py index 4a42625..466e373 100644 --- a/tests/test_manifest_seam.py +++ b/tests/test_manifest_seam.py @@ -5,8 +5,8 @@ (any OS) and the price producer (Windows for Norgate). Splitting the manifest by half means they never touch the same file. -Dual-write during the transition: the per-half file is what current readers prefer, the -legacy aggregate keeps consumers pinned to an older cotdata working. +Nothing writes the legacy aggregate any more. It is read only as a per-half fallback +for a store that has not run `--migrate-manifests` yet. """ import json @@ -54,17 +54,7 @@ def test_price_and_cot_domains_land_on_opposite_sides(): assert store.half_for("cot_tff") == "cot" -# ── dual write ──────────────────────────────────────────────────────────── -def test_write_lands_in_both_the_half_and_the_legacy_aggregate(store_env): - from cotdata import config, store - store.write_prices("ES", "backadj", _prices(), source="test") - - half = json.loads(config.manifest_path_for("prices").read_text()) - assert half["prices"]["ES_backadj"]["n_rows"] == 3 - legacy = json.loads(config.manifest_path().read_text()) - assert legacy["prices"]["ES_backadj"]["n_rows"] == 3 - - +# ── writes ──────────────────────────────────────────────────────────────── def test_the_two_halves_do_not_share_a_file(store_env): from cotdata import config, store store.write_prices("ES", "backadj", _prices(), source="test") @@ -76,9 +66,10 @@ def test_the_two_halves_do_not_share_a_file(store_env): assert "cot_legacy" in cot_half and "prices" not in cot_half -def test_a_lost_legacy_write_does_not_lose_the_entry(store_env): - """The hazard, simulated: the shared aggregate is clobbered by a concurrent - producer. Readers must still see both entries, because they prefer the halves.""" +def test_a_clobbered_legacy_file_cannot_lose_an_entry(store_env): + """The hazard, simulated. A concurrent producer (or a file-level sync between two + stores) overwrites the shared aggregate wholesale. Both entries must survive, + because a migrated store never consults that file.""" from cotdata import config, store store.write_prices("ES", "backadj", _prices(), source="test") store.write_cot_legacy("ES_13874A", _cot(), source="test") @@ -198,3 +189,93 @@ def test_every_action_flag_is_assigned_to_a_half(): actions = {"prices", "metadata", "prices_yahoo", "ingest_databento", "build_databento", "cot_legacy", "cot_disagg", "cot_tff", "cot_all"} assert actions == assigned + + +# ── dropping the legacy aggregate ───────────────────────────────────────── +def test_writes_no_longer_touch_the_legacy_aggregate(store_env): + """A single file holding both halves is unsafe two ways: concurrent producers + lose each other's entries, and a file-level sync between two stores resolves it + last-writer-wins. Nothing writes it any more.""" + from cotdata import config, store + store.write_prices("ES", "backadj", _prices(), source="test") + assert config.manifest_path_for("prices").exists() + assert not config.manifest_path().exists() + + +def test_migrate_splits_a_legacy_manifest(store_env): + from cotdata import config, store + config.manifest_path().parent.mkdir(parents=True, exist_ok=True) + config.manifest_path().write_text(json.dumps({ + "schema_version": 2, + "prices": {"ES_backadj": {"n_rows": 3}}, + "metadata": {"contract_specs": {"n_rows": 47}}, + "cot_legacy": {"GC_088691": {"n_rows": 5}}, + "cot_tff": {"ES_13874A": {"n_rows": 9}}, + })) + added = store.migrate_manifests() + assert added == {"prices": 2, "cot": 2} + + prices = json.loads(config.manifest_path_for("prices").read_text()) + cot = json.loads(config.manifest_path_for("cot").read_text()) + assert set(prices) == {"prices", "metadata", "schema_version"} + assert set(cot) == {"cot_legacy", "cot_tff", "schema_version"} + + +def test_migrate_is_idempotent_and_never_resurrects_stale_entries(store_env): + from cotdata import config, store + store.write_prices("ES", "backadj", _prices(), source="test") # 3 rows, current + config.manifest_path().write_text(json.dumps( + {"schema_version": 2, "prices": {"ES_backadj": {"n_rows": 999}}})) # stale + + assert store.migrate_manifests() == {"prices": 0, "cot": 0} + assert store.load_manifest()["prices"]["ES_backadj"]["n_rows"] == 3 + assert store.migrate_manifests() == {"prices": 0, "cot": 0} # re-run is a no-op + + +def test_an_incomplete_half_file_still_falls_back_per_domain(store_env): + """Found on a real store. manifests/prices.json held `prices` but not `metadata`, + because the price producer had run on the new code while the metadata producer + had not. A per-HALF fallback rule hid `metadata` entirely until the migration + ran, so the fallback is per DOMAIN.""" + from cotdata import config, store + config.manifest_path().parent.mkdir(parents=True, exist_ok=True) + config.manifest_path().write_text(json.dumps({ + "schema_version": 2, + "prices": {"OLD_backadj": {"n_rows": 1}}, + "metadata": {"contract_specs": {"n_rows": 47}}, + })) + store.write_prices("ES", "backadj", _prices(), source="test") # prices half only + + m = store.load_manifest() + assert "ES_backadj" in m["prices"] # from the half file + assert m["metadata"]["contract_specs"]["n_rows"] == 47 # still from legacy + + +def test_legacy_is_read_only_for_a_domain_that_has_not_migrated(store_env): + """A store can be part migrated: prices written by the new code, COT still only + in the aggregate.""" + from cotdata import config, store + config.manifest_path().parent.mkdir(parents=True, exist_ok=True) + config.manifest_path().write_text(json.dumps({ + "schema_version": 2, + "prices": {"OLD_backadj": {"n_rows": 1}}, + "cot_legacy": {"GC_088691": {"n_rows": 5}}, + })) + store.write_prices("ES", "backadj", _prices(), source="test") + + m = store.load_manifest() + assert m["cot_legacy"]["GC_088691"]["n_rows"] == 5 # from legacy, cot unmigrated + assert "ES_backadj" in m["prices"] # from the migrated half + assert "OLD_backadj" not in m["prices"] # legacy prices NOT consulted + + +def test_fully_migrated_store_ignores_the_legacy_file(store_env): + from cotdata import config, store + store.write_prices("ES", "backadj", _prices(), source="test") + store.write_cot_legacy("GC_088691", _cot(), source="test") + config.manifest_path().write_text(json.dumps( + {"schema_version": 2, "prices": {"GHOST": {"n_rows": 1}}, + "cot_legacy": {"GHOST": {"n_rows": 1}}})) + + m = store.load_manifest() + assert "GHOST" not in m["prices"] and "GHOST" not in m["cot_legacy"]