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/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/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 ", diff --git a/tests/test_findability.py b/tests/test_findability.py index a6272ca..dc4863f 100644 --- a/tests/test_findability.py +++ b/tests/test_findability.py @@ -187,14 +187,36 @@ 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_latest=explodes, list_documents=explodes, + ) + + # 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=[{"name": "run__sidecar.parquet", "file_id": "a", + "doc_type": "sampled_forecast_sidecar", "category": "forecast"}], + resolve_latest=explodes, list_documents=explodes, + ) + assert "sampled_forecast_sidecar" in str(caught.value), ( + "an unverified scope must name which artefact type could not be checked" ) @@ -222,3 +244,260 @@ 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 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 _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(): + """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 = _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={}, 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 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_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 + selection check is satisfied by the very run that is unservable. + """ + findability.verify( + 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 + 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 = _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_latest=_never, list_documents=_store({}), + ) + 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(): + """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=objects, + resolve_latest=_never, + list_documents=_store({"run2__sidecar.parquet": "id-from-run1"}), + ) + message = str(caught.value) + assert "WRONG DOCUMENT" 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 = _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": "id-m.json"}, objects=objects, + resolve_latest=lambda _f: "id-m.json", list_documents=_store(landed), + ) + + +@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" + ) + + +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 + + +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/tests/test_store_port.py b/tests/test_store_port.py index 46c82a4..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 @@ -63,6 +65,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 +76,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 +301,59 @@ 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 + + +@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/contract/wire/sink.py b/views_postprocessing/contract/wire/sink.py index 3c88b59..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 ) @@ -251,5 +251,12 @@ 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, "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 d5f9982..947a8af 100644 --- a/views_postprocessing/crafd/managers/crafd.py +++ b/views_postprocessing/crafd/managers/crafd.py @@ -124,23 +124,16 @@ 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. - - 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. - """ - 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) +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)) + findability.verify( + consumer_name=consumer_name, legs=legs, objects=objects, + resolve_latest=port.latest_file_id, list_documents=port.documents, + ) class CRAFDPostProcessorManager(PostprocessorManager, ForecastingModelManager): @@ -394,6 +387,8 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, + [*summary["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..a13e4cb 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 @@ -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: @@ -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 508dea3..91300a1 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 @@ -27,6 +30,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.""" @@ -36,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 @@ -48,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." ) @@ -105,3 +112,162 @@ 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_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: + 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 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 + 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 {sorted(found)})" + for name, (exp, found) in resolved.items() + if found and exp not in found + ) + 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"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." + ) + 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: list, resolve_latest, list_documents +) -> 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: 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. + + **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 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_latest({"name": consumer_name, "category": category}) + except Exception as exc: # could not ask != asked and got nothing (C-99, C-103) + raise unverified(f"the {category!r} leg", exc) from exc + assert_findable( + found, expected_file_id=expected, consumer_name=consumer_name, category=category + ) + + scopes = {(o["category"], o["doc_type"]) for o in objects} + # 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: + 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 + # `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"], 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.", + len(legs), len(objects), consumer_name, + ) diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index 8720230..fd047d2 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -124,23 +124,16 @@ 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. - - 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. - """ - 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) +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)) + findability.verify( + consumer_name=consumer_name, legs=legs, objects=objects, + resolve_latest=port.latest_file_id, list_documents=port.documents, + ) class UNFAOPostProcessorManager(PostprocessorManager, ForecastingModelManager): @@ -394,6 +387,8 @@ def _save_contract(self) -> dict: self._model_path, product.CONSUMER_DOCUMENT_NAME, {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, + [*summary["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..a13e4cb 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 @@ -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: @@ -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))