diff --git a/docs/CICs/UNFAOPostProcessorManager.md b/docs/CICs/UNFAOPostProcessorManager.md index 15cab32..cb6da57 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-17 +**Last reviewed:** 2026-08-18 **Related ADRs:** ADR-001, ADR-002, ADR-008, ADR-009 --- @@ -90,10 +90,11 @@ 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` +- **Delivery invisible to the consumer:** raises `delivery.findability.DeliveryNotFindableError` (C-94, added 2026-08-18). After both legs are uploaded, the manager queries the partner store as the consumer does — `name == product.CONSUMER_DOCUMENT_NAME`, per category — and refuses a falsy answer. This is the one failure mode where every upload reports success and the consumer still sees nothing; run-0's historical leg stranded exactly that way (C-79). It runs only inside the §11.4 interlock, and queries through a store with pipeline-core's automatic `name == model_name` filter suppressed, so it verifies the declared name rather than the views-models directory name that happens to match (C-77). **It does not detect a delivery that never ran, or stale data served from the consumer's cache** — both recorded as gaps in C-94 - **Wrong forecast selected:** structurally impossible since #149. Selection is by **run manifest** — a commit marker whose contents are hash-verified — not by scanning the bucket for the newest `category="forecast"` upload. Declared identity is additionally checked **per shard header** against the launched ensemble inside `TargetLease.load()` (`contract/wire/source_selection.py:73-81`), so identity comes from the artifact's own content. The metadata-field check this bullet used to describe (`delivery/identity.py`) was retired in #150 and the legacy reader it served in #149; register C-25 is closed as *superseded by mechanism* - **Launch config incomplete:** raises `LaunchConfigError` naming the missing key. A launcher that omits `wire_contract` or declares a `data_format` other than `feature_frame` is **refused**, never quietly routed into a fallback (ADR-003, register C-63) - **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) -- **Fabricated historical tail:** `_read_historical_frame()` drops months beyond the producer's `last_valid_month_id` at the read (`_clip_observed_history` was the pandas equivalent, retired with that path in #149) so unobserved zero-padding is not shipped as observed history (S2/C-26). Two outcomes when the boundary is unavailable, and they are different on purpose (C-103, 2026-08-17): if the producer simply publishes no boundary — or the read fails — it **degrades open**, skipping the clip with a WARNING that states the unobserved tail will ship; if the producer client cannot be imported at all it **refuses** (`source_metadata.ProducerClientMissing`), because a broken environment is not a producer fact +- **Fabricated historical tail:** `_read_historical_frame()` drops months beyond the producer's `last_valid_month_id` at the read (`_clip_observed_history` was the pandas equivalent, retired with that path in #149) so unobserved zero-padding is not shipped as observed history (S2/C-26). Two outcomes when the boundary is unavailable, and they are different on purpose (C-103, 2026-08-17): if the producer simply publishes no boundary — or the read fails — it **degrades open**, skipping the clip with a WARNING that states the unobserved tail will ship; if the producer client cannot be imported at all it **refuses** (`source_metadata.ProducerClientUnavailable`), because a broken environment is not a producer fact - **Upload provenance:** the historical artifact's `description` carries structured provenance (lookup version, region, expected/actual cell counts, unmapped count) built by `delivery/provenance.py` (`build_provenance` → `compact_description`) via the manager's `_historical_frame_description()` (S5/C-15). The **forecast** side carries no such description: its guarantee is the wire's verified chain — per-shard content hashes recorded in the §4.2 run manifest, header asserts on load, and manifest-last commit ordering. That is identity and integrity, not the C-15 provenance field set; the §4.2 manifest's keys are exactly `contract_version`, `run_id`, `targets`, `shards`, `expected_months`, `expected_cell_count`, `sidecar` — and it carries **no** `lookup_version`, `region` or `unmapped_count`. `_delivery_description()` was the pandas-path equivalent and was deleted with it in #149 The following **must never** fail silently: diff --git a/reports/technical_risk_register.md b/reports/technical_risk_register.md index 0db40c2..6aba7a5 100644 --- a/reports/technical_risk_register.md +++ b/reports/technical_risk_register.md @@ -275,7 +275,7 @@ Cross-refs: **C-94**, **C-96**, þing-01 `orð_dómr.md` D2, issue #249. | ID | C-94 | | Tier | 2 — the failure mode is invisible by construction and lands on the live FAO path: upload succeeds, storage is billed, the consumer's endpoint returns empty, nothing raises anywhere. ADR-013 §4.1a's *"invisible to the consumer, not merely degraded."* | | Source | `/expert-code-review` of the standing decisions, 2026-08-12 | -| Trigger | **Re-specified twice on 2026-08-13; the first attempt was not exclusive and its own worked example matched two arms.** (a) A delivery is reported empty **and an upload occurred after the bucket reached the state under investigation** — that is what the preflight below would catch, and the time bound is what the first attempt omitted. (b) `APPWRITE_READ_API_KEY` is provisioned, at which point the deferral has no remaining cost. *(A third arm — "reported empty with no upload since" — was drafted and withdrawn: it is not observable from this repository, which the amendment says four lines on, and ADR-014 §4 requires a trigger someone can notice. It is a gap, and is stated as one below rather than dressed as a trigger.)* | +| Trigger | **Re-specified twice on 2026-08-13; the first attempt was not exclusive and its own worked example matched two arms.** (a) A delivery is reported empty **and an upload occurred after the bucket reached the state under investigation** — that is what the preflight below would catch, and the time bound is what the first attempt omitted. (b) `APPWRITE_READ_API_KEY` is provisioned, at which point the deferral has no remaining cost. *(A third arm — "reported empty with no upload since" — was drafted and withdrawn: it is not observable from this repository, which the amendment says four lines on, and ADR-014 §4 requires a trigger someone can notice. It is a gap, and is stated as one below rather than dressed as a trigger.)* **Rewritten 2026-08-18, because arm (b) expired without firing** (register conventions: a trigger whose event has already occurred reads identically to a pending one). The preflight was built *without* `APPWRITE_READ_API_KEY`, so "it is provisioned" can no longer arm anything. What remains live is arm (a), now narrowed: **a delivery is reported empty, an upload occurred since, and the preflight did NOT raise** — that combination means the check is looking in the wrong place, and it is the only arm this repository can still be surprised by. | | Owner | This repository, for the mechanism. The credential is the operator's. | | Location | `views_postprocessing/contract/wire/sink.py` (the upload path, where nothing verifies); `views_postprocessing/delivery/`. | @@ -293,6 +293,21 @@ FAO emailed at **09:15 UTC** that `faoapi.viewsforecasting.org` returned no data **So the trigger was mis-specified, not the mechanism.** "A delivery is reported empty" names a symptom with at least two causes, and this entry's preflight addresses only one of them. +**Partial mitigation, 2026-08-18 — the preflight is built.** After both legs are uploaded, each manager asks the partner store the question the consumer asks — `name == product.CONSUMER_DOCUMENT_NAME`, `category ∈ {forecast, historical}` — and refuses a falsy answer (`delivery/findability.py`, `DeliveryNotFindableError`). The two legs are checked separately on purpose: a run whose forecast landed and whose historical did not is invisible in exactly one half, and the historical leg is the one that actually stranded in run-0 (C-79). + +**Two decisions inside it that a later reader should not have to re-derive:** + +1. **It runs on the existing key, not the registry's `APPWRITE_READ_API_KEY` slot.** Verified in the Appwrite console 2026-08-18: the live `VIEWS Pipeline Core` key already carries `documents.read`, `rows.read`, `buckets.read` and `files.read`. A separate read credential would buy no isolation here, because the preflight runs *inside the delivery process*, which already holds the write key it just uploaded with. C-96's permission is about the operation being read-only, and it is. The registry slot stays `planned` for a preflight that runs **outside** the delivery, where the isolation would be real. +2. **It queries through a store with pipeline-core's automatic `name == model_name` filter suppressed** (`_build_partner_read_store`). `get_latest_file_id` delegates to `get_predictions_by_metadata`, which merges the path manager's model name into every query — so without the suppression the check would verify the views-models *directory* name, which equals the declared consumer name only by coincidence (**C-77**). Verifying the coincidence rather than the contract would leave this green while a rename took the delivery dark, which is the precise failure it exists to see. + +**The read-back is scoped to THIS run, and that is the whole guard.** The first implementation asked *"is there any document under the consumer's name for this category"* — a question the **previous** delivery already answers yes to. Documents accumulate across runs (that is what makes "latest" meaningful to the consumer), so from delivery 2 onward the check could never fail: run-2's upload reports success, its metadata document is never created — the exact C-79 shape — the query returns run-1's document, and the preflight logs *"passed"* while the consumer goes on serving run-1. Caught by `/code-review high` before merge. `_ContractStorePort.upload` now returns the uploaded `file_id` (it was discarding it), the sink carries the manifest's id out — it is uploaded last, so it is the newest `category="forecast"` document — and the check asserts the newest document the consumer would find **is the one this run put there**. The refusal distinguishes "nothing found" from "found the previous run's", because those are different operator situations. + +**A failed read-back is not an invisible delivery.** `findability.unverified` names that separately (`FindabilityUnverifiedError`): a store error after a successful upload means the delivery is UNVERIFIED, not known invisible, and quarantining on it would be an outage the guard manufactured. Same distinction C-103 draws between a missing producer client and a producer that publishes no boundary, and C-99 between an unrecognised store result and a real one — three instances now of the same rule, that *could not ask* and *asked and got nothing* call for different operator actions. + +**A note on which pipeline-core you read, because it changed a review's conclusion.** The same review reported that `unverified()` was dead code, on the grounds that `get_predictions_by_metadata` swallows a failed search and returns `[]` — so a store error would arrive as `None` and be reported as an invisible delivery. That is true of **2.3.0**, which is what the drifted developer venv holds (C-104). It is false of **3.0.1**, which `poetry.lock` pins and CI installs: there the method **raises `MetadataSearchIncomplete`**, with a comment in pipeline-core saying why — *"Returning [] here would tell every caller 'no predictions match', which is a statement about the shelf rather than about the lookup… a false negative to an external counterparty"* (views-pipeline-core C-241, its Cluster J). Verified 2026-08-19 by reading `3.0.1` from the views-pipeline-core checkout rather than the installed package. So the split holds where it runs. **This is C-104's hazard in its most expensive form yet**: not a wall of red, but a confident and wrong conclusion about production drawn from a stale environment. + +**The tier does not move, and the reason it does not is the point.** The Tier 2 rationale was *"upload succeeds, storage is billed, the consumer's endpoint returns empty, nothing raises anywhere."* For that cause, something now raises. What holds the entry at 2 is the two causes below, which this does not touch and which remain invisible — the tier now rests on the gaps rather than on the mechanism. + **Two uncovered causes, stated as gaps rather than dressed as triggers.** Neither is observable from here, so neither can be a trigger under ADR-014 §4 — a trigger nobody can notice is a wish: 1. *The bucket is empty because nothing was delivered.* This repository is not told when a delivery is due and has no view of whether the last one is still present. That is the 2026-08-12 case. @@ -1081,7 +1096,7 @@ So the scenario this entry was filed on — a missing dependency silently shippi **What is genuinely left, and it is why this stays open at Tier 3.** The `except Exception` was never only about the import. A network failure, an auth error, a reshaped `.zattrs`, a timeout — all still leave through one branch, log one WARNING, and deliver the unobserved tail as observed history. That is the recorded C-26 degrade-open decision applied far more broadly than C-26 argued for. -**Partial mitigation, 2026-08-17.** `source_metadata` now raises `ProducerClientMissing` instead of letting the import failure fall into the caller's broad `except`, and both managers re-raise it rather than degrading. Defence in depth for the paths `assert_queryset_was_importable` does not cover — a direct caller, a future launcher, a partner not going through the same queryset. The module also gained its first tests (`tests/test_source_metadata.py`, 7, mutation-proven against both the pre-fix return-`None` and a manager that collapses the two branches back into one); it had **none** before, which is how "return None like everything else" ever looked reasonable. The broad degrade-open is deliberately unchanged — narrowing it is a decision about what to tell the partner, not a refactor. +**Partial mitigation, 2026-08-17.** `source_metadata` now raises `ProducerClientUnavailable` instead of letting the import failure fall into the caller's broad `except`, and both managers re-raise it rather than degrading. Defence in depth for the paths `assert_queryset_was_importable` does not cover — a direct caller, a future launcher, a partner not going through the same queryset. The module also gained its first tests (`tests/test_source_metadata.py`, 7, mutation-proven against both the pre-fix return-`None` and a manager that collapses the two branches back into one); it had **none** before, which is how "return None like everything else" ever looked reasonable. The broad degrade-open is deliberately unchanged — narrowing it is a decision about what to tell the partner, not a refactor. **C-60 is this shape, and it was resolved by deleting the degradation.** There, a provenance stamp reached into the producer's ledger schema inside a bare `except … pass` and returned `"unknown"`; the fix was to raise. The difference here is that the degradation is deliberate and documented ("degrade-open, C-26") — which makes the question *whether the open side is still the right one*, not whether someone forgot. diff --git a/tests/test_clone_readiness.py b/tests/test_clone_readiness.py index 79413fe..54784d2 100644 --- a/tests/test_clone_readiness.py +++ b/tests/test_clone_readiness.py @@ -49,6 +49,7 @@ "views_postprocessing.delivery.coverage", "views_postprocessing.delivery.draws", "views_postprocessing.delivery.parity", + "views_postprocessing.delivery.findability", "views_postprocessing.delivery.observed_range", "views_postprocessing.delivery.provenance", "views_postprocessing.contract.wire.sink", diff --git a/tests/test_findability.py b/tests/test_findability.py new file mode 100644 index 0000000..a6272ca --- /dev/null +++ b/tests/test_findability.py @@ -0,0 +1,224 @@ +"""The C-94 findability preflight: does the consumer's own query find the delivery? + +Two layers, matching how the repo tests every other delivery invariant: the rule on +primitives here, and the wiring — that the manager actually calls it, and calls it +through a store whose injected name filter is suppressed — as declaration checks. + +The wiring checks are source reads rather than behavioural ones. Constructing a manager +needs pipeline-core, a views-models path manager and a live Appwrite environment (C-40), +and the properties worth holding are two lines: that the preflight runs only when +something was uploaded, and that it asks under the DECLARED consumer name rather than +the path manager's (C-77). +""" + +from __future__ import annotations + +import ast +from pathlib import Path + +import pytest + +from tests.conftest import PARTNER_PACKAGES +from views_postprocessing.delivery import findability + +_REPO = Path(__file__).resolve().parent.parent + + +def _calls_body(body) -> set[str]: + """Names called anywhere inside a list of statements.""" + found: set[str] = set() + for stmt in body: + for c in ast.walk(stmt): + if isinstance(c, ast.Call) and isinstance(c.func, ast.Name): + found.add(c.func.id) + return found + + +def _function_source(source: str, name: str) -> str: + """Exactly one function's source. + + Slicing to end-of-file instead would let any later occurrence in the module satisfy + these checks — the assertion would pass for text that is not in the function at all. + """ + fn = next( + n for n in ast.walk(ast.parse(source)) + if isinstance(n, ast.FunctionDef) and n.name == name + ) + return "\n".join(source.splitlines()[fn.lineno - 1:fn.end_lineno]) + + +def test_a_found_document_passes(): + findability.assert_findable( + "file-abc", expected_file_id="file-abc", consumer_name="un_fao", category="forecast" + ) + + +def test_nothing_found_is_refused(): + with pytest.raises(findability.DeliveryNotFindableError): + findability.assert_findable( + None, expected_file_id="x", consumer_name="un_fao", category="forecast" + ) + + +def test_the_refusal_names_the_query_that_found_nothing(): + with pytest.raises(findability.DeliveryNotFindableError) as excinfo: + findability.assert_findable( + None, expected_file_id="x", consumer_name="un_fao", category="historical" + ) + message = str(excinfo.value) + assert "un_fao" in message, "the refusal must name the consumer name it queried by" + assert "historical" in message, "and which leg was invisible" + assert "INVISIBLE" in message, ( + "the message must say the delivery is invisible rather than degraded — that " + "distinction is ADR-013 §4.1a and it is what tells an operator to quarantine" + ) + + +def test_an_empty_string_file_id_is_not_treated_as_found(): + """`None` is the documented 'no match', but a store returning '' is not a find. + + pipeline-core's `get_latest_file_id` warns and returns None on no match, and warns + again if a match is missing its `fileId` field — in which case it returns whatever + `.get("fileId", None)` produced. A falsy id is not something to deliver on. + """ + with pytest.raises(findability.DeliveryNotFindableError): + findability.assert_findable( + "", expected_file_id="x", consumer_name="un_fao", category="forecast" + ) + + +@pytest.mark.parametrize("partner", PARTNER_PACKAGES) +def test_the_manager_verifies_only_when_something_was_uploaded(partner): + source = (_REPO / "views_postprocessing" / partner / "managers" / f"{partner}.py").read_text() + tree = ast.parse(source) + save = next( + n for n in ast.walk(tree) + if isinstance(n, ast.FunctionDef) and n.name == "_save_contract" + ) + + def _calls(node): + return { + c.func.id + for c in ast.walk(node) + if isinstance(c, ast.Call) and isinstance(c.func, ast.Name) + } + + assert "_assert_delivery_is_findable" in _calls(save), ( + f"{partner}'s _save_contract no longer runs the C-94 preflight, so an upload " + "that lands somewhere the consumer cannot see it reports success (C-94)." + ) + # Structural, not textual: an ordering check on substrings stays green if the call + # is moved into the `else` branch or dedented out of the guard entirely — which is + # the regression this message claims to prevent. + guarded = [ + n for n in ast.walk(save) + if isinstance(n, ast.If) + and isinstance(n.test, ast.Name) + and n.test.id == "upload_enabled" + and "_assert_delivery_is_findable" in _calls_body(n.body) + ] + assert guarded, ( + f"{partner} runs the findability preflight outside the `if upload_enabled:` " + "body. With the interlock holding nothing was uploaded, so the check would " + "refuse every staged run for the absence of a delivery nobody made." + ) + + +@pytest.mark.parametrize("partner", PARTNER_PACKAGES) +def test_the_preflight_queries_the_declared_name_not_the_path_managers(partner): + """C-77: the two are equal today by coincidence, and only one is a declaration.""" + source = (_REPO / "views_postprocessing" / partner / "managers" / f"{partner}.py").read_text() + + assert "_build_partner_read_store" in source, ( + f"{partner} lost the read-back store builder; without it " + "`get_latest_file_id` merges the path manager's model name into the query" + ) + builder = source[source.index("def _build_partner_read_store"):source.index("class ")] + assert "store.model_path = None" in builder, ( + f"{partner}'s read-back store no longer suppresses pipeline-core's automatic " + "`name == model_name` filter, so the preflight verifies the views-models " + "directory name instead of the declared consumer name. A rename there would " + "leave this check green while the delivery went dark (C-77)." + ) + + preflight = _function_source(source, "_assert_delivery_is_findable") + assert "_build_partner_read_store" in preflight, ( + f"{partner}'s preflight must use the suppressed-injection store, not the " + "upload store, or it asks the wrong question" + ) + + # The declared name and the two legs are supplied by the CALLER, so that is where + # they must be asserted. Looking for them in the preflight would be looking in the + # wrong function — which the end-of-file slice used to hide. + save = _function_source(source, "_save_contract") + assert "product.CONSUMER_DOCUMENT_NAME" in save, ( + f"{partner} must pass the DECLARED consumer name to the preflight, never the " + "path manager's model name that happens to equal it (C-77)" + ) + for category in ('"forecast"', '"historical"'): + assert category in save, ( + f"{partner} no longer verifies {category}. Both legs are checked separately " + "because a run with one leg missing is invisible in half, and the historical " + "leg is the one that stranded in run-0 (C-79)." + ) + assert "manifest_file_id" in save, ( + f"{partner} no longer scopes the forecast read-back to this run's manifest, so " + "the previous delivery's document satisfies the check (C-94)" + ) + + +def test_unverified_is_not_the_same_refusal_as_not_findable(): + """A store that could not be asked is a different event from an empty answer. + + Quarantining a delivery because the *check* failed would be an outage the guard + manufactured. Same distinction as C-103 (a missing producer client vs a producer + publishing no boundary) and C-99 (an unrecognised store result vs a real one). + """ + exc = findability.unverified("forecast", TimeoutError("read timed out")) + assert isinstance(exc, findability.FindabilityUnverifiedError) + assert not isinstance(exc, findability.DeliveryNotFindableError), ( + "the two must not share a type, or a caller cannot act differently on them" + ) + message = str(exc) + assert "UNVERIFIED, not known invisible" in message + assert "TimeoutError" in message and "read timed out" in message, ( + "the refusal must quote what actually stopped the check" + ) + 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_previous_runs_document_does_not_satisfy_this_run(): + """The finding that made the guard worth having: without run-scoping it passes + from delivery 2 onward exactly when a C-79 orphan appears. + + Run-1 delivered, so a document exists. Run-2's upload reports success but its + metadata document is never created. The consumer's query returns run-1's id. Asking + "is anything there" answers yes; asking "is what I just uploaded there" answers no. + """ + with pytest.raises(findability.DeliveryNotFindableError) as excinfo: + findability.assert_findable( + "run-1-document", + expected_file_id="run-2-document", + consumer_name="un_fao", + category="historical", + ) + message = str(excinfo.value) + assert "run-1-document" in message and "run-2-document" in message, ( + "the refusal must name both what the consumer will find and what this run " + "uploaded, or an operator cannot tell staleness from absence" + ) + assert "PREVIOUS" in message, ( + "the consequence — the consumer goes on serving the previous delivery — is the " + "part that distinguishes this from an empty bucket" + ) diff --git a/tests/test_store_construction.py b/tests/test_store_construction.py index a7d1895..9509cdb 100644 --- a/tests/test_store_construction.py +++ b/tests/test_store_construction.py @@ -69,7 +69,7 @@ def test_the_builders_are_functions_not_methods(partner): module = _managers(partner) for name in ("_build_prod_forecasts_store", "_build_partner_store", - "_partner_appwrite_config"): + "_build_partner_read_store", "_partner_appwrite_config"): fn = getattr(module, name, None) assert fn is not None, ( f"[{partner}] {name} is gone from the module namespace. If it moved back " @@ -163,7 +163,8 @@ def test_no_coordinate_value_is_baked_into_the_builders(partner): module = _managers(partner) source = "".join( inspect.getsource(getattr(module, name)) - for name in ("_build_prod_forecasts_store", "_partner_appwrite_config") + for name in ("_build_prod_forecasts_store", "_partner_appwrite_config", + "_build_partner_read_store") ) for line in source.splitlines(): if "os.getenv(" in line: diff --git a/views_postprocessing/contract/wire/sink.py b/views_postprocessing/contract/wire/sink.py index 251d921..8acb5cb 100644 --- a/views_postprocessing/contract/wire/sink.py +++ b/views_postprocessing/contract/wire/sink.py @@ -160,13 +160,19 @@ def deliver_run( common = {"name": consumer_name, "category": "forecast", "loa": "pgm"} - def _upload(file_name: str, doc_type: str, targets: list) -> None: - store.upload(staging / file_name, filename=file_name, doc_type=doc_type, targets=targets, **common) + def _upload(file_name: str, doc_type: str, targets: list): + file_id = store.upload( + staging / file_name, filename=file_name, doc_type=doc_type, targets=targets, **common + ) logger.info("uploaded %s (type=%s, run=%s)", file_name, doc_type, run_id) # the ledger + return file_id for record in shard_records: _upload(record["name"], SHARD_DOC_TYPE, [record["target"]]) _upload(sidecar_file, SIDECAR_DOC_TYPE, list(per_target)) - _upload(manifest_file, MANIFEST_DOC_TYPE, list(per_target)) # the commit marker + # The manifest is uploaded LAST, so it is the newest `category="forecast"` document + # 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)) summary["uploaded"] = True return summary diff --git a/views_postprocessing/crafd/managers/crafd.py b/views_postprocessing/crafd/managers/crafd.py index c60996f..5201056 100644 --- a/views_postprocessing/crafd/managers/crafd.py +++ b/views_postprocessing/crafd/managers/crafd.py @@ -17,7 +17,7 @@ from views_postprocessing.crafd.store_port import _ContractStorePort from views_postprocessing.contract.wire import sink as wire_sink from views_postprocessing.contract.wire import source_selection -from views_postprocessing.delivery import coverage, observed_range, provenance +from views_postprocessing.delivery import coverage, findability, observed_range, provenance from pathlib import Path logger = logging.getLogger(__name__) @@ -113,6 +113,36 @@ def _partner_appwrite_config(model_path) -> AppwriteConfig: ) +def _build_partner_read_store(model_path) -> DatastoreModule: + """The partner store for the C-94 read-back, with pipeline-core's automatic + ``name == model_name`` filter suppressed — otherwise the preflight would verify the + views-models directory name, which equals the declared consumer name only by + coincidence (**C-77**). C-94 records why it reuses the write key. + """ + store = DatastoreModule(appwrite_file_manager_config=_partner_appwrite_config(model_path)) + store.model_path = None + 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) + + class CRAFDPostProcessorManager(PostprocessorManager, ForecastingModelManager): def __init__( self, @@ -332,7 +362,7 @@ def _save_contract(self) -> dict: Path(summary["staging_dir"]), lookup ) if upload_enabled: - store.upload( + hist_file_id = store.upload( hist_path, filename=hist_path.name, # The DECLARED consumer name, not `self._model_path.model_name` @@ -351,6 +381,13 @@ def _save_contract(self) -> dict: description=hist_description, ) logger.info("uploaded %s (historical, run %s)", hist_path.name, summary["run_id"]) + # C-94: nothing above observes the OUTCOME of an upload. Every call + # reported success in run-0 too, and the historical leg still stranded. + _assert_delivery_is_findable( + self._model_path, + product.CONSUMER_DOCUMENT_NAME, + {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, + ) else: logger.info( "Interlock holding: historical artifact staged at %s (no store calls).", diff --git a/views_postprocessing/crafd/store_port.py b/views_postprocessing/crafd/store_port.py index cda03ff..b0ee757 100644 --- a/views_postprocessing/crafd/store_port.py +++ b/views_postprocessing/crafd/store_port.py @@ -63,7 +63,7 @@ def download(self, file_id: str) -> bytes: "these by name and cannot tell a failed download from an empty artifact." ) - def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, description=None) -> None: + def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, description=None) -> str | None: result = self._dsm.upload_data( file=file_path, filename=filename, @@ -98,3 +98,9 @@ def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, f"without a metadata document): {error}. The store reported " f"success={success!r} (result type {type(result).__name__})." ) + # The uploaded file's id, so the C-94 read-back can be scoped to THIS run. + # Discarding it (as this did until 2026-08-19) makes the only available check + # "is there any document under the consumer's name", which the previous run + # already satisfies — so the guard could never fail from delivery 2 onward. + data = getattr(result, "data", None) + return data.get("file_id") if isinstance(data, dict) else None diff --git a/views_postprocessing/delivery/findability.py b/views_postprocessing/delivery/findability.py new file mode 100644 index 0000000..508dea3 --- /dev/null +++ b/views_postprocessing/delivery/findability.py @@ -0,0 +1,107 @@ +"""Findability invariant: a delivered artifact must be retrievable under the name the +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. + +**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 +calls it *"invisible to the consumer, not merely degraded"*, and it has happened: run-0's +historical artifact was stranded on 2026-07-27 as a file with no metadata document +(register C-79). Nothing in this repository observed it. Every other mechanism the +platform aims at this is a **CI-time proxy** — we check our label against the registry, +the consumer checks theirs — and none of them observes the outcome of a real upload. + +**Scope, stated plainly, because this entry has already been over-claimed once.** This +catches *the delivery ran and the consumer cannot see it*. It does **not** catch: + +* *no delivery happened at all* — the 2026-08-12 empty bucket, which was an upstream + destructive migration with no run since. A post-upload check observes nothing when + there was no upload. This repository is not told when a delivery is due. +* *the bucket is fine but what is served is stale* — faoapi's warm per-key cache can + serve stale historical over an emptied bucket. + +Both are recorded as gaps in C-94 rather than dressed as things this covers. +""" + +from __future__ import annotations + + +class DeliveryNotFindableError(RuntimeError): + """An upload succeeded, and the consumer's own query cannot find it.""" + + +class FindabilityUnverifiedError(RuntimeError): + """The read-back could not be performed, so findability is unknown.""" + + +def unverified(category: 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 + producer client from a producer that publishes no boundary (C-103), and for the + same reason ``_ContractStorePort.download`` refuses an unrecognised result rather + than adapting to it (C-99): the two conditions call for different operator actions. + A delivery that cannot be found is quarantined. A delivery that could not be + *checked* may be perfectly fine, and quarantining it on a transient store error + 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"{type(exc).__name__}: {exc}. The delivery is UNVERIFIED, not known invisible — " + "re-run the check before quarantining anything." + ) + + +def assert_findable(file_id, *, expected_file_id, consumer_name: str, category: str) -> None: + """Raise unless the store returned something for the consumer's own query. + + Args: + file_id: the result of querying the partner store for the newest document + matching the consumer's filters. ``None`` is pipeline-core's documented + "no match", but **any falsy value is refused**: `get_latest_file_id` also + warns and returns ``.get("fileId", None)`` when it finds a document that + is missing that field, so an empty id means "found something unusable" + rather than "found". Same polarity as ``_ContractStorePort.download``, + which refuses zero bytes for the reason C-99 records — an unrecognised + result is refused and named, not adapted to silently. + expected_file_id: the id THIS run uploaded for this leg — the manifest for the + forecast leg (uploaded last, so it is the newest such document) and the + historical artifact for its own. Without it the only available question is + *"does any document exist under the consumer's name"*, which the previous + delivery already answered yes to — so the check would pass on every run + after the first, precisely when a C-79-shaped orphan appeared. + consumer_name: the DECLARED store-document ``name`` the consumer filters on + (``product.CONSUMER_DOCUMENT_NAME``), never a path-manager or directory + name that happens to equal it (C-77). + category: the delivery leg being verified — ``"forecast"`` or ``"historical"``. + Checked separately on purpose: a run whose forecast landed and whose + historical did not is invisible in exactly one half, and a single + whole-delivery check would pass on it. + + Raises: + DeliveryNotFindableError: naming the query that found nothing. + """ + if file_id and file_id == expected_file_id: + return + if file_id: + raise DeliveryNotFindableError( + f"delivery is INVISIBLE to the consumer: the newest {category!r} document " + f"under name == {consumer_name!r} is {file_id!r}, but this run uploaded " + f"{expected_file_id!r}. The consumer will go on serving the PREVIOUS " + "delivery while this one reports success — which is why the check is scoped " + "to this run rather than asking whether any document exists (a question the " + "previous run already answered). Quarantine and inspect the partner bucket." + ) + raise DeliveryNotFindableError( + f"delivery is INVISIBLE to the consumer: uploads for category {category!r} " + f"reported success, but querying the partner store as the consumer does — " + f"name == {consumer_name!r}, category == {category!r} — returns nothing. " + "The artifacts may exist as files while carrying no metadata document, which " + "is how run-0's historical leg was stranded (C-79); a document under any other " + "name is equally invisible, because the consumer filters on this one " + "unconditionally (ADR-013 §4.1a). Quarantine the run and check the partner " + "bucket before re-delivering — the contract has no retraction primitive, so a " + "correction is a new complete run." + ) diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index 4f98508..4d6ddf6 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -17,7 +17,7 @@ from views_postprocessing.unfao.store_port import _ContractStorePort from views_postprocessing.contract.wire import sink as wire_sink from views_postprocessing.contract.wire import source_selection -from views_postprocessing.delivery import coverage, observed_range, provenance +from views_postprocessing.delivery import coverage, findability, observed_range, provenance from pathlib import Path logger = logging.getLogger(__name__) @@ -113,6 +113,36 @@ def _partner_appwrite_config(model_path) -> AppwriteConfig: ) +def _build_partner_read_store(model_path) -> DatastoreModule: + """The partner store for the C-94 read-back, with pipeline-core's automatic + ``name == model_name`` filter suppressed — otherwise the preflight would verify the + views-models directory name, which equals the declared consumer name only by + coincidence (**C-77**). C-94 records why it reuses the write key. + """ + store = DatastoreModule(appwrite_file_manager_config=_partner_appwrite_config(model_path)) + store.model_path = None + 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) + + class UNFAOPostProcessorManager(PostprocessorManager, ForecastingModelManager): def __init__( self, @@ -332,7 +362,7 @@ def _save_contract(self) -> dict: Path(summary["staging_dir"]), lookup ) if upload_enabled: - store.upload( + hist_file_id = store.upload( hist_path, filename=hist_path.name, # The DECLARED consumer name, not `self._model_path.model_name` @@ -351,6 +381,13 @@ def _save_contract(self) -> dict: description=hist_description, ) logger.info("uploaded %s (historical, run %s)", hist_path.name, summary["run_id"]) + # C-94: nothing above observes the OUTCOME of an upload. Every call + # reported success in run-0 too, and the historical leg still stranded. + _assert_delivery_is_findable( + self._model_path, + product.CONSUMER_DOCUMENT_NAME, + {"forecast": summary["manifest_file_id"], "historical": hist_file_id}, + ) else: logger.info( "Interlock holding: historical artifact staged at %s (no store calls).", diff --git a/views_postprocessing/unfao/store_port.py b/views_postprocessing/unfao/store_port.py index cda03ff..b0ee757 100644 --- a/views_postprocessing/unfao/store_port.py +++ b/views_postprocessing/unfao/store_port.py @@ -63,7 +63,7 @@ def download(self, file_id: str) -> bytes: "these by name and cannot tell a failed download from an empty artifact." ) - def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, description=None) -> None: + def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, description=None) -> str | None: result = self._dsm.upload_data( file=file_path, filename=filename, @@ -98,3 +98,9 @@ def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, f"without a metadata document): {error}. The store reported " f"success={success!r} (result type {type(result).__name__})." ) + # The uploaded file's id, so the C-94 read-back can be scoped to THIS run. + # Discarding it (as this did until 2026-08-19) makes the only available check + # "is there any document under the consumer's name", which the previous run + # already satisfies — so the guard could never fail from delivery 2 onward. + data = getattr(result, "data", None) + return data.get("file_id") if isinstance(data, dict) else None