From 86e9bac6b870ee4d7c7ca892481205d8bf654ffb Mon Sep 17 00:00:00 2001 From: Polichinl Date: Mon, 27 Jul 2026 19:14:07 +0200 Subject: [PATCH] =?UTF-8?q?feat(historical):=20B2/#126=20=E2=80=94=20frame?= =?UTF-8?q?-native=20historical=20path;=20pandas=20leaves=20the=20construc?= =?UTF-8?q?tion,=20the=20artifact=20stays=20reader-identical?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _read_historical_data dispatches on declared_data_format (pipeline-core's contract fn — the one gate, never re-implemented): feature_frame descriptors fetch via get_feature_frame → FeatureFrame, making the FAO product the frame path's first production consumer (C-40 gate lifted). Frame-native month clip (frame_extraction.drop_months_above; same producer-sourced boundary, degrade-open). unfao/historical.py builds the SAME artifact faoapi's deployed reader ingests — month_id/priogrid_id/targets/9-GAUL-cols plain parquet — pandas-free (searchsorted gather from the injected ADR-011 lookup; sidecar cast discipline: plain strings, float64 codes); null-gate + unmapped-count carried over fail-loud. Deliberate divergences, both reader-safe and pinned by tests: junk row/col columns DROPPED (the legacy chain would have shipped them as bogus auto-detected targets) and values ship as SCALARS, not the C-40 single-element lists (faoapi's _validate_feature_samples normalizes both identically — verified in their code). Contract-mode gap closed: _save_contract now stages+uploads the historical artifact alongside the wire, behind the same §11.4 interlock (S7 had scoped it to the wire only — run 0 would have shipped forecasts with no actuals). Reader-level parity with the B1 characterization golden: green (same slice, frame chain vs legacy chain, compared through faoapi's reader semantics). Closes #126 (with #113-era amendments). Refs #92, #89. Co-Authored-By: Claude Fable 5 --- .../historical_golden/raw_slice_input.parquet | Bin 0 -> 5354 bytes tests/test_historical_builder.py | 123 +++++++++++++++ .../unfao/frame_extraction.py | 15 ++ views_postprocessing/unfao/historical.py | 115 ++++++++++++++ views_postprocessing/unfao/managers/unfao.py | 147 ++++++++++++++++-- 5 files changed, 386 insertions(+), 14 deletions(-) create mode 100644 tests/fixtures/historical_golden/raw_slice_input.parquet create mode 100644 tests/test_historical_builder.py create mode 100644 views_postprocessing/unfao/historical.py diff --git a/tests/fixtures/historical_golden/raw_slice_input.parquet b/tests/fixtures/historical_golden/raw_slice_input.parquet new file mode 100644 index 0000000000000000000000000000000000000000..f128c256fc5f790e214b10b9e3ec9d96833d5828 GIT binary patch literal 5354 zcmd6r3v^TU9mj7Ugair}$L$i=Va*_*mXf3m5CpfI(4+}%X@kuRu;$)8+k0<99&J;` z#$Jl}s;DSF@r{bly%Z4<1=-vfPUm#%rta(F3`3 z{{R2)_y7HV|KI=r&uI_iGSIbjJw3mZuBNH8DN5N$*os%o!-(w(Skq-SYNl%igvR`JRHNGRB!V(IFJy^4lIRSn0g zbSrhrGnqis47Vw|n!=+VFN&IuFve@k!w^%&VH~T`?XLO%59Ixs<$=UY0(n!R+jmAw zUU5a5=!KnRp-@=$tFyWC8O^=1bd7XVNrC0N)XMx|Qefq+urRkRUtnW3)L2cGZujg) zx;9s1N+B+$qy4dTsQM!C`7f z6?N_$CDjDZ1KK&t+Fu`;vACA{4W!?K-+}kR@4*M)58#jBL+~eze)_1&OFaYO0C*N0 z1kZs(;CXNuya0}X7r{&5C^!aQ2FJm7VCeGWsxj&c2s^-)U?;c=Tn(-P*MjT77`Ptn z0ylsgq1^Y1YJ_q_@BlCH0YBIXHi6CHQqT`B0|CH7ap4J7f?5PY2Nr`R;CyfaxDYgh zrJw~Y19~WZcv2N5!Bc+*e*u35e*=F9r@%kJN8q2}WAHC1QLn0^wMSl4O_>GiyBPfn z_#XH^I01eDPQt>jHxOhlV(KO|?glr5TfiQ0E4U5Z4({ICug)2_swHQ|(xtd{5r_soWa} zz*7qPa81b(kndEJjU|1^?ksOb*PfaWPUa2)Wr(V(uBn}&n5k6Fnq4<%?pgKo&i=%F zy5W;*%{iZH{Pcozo6ggIX5k{;;w9%_aAEV(mSuW_v9-;#e8o!hs*6^yx%ji6V=UIS z?KXRd!?~`rt9yMJdz#V z`uS~N_~Q02efcY2{o3VM?6`90Raakg?R8_<@4Df}n|9xP%br_ryZw$k@7jC!J@?*s z{{s*1d+6au_CNaA<4-*K)Yrf9&2N4C>1Pf+d+@nK&mVr_$crx>J@)d)$0_-5-f`W* z1rtZw>_jpak0fKkiBqkqVSJwMBaQBzI-F|d89GhgDxGQ?-bAStDJ~ik0`x43p0%jD z=$#Z;*rhg^&Cg)>(ci~Xf=&ubbro2KJ;fNokYY2CEkpwn`vO;?pQ{6dE^&?y^s zs^?Rbrdh2>1jJwYm!f4J|0LfZ!7nysH4a! z6705S_?8PaZyPB*-L!Jo|7P0w>8F*u{>#%Amq=-ji+d(b8WSrfxX?URdW-9@EWwjH z|HEme%1VX;>3CAAK0d}K1Ia8YzC<>VOvZ=F5)XxvsaRYjt%l_-h8DdR<)FE-*=@+BZM4`& zzc4c7^Y$9Mtqu|=a$y(_xy|~J*TR!|9fCi_8g2USfnI4I*5lUuJg#Wa-d@OEc*8J> zC*$)NaylbfuQAB`ZRQN}P3#H*k2^)=DaKUFvB28R(coHNE{Ekgp5%YpkNhG=y|1u- zBj!#ksiXXQV_iC&VZYIhIx&yRYvWJ&JtL5(7$aJutl)}rhADC*y34y<)Y$=aq_qe&mInG=3_O)BKBtoAGT#%HF=8}`GQt?A_bMfRLfo;R+;nP81U z>_r#$oLn=N>aJw|^1LNju~3s0*VK!ea|qp8zsII$TOFhpvBo0_pXiKo;-Iwuuon34 z916-cX!LuXqc|fnQe6+U8#@QuGhNtkWdBU~=;t}H8!@y=Ho7*n89E2LVqy6t#QDcC zX7cDH56epO>>-cbAxDl*2o@$PHFD&M^Dh69v?k?<8SR0RbUaGt8Gl~xbx5;1CGGrC z%rI6W+b3#3HXb1dF3wX<~wMadyHf5}t$6pP7C} z((jtk$EDCm&U@S+m|6<7cOoh6Ysp2Eb6WaJM1~mzmf~}c=p#NKdCQ*&vxZVPB?8GU lnNWHumq;v?F2SbUCD$Zhf0|AB_hs=PHC{(iJ@_x@zX5mU3c>&Y literal 0 HcmV?d00001 diff --git a/tests/test_historical_builder.py b/tests/test_historical_builder.py new file mode 100644 index 0000000..12b7a08 --- /dev/null +++ b/tests/test_historical_builder.py @@ -0,0 +1,123 @@ +"""The frame-native historical artifact builder (unfao/historical.py, #126) — +reader-level parity with the legacy characterization golden, plus its own +fail-loud properties.""" + +from pathlib import Path + +import numpy as np +import pyarrow as pa +import pyarrow.parquet as pq +import pytest +from views_frames import FeatureFrame, SpatialLevel, SpatioTemporalIndex + +from views_postprocessing.unfao import historical +from views_postprocessing.unfao.frame_extraction import drop_months_above +from views_postprocessing.unfao.gaul_schema import METADATA_COLS + +import sys + +sys.path.insert(0, str(Path(__file__).resolve().parent)) +from test_historical_parity import GOLDEN, REAL_TARGETS, reader_view # noqa: E402 + +_FIXDIR = Path(__file__).resolve().parent / "fixtures" / "historical_golden" +_LOOKUP = Path(__file__).resolve().parents[1] / "views_postprocessing" / "data" / "gaul_lookup.parquet" + + +def _slice_frame() -> FeatureFrame: + """The committed raw input slice as the FeatureFrame the loader would return.""" + table = pq.read_table(_FIXDIR / "raw_slice_input.parquet") + time = np.asarray(table.column("month_id"), dtype=np.int64) + unit = np.asarray(table.column("priogrid_id"), dtype=np.int64) + values_2d = np.column_stack( + [np.asarray(table.column(t), dtype=np.float32) for t in REAL_TARGETS] + ) + index = SpatioTemporalIndex(time=time, unit=unit, level=SpatialLevel.PGM) + return FeatureFrame.from_2d(values_2d, index, feature_names=list(REAL_TARGETS)) + + +def test_reader_level_parity_with_the_legacy_golden(tmp_path): + # The #126 acceptance oracle: same slice, frame-native chain, and faoapi's + # reader must see the same frame the legacy chain produced (minus junk). + table = historical.build_historical_table(_slice_frame(), pq.read_table(_LOOKUP)) + historical.assert_metadata_complete(table) + name, _ = historical.write_historical_artifact(table, tmp_path / "new.parquet") + new = reader_view(tmp_path / "new.parquet") + old = reader_view(GOLDEN) + assert set(new["targets"]) == set(REAL_TARGETS) # junk row/col GONE by design + assert new["rows"].keys() == old["rows"].keys() # identical (month, cell) set + def norm(v): + # faoapi's reader semantics: pandas collapses parquet null/NaN to NaN, + # and _validate_feature_samples (handlers.py:374-376) normalizes scalar + # and single-element-list values to the same 1-array — the legacy chain + # shipped [v] (the C-40 list-in-cell disease), the frame builder ships v; + # both reach faoapi identical. The oracle compares post-normalization. + if isinstance(v, list) and len(v) == 1: + v = v[0] + return float("nan") if v is None else v + + for key, new_cols in new["rows"].items(): + old_cols = old["rows"][key] + for col in (*REAL_TARGETS, *METADATA_COLS): + a, b = norm(new_cols[col]), norm(old_cols[col]) + if isinstance(a, float) and isinstance(b, float): + if a != a and b != b: # both NaN — equal through the reader + continue + assert a == pytest.approx(b, rel=1e-6), (key, col) + else: + assert a == b, (key, col) + + +def test_junk_columns_are_dropped(): + table = historical.build_historical_table(_slice_frame(), pq.read_table(_LOOKUP)) + assert "row" not in table.column_names and "col" not in table.column_names + assert table.column_names == ["month_id", "priogrid_id", *REAL_TARGETS, *METADATA_COLS] + + +def test_strings_are_plain_and_codes_float64(): + table = historical.build_historical_table(_slice_frame(), pq.read_table(_LOOKUP)) + assert table.schema.field("country_iso_a3").type == pa.string() + assert table.schema.field("admin1_gaul0_code").type == pa.float64() + + +def test_cell_absent_from_lookup_fails_loud(): + frame = _slice_frame() + bad_index = SpatioTemporalIndex( + time=np.asarray(frame.index.time), + unit=np.where(np.arange(frame.n_rows) == 0, 999_999, np.asarray(frame.index.unit)), + level=SpatialLevel.PGM, + ) + bad = FeatureFrame(frame.values, bad_index, feature_names=list(frame.feature_names)) + with pytest.raises(historical.HistoricalArtifactError, match="999999"): + historical.build_historical_table(bad, pq.read_table(_LOOKUP)) + + +def test_multi_sample_historical_refused(): + frame = _slice_frame() + fat = FeatureFrame( + np.repeat(frame.values, 2, axis=2), frame.index, feature_names=list(frame.feature_names) + ) + with pytest.raises(historical.HistoricalArtifactError, match="single-valued"): + historical.build_historical_table(fat, pq.read_table(_LOOKUP)) + + +def test_unmapped_cell_count_counts_distinct_gids(): + table = historical.build_historical_table(_slice_frame(), pq.read_table(_LOOKUP)) + assert historical.unmapped_cell_count(table) == 0 + # poke one null into a metadata column across both months of one cell + poked = table.set_column( + table.column_names.index("country_iso_a3"), + "country_iso_a3", + pa.array( + [None if i < 2 else v for i, v in enumerate(table.column("country_iso_a3").to_pylist())], + pa.string(), + ), + ) + assert historical.unmapped_cell_count(poked) >= 1 + + +def test_frame_native_month_clip(): + frame = _slice_frame() # months 121, 122 + clipped = drop_months_above(frame, 121) + months = sorted(set(np.asarray(clipped.index.time).tolist())) + assert months == [121] + assert drop_months_above(frame, 122) is frame # nothing to clip → identity diff --git a/views_postprocessing/unfao/frame_extraction.py b/views_postprocessing/unfao/frame_extraction.py index c6aaaab..2e4c45e 100644 --- a/views_postprocessing/unfao/frame_extraction.py +++ b/views_postprocessing/unfao/frame_extraction.py @@ -39,6 +39,21 @@ def months_of(frame: PredictionFrame) -> NDArray[np.int64]: return np.unique(np.asarray(frame.index.time, dtype=np.int64)) +def drop_months_above(frame, last_valid_month_id: int): + """The frame clipped to observed months (``time <= last_valid_month_id``). + + Frame-native sibling of ``extraction.drop_months_above`` (S7/#92): works on + any frame sharing the ``SpatioTemporalIndex`` surface (PredictionFrame, + FeatureFrame) via ``.select`` on a boolean row mask. The boundary value is + declared by the caller (producer-sourced, C-26) — never inferred here. + """ + time = np.asarray(frame.index.time, dtype=np.int64) + mask = time <= int(last_valid_month_id) + if mask.all(): + return frame + return frame.select(mask) + + def drop_units(frame: PredictionFrame, excluded: frozenset) -> PredictionFrame: """The frame without the DECLARED excluded cells (row filter on ``unit``). diff --git a/views_postprocessing/unfao/historical.py b/views_postprocessing/unfao/historical.py new file mode 100644 index 0000000..5d8f5cc --- /dev/null +++ b/views_postprocessing/unfao/historical.py @@ -0,0 +1,115 @@ +"""The FAO historical (actuals) artifact, built pandas-free (#126). + +One concept: turn the historical ``views_frames.FeatureFrame`` plus the injected +ADR-011 GAUL lookup into the SAME artifact faoapi's deployed reader already +ingests — plain-column parquet with ``month_id, priogrid_id, , +``. The wire artifact does not change shape +(ADR-013 §4.1/§8); only its construction leaves pandas. Reader-level parity with +the legacy chain is pinned by ``tests/test_historical_parity.py``'s golden. + +Rules carried over from the sidecar builder (§5.1 discipline): +- lookup injected (DIP) — production passes ``data/gaul_lookup.parquet``; +- a cell absent from the lookup FAILS LOUD (geography never silently vanishes); +- string columns plain (never dictionary/categorical), code columns float64; +- deliberately DROPPED: any non-feature payload columns (the legacy chain would + have shipped datafactory's ``row``/``col`` grid coordinates as bogus + auto-detected "targets" — pinned in the characterization golden). + +The null-gate (every metadata value present) is this module's fail-loud +equivalent of the legacy ``_validate`` historical leg; ``unmapped_cell_count`` +feeds the same provenance field it always did (C-15). +""" + +from __future__ import annotations + +import hashlib +from pathlib import Path + +import numpy as np +import pyarrow as pa +import pyarrow.compute as pc +import pyarrow.parquet as pq + +from views_postprocessing.unfao.gaul_schema import CODE_COLS, COORD_COLS, METADATA_COLS + + +class HistoricalArtifactError(ValueError): + """The historical artifact cannot be built as declared.""" + + +def build_historical_table(frame, lookup: pa.Table) -> pa.Table: + """FeatureFrame + GAUL lookup → the reader-compatible historical table. + + ``frame.values`` is ``(N, F, S)``; actuals are single-valued (S=1) — a + multi-sample historical is a declaration error, refused loud. + """ + if frame.sample_count != 1: + raise HistoricalArtifactError( + f"historical actuals must be single-valued; got sample_count={frame.sample_count}." + ) + missing_cols = [c for c in ("priogrid_gid", *METADATA_COLS) if c not in lookup.column_names] + if missing_cols: + raise HistoricalArtifactError(f"lookup lacks required columns: {missing_cols}.") + + unit = np.asarray(frame.index.unit, dtype=np.int64) + time = np.asarray(frame.index.time, dtype=np.int64) + + lookup_gids = lookup.column("priogrid_gid").to_numpy(zero_copy_only=False) + order = np.argsort(lookup_gids) + sorted_gids = lookup_gids[order] + pos = np.searchsorted(sorted_gids, unit) + pos_valid = (pos < len(sorted_gids)) + if not pos_valid.all() or not (sorted_gids[np.clip(pos, 0, len(sorted_gids) - 1)] == unit).all(): + absent = np.unique(unit[~pos_valid | (sorted_gids[np.clip(pos, 0, len(sorted_gids) - 1)] != unit)]) + raise HistoricalArtifactError( + f"{absent.size} historical cell(s) absent from the GAUL lookup — geography " + f"must never silently vanish. First missing gids: {absent[:5].tolist()}." + ) + row_indices = pa.array(order[pos]) + + columns: dict = { + "month_id": pa.array(time, pa.int64()), + "priogrid_id": pa.array(unit, pa.int64()), + } + for i, name in enumerate(frame.feature_names): + columns[name] = pa.array(frame.values[:, i, 0]) + for col in METADATA_COLS: + array = lookup.column(col).combine_chunks().take(row_indices) + if col in CODE_COLS: + array = pc.fill_null(array.cast(pa.float64()), float("nan")) + elif col in COORD_COLS: + array = array.cast(pa.float64()) + else: + array = array.cast(pa.string()) # plain, never dictionary (§5.1 discipline) + columns[col] = array + return pa.table(columns) + + +def assert_metadata_complete(table: pa.Table) -> None: + """Fail loud on any missing metadata value (the legacy null-gate, kept).""" + for col in METADATA_COLS: + nulls = table.column(col).null_count + if nulls: + raise HistoricalArtifactError( + f"historical artifact has {nulls} null value(s) in required metadata " + f"column {col!r} ({nulls}/{table.num_rows} rows)." + ) + + +def unmapped_cell_count(table: pa.Table) -> int: + """Distinct cells with a null in ANY metadata column (the C-15 provenance count).""" + mask = None + for col in METADATA_COLS: + col_null = pc.is_null(table.column(col)) + mask = col_null if mask is None else pc.or_(mask, col_null) + if not pc.any(mask).as_py(): + return 0 + gids = pc.filter(table.column("priogrid_id"), mask) + return len(pc.unique(gids)) + + +def write_historical_artifact(table: pa.Table, path: Path) -> tuple[str, str]: + """Serialize; return ``(file_name, sha256_of_bytes)``.""" + path = Path(path) + pq.write_table(table, path) + return path.name, hashlib.sha256(path.read_bytes()).hexdigest() diff --git a/views_postprocessing/unfao/managers/unfao.py b/views_postprocessing/unfao/managers/unfao.py index c75db4b..92326ad 100644 --- a/views_postprocessing/unfao/managers/unfao.py +++ b/views_postprocessing/unfao/managers/unfao.py @@ -19,7 +19,8 @@ from dotenv import load_dotenv from views_postprocessing.unfao.enrichment import _DEFAULT_LOOKUP, GaulLookupEnricher from views_postprocessing.unfao.gaul_schema import METADATA_COLS -from views_postprocessing.unfao import extraction, product, source_metadata +from views_pipeline_core.modules.dataloaders.datafactory_contract import declared_data_format +from views_postprocessing.unfao import extraction, frame_extraction, historical, product, source_metadata from views_postprocessing.unfao.wire import sink as wire_sink from views_postprocessing.unfao.wire import source_selection from views_postprocessing.delivery import coverage, identity, observed_range, provenance @@ -53,7 +54,7 @@ def download(self, file_id: str) -> bytes: self._dsm.download_prediction(file_id).to_dict().get("data", {}).get("file_bytes", None) ) - def upload(self, file_path, *, filename, name, doc_type, category, loa, targets) -> None: + def upload(self, file_path, *, filename, name, doc_type, category, loa, targets, description=None) -> None: self._dsm.upload_data( file=file_path, filename=filename, @@ -62,6 +63,7 @@ def upload(self, file_path, *, filename, name, doc_type, category, loa, targets) category=category, loa=loa, targets=targets, + description=description, ) @@ -82,12 +84,42 @@ def __init__( self._historical_dataset = None self._forecast_dataset = None self._forecast_resolution = None # contract mode: {target: TargetLease} + self._historical_frame = None # frame-native historical (#126) self._enricher = GaulLookupEnricher() self.ensemble_path_manager = None + def _read_historical_frame(self): + """#126: historical actuals as a views_frames.FeatureFrame — the first + production consumer of pipeline-core's frame-native fetch. Same clip + policy as the legacy path (producer-sourced boundary, degrade-open).""" + frame = self._data_loader.get_feature_frame( + partition="forecasting", use_saved=False, level="pgm", validate=True + ) + try: + lv = source_metadata.last_valid_month_id(self.configs.get("zarr_url")) + except Exception: + logger.warning("last_valid_month_id unavailable; skipping clip (degrade-open, C-26).", exc_info=True) + lv = None + if lv is None: + self._historical_frame = frame + return + fabricated = observed_range.fabricated_months(frame_extraction.months_of(frame), lv) + if len(fabricated): + logger.warning( + "Dropping %d fabricated (unobserved) month(s) above last_valid_month_id=%d.", + len(fabricated), lv, + ) + frame = frame_extraction.drop_months_above(frame, lv) + self._historical_frame = frame + def _read_historical_data(self): self._initialize_data_loader() run_type = "forecasting" + # Declared dispatch (never inferred): the queryset descriptor's + # data_format decides — pipeline-core's contract fn is the one gate. + if declared_data_format(self._model_path.get_queryset()) == "feature_frame": + self._read_historical_frame() + return self._data_loader.get_data( use_saved=False, @@ -244,7 +276,9 @@ def _append_metadata(self, dataset: PGMDataset) -> pd.DataFrame: def _transform( self, ) -> list: - self._historical_dataframe = self._append_metadata(self._historical_dataset) + if self._historical_frame is None: + self._historical_dataframe = self._append_metadata(self._historical_dataset) + # frame-native historical attaches geography at artifact build (historical.py) if self.configs.get("wire_contract"): # Contract mode: the forecast is frames, not a dataframe; geography ships # as the §5 sidecar (built at _save), never joined into the payload. @@ -254,7 +288,15 @@ def _transform( def _validate(self) -> pd.DataFrame: _necessary_metadata_cols = METADATA_COLS - for col in _necessary_metadata_cols: + if self._historical_frame is not None: + logger.info( + "Historical is frame-native: the metadata null-gate is enforced at " + "artifact build (historical.assert_metadata_complete)." + ) + _historical_cols_to_check = [] + else: + _historical_cols_to_check = _necessary_metadata_cols + for col in _historical_cols_to_check: if col not in self._historical_dataframe.columns: err_msg = f"Historical dataframe is missing required metadata column: {col}. Found columns: {self._historical_dataframe.columns.tolist()}" logger.error(err_msg) @@ -356,12 +398,10 @@ def _check_coverage(self) -> None: # Streaming mode: forecast coverage is asserted inside each lease's # load() (where the frame exists — see wire/source_selection); only # the historical (still pandas until #126) is checked here. - sources = [ - ("historical", extraction.cells_of(self._historical_dataframe), len(self._historical_dataframe)), - ] + sources = [self._historical_coverage_source()] else: sources = [ - ("historical", extraction.cells_of(self._historical_dataframe), len(self._historical_dataframe)), + self._historical_coverage_source(), ("forecast", extraction.cells_of(self._forecast_dataframe), len(self._forecast_dataframe)), ] for label, cells, n_rows in sources: @@ -403,13 +443,42 @@ def _save_contract(self) -> dict: ) upload_enabled = bool(self.configs.get("wire_upload_enabled", product.UPLOAD_ENABLED)) store = _ContractStorePort(self._unfao_datastore()) if upload_enabled else None - return wire_sink.deliver_run( + summary = wire_sink.deliver_run( self._forecast_resolution, lookup=pq.read_table(_DEFAULT_LOOKUP), staging_dir=Path(self._model_path.data_generated) / "wire_contract", store=store, upload_enabled=upload_enabled, ) + # Historical leg (#126): the FAO product ships actuals alongside the wire — + # frame-built, same artifact shape faoapi already ingests, same interlock. + if self._historical_frame is None: + raise ValueError( + "contract _save: no historical frame — the un_fao descriptor must " + "declare data_format: feature_frame (#126)." + ) + hist_path, hist_description, _ = self._build_historical_artifact( + Path(summary["staging_dir"]) + ) + if upload_enabled: + store.upload( + hist_path, + filename=hist_path.name, + name=self._model_path.model_name, + doc_type="model", + category="historical", + loa="pgm", + targets=list(self.configs.get("targets", [])), + description=hist_description, + ) + logger.info("uploaded %s (historical, run %s)", hist_path.name, summary["run_id"]) + else: + logger.info( + "Interlock holding: historical artifact staged at %s (no store calls).", + hist_path, + ) + summary["historical"] = hist_path.name + return summary def _unfao_datastore(self) -> DatastoreModule: """The FAO-facing store (`unfao_bucket`).""" @@ -458,18 +527,24 @@ def _save(self) -> list: dsm = DatastoreModule(appwrite_file_manager_config=unfao_appwrite_config) timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") - historical_file_path = self._model_path.data_generated / f"historical_dataset_{timestamp}.parquet" forecast_file_path = self._model_path.data_generated / f"forecast_dataset_{timestamp}.parquet" - self._historical_dataframe.to_parquet( - historical_file_path - ) + if self._historical_frame is not None: + historical_file_path, hist_description, timestamp = self._build_historical_artifact( + self._model_path.data_generated + ) + else: + historical_file_path = self._model_path.data_generated / f"historical_dataset_{timestamp}.parquet" + self._historical_dataframe.to_parquet( + historical_file_path + ) + hist_description = self._delivery_description(self._historical_dataframe, timestamp) dsm.upload_data(file=historical_file_path, filename=Path(historical_file_path).name, name=self._model_path.model_name, loa="pgm", type="model", targets=self.configs.get("targets", []), - description=self._delivery_description(self._historical_dataframe, timestamp), + description=hist_description, category="historical") self._forecast_dataframe.to_parquet( @@ -483,6 +558,50 @@ def _save(self) -> list: description=self._delivery_description(self._forecast_dataframe, timestamp), category="forecast") + def _historical_coverage_source(self): + """(label, cells, n_rows) for whichever historical representation is live.""" + if self._historical_frame is not None: + return ( + "historical", + frame_extraction.cells_of(self._historical_frame), + self._historical_frame.n_rows, + ) + return ( + "historical", + extraction.cells_of(self._historical_dataframe), + len(self._historical_dataframe), + ) + + def _historical_frame_description(self, table, timestamp: str) -> str: + """The C-15 provenance description for the frame-built historical artifact.""" + region = self.configs.get("region") + prov = provenance.build_provenance( + lookup_version=self._enricher.lookup_version, + region=region, + expected_cell_count=coverage.expected_for(region), + actual_cell_count=len(frame_extraction.cells_of(self._historical_frame)), + unmapped_count=historical.unmapped_cell_count(table), + ) + return ( + f"Enriched with geographic metadata on {timestamp} using precomputed GAUL " + f"lookup (ADR-011, version={self._enricher.lookup_version}). " + f"provenance={json.dumps(prov, separators=(',', ':'))}" + ) + + def _build_historical_artifact(self, directory) -> tuple: + """Frame-built historical artifact staged into ``directory``; returns + (path, description, timestamp). Null-gate enforced here (fail loud).""" + import pyarrow.parquet as pq + + timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") + table = historical.build_historical_table( + self._historical_frame, pq.read_table(_DEFAULT_LOOKUP) + ) + historical.assert_metadata_complete(table) + path = Path(directory) / f"historical_dataset_{timestamp}.parquet" + historical.write_historical_artifact(table, path) + return path, self._historical_frame_description(table, timestamp), timestamp + def _delivery_description(self, df: pd.DataFrame, timestamp: str) -> str: """Human prefix + structured provenance (S5/C-15) for an upload's metadata.