#!/usr/bin/env python3
"""Build a verified history from four WebTruffle federal-award releases.

The recipe pins and verifies the manifests before it trusts their asset
declarations, verifies the 66-field award schema and CSV headers, and loads all
source fields as text with DuckDB. It retains three distinct analytical grains:

* release membership: every award row present in every pinned release;
* content state: one row per distinct ``award_id`` and ``content_hash``; and
* latest observed state: one row per ``award_id`` in this release corpus.

Change status is independently reconstructed against persistent last-observed
state. It is not calculated by joining only adjacent files. Cumulative award
balances are checked at release-membership and latest-observed grain so an
appended balance cannot be mislabeled as period spending.

Requires Python 3.10+ and duckdb==1.5.5.

SPDX-License-Identifier: MIT
"""

from __future__ import annotations

import argparse
import csv
import hashlib
import json
import os
import platform
import re
import sys
import tempfile
from datetime import datetime, timezone
from decimal import Decimal, InvalidOperation
from pathlib import Path
from typing import Any
from urllib.request import Request, urlopen

import duckdb

REPOSITORY = "webtruffle/us-federal-contract-awards"
DATASET_ID = "us-federal-contract-awards"
EXPECTED_SCHEMA_VERSION = "1.0"
AWARDS_FILE = "us-federal-contract-awards.csv"
SCHEMA_FILE = "schema.json"
EXPECTED_AWARD_HEADER_COUNT = 66
EXPECTED_SCHEMA_BYTES = 19_271
EXPECTED_SCHEMA_SHA256 = (
    "12d2718f6f057d6a64dc09574652cfdb17625f7ab8434d6424f58013db73f1f6"
)
EXPECTED_AWARD_GRAIN = "one prime contract award summary per generated award identifier"
EXPECTED_FILE_GRAIN = "one award holder"
USER_AGENT = (
    "WebTruffle-US-federal-awards-history-Python-recipe/1.0 "
    "(+https://www.webtruffle.com/)"
)

RELEASES: tuple[dict[str, Any], ...] = (
    {
        "sequence": 1,
        "tag": "2026-08-25",
        "window_start": "2026-08-23",
        "window_end": "2026-08-25",
        "generated_at": "2026-08-26T14:01:08Z",
        "manifest_bytes": 11_153,
        "manifest_sha256": (
            "fcfe470fc3ef678319f88471849f8f4d448ceaafa406bd40584b876216443c7e"
        ),
        "awards_bytes": 20_403_764,
        "awards_sha256": (
            "bcf2f1131dcd7de8df89534dae3d5f4ef8c700418bd7b60d23a81aeb5e6a69fa"
        ),
        "record_count": 17_945,
        "supplier_count": 3_145,
        "relationship_count": 17_945,
        "change_counts": {"new": 17_945, "updated": 0, "unchanged": 0},
    },
    {
        "sequence": 2,
        "tag": "2026-08-26",
        "window_start": "2026-08-24",
        "window_end": "2026-08-26",
        "generated_at": "2026-08-27T14:02:44Z",
        "manifest_bytes": 11_160,
        "manifest_sha256": (
            "6d7e298d6982ad0f98c47e5ba73e7eed4bee011ea42ad88ad611b7c438c51728"
        ),
        "awards_bytes": 19_340_326,
        "awards_sha256": (
            "cb0f86800b523a62b1ed22bc9edf88e813b425d27b823c1e681f1cf6d6b0c632"
        ),
        "record_count": 17_920,
        "supplier_count": 4_362,
        "relationship_count": 17_920,
        "change_counts": {"new": 13_003, "updated": 81, "unchanged": 4_836},
    },
    {
        "sequence": 3,
        "tag": "2026-08-27",
        "window_start": "2026-08-25",
        "window_end": "2026-08-27",
        "generated_at": "2026-08-28T14:06:05Z",
        "manifest_bytes": 11_168,
        "manifest_sha256": (
            "aeef931ffd405bbff47bae6c10e94a1ce741eb51736978d69d86ecdac2f91ffc"
        ),
        "awards_bytes": 36_730_034,
        "awards_sha256": (
            "fbae9277a431631656d16ab806d64e5230eea1176ec218db37c19571da01ac5e"
        ),
        "record_count": 33_211,
        "supplier_count": 4_224,
        "relationship_count": 33_211,
        "change_counts": {"new": 19_900, "updated": 365, "unchanged": 12_946},
    },
    {
        "sequence": 4,
        "tag": "2026-08-29",
        "window_start": "2026-08-27",
        "window_end": "2026-08-29",
        "generated_at": "2026-08-30T13:59:47Z",
        "manifest_bytes": 11_162,
        "manifest_sha256": (
            "233229a8cb665da4c3b849f1b484cad8a789b935bbb70e19e4a2f30f354da28e"
        ),
        "awards_bytes": 36_471_909,
        "awards_sha256": (
            "657fb4c8629a1e865942c904d86183ad921381cdba9e4674546f0569331f67b8"
        ),
        "record_count": 32_313,
        "supplier_count": 5_788,
        "relationship_count": 32_313,
        "change_counts": {"new": 31_895, "updated": 418, "unchanged": 0},
    },
)

OUTPUT_NAMES = (
    "award-history-release-checkpoints.csv",
    "award-history-state-reconciliation.csv",
    "award-history-release-overlap.csv",
    "award-history-obligation-check.csv",
    "award-history-field-change-counts.csv",
    "award-history-update-review.csv",
    "award-history-latest-observed.csv",
)
PROVENANCE_NAME = "award-history-provenance.json"

EXPECTED_STATE_COUNTS = {
    "release_membership_rows": 101_389,
    "unique_award_ids": 82_743,
    "repeated_release_appearances": 18_646,
    "unique_content_states": 83_607,
    "unchanged_content_observations": 17_782,
    "update_observations": 864,
    "distinct_awards_ever_updated": 849,
    "first_observations": 82_743,
    "reconstructed_new": 82_743,
    "reconstructed_unchanged": 17_782,
    "reconstructed_updated": 864,
    "change_type_mismatches": 0,
    "updates_present_in_immediately_previous_release": 531,
    "updates_not_in_immediately_previous_release": 333,
    "missing_award_ids": 0,
    "missing_content_hashes": 0,
}

EXPECTED_OVERLAPS = {
    ("2026-08-25", "2026-08-26"): (4_917, 4_836, 81),
    ("2026-08-25", "2026-08-27"): (331, 0, 331),
    ("2026-08-25", "2026-08-29"): (145, 0, 145),
    ("2026-08-26", "2026-08-27"): (13_112, 12_946, 166),
    ("2026-08-26", "2026-08-29"): (199, 0, 199),
    ("2026-08-27", "2026-08-29"): (284, 0, 284),
}

MONEY_FIELDS = (
    "total_obligated_amount_usd",
    "total_outlay_amount_usd",
    "current_total_value_of_award_usd",
    "potential_total_value_of_award_usd",
)
EXPECTED_AMOUNT_CHECKS = {
    "total_obligated_amount_usd": {
        "membership_present": 101_389,
        "membership_total": Decimal("202845310918.86"),
        "latest_present": 82_743,
        "latest_total": Decimal("173612813504.62"),
        "repeated_balance": Decimal("29232497414.24"),
    },
    "total_outlay_amount_usd": {
        "membership_present": 24_555,
        "membership_total": Decimal("39392096509.13"),
        "latest_present": 21_412,
        "latest_total": Decimal("27666542124.45"),
        "repeated_balance": Decimal("11725554384.68"),
    },
    "current_total_value_of_award_usd": {
        "membership_present": 101_389,
        "membership_total": Decimal("235950561282.42"),
        "latest_present": 82_743,
        "latest_total": Decimal("194881219348.78"),
        "repeated_balance": Decimal("41069341933.64"),
    },
    "potential_total_value_of_award_usd": {
        "membership_present": 101_389,
        "membership_total": Decimal("337088327295.66"),
        "latest_present": 82_743,
        "latest_total": Decimal("281944398994.80"),
        "repeated_balance": Decimal("55143928300.86"),
    },
}
EXPECTED_UPDATE_OBLIGATION_CHECK = {
    "updates_with_higher_obligation": 123,
    "updates_with_lower_obligation": 200,
    "updates_with_unchanged_obligation": 541,
    "net_summary_obligation_difference_usd": Decimal("33512099.66"),
}

EXCLUDED_CHANGE_FIELDS = {
    "edition_date",
    "first_seen_at",
    "last_seen_at",
    "content_hash",
    "change_type",
}
EXPECTED_NONZERO_FIELD_CHANGES = {
    "last_modified_at": 863,
    "latest_action_date": 380,
    "total_obligated_amount_usd": 323,
    "current_total_value_of_award_usd": 315,
    "potential_total_value_of_award_usd": 305,
    "performance_current_end_date": 117,
    "performance_potential_end_date": 96,
    "description": 18,
    "funding_office_code": 15,
    "funding_office_name": 15,
    "place_city_name": 13,
    "place_state_code": 11,
    "place_state_name": 11,
    "place_county_name": 11,
    "solicitation_date": 6,
    "solicitation_identifier": 5,
    "performance_start_date": 4,
    "set_aside_code": 4,
    "set_aside": 4,
    "awarding_office_code": 3,
    "awarding_office_name": 3,
    "base_action_date": 3,
    "extent_competed_code": 2,
    "extent_competed": 2,
    "naics_code": 2,
    "naics_description": 2,
    "product_or_service_description": 2,
    "awarding_sub_agency_code": 2,
    "awarding_sub_agency_name": 2,
    "supplier_id": 1,
    "recipient_uei": 1,
    "recipient_cage_code": 1,
    "recipient_name": 1,
    "solicitation_procedures_code": 1,
    "solicitation_procedures": 1,
    "funding_sub_agency_code": 1,
    "funding_sub_agency_name": 1,
    "business_size_code": 1,
    "business_size": 1,
    "product_or_service_code": 1,
    "number_of_actions": 1,
}

DECIMAL_TEXT = re.compile(r"^-?(?:0|[1-9]\d*)(?:\.\d+)?$")
SAFE_FIELD_NAME = re.compile(r"^[a-z][a-z0-9_]*$")


def positive_review_limit(value: str) -> int:
    try:
        parsed = int(value)
    except ValueError as error:
        raise argparse.ArgumentTypeError("must be an integer") from error
    if not 1 <= parsed <= 1_000:
        raise argparse.ArgumentTypeError("must be between 1 and 1000")
    return parsed


def parse_args() -> argparse.Namespace:
    parser = argparse.ArgumentParser(
        description=(
            "Verify four pinned US federal award releases and build persistent "
            "latest-observed and change-history evidence."
        )
    )
    parser.add_argument(
        "--data-dir",
        type=Path,
        default=Path("us-federal-contract-awards-history-2026-08-25-to-2026-08-29"),
        help="Directory for verified release inputs and deterministic outputs.",
    )
    parser.add_argument(
        "--update-review-limit",
        type=positive_review_limit,
        default=200,
        help="Maximum update-review rows to export (1-1000; default: 200).",
    )
    return parser.parse_args()


def sha256_and_size(path: Path) -> tuple[str, int]:
    digest = hashlib.sha256()
    size = 0
    with path.open("rb") as source:
        for chunk in iter(lambda: source.read(1024 * 1024), b""):
            digest.update(chunk)
            size += len(chunk)
    return digest.hexdigest(), size


def assert_file(
    path: Path,
    expected_sha256: str,
    expected_bytes: int,
) -> tuple[str, int]:
    actual_sha256, actual_bytes = sha256_and_size(path)
    if actual_bytes != expected_bytes or actual_sha256 != expected_sha256:
        raise RuntimeError(
            f"Verification failed for {path}: expected {expected_bytes} bytes / "
            f"{expected_sha256}, got {actual_bytes} bytes / {actual_sha256}. "
            "Move or delete the unexpected file before retrying."
        )
    return actual_sha256, actual_bytes


def download_missing(
    url: str,
    destination: Path,
    expected_sha256: str,
    expected_bytes: int,
) -> None:
    if destination.exists():
        assert_file(destination, expected_sha256, expected_bytes)
        print(f"verified existing  {destination.parent.name}/{destination.name}")
        return

    request = Request(url, headers={"User-Agent": USER_AGENT, "Accept": "*/*"})
    descriptor, temporary_name = tempfile.mkstemp(
        dir=destination.parent,
        prefix=f".{destination.name}.",
        suffix=".part",
    )
    partial = Path(temporary_name)
    try:
        downloaded_bytes = 0
        digest = hashlib.sha256()
        with (
            urlopen(request, timeout=120) as response,
            os.fdopen(descriptor, "wb") as target,
        ):
            while chunk := response.read(1024 * 1024):
                downloaded_bytes += len(chunk)
                if downloaded_bytes > expected_bytes:
                    raise RuntimeError(
                        f"Download exceeded {expected_bytes} declared bytes for "
                        f"{destination.name}."
                    )
                digest.update(chunk)
                target.write(chunk)
        if downloaded_bytes != expected_bytes or digest.hexdigest() != expected_sha256:
            raise RuntimeError(
                f"Verification failed for downloaded {destination.name}: expected "
                f"{expected_bytes} bytes / {expected_sha256}, got {downloaded_bytes} "
                f"bytes / {digest.hexdigest()}."
            )
        os.replace(partial, destination)
    finally:
        try:
            os.close(descriptor)
        except OSError:
            pass
        partial.unlink(missing_ok=True)
    print(f"downloaded + verified {destination.parent.name}/{destination.name}")


def asset_base_url(tag: str) -> str:
    return f"https://github.com/{REPOSITORY}/releases/download/{tag}"


def validate_manifest(
    release: dict[str, Any],
    manifest: dict[str, Any],
) -> None:
    tag = release["tag"]
    exact_values = {
        "schema_version": EXPECTED_SCHEMA_VERSION,
        "dataset_id": DATASET_ID,
        "target_date": tag,
        "edition_date": tag,
        "coverage_start_date": release["window_start"],
        "coverage_end_date": release["window_end"],
        "coverage_date_type": "last_modified_date",
        "record_count": release["record_count"],
        "supplier_count": release["supplier_count"],
        "relationship_count": release["relationship_count"],
        "generated_at": release["generated_at"],
    }
    observed_values = {key: manifest.get(key) for key in exact_values}
    if observed_values != exact_values:
        raise RuntimeError(
            f"Pinned manifest values changed for {tag}: expected {exact_values}, "
            f"got {observed_values}."
        )
    if manifest.get("change_counts") != release["change_counts"]:
        raise RuntimeError(
            f"Pinned change counts changed for {tag}: "
            f"{manifest.get('change_counts')!r}."
        )

    product = manifest.get("products", {}).get("awards", {})
    if int(product.get("record_count", -1)) != release["record_count"]:
        raise RuntimeError(f"Pinned awards product count changed for {tag}.")
    if product.get("grain") != EXPECTED_AWARD_GRAIN:
        raise RuntimeError(
            f"Pinned awards product grain changed for {tag}: {product.get('grain')!r}."
        )

    awards_declaration = manifest.get("files", {}).get(AWARDS_FILE, {})
    expected_awards_declaration = {
        "bytes": release["awards_bytes"],
        "sha256": release["awards_sha256"],
        "record_count": release["record_count"],
        "grain": EXPECTED_FILE_GRAIN,
    }
    observed_awards_declaration = {
        key: awards_declaration.get(key) for key in expected_awards_declaration
    }
    if observed_awards_declaration != expected_awards_declaration:
        raise RuntimeError(
            f"Pinned awards declaration changed for {tag}: "
            f"{observed_awards_declaration}."
        )

    schema_declaration = manifest.get("files", {}).get(SCHEMA_FILE, {})
    expected_schema_declaration = {
        "bytes": EXPECTED_SCHEMA_BYTES,
        "sha256": EXPECTED_SCHEMA_SHA256,
    }
    observed_schema_declaration = {
        key: schema_declaration.get(key) for key in expected_schema_declaration
    }
    if observed_schema_declaration != expected_schema_declaration:
        raise RuntimeError(
            f"Pinned schema declaration changed for {tag}: "
            f"{observed_schema_declaration}."
        )

    sources = manifest.get("sources", [])
    if len(sources) != 1 or sources[0].get("source") != "usaspending":
        raise RuntimeError(f"Unexpected source declaration for {tag}: {sources!r}.")
    if sources[0].get("status") != "succeeded":
        raise RuntimeError(f"Pinned source did not succeed for {tag}: {sources!r}.")


def load_release_manifest(
    release_dir: Path,
    release: dict[str, Any],
) -> dict[str, Any]:
    tag = release["tag"]
    manifest_path = release_dir / "manifest.json"
    download_missing(
        f"{asset_base_url(tag)}/manifest.json",
        manifest_path,
        release["manifest_sha256"],
        release["manifest_bytes"],
    )
    manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
    validate_manifest(release, manifest)
    return manifest


def download_declared_release_assets(
    release_dir: Path,
    release: dict[str, Any],
    manifest: dict[str, Any],
) -> None:
    tag = release["tag"]
    for file_name in (SCHEMA_FILE, AWARDS_FILE):
        declaration = manifest["files"][file_name]
        download_missing(
            f"{asset_base_url(tag)}/{file_name}",
            release_dir / file_name,
            declaration["sha256"],
            int(declaration["bytes"]),
        )


def read_csv_header(path: Path) -> list[str]:
    with path.open("r", encoding="utf-8-sig", newline="") as source:
        return next(csv.reader(source))


def validate_schema_and_header(
    release_dir: Path,
    manifest: dict[str, Any],
) -> list[str]:
    schema = json.loads((release_dir / SCHEMA_FILE).read_text(encoding="utf-8"))
    product_refs = schema.get("x-webtruffle-products", {})
    expected_refs = {
        "awards": {"$ref": "#/$defs/award"},
        "award-suppliers": {"$ref": "#/$defs/award-supplier"},
        "suppliers": {"$ref": "#/$defs/supplier"},
    }
    if product_refs != expected_refs:
        raise RuntimeError(f"Unexpected schema product references: {product_refs!r}.")

    award_definition = schema.get("$defs", {}).get("award", {})
    headers = list(award_definition.get("properties", {}))
    if len(headers) != EXPECTED_AWARD_HEADER_COUNT or len(headers) != len(set(headers)):
        raise RuntimeError(
            f"Expected {EXPECTED_AWARD_HEADER_COUNT} unique award schema fields, "
            f"got {len(headers)}."
        )
    if not all(SAFE_FIELD_NAME.fullmatch(field) for field in headers):
        raise RuntimeError("The pinned schema contains an unsafe SQL field name.")
    required = set(award_definition.get("required", []))
    if not required <= set(headers):
        raise RuntimeError(
            f"Award schema required fields are missing: {sorted(required - set(headers))}."
        )

    csv_headers = read_csv_header(release_dir / AWARDS_FILE)
    if csv_headers != headers:
        raise RuntimeError(
            "Award CSV headers do not exactly match the verified schema order: "
            f"missing={sorted(set(headers) - set(csv_headers))}, "
            f"extra={sorted(set(csv_headers) - set(headers))}."
        )
    if manifest.get("record_fields") != headers:
        raise RuntimeError("Manifest record_fields do not match the verified header.")
    return headers


def validate_csv_rows(
    path: Path,
    release: dict[str, Any],
    headers: list[str],
) -> dict[str, Any]:
    tag = release["tag"]
    award_ids: set[str] = set()
    change_counts = {"new": 0, "updated": 0, "unchanged": 0}
    money_evidence = {
        field: {"present": 0, "missing": 0, "maximum_scale": 0}
        for field in MONEY_FIELDS
    }
    row_count = 0
    with path.open("r", encoding="utf-8-sig", newline="") as source:
        reader = csv.DictReader(source)
        if reader.fieldnames != headers:
            raise RuntimeError(f"Header changed while reading {path}.")
        for row_number, row in enumerate(reader, start=2):
            row_count += 1
            award_id = row["award_id"]
            if not award_id:
                raise RuntimeError(f"Missing award_id at {path}:{row_number}.")
            if award_id in award_ids:
                raise RuntimeError(f"Duplicate award_id {award_id!r} in {path}.")
            award_ids.add(award_id)
            if not row["content_hash"]:
                raise RuntimeError(f"Missing content_hash at {path}:{row_number}.")
            if row["edition_date"] != tag or row["source"] != "usaspending":
                raise RuntimeError(
                    f"Wrong edition/source at {path}:{row_number}: "
                    f"{row['edition_date']!r}/{row['source']!r}."
                )
            if not row["last_modified_at"]:
                raise RuntimeError(f"Missing last_modified_at at {path}:{row_number}.")
            change_type = row["change_type"]
            if change_type not in change_counts:
                raise RuntimeError(
                    f"Unexpected change_type at {path}:{row_number}: {change_type!r}."
                )
            change_counts[change_type] += 1

            for field in MONEY_FIELDS:
                text = row[field]
                if text == "":
                    money_evidence[field]["missing"] += 1
                    continue
                if not DECIMAL_TEXT.fullmatch(text):
                    raise RuntimeError(
                        f"Non-canonical decimal at {path}:{row_number} "
                        f"field {field}: {text!r}."
                    )
                try:
                    value = Decimal(text)
                except InvalidOperation as error:
                    raise RuntimeError(
                        f"Invalid decimal at {path}:{row_number} field {field}."
                    ) from error
                if not value.is_finite():
                    raise RuntimeError(
                        f"Non-finite decimal at {path}:{row_number} field {field}."
                    )
                scale = max(0, -value.as_tuple().exponent)
                if scale > 2:
                    raise RuntimeError(
                        f"Money exceeds two decimal places at {path}:{row_number} "
                        f"field {field}: {text!r}."
                    )
                money_evidence[field]["present"] += 1
                money_evidence[field]["maximum_scale"] = max(
                    money_evidence[field]["maximum_scale"], scale
                )

    if row_count != release["record_count"]:
        raise RuntimeError(
            f"Pinned row count changed for {tag}: expected {release['record_count']}, "
            f"got {row_count}."
        )
    if change_counts != release["change_counts"]:
        raise RuntimeError(
            f"CSV change counts changed for {tag}: expected "
            f"{release['change_counts']}, got {change_counts}."
        )
    return {
        "rows": row_count,
        "unique_award_ids_within_release": len(award_ids),
        "change_counts": change_counts,
        "money_fields": money_evidence,
    }


def quote_identifier(identifier: str) -> str:
    if not SAFE_FIELD_NAME.fullmatch(identifier):
        raise ValueError(f"Unsafe field name: {identifier!r}.")
    return f'"{identifier}"'


def register_release_membership(
    connection: duckdb.DuckDBPyConnection,
    releases_dir: Path,
    headers: list[str],
) -> dict[str, Any]:
    first = True
    for release in RELEASES:
        tag = release["tag"]
        source_view = f"release_source_{release['sequence']}"
        connection.read_csv(
            str(releases_dir / tag / AWARDS_FILE),
            header=True,
            all_varchar=True,
        ).create_view(source_view)
        source_info = connection.execute(
            f"PRAGMA table_info('{source_view}')"
        ).fetchall()
        if [row[1] for row in source_info] != headers:
            raise RuntimeError(f"DuckDB column order changed for release {tag}.")
        source_types = {row[2] for row in source_info}
        if source_types != {"VARCHAR"}:
            raise RuntimeError(
                f"Release {tag} contains inferred non-VARCHAR source fields: "
                f"{sorted(source_types)}."
            )

        select_sql = f"""
        SELECT
          {release["sequence"]}::BIGINT AS release_sequence,
          '{tag}'::VARCHAR AS release_tag,
          '{release["generated_at"]}'::VARCHAR AS release_generated_at,
          *
        FROM {source_view}
        """
        if first:
            connection.execute(f"CREATE TEMP TABLE release_membership AS {select_sql}")
            first = False
        else:
            connection.execute(f"INSERT INTO release_membership {select_sql}")
        connection.execute(f"DROP VIEW {source_view}")

    membership_info = connection.execute(
        "PRAGMA table_info('release_membership')"
    ).fetchall()
    membership_types = {row[1]: row[2] for row in membership_info}
    non_text_source_fields = {
        field: membership_types.get(field)
        for field in headers
        if membership_types.get(field) != "VARCHAR"
    }
    if non_text_source_fields:
        raise RuntimeError(
            f"Materialized source fields are not all VARCHAR: {non_text_source_fields}."
        )
    if membership_types.get("release_sequence") != "BIGINT":
        raise RuntimeError("release_sequence was not materialized as BIGINT.")
    return {
        "membership_columns": len(membership_info),
        "source_field_count": len(headers),
        "source_field_type": "VARCHAR",
        "release_sequence_type": "BIGINT",
    }


def create_release_metadata_table(connection: duckdb.DuckDBPyConnection) -> None:
    connection.execute(
        """
        CREATE TEMP TABLE release_metadata (
          release_sequence BIGINT,
          release_tag VARCHAR,
          release_url VARCHAR,
          coverage_start_date VARCHAR,
          coverage_end_date VARCHAR,
          generated_at VARCHAR,
          manifest_bytes BIGINT,
          manifest_sha256 VARCHAR,
          awards_bytes BIGINT,
          awards_sha256 VARCHAR,
          schema_bytes BIGINT,
          schema_sha256 VARCHAR,
          declared_award_rows BIGINT,
          declared_supplier_rows BIGINT,
          declared_relationship_rows BIGINT,
          declared_new_rows BIGINT,
          declared_updated_rows BIGINT,
          declared_unchanged_rows BIGINT
        )
        """
    )
    rows = []
    for release in RELEASES:
        tag = release["tag"]
        rows.append(
            (
                release["sequence"],
                tag,
                f"https://github.com/{REPOSITORY}/releases/tag/{tag}",
                release["window_start"],
                release["window_end"],
                release["generated_at"],
                release["manifest_bytes"],
                release["manifest_sha256"],
                release["awards_bytes"],
                release["awards_sha256"],
                EXPECTED_SCHEMA_BYTES,
                EXPECTED_SCHEMA_SHA256,
                release["record_count"],
                release["supplier_count"],
                release["relationship_count"],
                release["change_counts"]["new"],
                release["change_counts"]["updated"],
                release["change_counts"]["unchanged"],
            )
        )
    placeholders = ", ".join("?" for _ in range(18))
    connection.executemany(
        f"INSERT INTO release_metadata VALUES ({placeholders})",
        rows,
    )


def create_state_tables(
    connection: duckdb.DuckDBPyConnection,
    headers: list[str],
) -> None:
    source_columns = ", ".join(quote_identifier(field) for field in headers)
    connection.execute(
        """
        CREATE TEMP TABLE membership_reconstruction AS
        WITH ordered AS (
          SELECT
            *,
            row_number() OVER (
              PARTITION BY award_id ORDER BY release_sequence
            ) AS observation_number,
            lag(release_sequence) OVER (
              PARTITION BY award_id ORDER BY release_sequence
            ) AS previous_release_sequence,
            lag(content_hash) OVER (
              PARTITION BY award_id ORDER BY release_sequence
            ) AS previous_content_hash
          FROM release_membership
        )
        SELECT
          *,
          CASE
            WHEN observation_number = 1 THEN 'new'
            WHEN content_hash = previous_content_hash THEN 'unchanged'
            ELSE 'updated'
          END AS reconstructed_change_type
        FROM ordered
        """
    )
    connection.execute(
        f"""
        CREATE TEMP TABLE content_states AS
        SELECT
          release_sequence,
          release_tag,
          release_generated_at,
          {source_columns}
        FROM membership_reconstruction
        QUALIFY row_number() OVER (
          PARTITION BY award_id, content_hash ORDER BY release_sequence
        ) = 1
        """
    )
    connection.execute(
        f"""
        CREATE TEMP TABLE latest_observed AS
        SELECT
          release_sequence,
          release_tag,
          release_generated_at,
          {source_columns}
        FROM membership_reconstruction
        QUALIFY row_number() OVER (
          PARTITION BY award_id ORDER BY release_sequence DESC
        ) = 1
        """
    )
    connection.execute(
        """
        CREATE TEMP TABLE award_state_summary AS
        SELECT
          award_id,
          count(*) AS release_observation_count,
          count(DISTINCT content_hash) AS content_state_count,
          min(release_sequence) AS first_release_sequence,
          max(release_sequence) AS last_release_sequence,
          arg_min(release_tag, release_sequence) AS first_observed_release_tag,
          arg_max(release_tag, release_sequence) AS last_observed_release_tag
        FROM release_membership
        GROUP BY award_id
        """
    )


def create_update_pairs(
    connection: duckdb.DuckDBPyConnection,
    comparison_fields: list[str],
) -> None:
    paired_fields = []
    for field in comparison_fields:
        quoted = quote_identifier(field)
        paired_fields.extend(
            (
                f"previous.{quoted} AS previous__{field}",
                f"current.{quoted} AS current__{field}",
            )
        )
    paired_sql = ",\n          ".join(paired_fields)
    connection.execute(
        f"""
        CREATE TEMP TABLE update_pairs AS
        SELECT
          current.release_sequence AS current_release_sequence,
          current.release_tag AS current_release_tag,
          previous.release_sequence AS previous_release_sequence,
          previous.release_tag AS previous_release_tag,
          {paired_sql}
        FROM membership_reconstruction AS current
        JOIN release_membership AS previous
          ON previous.award_id = current.award_id
         AND previous.release_sequence = current.previous_release_sequence
        WHERE current.reconstructed_change_type = 'updated'
        """
    )


def fetch_named_row(
    connection: duckdb.DuckDBPyConnection,
    sql: str,
) -> dict[str, Any]:
    cursor = connection.execute(sql)
    row = cursor.fetchone()
    if row is None:
        raise RuntimeError("Expected one validation row, got none.")
    return dict(zip((column[0] for column in cursor.description), row))


STATE_RECONCILIATION_SQL = """
SELECT
  (SELECT count(*) FROM release_membership) AS release_membership_rows,
  (SELECT count(DISTINCT award_id) FROM release_membership) AS unique_award_ids,
  (SELECT count(*) - count(DISTINCT award_id) FROM release_membership)
    AS repeated_release_appearances,
  (SELECT count(*) FROM content_states) AS unique_content_states,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'unchanged')
    AS unchanged_content_observations,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'updated') AS update_observations,
  (SELECT count(DISTINCT award_id) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'updated')
    AS distinct_awards_ever_updated,
  (SELECT count(*) FROM membership_reconstruction
    WHERE observation_number = 1) AS first_observations,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'new') AS reconstructed_new,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'unchanged') AS reconstructed_unchanged,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'updated') AS reconstructed_updated,
  (SELECT count(*) FROM membership_reconstruction
    WHERE change_type IS DISTINCT FROM reconstructed_change_type)
    AS change_type_mismatches,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'updated'
      AND previous_release_sequence = release_sequence - 1)
    AS updates_present_in_immediately_previous_release,
  (SELECT count(*) FROM membership_reconstruction
    WHERE reconstructed_change_type = 'updated'
      AND previous_release_sequence < release_sequence - 1)
    AS updates_not_in_immediately_previous_release,
  (SELECT count(*) FROM release_membership
    WHERE nullif(trim(award_id), '') IS NULL) AS missing_award_ids,
  (SELECT count(*) FROM release_membership
    WHERE nullif(trim(content_hash), '') IS NULL) AS missing_content_hashes,
  'release membership' AS release_membership_grain,
  'award_id + content_hash' AS content_state_grain,
  'latest observed release row per award_id' AS latest_observed_grain
""".strip()


RELEASE_CHECKPOINTS_SQL = """
SELECT
  metadata.release_sequence,
  metadata.release_tag,
  metadata.release_url,
  metadata.coverage_start_date,
  metadata.coverage_end_date,
  metadata.generated_at,
  metadata.manifest_bytes,
  metadata.manifest_sha256,
  metadata.awards_bytes,
  metadata.awards_sha256,
  metadata.schema_bytes,
  metadata.schema_sha256,
  metadata.declared_award_rows,
  count(membership.award_id) AS observed_award_rows,
  count(DISTINCT membership.award_id) AS observed_unique_award_ids,
  metadata.declared_supplier_rows,
  metadata.declared_relationship_rows,
  metadata.declared_new_rows,
  count_if(membership.change_type = 'new') AS observed_new_rows,
  metadata.declared_updated_rows,
  count_if(membership.change_type = 'updated') AS observed_updated_rows,
  metadata.declared_unchanged_rows,
  count_if(membership.change_type = 'unchanged') AS observed_unchanged_rows
FROM release_metadata AS metadata
JOIN release_membership AS membership USING (release_sequence, release_tag)
GROUP BY ALL
ORDER BY metadata.release_sequence
""".strip()


RELEASE_OVERLAP_SQL = """
WITH release_counts AS (
  SELECT release_sequence, release_tag, count(*) AS release_rows
  FROM release_membership
  GROUP BY release_sequence, release_tag
)
SELECT
  earlier.release_sequence AS earlier_release_sequence,
  earlier.release_tag AS earlier_release_tag,
  later.release_sequence AS later_release_sequence,
  later.release_tag AS later_release_tag,
  earlier_count.release_rows AS earlier_release_rows,
  later_count.release_rows AS later_release_rows,
  count(*) AS shared_award_ids,
  count_if(earlier.content_hash = later.content_hash) AS same_content_hash,
  count_if(earlier.content_hash <> later.content_hash) AS changed_content_hash,
  round(count(*) * 100.0 / later_count.release_rows, 2)
    AS shared_share_of_later_release_percent,
  later.release_sequence = earlier.release_sequence + 1
    AS adjacent_published_releases
FROM release_membership AS earlier
JOIN release_membership AS later
  ON later.award_id = earlier.award_id
 AND later.release_sequence > earlier.release_sequence
JOIN release_counts AS earlier_count
  ON earlier_count.release_sequence = earlier.release_sequence
JOIN release_counts AS later_count
  ON later_count.release_sequence = later.release_sequence
GROUP BY
  earlier.release_sequence,
  earlier.release_tag,
  later.release_sequence,
  later.release_tag,
  earlier_count.release_rows,
  later_count.release_rows
ORDER BY earlier.release_sequence, later.release_sequence
""".strip()


OBLIGATION_CHECK_SQL = """
SELECT
  (SELECT count(*) FROM release_membership) AS release_membership_rows,
  (SELECT count(*) FROM release_membership
    WHERE nullif(total_obligated_amount_usd, '') IS NOT NULL)
    AS release_membership_obligation_present,
  (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
    FROM release_membership) AS release_membership_cumulative_obligation_usd,
  (SELECT count(*) FROM latest_observed) AS latest_observed_rows,
  (SELECT count(*) FROM latest_observed
    WHERE nullif(total_obligated_amount_usd, '') IS NOT NULL)
    AS latest_observed_obligation_present,
  (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
    FROM latest_observed) AS latest_observed_cumulative_obligation_usd,
  (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
    FROM release_membership)
    -
  (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
    FROM latest_observed) AS repeated_older_snapshot_balance_usd,
  round(
    (
      (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
        FROM release_membership)
      -
      (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
        FROM latest_observed)
    ) * 100.0
    /
    (SELECT sum(try_cast(total_obligated_amount_usd AS DECIMAL(38, 2)))
      FROM latest_observed),
    2
  ) AS raw_append_above_latest_observed_percent,
  (SELECT count(*) FROM update_pairs
    WHERE try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
      > try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2)))
    AS updates_with_higher_obligation,
  (SELECT count(*) FROM update_pairs
    WHERE try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
      < try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2)))
    AS updates_with_lower_obligation,
  (SELECT count(*) FROM update_pairs
    WHERE try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
      = try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2)))
    AS updates_with_unchanged_obligation,
  (SELECT sum(
      try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
      - try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2))
    ) FROM update_pairs) AS net_summary_obligation_difference_usd,
  'Cumulative award-summary balances; not period obligations or cash paid'
    AS interpretation
""".strip()


def build_field_change_counts_sql(comparison_fields: list[str]) -> str:
    parts = []
    for order, field in enumerate(comparison_fields, start=1):
        parts.append(
            f"""
            SELECT
              {order} AS schema_field_order,
              '{field}' AS field_name,
              count_if(previous__{field} IS DISTINCT FROM current__{field})
                AS changed_update_observations,
              count(*) AS total_update_observations,
              round(
                count_if(previous__{field} IS DISTINCT FROM current__{field})
                * 100.0 / count(*),
                2
              ) AS changed_share_of_updates_percent
            FROM update_pairs
            """.strip()
        )
    return (
        "SELECT * FROM (\n"
        + "\nUNION ALL\n".join(parts)
        + "\n) AS field_counts\nORDER BY schema_field_order"
    )


def build_update_review_sql(
    comparison_fields: list[str],
    review_limit: int,
) -> str:
    changed_names = ",\n      ".join(
        f"CASE WHEN previous__{field} IS DISTINCT FROM current__{field} "
        f"THEN '{field}' END"
        for field in comparison_fields
    )
    changed_count = " + ".join(
        f"(previous__{field} IS DISTINCT FROM current__{field})::INTEGER"
        for field in comparison_fields
    )
    return f"""
    SELECT
      current_release_tag,
      previous_release_tag,
      current_release_sequence - previous_release_sequence
        AS published_release_sequence_gap,
      current__award_id AS award_id,
      current__piid AS piid,
      current__recipient_name AS recipient_name,
      current__awarding_agency_name AS awarding_agency_name,
      previous__last_modified_at AS previous_last_modified_at,
      current__last_modified_at AS current_last_modified_at,
      ({changed_count}) AS changed_field_count,
      concat_ws('|',
        {changed_names}
      ) AS changed_fields,
      previous__total_obligated_amount_usd
        AS previous_total_obligated_amount_usd,
      current__total_obligated_amount_usd
        AS current_total_obligated_amount_usd,
      try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
        - try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2))
        AS net_summary_obligation_difference_usd,
      previous__latest_action_date AS previous_latest_action_date,
      current__latest_action_date AS current_latest_action_date,
      previous__performance_current_end_date
        AS previous_performance_current_end_date,
      current__performance_current_end_date AS current_performance_current_end_date,
      current__naics_code AS current_naics_code,
      current__product_or_service_code AS current_product_or_service_code,
      current__description AS current_description,
      current__source_url AS official_source_url
    FROM update_pairs
    ORDER BY
      current_release_sequence DESC,
      current__last_modified_at DESC,
      current__award_id
    LIMIT {review_limit}
    """.strip()


def build_latest_observed_sql(headers: list[str]) -> str:
    source_fields = ",\n  ".join(
        f"latest.{quote_identifier(field)}" for field in headers
    )
    return f"""
    SELECT
      summary.first_observed_release_tag,
      latest.release_tag AS last_observed_release_tag,
      summary.release_observation_count,
      summary.content_state_count,
      {source_fields}
    FROM latest_observed AS latest
    JOIN award_state_summary AS summary USING (award_id)
    ORDER BY latest.award_id
    """.strip()


def assert_state_counts(connection: duckdb.DuckDBPyConnection) -> dict[str, Any]:
    checks = fetch_named_row(connection, STATE_RECONCILIATION_SQL)
    numeric_checks = {key: checks[key] for key in EXPECTED_STATE_COUNTS}
    if numeric_checks != EXPECTED_STATE_COUNTS:
        raise RuntimeError(
            f"Persistent-state checkpoints failed: expected {EXPECTED_STATE_COUNTS}, "
            f"got {numeric_checks}."
        )
    latest_rows = connection.execute("SELECT count(*) FROM latest_observed").fetchone()[
        0
    ]
    if latest_rows != EXPECTED_STATE_COUNTS["unique_award_ids"]:
        raise RuntimeError(
            f"Latest-observed grain has {latest_rows} rows, expected 82,743."
        )
    return checks


def assert_release_overlaps(
    connection: duckdb.DuckDBPyConnection,
) -> dict[str, dict[str, int]]:
    cursor = connection.execute(RELEASE_OVERLAP_SQL)
    columns = [column[0] for column in cursor.description]
    observed: dict[tuple[str, str], tuple[int, int, int]] = {}
    evidence: dict[str, dict[str, int]] = {}
    for raw_row in cursor.fetchall():
        row = dict(zip(columns, raw_row))
        key = (row["earlier_release_tag"], row["later_release_tag"])
        values = (
            row["shared_award_ids"],
            row["same_content_hash"],
            row["changed_content_hash"],
        )
        observed[key] = values
        evidence[f"{key[0]}->{key[1]}"] = {
            "shared_award_ids": values[0],
            "same_content_hash": values[1],
            "changed_content_hash": values[2],
        }
    if observed != EXPECTED_OVERLAPS:
        raise RuntimeError(
            f"Release-overlap checkpoints failed: expected {EXPECTED_OVERLAPS}, "
            f"got {observed}."
        )
    return evidence


