diff --git a/docs/CICs/UNFAOPostProcessorManager.md b/docs/CICs/UNFAOPostProcessorManager.md index cb6da57..cd589f7 100644 --- a/docs/CICs/UNFAOPostProcessorManager.md +++ b/docs/CICs/UNFAOPostProcessorManager.md @@ -3,7 +3,7 @@ **Status:** Active **Owner:** PRIO MD&D Team -**Last reviewed:** 2026-08-18 +**Last reviewed:** 2026-08-19 **Related ADRs:** ADR-001, ADR-002, ADR-008, ADR-009 --- @@ -89,7 +89,7 @@ Assumptions that are not met **must cause failure**, not fallback behavior. The - **Missing required metadata columns after enrichment:** Raises `ValueError` listing missing columns - **Null values in required metadata columns:** Raises `ValueError` with null count and affected column name (C-01 resolved — validation active) - **Dataset initialization failure:** Raises `ValueError` in `_save()` if datasets are None -- **Appwrite upload failure:** Propagates exception from `DatastoreModule` +- **Appwrite upload failure:** Propagates exception from `DatastoreModule`. **Inside the wire leg** it is wrapped as `contract.wire.sink.TornRunError` (C-105, 2026-08-19), naming the run, how many of how many objects landed, and their names and file ids — the consumer cannot see a torn run (the manifest is the commit marker and never landed), the listed objects are **not** removed, and a re-run uploads all of them again under the same names. **The historical leg is not covered by that wrapper**: it uploads after the wire run is committed, so a failure there leaves a visible, complete forecast run alongside the *previous* run's historical artifact. Recorded as remaining scope in C-105 - **Delivery invisible to the consumer:** raises `delivery.findability.DeliveryNotFindableError` (C-94, added 2026-08-18). After both legs are uploaded, the manager queries the partner store as the consumer does — `name == product.CONSUMER_DOCUMENT_NAME`, per category — and refuses a falsy answer. This is the one failure mode where every upload reports success and the consumer still sees nothing; run-0's historical leg stranded exactly that way (C-79). It runs only inside the §11.4 interlock, and queries through a store with pipeline-core's automatic `name == model_name` filter suppressed, so it verifies the declared name rather than the views-models directory name that happens to match (C-77). **It does not detect a delivery that never ran, or stale data served from the consumer's cache** — both recorded as gaps in C-94 - **Wrong forecast selected:** structurally impossible since #149. Selection is by **run manifest** — a commit marker whose contents are hash-verified — not by scanning the bucket for the newest `category="forecast"` upload. Declared identity is additionally checked **per shard header** against the launched ensemble inside `TargetLease.load()` (`contract/wire/source_selection.py:73-81`), so identity comes from the artifact's own content. The metadata-field check this bullet used to describe (`delivery/identity.py`) was retired in #150 and the legacy reader it served in #149; register C-25 is closed as *superseded by mechanism* - **Launch config incomplete:** raises `LaunchConfigError` naming the missing key. A launcher that omits `wire_contract` or declares a `data_format` other than `feature_frame` is **refused**, never quietly routed into a fallback (ADR-003, register C-63) diff --git a/reports/technical_risk_register.md b/reports/technical_risk_register.md index 6aba7a5..66ba62c 100644 --- a/reports/technical_risk_register.md +++ b/reports/technical_risk_register.md @@ -1146,7 +1146,17 @@ Cross-refs: **C-72** (owns the pyarrow half — that half is not re-registered h `docs/operations/correction_procedure.md` covers the *wrong value* case — the contract has no retraction primitive, so a correction is a new complete run, manifest last. A torn attempt is a different case and is not covered by it. -Cross-refs: **C-94** (nothing observes the outcome of an upload at the time it happens), **C-79** (the single-file orphan this generalises), **D-12** (the unnamed retention owner this compounds with). Part of causal cluster: **Cluster J — Delivery aftercare has no mechanism**. +**Partial mitigation, 2026-08-19 — the tear is now documented, not removed.** `deliver_run` keeps an in-memory ledger of what it has uploaded, and a failure anywhere in the upload phase raises `TornRunError` naming the run, how many of how many objects were *confirmed* uploaded, and what is true about the consumer. Tested in `tests/test_torn_run.py` (8), mutation-proven three ways. + +**Three corrections `/code-review high` made to the first draft, each of which would have sent an operator the wrong way.** (a) The object that FAILED was omitted, and it is the likeliest orphan of the whole run: `_ContractStorePort.upload` raises precisely when the store returns failure *with the file already uploaded* (the C-79 shape), so the refusal now names it as a separate thing to go and look for. (b) File ids were truncated to five in the message and recorded nowhere else, so at run-0 scale ~104 ids existed only in a string nobody kept — the log ledger now carries `file_id` per upload and is the persistent record. (c) A failure on the *first* upload printed an empty list and a dangling period while telling the operator to audit a bucket. The message also no longer asserts categorically that the consumer cannot see the run: a manifest upload can fail after the store committed the document, and steering a re-run on a false certainty duplicates every object. + +The refusal says three things an operator otherwise has to establish by hand: the consumer **cannot see this run** (the manifest is the commit marker and never landed, so nothing partial is being served — §4.2 working as designed); the objects listed are **still there and were NOT removed**; and a re-run will upload all of them again under the same names, with supersede-or-duplicate being a store semantic this repository does not assert. + +**Remaining scope, found by `/review-diff` on the fix itself: the historical leg is not covered.** `TornRunError` wraps the upload phase inside `deliver_run`. The historical artifact uploads *after* the wire run is committed, from the manager, so a failure there raises unwrapped — and its consequence is different rather than smaller: the manifest already landed, so the consumer sees a **complete, visible forecast run** sitting next to the *previous* run's historical artifact. Not corrupt (the historical is a full snapshot, so the older one is valid, just one run stale) and the delivery does report failure — but it is the one tear where "the consumer cannot see this run" is false, and the wrapper's message would be wrong if it fired there. It does not fire there. Left uncovered deliberately rather than widening this change; the manager is at 435/450 of its line budget and the fix belongs with whoever takes the deletion decision below. + +**What is deliberately NOT done: deletion.** Removing objects from a partner bucket is irreversible and an operator decision rather than a delivery-path one, and the neighbouring delete surface is its own open question (**C-58**, views-pipeline-core #333, blocked on a test key). So this entry stays open: the mess is now legible, and it is still a mess. Closing it needs a decision about who cleans up and whether the store supersedes — neither of which is engineering work here. + +Cross-refs: **C-94** (nothing observes the outcome of an upload at the time it happens), **C-79** (the single-file orphan this generalises), **C-58** (the delete surface deletion would have to go through), **D-12** (the unnamed retention owner this compounds with). Part of causal cluster: **Cluster J — Delivery aftercare has no mechanism**. --- diff --git a/tests/test_torn_run.py b/tests/test_torn_run.py new file mode 100644 index 0000000..2188404 --- /dev/null +++ b/tests/test_torn_run.py @@ -0,0 +1,198 @@ +"""A run that dies mid-upload must say what it left behind (register C-105). + +The contract already handles the **consumer's** side of a torn run correctly and by +design: the manifest is uploaded last, so an attempt that dies before it has no commit +marker and is invisible rather than half-visible (ADR-013 §4.2). Nothing partial is +served. + +What was missing is our side. The objects that *did* land stay in the partner store, +and until 2026-08-19 nothing recorded that they had — an operator was left to diff the +bucket by hand. At run-0 scale a retry adds ~110 more under the same names. + +**Nothing is deleted, deliberately.** Removing objects from a partner bucket is +irreversible and an operator decision; the neighbouring delete surface is its own open +question (C-58, views-pipeline-core #333, blocked on a test key). These tests pin the +part this repository can honestly own: turning an invisible mess into a documented one. +""" + +from __future__ import annotations + +from pathlib import Path + +import numpy as np +import pytest + +import pyarrow as pa + +from views_postprocessing.contract.frames import build_prediction_frame +from views_postprocessing.contract.gaul_schema import CODE_COLS, COORD_COLS, METADATA_COLS +from views_postprocessing.contract.wire import sink +from views_postprocessing.unfao import product + +_GIDS = [100001, 100002, 100003, 100004, 100005, 100006] + + +class FakeLease: + """Preloaded values — this file tests the upload phase, not the inbound chain.""" + + def __init__(self, run_id, frame, headers): + self.run_id = run_id + self._value = (frame, headers) + + def load(self): + return self._value + + +def _synthetic_lookup() -> pa.Table: + """A lookup covering exactly `_GIDS`, built from the declared schema. + + Derived from `gaul_schema` rather than copied as a literal table: the columns are + the contract's, and a table hand-written here would drift from it silently. Built + locally rather than imported from another test module — reaching into a sibling + test's private helper couples two files that should be able to change apart. + """ + columns = {"priogrid_gid": pa.array(_GIDS, pa.int64())} + for col in METADATA_COLS: + if col in COORD_COLS: + columns[col] = pa.array([10.25 + i for i in range(len(_GIDS))], pa.float64()) + elif col in CODE_COLS: + columns[col] = pa.array(list(range(1, len(_GIDS) + 1)), pa.int64()) + else: + columns[col] = pa.array([f"{col}-{i}" for i in range(len(_GIDS))]) + return pa.table(columns) +_PRODUCT = {"consumer_name": product.CONSUMER_DOCUMENT_NAME, "s_min": product.S_MIN} + + +class FailAfter: + """Uploads `n` objects, then refuses — the C-79 shape mid-run.""" + + def __init__(self, n: int): + self.n = n + self.calls: list[str] = [] + + def upload(self, file_path, **kwargs): + if len(self.calls) >= self.n: + raise RuntimeError("store said no") + self.calls.append(Path(file_path).name) + return f"id-{len(self.calls)}" + + +def _one_target_leases(): + values = np.tile(np.array([[1.0, 2.0, 3.0, 4.0]], dtype=np.float32), (6, 1)) + time = np.full(6, 543, dtype=np.int64) + unit = np.array(_GIDS, dtype=np.int64) + frame = build_prediction_frame(values, time, unit) + headers = [{ + "run_id": "fixture_run_0", "target": "lr_ged_sb", "time_id": 543, + "sample_count": 4, "provenance": {"ensemble": "fixture_ensemble"}, + }] + return {"lr_ged_sb": FakeLease("fixture_run_0", frame, headers)} + + +def _deliver(store, tmp_path): + return sink.deliver_run( + _one_target_leases(), lookup=_synthetic_lookup(), staging_dir=tmp_path, + **_PRODUCT, store=store, upload_enabled=True, + ) + + +def test_a_torn_run_refuses_and_names_what_already_landed(tmp_path): + store = FailAfter(1) # the shard lands; the sidecar does not + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(store, tmp_path) + message = str(excinfo.value) + assert "TORN" in message + assert "fixture_run_0" in message, "the refusal must name the run" + assert store.calls[0] in message, ( + "the refusal must name the objects already in the partner store; without them " + "an operator has to diff the bucket by hand (C-105)" + ) + assert "id-1" in message, "and their file ids, so they can be found again" + + +def test_it_says_how_many_of_how_many(tmp_path): + store = FailAfter(1) + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(store, tmp_path) + # one shard + sidecar + manifest = 3 objects for this single-target run + assert "1 of 3" in str(excinfo.value) + + +def test_it_says_the_consumer_is_unaffected_and_nothing_was_removed(tmp_path): + """Both halves matter: no partial serve, and no silent cleanup either.""" + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(FailAfter(1), tmp_path) + message = str(excinfo.value) + assert "almost certainly" in message and "cannot see this run" in message, ( + "an operator's first question is whether the partner is being served garbage. " + "The §4.2 commit marker means almost certainly not — but the code cannot KNOW " + "it, because a manifest upload can fail after the store committed the document. " + "Asserting it categorically would steer a re-run that duplicates every object." + ) + assert "NOT removed" in message, ( + "the refusal must be explicit that it deleted nothing — a reader who assumes " + "cleanup happened will not go looking" + ) + assert "correction_procedure.md" in message and "C-105" in message + + +def test_a_failure_on_the_manifest_is_still_torn(tmp_path): + """The last upload is the commit marker; losing it is the canonical torn run.""" + store = FailAfter(2) # shard + sidecar land, manifest does not + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(store, tmp_path) + assert "2 of 3" in str(excinfo.value) + + +def test_a_clean_run_reports_no_tear(tmp_path): + store = FailAfter(99) + summary = _deliver(store, tmp_path) + assert summary["uploaded"] is True + assert summary["manifest_file_id"] == "id-3" + assert len(store.calls) == 3 + + +def test_the_object_that_failed_is_named_as_the_likeliest_orphan(tmp_path): + """The failing object is the one most likely to be an orphan, and it is not in the + confirmed list — because the list only holds uploads that returned. + + `_ContractStorePort.upload` raises precisely in the case its own comment documents: + the store "RETURNS success=False with the file already uploaded". So the object that + failed is the C-79 shape, sitting in the bucket with no metadata document. Listing + only the successes and calling it what remains would send an operator past the very + orphan this exists to surface. + """ + store = FailAfter(1) + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(store, tmp_path) + message = str(excinfo.value) + assert "MAY ALSO HAVE LANDED" in message + assert "C-79" in message, "and say which shape to look for" + # the sidecar is what failed here; it must appear even though it is not "confirmed" + assert "sidecar" in message + + +def test_a_failure_on_the_very_first_upload_does_not_claim_an_empty_list(tmp_path): + """The commonest infrastructure failure: credentials expire, upload #1 refuses. + + The first draft printed "Already in the partner store, and NOT removed: ." — an + empty list with a dangling period, presented as a bucket to audit. + """ + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(FailAfter(0), tmp_path) + message = str(excinfo.value) + assert "Nothing is confirmed in the partner store" in message + assert "NOT removed: ." not in message, "no dangling empty list" + # even here the failing object may have landed, so the caveat must still be present + assert "MAY ALSO HAVE LANDED" in message + + +def test_a_tear_is_not_a_malformed_run(tmp_path): + """Opposite retry semantics: SinkError means do-not-retry, a tear is transient.""" + with pytest.raises(sink.TornRunError) as excinfo: + _deliver(FailAfter(1), tmp_path) + assert not isinstance(excinfo.value, sink.SinkError), ( + "TornRunError must not share a base with the malformed-run family, or an " + "orchestration layer treating SinkError as do-not-retry would silently swallow " + "store outages — and test_hop_b_sink_e2e already asserts SinkError for malformed" + ) diff --git a/views_postprocessing/contract/wire/sink.py b/views_postprocessing/contract/wire/sink.py index 8acb5cb..cd24c24 100644 --- a/views_postprocessing/contract/wire/sink.py +++ b/views_postprocessing/contract/wire/sink.py @@ -54,6 +54,57 @@ class SinkError(ValueError): """The assembled run cannot be delivered as declared.""" +class TornRunError(RuntimeError): + """An upload failed partway, leaving objects in the partner store. + + Deliberately **not** a ``SinkError``. That family means "the assembled run cannot be + delivered as declared" — malformed input, where retrying is pointless; a tear is a + transient infrastructure failure with the opposite retry semantics, and + ``tests/test_hop_b_sink_e2e`` already uses ``pytest.raises(SinkError)`` as the + malformed-run assertion. Sharing a base would let any future orchestration layer + that treats ``SinkError`` as do-not-retry silently swallow store outages. Same + reasoning, same base, as ``delivery.findability``'s two error types. + """ + + +def _torn_run_error(run_id, failed_on, uploaded, total, exc) -> TornRunError: + """The refusal for a run that died mid-upload, naming what may be left (C-105). + + The contract handles the *consumer's* side correctly and by design: the run manifest + is uploaded last, so a torn attempt has no commit marker and is invisible rather + than half-visible (§4.2). What it does not handle is our side — the objects that did + land stay there, and until this existed nothing recorded that they had. At run-0 + scale a retry adds ~110 more under the same names. + + **Nothing is deleted here, deliberately.** Removing objects from a partner bucket is + irreversible and an operator decision, not a delivery-path one; the neighbouring + delete surface is its own open question (C-58, views-pipeline-core #333, blocked on + a test key). This turns an invisible mess into a documented one, which is the part + this repository can honestly own. + """ + landed = ", ".join(f"{u['name']}#{u['file_id']}" for u in uploaded[:5]) + more = "" if len(uploaded) <= 5 else f" (+{len(uploaded) - 5} more; every id is in the run log)" + confirmed = ( + f"Confirmed in the partner store, and NOT removed: {landed}{more}." + if uploaded + else "Nothing is confirmed in the partner store — this was the first upload." + ) + return TornRunError( + f"run {run_id!r} is TORN: {len(uploaded)} of {total} objects were confirmed " + f"uploaded before {failed_on!r} failed ({type(exc).__name__}: {exc}).\n" + f"{confirmed}\n" + f"{failed_on!r} MAY ALSO HAVE LANDED, as a file carrying no metadata document — " + "the store reports failure *after* uploading the file when the document write " + "fails, which is the C-79 orphan shape. Check for it as well as anything listed above; " + "it is the likeliest orphan of the whole run.\n" + "The manifest upload did not report success, so the consumer almost certainly " + "cannot see this run (§4.2 — the manifest is the commit marker). Verify that " + "before re-running: a re-run uploads every object again under the same names, " + "and whether the store supersedes or duplicates is not something this " + "repository asserts. See docs/operations/correction_procedure.md and C-105." + ) + + def deliver_run( per_target: dict, *, @@ -160,11 +211,22 @@ def deliver_run( common = {"name": consumer_name, "category": "forecast", "loa": "pgm"} + # The ledger is kept in memory as well as in the log, so a torn run can say what + # it left behind rather than leaving an operator to diff the bucket (C-105). + uploaded: list[dict] = [] + total = len(shard_records) + 2 # shards + sidecar + manifest + def _upload(file_name: str, doc_type: str, targets: list): - file_id = store.upload( - staging / file_name, filename=file_name, doc_type=doc_type, targets=targets, **common + try: + file_id = store.upload( + staging / file_name, filename=file_name, doc_type=doc_type, targets=targets, **common + ) + except Exception as exc: + raise _torn_run_error(run_id, file_name, uploaded, total, exc) from exc + uploaded.append({"name": file_name, "file_id": file_id}) + logger.info( # the ledger — file_id included, it is the only persistent record + "uploaded %s (type=%s, run=%s, file_id=%s)", file_name, doc_type, run_id, file_id ) - logger.info("uploaded %s (type=%s, run=%s)", file_name, doc_type, run_id) # the ledger return file_id for record in shard_records: