Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog.d/523.fixed.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Load certified US annual projection files for the requested calculation years, require complete annual coverage from each family's base, retain earlier annual inputs for formula lookbacks, and reject periods outside the selected coverage. Separate source and derived input caches by artifact, model, and SPM identity.
90 changes: 90 additions & 0 deletions docs/engineering/skills/data-certification.md
Original file line number Diff line number Diff line change
Expand Up @@ -104,3 +104,93 @@ before they can be updated through this path.
The retired `policyengine-bundles` flow (candidates → generated bundle →
archive import) is preserved read-only in that repo's history; bundles
4.15.x–4.16.x remain the historical record of earlier certifications.

## Annual national artifact families

A producer can advertise exact annual inputs through release metadata:

```json
"metadata": {
"dataset_years": {
"populace_us_2024": {
"2024": "populace_us_2024",
"2025": "populace_us_2025"
}
}
}
```

Each value names an ordinary `artifacts` entry with an H5 path, explicit revision
and SHA256. Each family must map its earliest year to the family/base artifact
and include every year through its projection horizon. Certification rejects
gaps and omitted bases so formula lookbacks retain the complete input history.
Do not use path templates or implicit extension for missing years. Keep the source base
release and content identity in producer provenance separately. Certification
validates this mapping and copies it into `data_releases.us.dataset_years`.
It does not invent artifact pins or certify an unpublished candidate.

Annual files use the existing single-year entity tables plus `_time_period`.
The runtime requires that stored year to match the selected manifest year.
`pe.us.ensure_datasets(years=[2025])` fetches only the requested annual file and
loads its native inputs; it does not calculate or uprate those inputs again.
Return keys retain the requested family name, such as
`populace_us_2024_2025`. Explicit annual artifact names work too, with coverage
limited to that artifact's year. Families without this metadata retain engine
extension behavior.

`pe.us.managed_microsimulation(years=[2025, 2030])` permits external calculations
for those two years. Without `years`, the wrapper selects the years in an explicit
`default_calculation_period`, or the current calendar year. With `years` and no
explicit default period, calculations default to the first selected year. The
wrapper checks coverage before downloading any files and preserves the chosen
default after the engine's initialization. External periods outside the selected
years raise `ValueError`, including years available only for internal lookbacks.
Legacy families without annual metadata do not accept `years`; their existing
country-model period behavior remains unchanged.

The country `USMultiYearDataset` receives the selected years and every advertised
earlier year through the last selected year, with no future inputs. This history
prefix supplies actual prior-year income for Medicare IRMAA, state tax provisions,
and recursive employment-income formulas. Selecting an individual annual artifact
also loads its parent family's earlier inputs. Internal formula lookbacks keep
their normal behavior; the wrapper supplies no history before the family's first
year, so earlier lookbacks retain the country's existing assumptions. The loader
checks schema, row counts and IDs across files, and preserves each year's source
URI, revision and digest in `policyengine_bundle["annual_datasets"]`.

The ordinary `Simulation(dataset=ensure_datasets(...)[...])` route uses the same
annual history. `ensure_datasets` records the required source references in the
dataset's JSON metadata without fetching them; `Simulation.run()` fetches earlier
files when needed. Both the baseline and reform use the country multi-year loader,
including its income-response normalization. Regional runs retain the same current
household IDs in every earlier input year. The wrapper rejects missing or changed
history pins, and the output records those references in `annual_input_sources`.

For a 2024–2035 family in calendar year 2026, default construction loads three
files (2024–2026). Selecting 2035 explicitly loads all 12 files. A local candidate
built during this integration measured 1,754,538,463 bytes (1.755 GB) for the
three-year prefix and 5,928,824,315 bytes (5.929 GB) for all 12 files. Its base
file occupied 826,917,837 bytes; each projected file occupied about 463.81 MB.
These measurements describe candidate artifacts; release certification remains a
separate step. `annual_selected_years`, `annual_loaded_years`,
`annual_input_bytes_by_year`, and `annual_input_bytes` record the actual selection
and file sizes. The total describes input storage and the maximum download size;
cached files do not require another download. Memory also includes decoded tables
and engine arrays, so production adoption still needs a population-scale memory
check. The wrapper performs one full entity load per annual file; its preliminary
`_time_period` check reads only that small entry. `ensure_datasets` materializes
only the requested files and does not load this history prefix.

Source caches live beneath `.policyengine/sources/` and include the artifact's
repository, path, revision, content hash and metadata hash. Derived input caches
beneath `.policyengine/derived/` also include the installed model source/version,
core/wrapper versions, SPM selection and requested year. Legacy basename-only
caches cannot establish that identity and are not reused. Loading never rewrites
annual native files; no model-specific derived input file is necessary.
Annual output cache keys also include the selected dataset, frozen history pins,
and runtime identity, even when a caller reuses an explicit simulation ID.

US regional analysis keeps row filters over the annual national dataset.
Positional weight replacement cannot establish annual year/ID alignment, so the
wrapper rejects it for annual inputs. A separate certified alignment contract
would be necessary before supporting such an overlay.
19 changes: 17 additions & 2 deletions src/policyengine/core/simulation.py
Original file line number Diff line number Diff line change
Expand Up @@ -125,11 +125,26 @@ def spm_provenance(self) -> Optional[dict[str, Any]]:
@property
def storage_id(self) -> str:
"""Include resolved SPM settings in cache and saved-result identity."""
identity = self.id
history = getattr(self.dataset, "metadata", {}).get("annual_input_sources")
if history:
from policyengine.tax_benefit_models.us.datasets import (
_derived_cache_identity,
)

annual = {
"dataset_id": self.dataset.id,
"year": self.dataset.year,
"history": history,
"runtime": _derived_cache_identity(),
}
encoded = json.dumps(annual, sort_keys=True, separators=(",", ":")).encode()
identity += f"-annual-{hashlib.sha256(encoded).hexdigest()}"
config = self.spm_config
if config is None:
return self.id
return identity
encoded = json.dumps(config, sort_keys=True, separators=(",", ":")).encode()
return f"{self.id}-spm-{hashlib.sha256(encoded).hexdigest()}"
return f"{identity}-spm-{hashlib.sha256(encoded).hexdigest()}"

@model_validator(mode="after")
def _compile_dict_reforms(self) -> "Simulation":
Expand Down
12 changes: 11 additions & 1 deletion src/policyengine/provenance/certification.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@

from policyengine.provenance.manifest import (
HF_REQUEST_TIMEOUT_SECONDS,
CountryReleaseManifest,
DataReleaseManifest,
_specifier_matches,
fetch_pypi_wheel_metadata,
Expand Down Expand Up @@ -561,7 +562,7 @@ def build_country_manifest_payload(
model_build.data_build_fingerprint
)

return {
payload = {
"schema_version": 1,
"bundle_id": f"{country}-{policyengine_version}",
"country_id": country,
Expand All @@ -583,10 +584,19 @@ def build_country_manifest_payload(
"default_dataset": default_dataset,
"datasets": datasets,
"region_datasets": region_datasets,
**(
{"dataset_years": manifest.metadata["dataset_years"]}
if "dataset_years" in manifest.metadata
else {}
),
"certified_data_artifact": certified_artifact,
"certification": certification,
}

# Annual metadata must survive certification only with complete artifact pins.
CountryReleaseManifest.model_validate(payload)
return payload


def build_bundle_data_release_payload(
*,
Expand Down
23 changes: 22 additions & 1 deletion src/policyengine/provenance/dataset_materialization.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

from __future__ import annotations

import hashlib
import json
import os
import tempfile
from dataclasses import dataclass
Expand All @@ -16,6 +18,7 @@
from .manifest import (
CountryReleaseManifest,
_artifact_revision,
_dataset_for_year,
build_hf_uri,
dataset_logical_name,
get_release_manifest,
Expand Down Expand Up @@ -104,6 +107,18 @@ def _resolve_bundle_dataset(
f"Managed dataset {dataset_name!r} is missing a certified sha256."
)

identity = {
"repo_id": reference.repo_id or country_manifest.data_package.repo_id,
"repo_type": reference.repo_type or country_manifest.data_package.repo_type,
"path": reference.path,
"revision": reference.revision
or _artifact_revision(country_manifest.data_package),
"sha256": reference.sha256,
"metadata_sha256": reference.metadata_sha256,
}
cache_key = hashlib.sha256(
json.dumps(identity, sort_keys=True).encode()
).hexdigest()
return _BundleDatasetSpec(
country_id=country_id,
dataset=dataset_name,
Expand All @@ -116,7 +131,11 @@ def _resolve_bundle_dataset(
revision=reference.revision
or _artifact_revision(country_manifest.data_package),
sha256=reference.sha256,
destination=data_dir / Path(reference.path).name,
destination=data_dir
/ ".policyengine"
/ "sources"
/ cache_key
/ Path(reference.path).name,
metadata_sha256=reference.metadata_sha256,
)

Expand All @@ -127,10 +146,12 @@ def materialize_dataset(
*,
allow_unmanaged: bool = False,
data_dir: Path = DEFAULT_DATA_DIR,
year: Optional[int] = None,
) -> DatasetSource:
"""Select a dataset source and return the local file used for calculation."""

manifest = get_release_manifest(country_id)
dataset = _dataset_for_year(manifest, dataset, year)
if dataset is None or dataset == manifest.default_dataset_uri:
return _use_bundle_dataset(
country_id,
Expand Down
91 changes: 90 additions & 1 deletion src/policyengine/provenance/manifest.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
from urllib.parse import quote

import requests
from pydantic import BaseModel, Field
from pydantic import BaseModel, Field, model_validator

HF_REQUEST_TIMEOUT_SECONDS = 30
PYPI_REQUEST_TIMEOUT_SECONDS = 30
Expand Down Expand Up @@ -171,11 +171,52 @@ class CountryReleaseManifest(BaseModel):
default_dataset: str
datasets: dict[str, ArtifactPathReference] = Field(default_factory=dict)
region_datasets: dict[str, ArtifactPathTemplate] = Field(default_factory=dict)
dataset_years: dict[str, dict[int, str]] = Field(default_factory=dict)
certified_data_artifact: Optional[CertifiedDataArtifact] = None
certification: Optional[DataCertification] = None
source_sha256: Optional[str] = Field(default=None, exclude=True)
"""Byte sha256 of the bundled manifest before Pydantic normalization."""

@model_validator(mode="after")
def validate_dataset_years(self):
if self.dataset_years and self.country_id != "us":
raise ValueError(
"Annual native dataset families currently support only the US"
)
for family, years in self.dataset_years.items():
if family not in self.datasets or not years:
raise ValueError(
f"Annual dataset family {family!r} must name a dataset and contain years"
)
for year, name in years.items():
reference = self.datasets.get(name)
if year < 1 or reference is None:
raise ValueError(
f"Unknown annual dataset {name!r} for {family!r}, year {year}"
)
if not reference.path.endswith(".h5") or not reference.revision:
raise ValueError(
f"Annual dataset {name!r} requires an H5 path and explicit revision"
)
digest = reference.sha256 or ""
if len(digest) != 64 or any(
c not in "0123456789abcdef" for c in digest
):
raise ValueError(
f"Annual dataset {name!r} requires a SHA256 digest"
)
first_year = min(years)
if len(years) != max(years) - first_year + 1:
raise ValueError(
f"Annual dataset family {family!r} must contain contiguous years"
)
if years[first_year] != family:
raise ValueError(
f"The earliest year of annual dataset family {family!r} "
"must map to the family dataset"
)
return self

@property
def default_dataset_uri(self) -> str:
if (
Expand Down Expand Up @@ -516,6 +557,54 @@ def certify_data_release_compatibility(
)


def _dataset_years(
manifest: CountryReleaseManifest,
dataset: Optional[str],
*,
include_history: bool = False,
) -> dict[int, str]:
"""Return calculation coverage, or the parent family's input history."""
name = dataset or manifest.default_dataset
for artifact, reference in manifest.datasets.items():
if name == build_hf_uri(
reference.repo_id or manifest.data_package.repo_id,
reference.path,
reference.revision or _artifact_revision(manifest.data_package),
):
name = artifact
break
if name in manifest.dataset_years:
return manifest.dataset_years[name]
families = [
years for years in manifest.dataset_years.values() if name in years.values()
]
if len(families) > 1:
raise ValueError(f"Annual artifact {name!r} belongs to multiple families")
if not families:
return {}
if include_history:
return families[0]
return {
year: artifact for year, artifact in families[0].items() if artifact == name
}


def _dataset_for_year(
manifest: CountryReleaseManifest, dataset: Optional[str], year: Optional[int]
) -> Optional[str]:
years = _dataset_years(manifest, dataset)
if not years:
return dataset
selected = min(years) if year is None else year
if selected not in years:
raise ValueError(
f"Year {selected} is outside annual dataset coverage for "
f"{dataset or manifest.default_dataset!r}: {sorted(years)}. "
"The managed family cannot fall back to engine uprating."
)
return years[selected]


def resolve_dataset_reference(country_id: str, dataset: str) -> str:
if "://" in dataset:
return dataset
Expand Down
24 changes: 24 additions & 0 deletions src/policyengine/tax_benefit_models/common/model_version.py
Original file line number Diff line number Diff line change
Expand Up @@ -367,10 +367,25 @@ def save(self, simulation: Simulation) -> None:
raise ValueError(
"SPM settings changed since this output was calculated; run again before saving"
)
annual_sources = getattr(simulation.dataset, "metadata", {}).get(
"annual_input_sources"
)
if (
simulation.output_dataset.metadata.get("annual_input_sources")
!= annual_sources
):
raise ValueError(
"Annual input pins changed since this output was calculated; run again before saving"
)
serialized_spm = json.dumps(
{
"config": simulation.spm_config,
"provenance": receipt.model_dump(mode="json"),
**(
{"annual_input_sources": annual_sources}
if annual_sources
else {}
),
},
sort_keys=True,
)
Expand Down Expand Up @@ -413,6 +428,11 @@ def load(self, simulation: Simulation) -> None:
recorded = json.loads(raw)
if recorded["config"] != simulation.spm_config:
raise ValueError("Saved US simulation uses different SPM settings")
annual_sources = getattr(simulation.dataset, "metadata", {}).get(
"annual_input_sources"
)
if recorded.get("annual_input_sources") != annual_sources:
raise ValueError("Saved US simulation uses different annual input pins")
receipt = SPMProvenance.model_validate(recorded["provenance"])

simulation.output_dataset = self._dataset_class(
Expand All @@ -429,6 +449,10 @@ def load(self, simulation: Simulation) -> None:
simulation.spm_receipt = receipt
simulation.spm = SPMSelection.model_validate(recorded["config"])
simulation.output_dataset.metadata["spm_config"] = recorded["config"]
if recorded.get("annual_input_sources"):
simulation.output_dataset.metadata["annual_input_sources"] = recorded[
"annual_input_sources"
]

if os.path.exists(filepath):
simulation.created_at = datetime.datetime.fromtimestamp(
Expand Down
Loading
Loading