def validate_amount_semantics(
    connection: duckdb.DuckDBPyConnection,
) -> dict[str, Any]:
    evidence: dict[str, Any] = {}
    for field in MONEY_FIELDS:
        row = connection.execute(
            f"""
            SELECT
              (SELECT count(*) FROM release_membership
                WHERE nullif({field}, '') IS NOT NULL),
              (SELECT count(*) FROM release_membership
                WHERE nullif({field}, '') IS NOT NULL
                  AND try_cast({field} AS DECIMAL(38, 2)) IS NULL),
              (SELECT sum(try_cast({field} AS DECIMAL(38, 2)))
                FROM release_membership),
              (SELECT count(*) FROM latest_observed
                WHERE nullif({field}, '') IS NOT NULL),
              (SELECT count(*) FROM latest_observed
                WHERE nullif({field}, '') IS NOT NULL
                  AND try_cast({field} AS DECIMAL(38, 2)) IS NULL),
              (SELECT sum(try_cast({field} AS DECIMAL(38, 2)))
                FROM latest_observed)
            """
        ).fetchone()
        if row is None:
            raise RuntimeError(f"No amount validation row for {field}.")
        (
            membership_present,
            membership_invalid,
            membership_total,
            latest_present,
            latest_invalid,
            latest_total,
        ) = row
        if not isinstance(membership_total, Decimal) or not isinstance(
            latest_total, Decimal
        ):
            raise TypeError(f"DuckDB did not return {field} totals as Decimal.")
        observed = {
            "membership_present": membership_present,
            "membership_total": membership_total,
            "latest_present": latest_present,
            "latest_total": latest_total,
            "repeated_balance": membership_total - latest_total,
        }
        if membership_invalid or latest_invalid:
            raise RuntimeError(
                f"DuckDB DECIMAL validation failed for {field}: "
                f"membership_invalid={membership_invalid}, "
                f"latest_invalid={latest_invalid}."
            )
        if observed != EXPECTED_AMOUNT_CHECKS[field]:
            raise RuntimeError(
                f"Pinned amount checkpoints failed for {field}: expected "
                f"{EXPECTED_AMOUNT_CHECKS[field]}, got {observed}."
            )
        evidence[field] = {
            key: str(value) if isinstance(value, Decimal) else value
            for key, value in observed.items()
        }
        evidence[field]["duckdb_type"] = "DECIMAL(38, 2)"
    update_row = fetch_named_row(
        connection,
        """
        SELECT
          count_if(
            try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
              > try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2))
          ) AS updates_with_higher_obligation,
          count_if(
            try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
              < try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2))
          ) AS updates_with_lower_obligation,
          count_if(
            try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
              = try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2))
          ) AS updates_with_unchanged_obligation,
          sum(
            try_cast(current__total_obligated_amount_usd AS DECIMAL(38, 2))
              - try_cast(previous__total_obligated_amount_usd AS DECIMAL(38, 2))
          ) AS net_summary_obligation_difference_usd
        FROM update_pairs
        """,
    )
    if update_row != EXPECTED_UPDATE_OBLIGATION_CHECK:
        raise RuntimeError(
            "Update-obligation checkpoints failed: expected "
            f"{EXPECTED_UPDATE_OBLIGATION_CHECK}, got {update_row}."
        )
    evidence["update_obligation_differences"] = {
        key: str(value) if isinstance(value, Decimal) else value
        for key, value in update_row.items()
    }
    evidence["update_obligation_differences"]["interpretation"] = (
        "Net differences between cumulative summaries; not transaction amounts."
    )
    return evidence


def assert_field_change_counts(
    connection: duckdb.DuckDBPyConnection,
    field_change_sql: str,
    comparison_fields: list[str],
) -> dict[str, int]:
    cursor = connection.execute(field_change_sql)
    field_name_index = next(
        index
        for index, column in enumerate(cursor.description)
        if column[0] == "field_name"
    )
    count_index = next(
        index
        for index, column in enumerate(cursor.description)
        if column[0] == "changed_update_observations"
    )
    observed = {row[field_name_index]: row[count_index] for row in cursor.fetchall()}
    expected = {
        field: EXPECTED_NONZERO_FIELD_CHANGES.get(field, 0)
        for field in comparison_fields
    }
    if observed != expected:
        differences = {
            field: {"expected": expected.get(field), "observed": observed.get(field)}
            for field in sorted(set(expected) | set(observed))
            if expected.get(field) != observed.get(field)
        }
        raise RuntimeError(f"Field-change checkpoints failed: {differences}.")
    return observed


def reassert_input_files(
    releases_dir: Path,
    manifests: dict[str, dict[str, Any]],
) -> dict[str, Any]:
    evidence: dict[str, Any] = {}
    for release in RELEASES:
        tag = release["tag"]
        release_dir = releases_dir / tag
        files: dict[str, Any] = {}
        for file_name, expected_sha, expected_bytes in (
            (
                "manifest.json",
                release["manifest_sha256"],
                release["manifest_bytes"],
            ),
            (
                SCHEMA_FILE,
                manifests[tag]["files"][SCHEMA_FILE]["sha256"],
                int(manifests[tag]["files"][SCHEMA_FILE]["bytes"]),
            ),
            (
                AWARDS_FILE,
                manifests[tag]["files"][AWARDS_FILE]["sha256"],
                int(manifests[tag]["files"][AWARDS_FILE]["bytes"]),
            ),
        ):
            digest, size = assert_file(
                release_dir / file_name,
                expected_sha,
                expected_bytes,
            )
            files[file_name] = {"bytes": size, "sha256": digest}
        evidence[tag] = files
    return evidence


def write_query(
    connection: duckdb.DuckDBPyConnection,
    sql: str,
    destination: Path,
) -> dict[str, Any]:
    rows = connection.execute(f"SELECT count(*) FROM ({sql}) AS result").fetchone()[0]
    connection.sql(sql).write_csv(
        str(destination),
        header=True,
        overwrite=True,
        use_tmp_file=True,
    )
    digest, size = sha256_and_size(destination)
    print(f"wrote {destination.name:<48} {rows:>7,} rows")
    return {"rows": rows, "bytes": size, "sha256": digest}


