From d1f1a34e74da411325344f87690b85fddd128f11 Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 21:14:00 +0200 Subject: [PATCH 1/7] fix(delivery): verify every uploaded artefact by name, not two legs by category (#312) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The first live UN-FAO delivery, 2026-09-29, uploaded 109 of 110 objects and reported success. views-faoapi refused it. The missing object was the §5 GAUL sidecar: its bytes matched the previous run's, the content-addressed store correctly declined a second copy, update_document ran against the OLD file, and _ContractStorePort.upload returned a real file id for the wrong document. The consumer resolves by filename, found nothing, refused. The C-94 guard did not catch it, and could not have. It checked two entries, {"forecast": manifest_id, "historical": hist_id}, via latest_file_id({"name": ..., "category": category}). Every wire object is uploaded with category="forecast" (sink.py:227), so the newest forecast document is always the manifest, uploaded last by design. Adding the sidecar to that dict would have resolved the manifest and passed. The lookup key was wrong, not the dict. So the guard now asks two different questions, because neither implies the other: legs — does the consumer's SELECTION land on this run? (unchanged) objects — does each artefact resolve by its own filename AT ALL? (new) Per object rather than by count, and the decisive reason is not the one I started with. Pooling upstream is deterministic — measured in views-models on 2026-09-29, 25 of 25 anchor cells byte-identical across an accidental re-pool — and naming.py embeds the run id in every filename. A re-run therefore writes NEW names over IDENTICAL bytes, so the dedup path that took the sidecar takes all 110 objects at once: the store creates the documents, the count is right, nothing is servable. A count-plus-spot-check passes that. Re-running is what C-105 and C-22 tell an operator to do after a torn run, which makes it the realistic case rather than the exotic one — our own remedy is the trigger for the worst version. The query loop moved OUT of both partner managers into delivery/findability.py. Not incidental: the partner packages are under a ratcheting line budget whose own comment says the response to it binding is to move code out of the package rather than raise the number, and crafd/ sat at 699 of 700. Delegating paid for the wider call site exactly — both partners are unchanged at 690 and 699. It also leaves one copy of the rule instead of two that can disagree (the C-75 shape). An unqueryable filename degrades to UNVERIFIED, never to "invisible". filename is a declared collection attribute (pipeline-core provisioning.py:101), but if a store cannot be queried on it the honest answer is "could not ask" — a guard that manufactured an outage on every delivery would be deleted within a week (ADR-014 §3). Five mutations, each caught: report only the first missing object; accept any id as evidence; pass only the two legs (the shipped bug); stop exporting the sink's ledger; swallow a failed query as not-found. Complementary to views-pipeline-core's separate fix for the upload that reports success having written nothing. Neither waits on the other: theirs stops the lie at the source, this one stops us believing a lie from any source. Noted, not acted on: this is the first bug requiring an identical hand-patch in both partner managers, which is one of the two triggers C-33 names for extracting the partner seam. Recorded on #312; extracting it here would be a different change. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- tests/test_findability.py | 181 ++++++++++++++++++- views_postprocessing/contract/wire/sink.py | 5 + views_postprocessing/crafd/managers/crafd.py | 28 +-- views_postprocessing/delivery/findability.py | 115 ++++++++++++ views_postprocessing/unfao/managers/unfao.py | 28 +-- 5 files changed, 321 insertions(+), 36 deletions(-) diff --git a/tests/test_findability.py b/tests/test_findability.py index a6272ca..b735c72 100644 --- a/tests/test_findability.py +++ b/tests/test_findability.py @@ -187,14 +187,35 @@ def test_unverified_is_not_the_same_refusal_as_not_findable(): assert "forecast" in message, "and say which leg is unverified" -@pytest.mark.parametrize("partner", PARTNER_PACKAGES) -def test_the_preflight_does_not_report_a_failed_query_as_an_invisible_delivery(partner): - source = (_REPO / "views_postprocessing" / partner / "managers" / f"{partner}.py").read_text() - preflight = _function_source(source, "_assert_delivery_is_findable") - assert "findability.unverified(" in preflight, ( - f"{partner}'s preflight no longer distinguishes a store error from an empty " - "answer, so a transient network failure after a successful delivery would be " - "reported as the delivery being invisible — and quarantined (C-94)." +def test_the_preflight_does_not_report_a_failed_query_as_an_invisible_delivery(): + """The distinction must survive wherever the query loop lives. + + It used to live in each partner manager and this test read it there. #312 moved it + into `delivery/findability.verify` — one copy instead of two, and out of the partner + packages, whose line budget says to move code out rather than raise it. The claim is + unchanged: a store that could not be asked must not be reported as an invisible + delivery. Only its address changed, so the test follows it rather than being deleted. + + Behavioural, not a source scan, now that the logic is reachable without a manager. + """ + def explodes(_filters): + raise TimeoutError("read timed out") + + with pytest.raises(findability.FindabilityUnverifiedError): + findability.verify( + consumer_name="un_fao", legs={"forecast": "a"}, objects={}, resolve=explodes + ) + + # and the same for the by-name half, which #312 added + with pytest.raises(findability.FindabilityUnverifiedError) as caught: + findability.verify( + consumer_name="un_fao", + legs={}, + objects={"run__sidecar.parquet": "a"}, + resolve=explodes, + ) + assert "run__sidecar.parquet" in str(caught.value), ( + "an unverified object must name which object could not be checked" ) @@ -222,3 +243,147 @@ def test_the_previous_runs_document_does_not_satisfy_this_run(): "the consequence — the consumer goes on serving the previous delivery — is the " "part that distinguishes this from an empty bucket" ) + + +# ── #312: every artefact, by name ──────────────────────────────────────────── + + +def _resolver(store: dict): + """A store that answers by whatever filter key it is given, as Appwrite does.""" + def resolve(filters): + if "filename" in filters: + return store.get(filters["filename"]) + return store.get(("category", filters.get("category"))) + return resolve + + +def test_the_2026_09_29_incident_is_caught(): + """The regression case: 109 of 110 uploaded, every call reported success. + + The sidecar's bytes matched the previous run's, the store declined a second copy, + `update_document` ran against the OLD file, and the port returned a real file id + for the wrong document. The consumer resolves by filename, found nothing, refused. + """ + objects = {f"run2__shard{i}.parquet": f"id{i}" for i in range(3)} + objects["run2__sidecar.parquet"] = "id-sidecar" + objects["run2__manifest.json"] = "id-manifest" + + store = {name: fid for name, fid in objects.items() if "sidecar" not in name} + store[("category", "forecast")] = "id-manifest" + + with pytest.raises(findability.DeliveryNotFindableError) as caught: + findability.verify( + consumer_name="un_fao", + legs={"forecast": "id-manifest"}, + objects=objects, + resolve=_resolver(store), + ) + message = str(caught.value) + assert "run2__sidecar.parquet" in message, "the refusal must name the missing object" + assert "1 of 5" in message + assert "NOTHING in this run is servable" not in message, ( + "one missing object is not a total failure; the message must not overstate" + ) + + +def test_the_category_check_alone_would_have_passed_the_incident(): + """Why the by-name half had to be added rather than the dict widened. + + Shards, sidecar and manifest all carry `category="forecast"`, and the manifest is + uploaded last — so the newest forecast document is always the manifest. The + selection check is satisfied by the very run that is unservable. + """ + findability.verify( + consumer_name="un_fao", + legs={"forecast": "id-manifest"}, + objects={}, + resolve=_resolver({("category", "forecast"): "id-manifest"}), + ) + + +def test_a_deterministic_rerun_losing_everything_says_so(): + """The worst case, and the one our own remedy triggers. + + Pooling is deterministic and `naming.py` embeds the run id in every filename, so a + re-run writes NEW names over IDENTICAL bytes and every object deduplicates at once. + The store creates the documents, the count is right, nothing is servable — and + re-running is what C-105 and C-22 tell an operator to do. + """ + objects = {f"run2__shard{i}.parquet": f"id{i}" for i in range(4)} + with pytest.raises(findability.DeliveryNotFindableError) as caught: + findability.verify( + consumer_name="un_fao", legs={}, objects=objects, resolve=_resolver({}) + ) + message = str(caught.value) + assert "4 of 4" in message + assert "NOTHING in this run is servable" in message + assert "re-running unchanged reproduces this" in message, ( + "the refusal must not send an operator at the remedy that reproduces the fault" + ) + + +def test_an_object_resolving_to_the_wrong_document_is_refused(): + """A real id for the wrong file is the exact shape of the incident, and a check + that only asked 'did I get an id back' would pass it.""" + with pytest.raises(findability.DeliveryNotFindableError) as caught: + findability.verify( + consumer_name="un_fao", + legs={}, + objects={"run2__sidecar.parquet": "id-new"}, + resolve=_resolver({"run2__sidecar.parquet": "id-from-run1"}), + ) + message = str(caught.value) + assert "WRONG DOCUMENT" in message + assert "id-new" in message and "id-from-run1" in message + + +def test_a_fully_landed_run_passes(): + objects = {"run2__sidecar.parquet": "a", "run2__manifest.json": "b"} + store = dict(objects) + store[("category", "forecast")] = "b" + findability.verify( + consumer_name="un_fao", + legs={"forecast": "b"}, + objects=objects, + resolve=_resolver(store), + ) + + +@pytest.mark.parametrize("partner", PARTNER_PACKAGES) +def test_the_manager_hands_the_preflight_every_uploaded_object(partner): + """The wiring half. The rule above is worth nothing if the call site passes two ids. + + That is exactly what shipped: `{"forecast": ..., "historical": ...}` — the commit + marker and the historical leg, while 108 shards and the sidecar went unasked (#312). + """ + source = (_REPO / "views_postprocessing" / partner / "managers" / f"{partner}.py").read_text() + call = next( + n for n in ast.walk(ast.parse(source)) + if isinstance(n, ast.Call) + and isinstance(n.func, ast.Name) + and n.func.id == "_assert_delivery_is_findable" + ) + assert len(call.args) == 4, ( + f"{partner} calls the preflight with {len(call.args)} arguments; it needs the " + "per-object map as well as the per-leg one, or the sidecar goes unasked again" + ) + passed = ast.unparse(call.args[3]) + assert "uploaded_objects" in passed, ( + f"{partner}'s object map is {passed!r} — it must come from the sink's upload " + "ledger, which is the only record of what this run actually put in the bucket" + ) + assert "hist" in passed, ( + f"{partner} omits the historical artefact from the per-object check; it is " + "uploaded by the manager, so the sink's ledger does not contain it" + ) + + +def test_the_sink_carries_its_upload_ledger_out(): + """The ledger already existed for the torn-run refusal (C-105) and stayed local. + The manager cannot verify what it is not told.""" + source = (_REPO / "views_postprocessing" / "contract" / "wire" / "sink.py").read_text() + deliver = _function_source(source, "deliver_run") + assert 'summary["uploaded_objects"]' in deliver, ( + "deliver_run no longer exports the upload ledger, so the #312 per-object " + "preflight has nothing to check against" + ) diff --git a/views_postprocessing/contract/wire/sink.py b/views_postprocessing/contract/wire/sink.py index 3c88b59..11a97f6 100644 --- a/views_postprocessing/contract/wire/sink.py +++ b/views_postprocessing/contract/wire/sink.py @@ -251,5 +251,10 @@ def _upload(file_name: str, doc_type: str, targets: list): # in the store — which is exactly what the consumer's query returns. Carried out so # the C-94 read-back can assert the consumer would find THIS run (register C-94). summary["manifest_file_id"] = _upload(manifest_file, MANIFEST_DOC_TYPE, list(per_target)) + # The ledger leaves with the summary so the manager can verify EVERY object by + # filename, not just the commit marker (C-94, #312). It already exists for the + # torn-run refusal; carrying it out costs nothing and is the only record of what + # this run actually put in the bucket. + summary["uploaded_objects"] = {u["name"]: u["file_id"] for u in uploaded} summary["uploaded"] = True return summary diff --git a/views_postprocessing/crafd/managers/crafd.py b/views_postprocessing/crafd/managers/crafd.py index d5f9982..06c9cd0 100644 --- a/views_postprocessing/crafd/managers/crafd.py +++ b/views_postprocessing/crafd/managers/crafd.py @@ -124,23 +124,22 @@ def _build_partner_read_store(model_path) -> DatastoreModule: return store -def _assert_delivery_is_findable(model_path, consumer_name: str, uploaded: dict) -> None: - """C-94: ask the store the question the consumer asks, and refuse silence. +def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects: dict) -> None: + """C-94: ask the store the questions the consumer asks, and refuse silence. A function, not a method (C-40 (a)) — its refusal is observable without a manager - or an Appwrite environment. ``uploaded`` maps each leg to the file id THIS run put - there; `delivery/findability.py` carries why that scoping is the whole guard. + or an Appwrite environment. The rule lives in `delivery/findability.py`; this owns + only the port. ``legs`` scopes the consumer's SELECTION to this run; ``objects`` + asks whether each uploaded artefact resolves by filename at all (#312). """ - for category, expected in uploaded.items(): - try: - port = _ContractStorePort(_build_partner_read_store(model_path)) - found = port.latest_file_id({"name": consumer_name, "category": category}) - except Exception as exc: # could not ask != asked and got nothing (C-99, C-103) - raise findability.unverified(category, exc) from exc - findability.assert_findable( - found, expected_file_id=expected, consumer_name=consumer_name, category=category - ) - logger.info("Findability preflight passed: both legs retrievable under %r.", consumer_name) + port = _ContractStorePort(_build_partner_read_store(model_path)) + findability.verify( + consumer_name=consumer_name, legs=legs, objects=objects, resolve=port.latest_file_id + ) + logger.info( + "Findability preflight passed: %d leg(s) and %d object(s) under %r.", + len(legs), len(objects), consumer_name, + ) class CRAFDPostProcessorManager(PostprocessorManager, ForecastingModelManager): @@ -394,6 +393,7 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, + {**summary.get("uploaded_objects", {}), hist_path.name: hist_file_id}, ) else: logger.info( diff --git a/views_postprocessing/delivery/findability.py b/views_postprocessing/delivery/findability.py index 508dea3..8206d49 100644 --- a/views_postprocessing/delivery/findability.py +++ b/views_postprocessing/delivery/findability.py @@ -105,3 +105,118 @@ def assert_findable(file_id, *, expected_file_id, consumer_name: str, category: "bucket before re-delivering — the contract has no retraction primitive, so a " "correction is a new complete run." ) + + +def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: + """Raise unless EVERY artefact this run uploaded resolves under its own filename. + + Args: + resolved: ``{filename: (expected_file_id, found_file_id)}`` — one entry per + object the run uploaded, ``found_file_id`` being what the store returns + when asked for that exact filename under the consumer's name. ``None`` + means the query found nothing. + consumer_name: the DECLARED store-document ``name`` (C-77). + + Raises: + DeliveryNotFindableError: naming every object that does not resolve, not the + first — an operator auditing a partner bucket needs the whole list, and + the whole list is the difference between "one file is missing" and + "nothing in this run is servable". + + **Why by filename and not by id (register C-94, #312).** The first live UN-FAO + delivery, 2026-09-29, uploaded 109 of 110 objects and reported success. The GAUL + sidecar's bytes were identical to the previous run's, the content-addressed store + correctly declined a second copy, ``update_document`` ran against the OLD file, and + the port returned a **real file id for the wrong document**. views-faoapi resolves + by ``filename``, found nothing, and refused the delivery. + + So an id is not evidence. `assert_findable` above asks whether the consumer's + *selection* lands on this run; this asks whether each artefact the manifest + references is *there at all*. Both are needed and neither implies the other. + + **Why every object and not a count.** Pooling upstream is deterministic (measured + in views-models, 2026-09-29: 25 of 25 anchor cells byte-identical across a re-pool), + and `contract/wire/naming.py` embeds the run id in every filename. A re-run + therefore writes NEW filenames over IDENTICAL bytes, so the dedup path that took the + sidecar takes **all 110 objects at once** — the store creates the documents, the + count is right, and nothing is servable. A count-plus-spot-check passes that. + Re-running is our own documented remedy for a torn run (C-105) and for a correction + (C-22), which is what makes this the realistic case rather than the exotic one. + """ + missing = sorted(name for name, (_, found) in resolved.items() if not found) + wrong = sorted( + f"{name} (expected {exp!r}, store has {found!r})" + for name, (exp, found) in resolved.items() + if found and found != exp + ) + if not missing and not wrong: + return + + total = len(resolved) + parts = [ + f"delivery is INVISIBLE to the consumer: {len(missing) + len(wrong)} of {total} " + f"uploaded object(s) do not resolve under name == {consumer_name!r}." + ] + if missing: + parts.append( + f"NOT FOUND by filename ({len(missing)}): {', '.join(missing)}. Every upload " + "reported success, so these exist as ids pointing at some other document — " + "the shape that stranded the 2026-09-29 sidecar when the store deduplicated " + "identical bytes and updated the PREVIOUS run's file instead." + ) + if wrong: + parts.append(f"RESOLVES TO THE WRONG DOCUMENT ({len(wrong)}): {'; '.join(wrong)}.") + if len(missing) + len(wrong) == total: + parts.append( + "NOTHING in this run is servable. If this was a re-run of an earlier one, " + "that is the expected shape: pooling is deterministic, so a re-run writes " + "new filenames over identical bytes and every object deduplicates at once." + ) + parts.append( + "The contract has no retraction primitive, so a correction is a new complete " + "run — but re-running unchanged reproduces this. Fix the upload path first." + ) + raise DeliveryNotFindableError(" ".join(parts)) + + +def verify(*, consumer_name: str, legs: dict, objects: dict, resolve) -> None: + """Run both findability questions against a partner store. The caller owns the port. + + Args: + consumer_name: the DECLARED store-document ``name`` (C-77). + legs: ``{category: expected_file_id}`` — does the consumer's own *selection* + land on this run? Answered per category because a run whose forecast + landed and whose historical did not is invisible in exactly one half. + objects: ``{filename: expected_file_id}`` — is each artefact *there at all*? + resolve: ``callable(filters: dict) -> file_id | None``, normally the port's + ``latest_file_id``. A callable rather than a store keeps this module free + of store types, which is what lets its refusals be tested without Appwrite. + + Raises: + DeliveryNotFindableError: the delivery, or part of it, is invisible. + FindabilityUnverifiedError: a query failed, so findability is UNKNOWN. This is + the deliberate soft edge (ADR-014 §3): ``filename`` is a declared collection + attribute, but if a store cannot be queried on it the honest answer is "could + not ask", not a manufactured outage on every delivery. + + Lives here rather than in the partner managers because the logic is identical in + both and the partner packages are under a ratcheting line budget whose stated + response to binding is to move code **out of the package**. It also means one copy + of the rule instead of two that can disagree — the C-75 shape. + """ + for category, expected in legs.items(): + try: + found = resolve({"name": consumer_name, "category": category}) + except Exception as exc: # could not ask != asked and got nothing (C-99, C-103) + raise unverified(category, exc) from exc + assert_findable( + found, expected_file_id=expected, consumer_name=consumer_name, category=category + ) + + resolved = {} + for filename, expected in objects.items(): + try: + resolved[filename] = (expected, resolve({"name": consumer_name, "filename": filename})) + except Exception as exc: + raise unverified(f"object {filename!r}", exc) from exc + assert_all_findable(resolved, consumer_name=consumer_name) diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index 8720230..1413365 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -124,23 +124,22 @@ def _build_partner_read_store(model_path) -> DatastoreModule: return store -def _assert_delivery_is_findable(model_path, consumer_name: str, uploaded: dict) -> None: - """C-94: ask the store the question the consumer asks, and refuse silence. +def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects: dict) -> None: + """C-94: ask the store the questions the consumer asks, and refuse silence. A function, not a method (C-40 (a)) — its refusal is observable without a manager - or an Appwrite environment. ``uploaded`` maps each leg to the file id THIS run put - there; `delivery/findability.py` carries why that scoping is the whole guard. + or an Appwrite environment. The rule lives in `delivery/findability.py`; this owns + only the port. ``legs`` scopes the consumer's SELECTION to this run; ``objects`` + asks whether each uploaded artefact resolves by filename at all (#312). """ - for category, expected in uploaded.items(): - try: - port = _ContractStorePort(_build_partner_read_store(model_path)) - found = port.latest_file_id({"name": consumer_name, "category": category}) - except Exception as exc: # could not ask != asked and got nothing (C-99, C-103) - raise findability.unverified(category, exc) from exc - findability.assert_findable( - found, expected_file_id=expected, consumer_name=consumer_name, category=category - ) - logger.info("Findability preflight passed: both legs retrievable under %r.", consumer_name) + port = _ContractStorePort(_build_partner_read_store(model_path)) + findability.verify( + consumer_name=consumer_name, legs=legs, objects=objects, resolve=port.latest_file_id + ) + logger.info( + "Findability preflight passed: %d leg(s) and %d object(s) under %r.", + len(legs), len(objects), consumer_name, + ) class UNFAOPostProcessorManager(PostprocessorManager, ForecastingModelManager): @@ -394,6 +393,7 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, + {**summary.get("uploaded_objects", {}), hist_path.name: hist_file_id}, ) else: logger.info( From 4169fc5ed03e7f31f4ef12f2470e547ef3b06016 Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 21:23:08 +0200 Subject: [PATCH 2/7] fix(delivery): use the consumer's proven query shape, not a filename lookup (#312 review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review finding from the views-models seat, verified here before acting. The per-object check queried `filename` directly. Three facts say it would not work: 1. `Query.equal("filename", ...)` appears NOWHERE on the platform. The hits in pipeline-core, faoapi and crafdapi are `Query.equal("name", filename)` against the storage BUCKET — a different call on a different object. 2. `filename` is declared in pipeline-core's collection schema (provisioning.py:101) and **no index is declared anywhere in that file**. 3. views-faoapi — the consumer being modelled — deliberately does not query it. `resolve_artifact_file_ids` (prediction/manager.py:277-295) runs a type-scoped query on {category, type} and matches `filename` in Python. Its docstring says so, and that shape ran against the real store during the 2026-09-29 incident. Left as written, the guard would return UNVERIFIED on every delivery forever while its presence read as coverage — a guard that structurally cannot fire (views-models C-155), which is worse than the bug it replaces. The failure direction was right and the query was wrong. So the check now uses faoapi's shape: scope to {name, category, type}, match filename in Python. Better on its own merits, index question aside — it asks the question the consumer actually asks, which is C-94's thesis rather than an approximation, and it costs one query per artefact TYPE rather than one per object (110 -> 3). Carrying it needed the sink's ledger to record doc_type, and the port to expose `documents()` alongside `latest_file_id`. Both partner packages came in UNDER where they started: unfao 690 -> 688, crafd 699 -> 697. The guard's docstring and success log moved into `findability.verify` with the rule they describe — one home per fact, and the budget's own instruction is to move code out of the package rather than raise the number. Moving `store_port.py` out was considered and rejected: it is byte-identical in both partners and would free 111 lines each, but it is referenced by six register entries, an anti-drift test and ADR-015, and that blast radius does not belong inside an incident fix. Six mutations, each caught: revert to a filename query; accept any id as evidence; report only the first missing object; pass no objects (the shipped bug); swallow a failed query as not-found; stop exporting the ledger. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- tests/test_findability.py | 136 ++++++++++++------- views_postprocessing/contract/wire/sink.py | 6 +- views_postprocessing/crafd/managers/crafd.py | 21 +-- views_postprocessing/crafd/store_port.py | 5 + views_postprocessing/delivery/findability.py | 64 ++++++--- views_postprocessing/unfao/managers/unfao.py | 21 +-- views_postprocessing/unfao/store_port.py | 5 + 7 files changed, 159 insertions(+), 99 deletions(-) diff --git a/tests/test_findability.py b/tests/test_findability.py index b735c72..9220578 100644 --- a/tests/test_findability.py +++ b/tests/test_findability.py @@ -203,19 +203,20 @@ def explodes(_filters): with pytest.raises(findability.FindabilityUnverifiedError): findability.verify( - consumer_name="un_fao", legs={"forecast": "a"}, objects={}, resolve=explodes + consumer_name="un_fao", legs={"forecast": "a"}, objects=[], + resolve_latest=explodes, list_documents=explodes, ) - # and the same for the by-name half, which #312 added + # and the same for the per-object half, which #312 added with pytest.raises(findability.FindabilityUnverifiedError) as caught: findability.verify( - consumer_name="un_fao", - legs={}, - objects={"run__sidecar.parquet": "a"}, - resolve=explodes, + consumer_name="un_fao", legs={}, + objects=[{"name": "run__sidecar.parquet", "file_id": "a", + "doc_type": "sampled_forecast_sidecar", "category": "forecast"}], + resolve_latest=explodes, list_documents=explodes, ) - assert "run__sidecar.parquet" in str(caught.value), ( - "an unverified object must name which object could not be checked" + assert "sampled_forecast_sidecar" in str(caught.value), ( + "an unverified scope must name which artefact type could not be checked" ) @@ -245,16 +246,26 @@ def test_the_previous_runs_document_does_not_satisfy_this_run(): ) -# ── #312: every artefact, by name ──────────────────────────────────────────── +# ── #312: every artefact, by the consumer's own resolution shape ───────────── + + +def _objects(names, doc_type="sampled_forecast_shard", category="forecast"): + return [ + {"name": n, "file_id": f"id-{n}", "doc_type": doc_type, "category": category} + for n in names + ] -def _resolver(store: dict): - """A store that answers by whatever filter key it is given, as Appwrite does.""" - def resolve(filters): - if "filename" in filters: - return store.get(filters["filename"]) - return store.get(("category", filters.get("category"))) - return resolve +def _store(landed: dict): + """A store answering the way views-faoapi queries it: type-scoped, filename in + the document. `landed` maps filename -> fileId for what is actually retrievable.""" + def list_documents(_filters): + return [{"filename": n, "fileId": fid} for n, fid in landed.items()] + return list_documents + + +def _never(_filters): + raise AssertionError("the selection query must not run when legs is empty") def test_the_2026_09_29_incident_is_caught(): @@ -264,55 +275,49 @@ def test_the_2026_09_29_incident_is_caught(): `update_document` ran against the OLD file, and the port returned a real file id for the wrong document. The consumer resolves by filename, found nothing, refused. """ - objects = {f"run2__shard{i}.parquet": f"id{i}" for i in range(3)} - objects["run2__sidecar.parquet"] = "id-sidecar" - objects["run2__manifest.json"] = "id-manifest" - - store = {name: fid for name, fid in objects.items() if "sidecar" not in name} - store[("category", "forecast")] = "id-manifest" + objects = _objects(["run2__shard0.parquet", "run2__shard1.parquet"]) + objects += _objects(["run2__sidecar.parquet"], "sampled_forecast_sidecar") + landed = {o["name"]: o["file_id"] for o in objects if "sidecar" not in o["name"]} with pytest.raises(findability.DeliveryNotFindableError) as caught: findability.verify( - consumer_name="un_fao", - legs={"forecast": "id-manifest"}, - objects=objects, - resolve=_resolver(store), + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, list_documents=_store(landed), ) message = str(caught.value) assert "run2__sidecar.parquet" in message, "the refusal must name the missing object" - assert "1 of 5" in message + assert "1 of 3" in message assert "NOTHING in this run is servable" not in message, ( "one missing object is not a total failure; the message must not overstate" ) -def test_the_category_check_alone_would_have_passed_the_incident(): - """Why the by-name half had to be added rather than the dict widened. +def test_the_selection_check_alone_would_have_passed_the_incident(): + """Why the per-object half had to be added rather than the dict widened. - Shards, sidecar and manifest all carry `category="forecast"`, and the manifest is - uploaded last — so the newest forecast document is always the manifest. The + Shards, sidecar and manifest all carry category="forecast" and the manifest is + uploaded last, so the newest forecast document is always the manifest. The selection check is satisfied by the very run that is unservable. """ findability.verify( - consumer_name="un_fao", - legs={"forecast": "id-manifest"}, - objects={}, - resolve=_resolver({("category", "forecast"): "id-manifest"}), + consumer_name="un_fao", legs={"forecast": "id-manifest"}, objects=[], + resolve_latest=lambda _f: "id-manifest", list_documents=_never, ) def test_a_deterministic_rerun_losing_everything_says_so(): """The worst case, and the one our own remedy triggers. - Pooling is deterministic and `naming.py` embeds the run id in every filename, so a + Pooling is deterministic and naming.py embeds the run id in every filename, so a re-run writes NEW names over IDENTICAL bytes and every object deduplicates at once. The store creates the documents, the count is right, nothing is servable — and re-running is what C-105 and C-22 tell an operator to do. """ - objects = {f"run2__shard{i}.parquet": f"id{i}" for i in range(4)} + objects = _objects([f"run2__shard{i}.parquet" for i in range(4)]) with pytest.raises(findability.DeliveryNotFindableError) as caught: findability.verify( - consumer_name="un_fao", legs={}, objects=objects, resolve=_resolver({}) + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, list_documents=_store({}), ) message = str(caught.value) assert "4 of 4" in message @@ -323,29 +328,56 @@ def test_a_deterministic_rerun_losing_everything_says_so(): def test_an_object_resolving_to_the_wrong_document_is_refused(): - """A real id for the wrong file is the exact shape of the incident, and a check - that only asked 'did I get an id back' would pass it.""" + """Content-hash dedup returns a REAL id for the WRONG document, so a check that + only asked "did I get an id back" passes the incident. This one does not.""" + objects = _objects(["run2__sidecar.parquet"], "sampled_forecast_sidecar") with pytest.raises(findability.DeliveryNotFindableError) as caught: findability.verify( - consumer_name="un_fao", - legs={}, - objects={"run2__sidecar.parquet": "id-new"}, - resolve=_resolver({"run2__sidecar.parquet": "id-from-run1"}), + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, + list_documents=_store({"run2__sidecar.parquet": "id-from-run1"}), ) message = str(caught.value) assert "WRONG DOCUMENT" in message - assert "id-new" in message and "id-from-run1" in message + assert "id-from-run1" in message + + +def test_the_query_is_type_scoped_and_costs_one_per_type_not_one_per_object(): + """The #312 review finding: `filename` is declared but NOT indexed, and nothing on + the platform queries it. This uses views-faoapi's proven shape instead — scope by + {category, type}, match filename in Python — which also collapses 110 lookups to + one per artefact type.""" + objects = _objects([f"s{i}.parquet" for i in range(50)]) + objects += _objects(["sidecar.parquet"], "sampled_forecast_sidecar") + objects += _objects(["hist.parquet"], "model", "historical") + seen = [] + + def list_documents(filters): + seen.append(filters) + assert "filename" not in filters, ( + "queried by filename — that attribute has no index and this deployment " + "enforces index requirements, so the guard would return UNVERIFIED forever" + ) + return [{"filename": o["name"], "fileId": o["file_id"]} for o in objects] + + findability.verify( + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, list_documents=list_documents, + ) + assert len(seen) == 3, f"expected one query per (category, type), got {len(seen)}" + assert {tuple(sorted(f.items())) for f in seen} == { + (("category", "forecast"), ("name", "un_fao"), ("type", "sampled_forecast_shard")), + (("category", "forecast"), ("name", "un_fao"), ("type", "sampled_forecast_sidecar")), + (("category", "historical"), ("name", "un_fao"), ("type", "model")), + } def test_a_fully_landed_run_passes(): - objects = {"run2__sidecar.parquet": "a", "run2__manifest.json": "b"} - store = dict(objects) - store[("category", "forecast")] = "b" + objects = _objects(["a.parquet"]) + _objects(["m.json"], "sampled_forecast_manifest") + landed = {o["name"]: o["file_id"] for o in objects} findability.verify( - consumer_name="un_fao", - legs={"forecast": "b"}, - objects=objects, - resolve=_resolver(store), + consumer_name="un_fao", legs={"forecast": "id-m.json"}, objects=objects, + resolve_latest=lambda _f: "id-m.json", list_documents=_store(landed), ) diff --git a/views_postprocessing/contract/wire/sink.py b/views_postprocessing/contract/wire/sink.py index 11a97f6..ed4e2d6 100644 --- a/views_postprocessing/contract/wire/sink.py +++ b/views_postprocessing/contract/wire/sink.py @@ -238,7 +238,7 @@ def _upload(file_name: str, doc_type: str, targets: list): ) 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}) + uploaded.append({"name": file_name, "file_id": file_id, "doc_type": doc_type}) 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 ) @@ -255,6 +255,8 @@ def _upload(file_name: str, doc_type: str, targets: list): # filename, not just the commit marker (C-94, #312). It already exists for the # torn-run refusal; carrying it out costs nothing and is the only record of what # this run actually put in the bucket. - summary["uploaded_objects"] = {u["name"]: u["file_id"] for u in uploaded} + summary["uploaded_objects"] = [ + {**u, "category": common["category"]} for u in uploaded + ] summary["uploaded"] = True return summary diff --git a/views_postprocessing/crafd/managers/crafd.py b/views_postprocessing/crafd/managers/crafd.py index 06c9cd0..52a63b2 100644 --- a/views_postprocessing/crafd/managers/crafd.py +++ b/views_postprocessing/crafd/managers/crafd.py @@ -124,21 +124,13 @@ def _build_partner_read_store(model_path) -> DatastoreModule: return store -def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects: dict) -> None: - """C-94: ask the store the questions the consumer asks, and refuse silence. - - A function, not a method (C-40 (a)) — its refusal is observable without a manager - or an Appwrite environment. The rule lives in `delivery/findability.py`; this owns - only the port. ``legs`` scopes the consumer's SELECTION to this run; ``objects`` - asks whether each uploaded artefact resolves by filename at all (#312). - """ +def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects) -> None: + """C-94/#312: ask the store the questions the consumer asks. The rule and its + refusals live in `delivery/findability.verify`; this owns only the port.""" port = _ContractStorePort(_build_partner_read_store(model_path)) findability.verify( - consumer_name=consumer_name, legs=legs, objects=objects, resolve=port.latest_file_id - ) - logger.info( - "Findability preflight passed: %d leg(s) and %d object(s) under %r.", - len(legs), len(objects), consumer_name, + consumer_name=consumer_name, legs=legs, objects=objects, + resolve_latest=port.latest_file_id, list_documents=port.documents, ) @@ -393,7 +385,8 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, - {**summary.get("uploaded_objects", {}), hist_path.name: hist_file_id}, + [*summary.get("uploaded_objects", []), {"name": hist_path.name, + "file_id": hist_file_id, "doc_type": "model", "category": "historical"}], ) else: logger.info( diff --git a/views_postprocessing/crafd/store_port.py b/views_postprocessing/crafd/store_port.py index b0ee757..8dc64ae 100644 --- a/views_postprocessing/crafd/store_port.py +++ b/views_postprocessing/crafd/store_port.py @@ -27,6 +27,11 @@ def __init__(self, datastore) -> None: def latest_file_id(self, filters: dict): return self._dsm.get_latest_file_id(filters=filters) + def documents(self, filters: dict) -> list: + """Documents matching ``filters`` — views-faoapi's own resolution shape; + `delivery/findability.verify` carries why it is not a filename query.""" + return self._dsm.get_predictions_by_metadata(filters=filters) + def file_metadata(self, file_id: str) -> dict: return store_metadata.file_metadata(self._dsm.get_file_metadata(file_id)) diff --git a/views_postprocessing/delivery/findability.py b/views_postprocessing/delivery/findability.py index 8206d49..f7a148a 100644 --- a/views_postprocessing/delivery/findability.py +++ b/views_postprocessing/delivery/findability.py @@ -27,6 +27,10 @@ from __future__ import annotations +import logging + +logger = logging.getLogger(__name__) + class DeliveryNotFindableError(RuntimeError): """An upload succeeded, and the consumer's own query cannot find it.""" @@ -179,7 +183,7 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: raise DeliveryNotFindableError(" ".join(parts)) -def verify(*, consumer_name: str, legs: dict, objects: dict, resolve) -> None: +def verify(*, consumer_name: str, legs: dict, objects, resolve_latest, list_documents) -> None: """Run both findability questions against a partner store. The caller owns the port. Args: @@ -187,36 +191,62 @@ def verify(*, consumer_name: str, legs: dict, objects: dict, resolve) -> None: legs: ``{category: expected_file_id}`` — does the consumer's own *selection* land on this run? Answered per category because a run whose forecast landed and whose historical did not is invisible in exactly one half. - objects: ``{filename: expected_file_id}`` — is each artefact *there at all*? - resolve: ``callable(filters: dict) -> file_id | None``, normally the port's - ``latest_file_id``. A callable rather than a store keeps this module free - of store types, which is what lets its refusals be tested without Appwrite. + objects: records of every artefact uploaded — ``name``, ``file_id``, + ``doc_type``, ``category``. Is each one *there at all*? + resolve_latest: ``callable(filters) -> file_id | None`` (the port's + ``latest_file_id``), for the selection question. + list_documents: ``callable(filters) -> list[dict]`` (the port's ``documents``), + for the per-object question. Raises: DeliveryNotFindableError: the delivery, or part of it, is invisible. - FindabilityUnverifiedError: a query failed, so findability is UNKNOWN. This is - the deliberate soft edge (ADR-014 §3): ``filename`` is a declared collection - attribute, but if a store cannot be queried on it the honest answer is "could - not ask", not a manufactured outage on every delivery. + FindabilityUnverifiedError: a query failed, so findability is UNKNOWN. + + **Why a type-scoped query and Python matching, not a filename query (#312 review).** + The obvious implementation asks the store for ``filename == X``. It would be the + first code on this platform to do so: ``filename`` is declared in pipeline-core's + collection schema (`provisioning.py:101`) but **no index is declared for it + anywhere**, and this deployment enforces index requirements on attribute lookups. + An unsupported query would make every delivery report UNVERIFIED — a guard that + logs "could not ask" forever while its mere presence reads as coverage. That is the + guard-that-cannot-fire shape (**C-155**), and it is worse than the bug it replaces. + + So this uses the shape views-faoapi already proves against this store + (`prediction/manager.py::resolve_artifact_file_ids`): query scoped to + ``{category, type}``, match ``filename`` in Python. It also means the check asks + the question **the consumer actually asks**, which is C-94's thesis rather than an + approximation of it — and it costs one query per artefact type, not one per object. Lives here rather than in the partner managers because the logic is identical in - both and the partner packages are under a ratcheting line budget whose stated - response to binding is to move code **out of the package**. It also means one copy - of the rule instead of two that can disagree — the C-75 shape. + both and must not diverge (the C-75 shape), and because the partner packages are + under a ratcheting line budget whose stated response to binding is to move code + **out of the package**. """ for category, expected in legs.items(): try: - found = resolve({"name": consumer_name, "category": category}) + found = resolve_latest({"name": consumer_name, "category": category}) except Exception as exc: # could not ask != asked and got nothing (C-99, C-103) raise unverified(category, exc) from exc assert_findable( found, expected_file_id=expected, consumer_name=consumer_name, category=category ) - resolved = {} - for filename, expected in objects.items(): + scopes = {(o["category"], o["doc_type"]) for o in objects} + seen: dict[str, str] = {} + for category, doc_type in sorted(scopes): + filters = {"name": consumer_name, "category": category, "type": doc_type} try: - resolved[filename] = (expected, resolve({"name": consumer_name, "filename": filename})) + docs = list_documents(filters) except Exception as exc: - raise unverified(f"object {filename!r}", exc) from exc + raise unverified(f"{category}/{doc_type} objects", exc) from exc + for doc in docs or (): + filename, file_id = doc.get("filename"), doc.get("fileId") + if filename and file_id and filename not in seen: + seen[filename] = file_id + + resolved = {o["name"]: (o["file_id"], seen.get(o["name"])) for o in objects} assert_all_findable(resolved, consumer_name=consumer_name) + logger.info( + "Findability preflight passed: %d leg(s) and %d object(s) resolvable under %r.", + len(legs), len(objects), consumer_name, + ) diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index 1413365..abdaff8 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -124,21 +124,13 @@ def _build_partner_read_store(model_path) -> DatastoreModule: return store -def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects: dict) -> None: - """C-94: ask the store the questions the consumer asks, and refuse silence. - - A function, not a method (C-40 (a)) — its refusal is observable without a manager - or an Appwrite environment. The rule lives in `delivery/findability.py`; this owns - only the port. ``legs`` scopes the consumer's SELECTION to this run; ``objects`` - asks whether each uploaded artefact resolves by filename at all (#312). - """ +def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects) -> None: + """C-94/#312: ask the store the questions the consumer asks. The rule and its + refusals live in `delivery/findability.verify`; this owns only the port.""" port = _ContractStorePort(_build_partner_read_store(model_path)) findability.verify( - consumer_name=consumer_name, legs=legs, objects=objects, resolve=port.latest_file_id - ) - logger.info( - "Findability preflight passed: %d leg(s) and %d object(s) under %r.", - len(legs), len(objects), consumer_name, + consumer_name=consumer_name, legs=legs, objects=objects, + resolve_latest=port.latest_file_id, list_documents=port.documents, ) @@ -393,7 +385,8 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, - {**summary.get("uploaded_objects", {}), hist_path.name: hist_file_id}, + [*summary.get("uploaded_objects", []), {"name": hist_path.name, + "file_id": hist_file_id, "doc_type": "model", "category": "historical"}], ) else: logger.info( diff --git a/views_postprocessing/unfao/store_port.py b/views_postprocessing/unfao/store_port.py index b0ee757..8dc64ae 100644 --- a/views_postprocessing/unfao/store_port.py +++ b/views_postprocessing/unfao/store_port.py @@ -27,6 +27,11 @@ def __init__(self, datastore) -> None: def latest_file_id(self, filters: dict): return self._dsm.get_latest_file_id(filters=filters) + def documents(self, filters: dict) -> list: + """Documents matching ``filters`` — views-faoapi's own resolution shape; + `delivery/findability.verify` carries why it is not a filename query.""" + return self._dsm.get_predictions_by_metadata(filters=filters) + def file_metadata(self, file_id: str) -> dict: return store_metadata.file_metadata(self._dsm.get_file_metadata(file_id)) From 04d2d80dfc7e59da9288feb8803a98edde120a84 Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 21:38:26 +0200 Subject: [PATCH 3/7] docs(review-diff): the prose still described the query the review removed MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three findings from /review-diff on this branch, two of them self-inflicted by the rework that answered the review. 1. WARNING — findability.py said "asked for that exact filename", "Why by filename and not by id", and "NOT FOUND by filename". After the rework the store is NOT queried on filename; it is a type-scoped query matched in Python, precisely because that attribute has no index. A reader would have concluded the opposite, and anyone simplifying would have restored the broken form. The test guards the filters; the prose was inviting the change. 2. WARNING — the CIC still said "After both legs are uploaded" and "Two distinct refusals". Both wrong: every artefact is checked now, and there are four refusals across two questions. Updated, with why the per-object half resolves the way faoapi does rather than by filename. Review date moved. 3. SUGGESTION — `objects` carried no annotation while `legs: dict` sat beside it. Annotated in verify and in both managers. Checked and NOT findings, recorded so they are not re-raised: blank-line hygiene is correct, `logger` remains used in both managers, and `assert_findable` is still reachable through the legs loop. Both partner packages remain under budget at 690 and 699. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- docs/CICs/UNFAOPostProcessorManager.md | 4 ++-- views_postprocessing/crafd/managers/crafd.py | 4 +++- views_postprocessing/delivery/findability.py | 18 +++++++++++------- views_postprocessing/unfao/managers/unfao.py | 4 +++- 4 files changed, 19 insertions(+), 11 deletions(-) diff --git a/docs/CICs/UNFAOPostProcessorManager.md b/docs/CICs/UNFAOPostProcessorManager.md index 84e715a..56cc4bb 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-26 +**Last reviewed:** 2026-09-29 **Related ADRs:** ADR-001, ADR-002, ADR-008, ADR-009 --- @@ -90,7 +90,7 @@ Assumptions that are not met **must cause failure**, not fallback behavior. The - **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`. **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 asserts the newest document it finds **is the one this run just uploaded**. Two distinct refusals: nothing found at all, and *found the previous run's* (*"the newest forecast document is X, but this run uploaded Y"*). The run-scoping is the whole guard — asking merely whether any document exists is a question the previous delivery already answers yes to, so the check could never fail from delivery 2 onward. 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 +- **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 asserts the newest document it finds **is the one this run just uploaded**. **Four distinct refusals, in two questions.** The *selection* question — does the consumer's own query land on this run — refuses when nothing is found at all, and when it finds the previous run's (*"the newest forecast document is X, but this run uploaded Y"*). The *per-object* question, added 2026-09-29 by #312, refuses when no document carries an uploaded artefact's filename, and when one does but carries a different file id. **Every artefact is checked, not the two legs**: the first live delivery uploaded 109 of 110 and reported success, because the sidecar's bytes matched the previous run's, the store declined a duplicate, and `update_document` ran against the old file — returning a real id for the wrong document. The per-object half resolves the way views-faoapi does (type-scoped query, filename matched in Python) because `filename` carries no index; a direct filename query would report UNVERIFIED forever. The run-scoping is the whole guard — asking merely whether any document exists is a question the previous delivery already answers yes to, so the check could never fail from delivery 2 onward. 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) - **Region coverage mismatch:** Raises `CoverageError` in `_check_coverage()` (called from `_validate()`) if a pinned region's delivered cell count is wrong (S1/C-34) or a GAUL-uncovered excluded cell leaks into the delivery (S4/C-30) diff --git a/views_postprocessing/crafd/managers/crafd.py b/views_postprocessing/crafd/managers/crafd.py index 52a63b2..299ce66 100644 --- a/views_postprocessing/crafd/managers/crafd.py +++ b/views_postprocessing/crafd/managers/crafd.py @@ -124,7 +124,9 @@ def _build_partner_read_store(model_path) -> DatastoreModule: return store -def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects) -> None: +def _assert_delivery_is_findable( + model_path, consumer_name: str, legs: dict, objects: list +) -> None: """C-94/#312: ask the store the questions the consumer asks. The rule and its refusals live in `delivery/findability.verify`; this owns only the port.""" port = _ContractStorePort(_build_partner_read_store(model_path)) diff --git a/views_postprocessing/delivery/findability.py b/views_postprocessing/delivery/findability.py index f7a148a..c26b203 100644 --- a/views_postprocessing/delivery/findability.py +++ b/views_postprocessing/delivery/findability.py @@ -116,9 +116,10 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: Args: resolved: ``{filename: (expected_file_id, found_file_id)}`` — one entry per - object the run uploaded, ``found_file_id`` being what the store returns - when asked for that exact filename under the consumer's name. ``None`` - means the query found nothing. + object the run uploaded. ``found_file_id`` is the id carried by the store + document whose ``filename`` field matches, as recovered by `verify` from a + type-scoped query (**not** a query on ``filename`` — see `verify`). ``None`` + means no document carried that filename. consumer_name: the DECLARED store-document ``name`` (C-77). Raises: @@ -127,7 +128,7 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: the whole list is the difference between "one file is missing" and "nothing in this run is servable". - **Why by filename and not by id (register C-94, #312).** The first live UN-FAO + **Why matched on filename rather than trusted by id (register C-94, #312).** The first live UN-FAO delivery, 2026-09-29, uploaded 109 of 110 objects and reported success. The GAUL sidecar's bytes were identical to the previous run's, the content-addressed store correctly declined a second copy, ``update_document`` ran against the OLD file, and @@ -163,8 +164,9 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: ] if missing: parts.append( - f"NOT FOUND by filename ({len(missing)}): {', '.join(missing)}. Every upload " - "reported success, so these exist as ids pointing at some other document — " + f"NO DOCUMENT CARRIES THIS FILENAME ({len(missing)}): {', '.join(missing)}. " + "Every upload reported success, so these exist as ids pointing at some " + "other document — " "the shape that stranded the 2026-09-29 sidecar when the store deduplicated " "identical bytes and updated the PREVIOUS run's file instead." ) @@ -183,7 +185,9 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: raise DeliveryNotFindableError(" ".join(parts)) -def verify(*, consumer_name: str, legs: dict, objects, resolve_latest, list_documents) -> None: +def verify( + *, consumer_name: str, legs: dict, objects: list, resolve_latest, list_documents +) -> None: """Run both findability questions against a partner store. The caller owns the port. Args: diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index abdaff8..b033ded 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -124,7 +124,9 @@ def _build_partner_read_store(model_path) -> DatastoreModule: return store -def _assert_delivery_is_findable(model_path, consumer_name: str, legs: dict, objects) -> None: +def _assert_delivery_is_findable( + model_path, consumer_name: str, legs: dict, objects: list +) -> None: """C-94/#312: ask the store the questions the consumer asks. The rule and its refusals live in `delivery/findability.verify`; this owns only the port.""" port = _ContractStorePort(_build_partner_read_store(model_path)) From d934007147c76f93335d0702ef0f28665ef55f2e Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 21:43:25 +0200 Subject: [PATCH 4/7] fix(delivery): a failed check must not quarantine a healthy delivery (/code-review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Six findings from /code-review medium on PR #314. It also traced the new query down through pipeline-core to Appwrite and confirmed the mechanism holds: the scope set matches faoapi's query exactly, newest-first is guaranteed by the $createdAt sort so first-wins reproduces the consumer's selection rather than approximating it, pagination is handled upstream (MAX_METADATA_PAGES 1000 x 100) so 110 objects are not truncated, and pipeline-core's own comment at file.py:1216 independently corroborates that `filename` carries no index. MEDIUM 1 — a polarity bug, and the worst kind. The parse sat OUTSIDE the try. If the store returned objects rather than dicts, or renamed its keys, `doc.get` raised, `seen` stayed empty, and EVERY object was reported not found — firing the full "NOTHING in this run is servable" refusal on a healthy delivery. A failure of the check must never quarantine the delivery; this module argues that twice, in C-99's and C-103's words, and then did the opposite. The parse is now inside the try. Mutation-proven: moving it back out fails the new test. MEDIUM 2 — silent fallback. `summary.get("uploaded_objects", [])` would have let the guard quietly shrink to the single historical object and log "preflight passed" while 108 shards, the sidecar and the manifest went unasked — the precise shape of the bug being fixed. The key is guaranteed on every path that reaches the guard (the interlock-closed path returns before it), so its absence is a defect and now reads as one (ADR-003). LOW 3 — `documents()` had no test anywhere. It is the only method the per-object check depends on, and the findability tests inject callables, so a pipeline-core rename would have surfaced during a live delivery. Now asserted at the port, which is the seam whose job is absorbing exactly that. LOW 4 — the port documented itself as four-method; `documents` makes it five. A datastore satisfying the documented contract would have failed with AttributeError mid-delivery. Also fixed, reported as not-a-finding: `unverified()` hardcoded the word "leg", so a scope rendered as "'forecast/sampled_forecast_shard objects' leg". It now takes the phrase. And the module docstring claimed "no store types" while verify reads two document KEYS — the boundary is now stated honestly rather than overclaimed. LOW 5 (SEARCH_INCOMPLETE noise growing ~110 documents per run) and LOW 6 (the budget-driven inline dict at the call site) are recorded on #312 rather than changed here: 5 needs a C-94 note and no code, and 6 costs lines neither partner package has. 522 passing. Both packages unchanged at 690 and 699. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- tests/test_findability.py | 41 ++++++++++++++++++++ tests/test_store_port.py | 27 +++++++++++++ views_postprocessing/crafd/managers/crafd.py | 2 +- views_postprocessing/crafd/store_port.py | 4 +- views_postprocessing/delivery/findability.py | 29 ++++++++------ views_postprocessing/unfao/managers/unfao.py | 2 +- views_postprocessing/unfao/store_port.py | 4 +- 7 files changed, 92 insertions(+), 17 deletions(-) diff --git a/tests/test_findability.py b/tests/test_findability.py index 9220578..5f26115 100644 --- a/tests/test_findability.py +++ b/tests/test_findability.py @@ -419,3 +419,44 @@ def test_the_sink_carries_its_upload_ledger_out(): "deliver_run no longer exports the upload ledger, so the #312 per-object " "preflight has nothing to check against" ) + + +def test_an_unreadable_store_answer_is_unverified_not_an_invisible_delivery(): + """#312 review, finding 1 — a polarity bug, and the worst kind. + + The parse used to sit OUTSIDE the try. If the store returned objects rather than + dicts, or renamed its keys, `doc.get` raised, `seen` stayed empty, and every object + was reported NOT FOUND — firing the full "NOTHING in this run is servable" refusal + on a perfectly healthy delivery. A failure of the CHECK must never quarantine the + DELIVERY; that is the whole of C-99 and C-103 and this module says so twice. + """ + class NotADict: + pass + + objects = _objects([f"s{i}.parquet" for i in range(3)]) + with pytest.raises(findability.FindabilityUnverifiedError) as caught: + findability.verify( + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, list_documents=lambda _f: [NotADict()], + ) + message = str(caught.value) + assert "UNVERIFIED, not known invisible" in message + assert "sampled_forecast_shard" in message, "and name the scope that could not be read" + + +def test_a_missing_upload_ledger_fails_loudly_rather_than_shrinking_the_check(): + """#312 review, finding 2 — the silent fallback. + + `summary.get("uploaded_objects", [])` would have let the per-object guard quietly + shrink to the single historical object and log "preflight passed", leaving the 108 + shards, the sidecar and the manifest unasked — the precise shape of the bug being + fixed. Declared access instead: the key is guaranteed on every path that reaches + the guard, so its absence is a defect and must read as one (ADR-003). + """ + for partner in PARTNER_PACKAGES: + source = (_REPO / "views_postprocessing" / partner / "managers" / f"{partner}.py").read_text() + assert 'summary.get("uploaded_objects"' not in source, ( + f"{partner} defaults the upload ledger; a sink that stopped exporting it " + "would silently reduce the guard to one object instead of failing" + ) + assert 'summary["uploaded_objects"]' in source diff --git a/tests/test_store_port.py b/tests/test_store_port.py index 46c82a4..c81a934 100644 --- a/tests/test_store_port.py +++ b/tests/test_store_port.py @@ -63,6 +63,8 @@ def __init__(self, result, downloaded=None): self.downloaded = downloaded self.calls = [] self.downloads = [] + self.metadata_queries = [] + self.documents_result = [{"filename": "a.parquet", "fileId": "id-a"}] def upload_data(self, **kwargs): self.calls.append(kwargs) @@ -72,6 +74,10 @@ def download_prediction(self, file_id): self.downloads.append(file_id) return self.downloaded + def get_predictions_by_metadata(self, filters=None): + self.metadata_queries.append(filters) + return self.documents_result + def _port(partner: str, result, downloaded=None): """The partner's port, wrapping a fake store. Needs no Appwrite environment.""" @@ -293,3 +299,24 @@ def test_the_two_partners_ports_have_not_drifted(): "twice. Apply it to both, or if the divergence is deliberate, say so in " "C-33 and replace this check with one that allows it." ) + + +@pytest.mark.parametrize("partner", PARTNER_PACKAGES) +def test_documents_forwards_the_filters_to_the_store_unchanged(partner): + """#312 review, finding 3: `documents()` had no test anywhere. + + It is the fifth method on the port and the only one the per-object findability + check depends on. The findability tests inject plain callables, so a rename or + signature change in pipeline-core's `get_predictions_by_metadata` would have + surfaced during a live delivery — after the upload had already happened. Absorbing + exactly that change is the port's job, so the port is where it must be asserted. + """ + port, store = _port(partner, None) + filters = {"name": "un_fao", "category": "forecast", "type": "sampled_forecast_shard"} + got = port.documents(filters) + + assert store.metadata_queries == [filters], ( + "documents() must forward the filters verbatim — the guard's whole premise is " + "that it asks the store the same question views-faoapi asks" + ) + assert got == store.documents_result diff --git a/views_postprocessing/crafd/managers/crafd.py b/views_postprocessing/crafd/managers/crafd.py index 299ce66..947a8af 100644 --- a/views_postprocessing/crafd/managers/crafd.py +++ b/views_postprocessing/crafd/managers/crafd.py @@ -387,7 +387,7 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, - [*summary.get("uploaded_objects", []), {"name": hist_path.name, + [*summary["uploaded_objects"], {"name": hist_path.name, "file_id": hist_file_id, "doc_type": "model", "category": "historical"}], ) else: diff --git a/views_postprocessing/crafd/store_port.py b/views_postprocessing/crafd/store_port.py index 8dc64ae..596b0d2 100644 --- a/views_postprocessing/crafd/store_port.py +++ b/views_postprocessing/crafd/store_port.py @@ -1,8 +1,8 @@ -"""The prediction store behind a four-method port — the DIP seam of ADR-013 epic #105. +"""The prediction store behind a five-method port — the DIP seam of ADR-013 epic #105. ``wire/source_selection`` and ``wire/sink`` drive the store through this object and never see the client's types. That is the whole point of the seam, so the constructor -takes **any** object carrying the four methods below rather than naming a concrete +takes **any** object carrying the five methods below rather than naming a concrete client class — the contract is the methods, not the type. Both refusals here are the same rule applied twice: *an unrecognised result should be diff --git a/views_postprocessing/delivery/findability.py b/views_postprocessing/delivery/findability.py index c26b203..d707664 100644 --- a/views_postprocessing/delivery/findability.py +++ b/views_postprocessing/delivery/findability.py @@ -2,8 +2,11 @@ consumer actually queries (register C-94). Representation-free — a lookup result and the declared identity it was looked up by. -No store types, no pandas, no frames. The caller performs the query (it owns the port); -the rule about what the answer means lives here. +No store *types*, no pandas, no frames; the caller owns the port and performs the query, +and the rule about what the answer means lives here. `verify` does read two document +**keys** (``filename``, ``fileId``) when recovering what landed, which is store +vocabulary rather than a store type — the honest boundary, stated because the sentence +above used to claim more than it delivered (#312 review). **The failure this exists for is invisible by construction.** Every upload reports success, storage is billed, and the consumer's endpoint returns empty. ADR-013 §4.1a @@ -40,7 +43,7 @@ class FindabilityUnverifiedError(RuntimeError): """The read-back could not be performed, so findability is unknown.""" -def unverified(category: str, exc: BaseException) -> FindabilityUnverifiedError: +def unverified(what: str, exc: BaseException) -> FindabilityUnverifiedError: """The refusal for *could not ask*, which is not *asked and got nothing*. Distinguished for the same reason ``source_metadata`` distinguishes a missing @@ -52,7 +55,7 @@ def unverified(category: str, exc: BaseException) -> FindabilityUnverifiedError: would be an outage manufactured by the guard. """ return FindabilityUnverifiedError( - f"the {category!r} leg uploaded, but the C-94 read-back could not be performed: " + f"{what} uploaded, but the C-94 read-back could not be performed: " f"{type(exc).__name__}: {exc}. The delivery is UNVERIFIED, not known invisible — " "re-run the check before quarantining anything." ) @@ -230,7 +233,7 @@ def verify( try: found = resolve_latest({"name": consumer_name, "category": category}) except Exception as exc: # could not ask != asked and got nothing (C-99, C-103) - raise unverified(category, exc) from exc + raise unverified(f"the {category!r} leg", exc) from exc assert_findable( found, expected_file_id=expected, consumer_name=consumer_name, category=category ) @@ -240,13 +243,17 @@ def verify( for category, doc_type in sorted(scopes): filters = {"name": consumer_name, "category": category, "type": doc_type} try: - docs = list_documents(filters) + for doc in list_documents(filters) or (): + filename, file_id = doc.get("filename"), doc.get("fileId") + if filename and file_id and filename not in seen: + seen[filename] = file_id except Exception as exc: - raise unverified(f"{category}/{doc_type} objects", exc) from exc - for doc in docs or (): - filename, file_id = doc.get("filename"), doc.get("fileId") - if filename and file_id and filename not in seen: - seen[filename] = file_id + # The PARSE is inside the try deliberately. If the store ever returns + # objects rather than dicts, or renames these keys, `doc.get` raises and + # `seen` is left empty — which would report EVERY object as missing and + # fire the "NOTHING is servable" refusal on a healthy delivery. A failure + # of the check must never quarantine the delivery (C-99, C-103). + raise unverified(f"{category}/{doc_type} artefacts", exc) from exc resolved = {o["name"]: (o["file_id"], seen.get(o["name"])) for o in objects} assert_all_findable(resolved, consumer_name=consumer_name) diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index b033ded..fd047d2 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -387,7 +387,7 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, - [*summary.get("uploaded_objects", []), {"name": hist_path.name, + [*summary["uploaded_objects"], {"name": hist_path.name, "file_id": hist_file_id, "doc_type": "model", "category": "historical"}], ) else: diff --git a/views_postprocessing/unfao/store_port.py b/views_postprocessing/unfao/store_port.py index 8dc64ae..596b0d2 100644 --- a/views_postprocessing/unfao/store_port.py +++ b/views_postprocessing/unfao/store_port.py @@ -1,8 +1,8 @@ -"""The prediction store behind a four-method port — the DIP seam of ADR-013 epic #105. +"""The prediction store behind a five-method port — the DIP seam of ADR-013 epic #105. ``wire/source_selection`` and ``wire/sink`` drive the store through this object and never see the client's types. That is the whole point of the seam, so the constructor -takes **any** object carrying the four methods below rather than naming a concrete +takes **any** object carrying the five methods below rather than naming a concrete client class — the contract is the methods, not the type. Both refusals here are the same rule applied twice: *an unrecognised result should be From 1d3cbf311472c63bcec29522e7af44784d507474 Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 22:02:51 +0200 Subject: [PATCH 5/7] fix(delivery): drop the ordering assumption instead of documenting it (#312 re-review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The views-models seat could not find the `$createdAt` sort I had twice asserted as fact. They were right and I was wrong, and the way I was wrong is the worse part: an automated review asserted it, I relayed it to them twice without checking, and it became a rationale I was about to put in the code. Verified here. `search_files_by_metadata` builds `queries = []` and appends only `Query.equal(attribute, value)` per filter, then `limit` + `offset` (pipeline-core `modules/appwrite/file.py:1045-1050`). No `order_desc`, no `order_asc`, no sort field. The only ordering calls in that module are at `2643-2645`, inside `list_files` — a STORAGE listing, not a collection query — where `order_field` defaults to None. So document order is unspecified and "first match wins" was a coin flip dressed as a tie-break. Rather than record the caveat, the assumption is gone. `seen` now maps filename to the SET of ids carried under it, and an object is findable when the id THIS run uploaded is among them. That is the question we actually have, and it needs no ordering guarantee at all. The refusal now shows every id the store holds under a name instead of one, which is strictly more useful for an operator auditing a bucket. Two tests: the same delivery passes under three different document orders including a stale same-filename document returned first, and order-independence does not become permissiveness — a name carrying only other runs' ids is still refused. Mutation-proven by restoring first-wins, which fails both. NOT fixed here, because they are not ours, and both are now on #312: pipeline-core's `get_latest_file_id` DOCUMENTS "the newest matching file based on creation timestamp" and implements `files_list[0]` over that unsorted result — a stated guarantee with no mechanism, and the C-94 leg check has depended on it since August. views-faoapi's `resolve_artifact_file_ids` carries the same "Newest-first" assumption in its docstring. 524 passing. Both partner packages unchanged at 690 and 699. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- tests/test_findability.py | 41 ++++++++++++++++++++ views_postprocessing/delivery/findability.py | 32 +++++++++------ 2 files changed, 62 insertions(+), 11 deletions(-) diff --git a/tests/test_findability.py b/tests/test_findability.py index 5f26115..dc4863f 100644 --- a/tests/test_findability.py +++ b/tests/test_findability.py @@ -460,3 +460,44 @@ def test_a_missing_upload_ledger_fails_loudly_rather_than_shrinking_the_check(): "would silently reduce the guard to one object instead of failing" ) assert 'summary["uploaded_objects"]' in source + + +def test_a_correct_delivery_passes_whatever_order_the_store_returns(): + """#312 re-review: document order is unspecified, so the check must not depend on it. + + `search_files_by_metadata` appends only `Query.equal` per filter and never an + `order_desc`/`order_asc` (pipeline-core `modules/appwrite/file.py:1045-1050`). A + "first match wins" rule would have been a coin flip presented as a tie-break. The + question we actually have — is the id THIS run uploaded present under the name — + needs no ordering at all, so the same delivery must pass in any permutation. + """ + objects = _objects(["a.parquet", "b.parquet"]) + docs = [{"filename": o["name"], "fileId": o["file_id"]} for o in objects] + # a stale document from a previous run, sharing a filename, returned FIRST + stale = [{"filename": "a.parquet", "fileId": "id-from-run1"}] + + for order in ([*stale, *docs], [*docs, *stale], [docs[1], stale[0], docs[0]]): + findability.verify( + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, list_documents=lambda _f, o=order: o, + ) + + +def test_a_name_carrying_only_another_runs_id_is_still_refused(): + """Order-independence must not become permissiveness. If the only documents under + a filename belong to some other run, this run's object did not land.""" + objects = _objects(["a.parquet"]) + with pytest.raises(findability.DeliveryNotFindableError) as caught: + findability.verify( + consumer_name="un_fao", legs={}, objects=objects, + resolve_latest=_never, + list_documents=lambda _f: [ + {"filename": "a.parquet", "fileId": "id-run1"}, + {"filename": "a.parquet", "fileId": "id-run0"}, + ], + ) + message = str(caught.value) + assert "WRONG DOCUMENT" in message + assert "id-run0" in message and "id-run1" in message, ( + "the refusal must show every id the store holds under that name, not one" + ) diff --git a/views_postprocessing/delivery/findability.py b/views_postprocessing/delivery/findability.py index d707664..91300a1 100644 --- a/views_postprocessing/delivery/findability.py +++ b/views_postprocessing/delivery/findability.py @@ -118,11 +118,14 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: """Raise unless EVERY artefact this run uploaded resolves under its own filename. Args: - resolved: ``{filename: (expected_file_id, found_file_id)}`` — one entry per - object the run uploaded. ``found_file_id`` is the id carried by the store - document whose ``filename`` field matches, as recovered by `verify` from a - type-scoped query (**not** a query on ``filename`` — see `verify`). ``None`` - means no document carried that filename. + resolved: ``{filename: (expected_file_id, found_file_ids)}`` — one entry per + object the run uploaded. ``found_file_ids`` is the **set** of ids carried by + store documents whose ``filename`` matches, as recovered by `verify` from a + type-scoped query (**not** a query on ``filename`` — see `verify`). A set + rather than one id because document order is unspecified, so "the first + match" is not a defined thing; an empty set means no document carried that + filename. The delivery is findable when the id THIS run uploaded is among + them, which needs no ordering guarantee. consumer_name: the DECLARED store-document ``name`` (C-77). Raises: @@ -153,9 +156,9 @@ def assert_all_findable(resolved: dict, *, consumer_name: str) -> None: """ missing = sorted(name for name, (_, found) in resolved.items() if not found) wrong = sorted( - f"{name} (expected {exp!r}, store has {found!r})" + f"{name} (expected {exp!r}, store has {sorted(found)})" for name, (exp, found) in resolved.items() - if found and found != exp + if found and exp not in found ) if not missing and not wrong: return @@ -239,14 +242,21 @@ def verify( ) scopes = {(o["category"], o["doc_type"]) for o in objects} - seen: dict[str, str] = {} + # filename -> EVERY id carried under it, because document order is unspecified. + # `search_files_by_metadata` appends only `Query.equal` per filter and never an + # `order_desc`/`order_asc` (pipeline-core `modules/appwrite/file.py:1045-1050`); + # the only ordering calls in that module are inside the storage `list_files`. + # So "first wins" would have been a coin flip dressed as a tie-break. Asking + # whether the id THIS run uploaded is present under the name needs no order at + # all, and is the question we actually have (#312 review, views-models seat). + seen: dict[str, set] = {} for category, doc_type in sorted(scopes): filters = {"name": consumer_name, "category": category, "type": doc_type} try: for doc in list_documents(filters) or (): filename, file_id = doc.get("filename"), doc.get("fileId") - if filename and file_id and filename not in seen: - seen[filename] = file_id + if filename and file_id: + seen.setdefault(filename, set()).add(file_id) except Exception as exc: # The PARSE is inside the try deliberately. If the store ever returns # objects rather than dicts, or renames these keys, `doc.get` raises and @@ -255,7 +265,7 @@ def verify( # of the check must never quarantine the delivery (C-99, C-103). raise unverified(f"{category}/{doc_type} artefacts", exc) from exc - resolved = {o["name"]: (o["file_id"], seen.get(o["name"])) for o in objects} + resolved = {o["name"]: (o["file_id"], seen.get(o["name"], set())) for o in objects} assert_all_findable(resolved, consumer_name=consumer_name) logger.info( "Findability preflight passed: %d leg(s) and %d object(s) resolvable under %r.", From 45c2984526d58d7ebd8799ab2f9337d89d65f7b8 Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 22:08:35 +0200 Subject: [PATCH 6/7] fix(port): the documented datastore contract omitted the method documents() calls MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Follow-up to #312. /code-review's finding 4 said the port documented itself as four-method when `documents` made it five. I fixed the module header and missed the thing that mattered: the CLASS docstring also enumerates the datastore methods a caller must supply, and `get_predictions_by_metadata` was absent from it. Measured rather than reasoned about — a double built to the documented contract: class DocumentedContract: def get_latest_file_id(self, filters): ... def get_file_metadata(self, fid): ... def download_prediction(self, fid): ... def upload_data(self, **kw): ... AttributeError: 'DocumentedContract' object has no attribute 'get_predictions_by_metadata' Mid-delivery, after the upload has happened — which is the one place this port exists to prevent a surprise. Found while verifying, for this repo rather than taking it on trust, that a pipeline-core fix to `get_latest_file_id`'s missing sort would actually reach us. It does: this repo vendors no copy of the Appwrite client, calls `search_files_by_metadata` nowhere, and the port takes any duck-typed datastore with `views_pipeline_core.modules.datastore.DatastoreModule` passed in production. So the ordering fix reaches us with a pin bump and no change here. faoapi and crafdapi each carry their own copy of that module and are not covered. The guard is derived from the source, not a hardcoded list: it collects every `self._dsm.(` the port calls and asserts each appears in the module or class docstring. Prose cannot drift from the code without failing, which is how this one rotted. Mutation-proven by removing the fifth name, which fails both partners. ADR-014 §2 — the guard also asserts it found calls at all, so a refactor that renames `_dsm` cannot leave it silently scanning nothing. 526 passing. Both partner packages unchanged at 690 and 699. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- tests/test_store_port.py | 37 ++++++++++++++++++++++++ views_postprocessing/crafd/store_port.py | 2 +- views_postprocessing/unfao/store_port.py | 2 +- 3 files changed, 39 insertions(+), 2 deletions(-) diff --git a/tests/test_store_port.py b/tests/test_store_port.py index c81a934..a0e4d41 100644 --- a/tests/test_store_port.py +++ b/tests/test_store_port.py @@ -24,6 +24,8 @@ from dataclasses import dataclass +import ast +import re import pytest from pathlib import Path @@ -320,3 +322,38 @@ def test_documents_forwards_the_filters_to_the_store_unchanged(partner): "that it asks the store the same question views-faoapi asks" ) assert got == store.documents_result + + +@pytest.mark.parametrize("partner", PARTNER_PACKAGES) +def test_the_documented_datastore_contract_lists_every_method_the_port_calls(partner): + """A double satisfying the docstring must actually work. + + `/code-review` caught that the module header said "four-method port" when + `documents` made it five. Fixing the header missed the thing that matters: the + class docstring also enumerates the datastore methods a caller must supply, and + `get_predictions_by_metadata` was absent from it. A test double built to the + documented contract raised AttributeError — mid-delivery, after the upload. + + Derived from the source rather than listed here, so the assertion cannot rot the + way the prose did (ADR-014 §2: prove the guard's inputs are real). + """ + source = (_PKG / partner / "store_port.py").read_text() + called = set(re.findall(r"self\._dsm\.(\w+)\(", source)) + assert called, "found no datastore calls — this guard is scanning the wrong thing" + + tree = ast.parse(source) + doc = "\n".join( + [ast.get_docstring(tree) or ""] + + [ + ast.get_docstring(n) or "" + for n in ast.walk(tree) + if isinstance(n, ast.ClassDef) + ] + ) + documented = set(re.findall(r"``(\w+)``", doc)) + missing = sorted(called - documented) + assert not missing, ( + f"{partner}/store_port.py calls {missing} on the datastore but its class " + "docstring does not list them. A caller building to the documented contract " + "gets an AttributeError during a delivery, after the upload has happened." + ) diff --git a/views_postprocessing/crafd/store_port.py b/views_postprocessing/crafd/store_port.py index 596b0d2..a13e4cb 100644 --- a/views_postprocessing/crafd/store_port.py +++ b/views_postprocessing/crafd/store_port.py @@ -18,7 +18,7 @@ class _ContractStorePort: """Adapts a prediction-store client to the wire ports. ``datastore`` is any object exposing ``get_latest_file_id``, ``get_file_metadata``, - ``download_prediction`` and ``upload_data``. + ``download_prediction``, ``upload_data`` and ``get_predictions_by_metadata``. """ def __init__(self, datastore) -> None: diff --git a/views_postprocessing/unfao/store_port.py b/views_postprocessing/unfao/store_port.py index 596b0d2..a13e4cb 100644 --- a/views_postprocessing/unfao/store_port.py +++ b/views_postprocessing/unfao/store_port.py @@ -18,7 +18,7 @@ class _ContractStorePort: """Adapts a prediction-store client to the wire ports. ``datastore`` is any object exposing ``get_latest_file_id``, ``get_file_metadata``, - ``download_prediction`` and ``upload_data``. + ``download_prediction``, ``upload_data`` and ``get_predictions_by_metadata``. """ def __init__(self, datastore) -> None: From e8f2b4af66db0c8faedbd0ef13fc8ec2d81c2964 Mon Sep 17 00:00:00 2001 From: Polichinl Date: Tue, 29 Sep 2026 22:20:31 +0200 Subject: [PATCH 7/7] =?UTF-8?q?release:=201.4.0=20=E2=80=94=20the=20findab?= =?UTF-8?q?ility=20fix,=20which=20is=20unreachable=20until=20it=20is=20tag?= =?UTF-8?q?ged?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bumps 1.3.0 -> 1.4.0. MINOR rather than PATCH by the same reasoning as 1.2.0: a previously-passing delivery can now stop. The guard refuses in four situations instead of two. That is a behaviour change a launcher must be told about, not a patch. Cut now rather than at convenience because the launcher installs from a git TAG (views-models tools/launcher/postprocessor.sh:57), not from PyPI and not from poetry.lock. Until this tag exists, the fix for today's unservable delivery is merged, reviewed, tested — and reachable by nothing. 1.2.0 sat unpinned for three weeks for exactly this reason. The changelog entry leads with the consequence for someone still on an older pin: a delivery can complete successfully and be unservable. Then what changed, then why taking it promptly is worth it — the deterministic-re-run case, where the dedup path that took one object takes all 110 and the documented remedy for a torn run is what triggers it. Recorded there as known-and-not-ours: pipeline-core's get_latest_file_id documents "newest by creation timestamp" over an unsorted result, which the SELECTION half of the guard has relied on since August. This release removes the equivalent assumption from the per-object half. The upstream half is filed there. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01ANY1CCy9Xo7zjMY4XJ69v9 --- CHANGELOG.md | 58 ++++++++++++++++++++++++++++++++++++++++++++++++++ pyproject.toml | 2 +- 2 files changed, 59 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 021c90b..ade3860 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,64 @@ This file exists because the version number was the only signal a consumer got (register C-111). Releases before 1.2.0 are summarised from their tags rather than reconstructed in detail. +## 1.4.0 — 2026-09-29 + +**Fixes an unservable delivery.** The first live UN-FAO run, earlier today, uploaded +109 of 110 objects, reported success, and was refused by views-faoapi at ingest. If you +are pinned below this version, a delivery can still complete successfully and be +unservable. + +### What was wrong + +The C-94 findability guard checked **two** things — the run manifest and the historical +artifact — by querying the newest document per *category*. Every wire object carries +`category="forecast"` and the manifest is uploaded last, so that query always returned +the manifest. The GAUL sidecar and the 108 shards were never asked about. + +The sidecar's bytes were identical to the previous run's, the content-addressed store +correctly declined a second copy, the metadata document was updated against the **old** +file, and the upload returned a real file id for the wrong document. Nothing observed it. + +### What changed for a consumer + +The guard now asks two questions instead of one, and **refuses in four situations rather +than two**. A delivery that previously completed can now stop: + +| | | +|---|---| +| *selection* | does the consumer's own query land on this run? — unchanged | +| *per-object* | does **every** uploaded artefact resolve under its own filename? — new | + +Both raise the existing `DeliveryNotFindableError`; no new exception types. A failed +*check* still reports `FindabilityUnverifiedError` rather than condemning the delivery. + +**If one of these fires after upgrading it is reporting a condition that was already +wrong and already invisible.** The refusal names every object that does not resolve, and +says so explicitly when nothing in the run is servable. + +### Why this is worth taking promptly + +Pooling upstream is deterministic, and every artefact's filename embeds the run id. So a +**re-run** writes new filenames over identical bytes, and the deduplication path that +took one object takes **all 110 at once** — the store creates the documents, the count is +right, and nothing is servable. Re-running is the documented remedy for a torn run, which +makes this the realistic case rather than the exotic one. + +### Also in this release + +- The store port gained a fifth method, `documents()`, and its documented duck-typing + contract now lists it. A datastore built to the previous docstring would have raised + `AttributeError` mid-delivery. +- `deliver_run`'s summary carries the upload ledger out, as `uploaded_objects`. + +### Known, and not ours + +`views-pipeline-core`'s `get_latest_file_id` documents *"the newest matching file based on +creation timestamp"* and takes the first element of an unsorted result. The *selection* +half of the guard has relied on that since August and can in principle raise a false +alarm. This release removes the equivalent assumption from the per-object half, which no +longer depends on document order at all. The upstream half is filed in pipeline-core. + ## 1.3.0 — 2026-09-19 **No new failure modes.** A launcher that ran 1.2.0 sees nothing new stop. This release diff --git a/pyproject.toml b/pyproject.toml index 04cc622..599d77a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "views-postprocessing" -version = "1.3.0" +version = "1.4.0" description = "" authors = [ "Dylan Pinheiro ",