def main() -> None:
    args = parse_args()
    data_dir = args.data_dir.expanduser().resolve()
    releases_dir = data_dir / "releases"
    releases_dir.mkdir(parents=True, exist_ok=True)

    manifests: dict[str, dict[str, Any]] = {}
    csv_validation: dict[str, Any] = {}
    canonical_headers: list[str] | None = None
    for release in RELEASES:
        tag = release["tag"]
        release_dir = releases_dir / tag
        release_dir.mkdir(parents=True, exist_ok=True)
        manifest = load_release_manifest(release_dir, release)
        download_declared_release_assets(release_dir, release, manifest)
        headers = validate_schema_and_header(release_dir, manifest)
        if canonical_headers is None:
            canonical_headers = headers
        elif headers != canonical_headers:
            raise RuntimeError(f"Award schema changed across releases at {tag}.")
        csv_validation[tag] = validate_csv_rows(
            release_dir / AWARDS_FILE,
            release,
            headers,
        )
        manifests[tag] = manifest

    if canonical_headers is None:
        raise RuntimeError("No pinned releases were configured.")
    comparison_fields = [
        field for field in canonical_headers if field not in EXCLUDED_CHANGE_FIELDS
    ]
    if len(comparison_fields) != 61:
        raise RuntimeError(
            f"Expected 61 comparison fields, got {len(comparison_fields)}."
        )

    connection = duckdb.connect()
    materialization_evidence = register_release_membership(
        connection,
        releases_dir,
        canonical_headers,
    )
    create_release_metadata_table(connection)
    create_state_tables(connection, canonical_headers)
    create_update_pairs(connection, comparison_fields)

    state_evidence = assert_state_counts(connection)
    overlap_evidence = assert_release_overlaps(connection)
    amount_evidence = validate_amount_semantics(connection)
    field_change_sql = build_field_change_counts_sql(comparison_fields)
    field_change_evidence = assert_field_change_counts(
        connection,
        field_change_sql,
        comparison_fields,
    )
    input_evidence = reassert_input_files(releases_dir, manifests)

    update_review_sql = build_update_review_sql(
        comparison_fields,
        args.update_review_limit,
    )
    latest_observed_sql = build_latest_observed_sql(canonical_headers)
    queries = {
        "award-history-release-checkpoints.csv": RELEASE_CHECKPOINTS_SQL,
        "award-history-state-reconciliation.csv": STATE_RECONCILIATION_SQL,
        "award-history-release-overlap.csv": RELEASE_OVERLAP_SQL,
        "award-history-obligation-check.csv": OBLIGATION_CHECK_SQL,
        "award-history-field-change-counts.csv": field_change_sql,
        "award-history-update-review.csv": update_review_sql,
        "award-history-latest-observed.csv": latest_observed_sql,
    }

    with tempfile.TemporaryDirectory(
        dir=data_dir,
        prefix=".award-history-results.",
    ) as temporary_directory:
        staging_dir = Path(temporary_directory)
        outputs = {
            file_name: write_query(
                connection,
                sql,
                staging_dir / file_name,
            )
            for file_name, sql in queries.items()
        }
        expected_output_rows = {
            "award-history-release-checkpoints.csv": 4,
            "award-history-state-reconciliation.csv": 1,
            "award-history-release-overlap.csv": 6,
            "award-history-obligation-check.csv": 1,
            "award-history-field-change-counts.csv": 61,
            "award-history-update-review.csv": min(args.update_review_limit, 864),
            "award-history-latest-observed.csv": 82_743,
        }
        observed_output_rows = {
            file_name: evidence["rows"] for file_name, evidence in outputs.items()
        }
        if observed_output_rows != expected_output_rows:
            raise RuntimeError(
                "Output row checkpoints failed: expected "
                f"{expected_output_rows}, got {observed_output_rows}."
            )

        recipe_path = Path(__file__).resolve()
        recipe_sha256, recipe_bytes = sha256_and_size(recipe_path)
        provenance = {
            "dataset_id": DATASET_ID,
            "recipe": "us-federal-contract-awards-history-python",
            "recipe_version": "1.0",
            "recipe_file": {
                "name": recipe_path.name,
                "bytes": recipe_bytes,
                "sha256": recipe_sha256,
            },
            "release_tags": [release["tag"] for release in RELEASES],
            "release_urls": [
                f"https://github.com/{REPOSITORY}/releases/tag/{release['tag']}"
                for release in RELEASES
            ],
            "schema_version": EXPECTED_SCHEMA_VERSION,
            "queried_at_utc": datetime.now(timezone.utc).isoformat(),
            "update_review_limit": args.update_review_limit,
            "duckdb_version": duckdb.__version__,
            "python_runtime": {
                "minimum_version": "3.10",
                "version": sys.version,
                "implementation": platform.python_implementation(),
                "platform": platform.platform(),
            },
            "input_files": input_evidence,
            "csv_validation": csv_validation,
            "materialized_types": materialization_evidence,
            "state_reconciliation": state_evidence,
            "release_overlap": overlap_evidence,
            "amount_checks": amount_evidence,
            "field_change_counts": field_change_evidence,
            "state_grains": {
                "release_membership": {
                    "key": ["release_tag", "award_id"],
                    "rows": 101_389,
                    "purpose": "Retain which award state appeared in each release.",
                },
                "content_state": {
                    "key": ["award_id", "content_hash"],
                    "rows": 83_607,
                    "purpose": "Retain each distinct normalized public award state.",
                },
                "latest_observed": {
                    "key": ["award_id"],
                    "rows": 82_743,
                    "purpose": (
                        "Select the latest row observed in this pinned release corpus; "
                        "not a complete active-contract inventory."
                    ),
                },
            },
            "outputs": outputs,
            "queries": queries,
            "interpretation": {
                "release_membership": (
                    "The same cumulative award summary can appear in multiple "
                    "overlapping releases. Membership rows are evidence, not additive "
                    "award balances."
                ),
                "persistent_state": (
                    "Change type is reconstructed against the last observed state for "
                    "each award_id. An adjacent-file diff would miss updates after an "
                    "award skipped one or more release memberships."
                ),
                "latest_observed": (
                    "Latest observed means latest within these four releases. Absence "
                    "from a later changed-record release does not prove deletion, "
                    "termination, inactivity, or lack of an award."
                ),
                "amounts": (
                    "Award obligations, outlays, current values, and potential values "
                    "are cumulative summary measures. Snapshot differences are net "
                    "summary differences, not individual transactions or period flow."
                ),
                "transactions": (
                    "Use the official USAspending transaction history for action dates, "
                    "modification numbers, and federal_action_obligation."
                ),
                "authority": (
                    "Verify decision-critical facts with each row's source_url and the "
                    "current authoritative federal record."
                ),
            },
        }
        staged_provenance = staging_dir / PROVENANCE_NAME
        staged_provenance.write_text(
            json.dumps(provenance, indent=2, sort_keys=True) + "\n",
            encoding="utf-8",
        )

        for file_name in OUTPUT_NAMES:
            os.replace(staging_dir / file_name, data_dir / file_name)
        os.replace(staged_provenance, data_dir / PROVENANCE_NAME)
        print(f"wrote {PROVENANCE_NAME}")

    connection.close()
    print(
        "verified checkpoint: 101,389 release rows; 82,743 award IDs; "
        "83,607 content states; 17,782 unchanged observations; 864 updates "
        "across 849 awards; $29,232,497,414.24 repeated older obligation balance"
    )


if __name__ == "__main__":
    main()
