"""Guarded one-shot projection backfill for the accepted 2026 term state.

This module deliberately writes only the durable Phase 1 integration outbox.
It never calls the Scheduling Engine and never changes a Timetabler activity,
allocation, week, Staff, or Location row.
"""

from __future__ import annotations

import datetime
import hashlib
import json
import time
import uuid
from dataclasses import dataclass
from types import SimpleNamespace
from typing import Any

from django.conf import settings
from django.db import connection, transaction

from api.models import (
    IntegrationAggregateVersion,
    IntegrationOutbox,
    IntegrationTransportCursor,
    TtAcademicTerm,
    TtActivity,
)
from api.services.integration.outbox import (
    IntegrationAtomicityError,
    append_activity_event,
    build_transport_envelope,
    canonical_hash,
    capture_activity_state,
    finalize_change_set,
    integration_source_scope,
    source_identity,
)
from api.services.integration.readiness import phase1_readiness


TERM_ID = 26
TERM_START = datetime.date(2026, 1, 5)
TERM_END = datetime.date(2026, 6, 26)
TERM_TIMEZONE = "Asia/Kuala_Lumpur"
TERM_ACTIVITY_IDS = tuple(range(672, 717))
TERM_ACTIVITY_IDS_SHA256 = "1b599dcd8709fcdb29b292b110ef1e5cb011929ba371011ef1ae055c9e408453"
ENTRY_SOURCE_WATERMARK = 2639
FINAL_SOURCE_WATERMARK = 2640
EXPECTED_OCCURRENCE_COUNT = 1125
EXPECTED_MISSING_STAFF_IDS = (686, 690)
EXPECTED_DATABASE_CONFIGURATION_SHA256 = (
    "13119e817a230f15dff68303013dc6f33d14110eba94bb54f977af001286716a"
)
EXPECTED_CURRENT_DATABASE_SHA256 = (
    "a7ba0f3444baa246334f8b32b2389ab564b395cd6aaebedfe164833163828aa8"
)
OPERATION_KEY = "phase1-term26-existing-state-backfill-v1"
OPERATION_ORIGIN = "term26_existing_source_state_backfill"
CHANGE_SET_ID = uuid.uuid5(uuid.NAMESPACE_URL, f"timetabler:{OPERATION_KEY}")
COMMITTED_AT_SIZE_SENTINEL = "9999-12-31T23:59:59.999999Z"


class TermProjectionBackfillError(RuntimeError):
    """A non-negotiable backfill precondition was not met."""


@dataclass(frozen=True)
class TermProjectionBackfillContract:
    term_id: int = TERM_ID
    term_start: datetime.date = TERM_START
    term_end: datetime.date = TERM_END
    timezone: str = TERM_TIMEZONE
    activity_ids: tuple[int, ...] = TERM_ACTIVITY_IDS
    activity_ids_sha256: str = TERM_ACTIVITY_IDS_SHA256
    entry_watermark: int = ENTRY_SOURCE_WATERMARK
    final_watermark: int = FINAL_SOURCE_WATERMARK
    expected_occurrence_count: int | None = EXPECTED_OCCURRENCE_COUNT
    expected_missing_staff_ids: tuple[int, ...] | None = EXPECTED_MISSING_STAFF_IDS
    expected_database_configuration_sha256: str | None = (
        EXPECTED_DATABASE_CONFIGURATION_SHA256
    )
    expected_current_database_sha256: str | None = EXPECTED_CURRENT_DATABASE_SHA256
    operation_key: str = OPERATION_KEY
    operation_origin: str = OPERATION_ORIGIN
    change_set_id: uuid.UUID = CHANGE_SET_ID


def _canonical_bytes(value: Any) -> bytes:
    return json.dumps(
        value,
        sort_keys=True,
        separators=(",", ":"),
        ensure_ascii=False,
        allow_nan=False,
        default=str,
    ).encode("utf-8")


def _sha256(value: Any) -> str:
    return hashlib.sha256(_canonical_bytes(value)).hexdigest()


def _database_attestation() -> dict[str, Any]:
    database = connection.settings_dict
    fingerprint = {
        "vendor": connection.vendor,
        "engine": str(database.get("ENGINE") or ""),
        "name": str(database.get("NAME") or ""),
        "host": str(database.get("HOST") or ""),
        "port": str(database.get("PORT") or ""),
    }
    with connection.cursor() as cursor:
        cursor.execute("SELECT current_database()")
        current_database = cursor.fetchone()[0]
    return {
        "alias": connection.alias,
        "vendor": connection.vendor,
        "configuration_sha256": _sha256(fingerprint),
        "current_database_sha256": hashlib.sha256(
            str(current_database).encode("utf-8")
        ).hexdigest(),
        "other_database_connection_used": False,
    }


def _configuration_attestation() -> dict[str, Any]:
    config = settings.RESOURCE_BOOKING_INTEGRATION
    return {
        "capture_enabled": bool(config["CAPTURE_ENABLED"]),
        "publish_enabled": bool(config["PUBLISH_ENABLED"]),
        "activation_approved": bool(config["ACTIVATION_APPROVED"]),
        "phase": str(config["PHASE"]),
        "reverse_delivery_enabled": bool(config["REVERSE_DELIVERY_ENABLED"]),
        "transport": str(config["TRANSPORT"]),
        "schema_version": str(config["SCHEMA_VERSION"]),
        "source_scope": integration_source_scope(),
        "topic_sha256": hashlib.sha256(
            str(config.get("KAFKA_TOPIC") or "").encode("utf-8")
        ).hexdigest(),
        "max_event_bytes": int(config["MAX_EVENT_BYTES"]),
        "max_events_per_change_set": int(config.get("MAX_EVENTS_PER_CHANGE_SET", 5000)),
    }


def _assert_configuration(configuration: dict[str, Any]) -> None:
    required = {
        "capture_enabled": True,
        "publish_enabled": True,
        "activation_approved": True,
        "phase": "phase1",
        "reverse_delivery_enabled": False,
        "transport": "kafka",
        "schema_version": "2",
    }
    mismatches = [key for key, value in required.items() if configuration.get(key) != value]
    if mismatches:
        raise TermProjectionBackfillError(
            "unsafe Phase 1 configuration: " + ",".join(sorted(mismatches))
        )


def _assert_readiness(expected_watermark: int) -> dict[str, Any]:
    readiness = phase1_readiness(require_enabled=True)
    mismatches: list[str] = []
    if not readiness["ready"]:
        mismatches.append("ready")
    if not readiness["publisher_live"]:
        mismatches.append("publisher_live")
    if readiness["publisher_status"] not in {"idle", "running"}:
        mismatches.append("publisher_status")
    if int(readiness["transport_watermark"]) != expected_watermark:
        mismatches.append("transport_watermark")
    if int(readiness["publisher_last_sequence"]) != expected_watermark:
        mismatches.append("publisher_last_sequence")
    if int(readiness["dead_letter_count"]) != 0:
        mismatches.append("dead_letter_count")
    pending = IntegrationOutbox.objects.filter(
        source_scope=readiness["source_scope"],
    ).exclude(
        status__in=(IntegrationOutbox.Status.PUBLISHED, IntegrationOutbox.Status.SUPERSEDED)
    )
    if pending.exists():
        mismatches.append("outbox_backlog")
    if mismatches:
        raise TermProjectionBackfillError(
            "source/publisher fence is not safe: " + ",".join(sorted(mismatches))
        )
    return {
        "transport_watermark": int(readiness["transport_watermark"]),
        "publisher_last_sequence": int(readiness["publisher_last_sequence"]),
        "publisher_live": bool(readiness["publisher_live"]),
        "publisher_status": str(readiness["publisher_status"]),
        "dead_letter_count": int(readiness["dead_letter_count"]),
        "outbox_backlog": 0,
    }


def _set_transaction_timeouts() -> None:
    if connection.vendor != "postgresql":
        raise TermProjectionBackfillError("atomic backfill requires PostgreSQL")
    lock_timeout_ms = int(
        settings.RESOURCE_BOOKING_INTEGRATION.get("LOCK_TIMEOUT_MS", 5000)
    )
    with connection.cursor() as cursor:
        cursor.execute("SET LOCAL lock_timeout = %s", [f"{lock_timeout_ms}ms"])
        cursor.execute("SET LOCAL statement_timeout = %s", ["120s"])


def _activity_queryset(*, locked: bool = False):
    queryset = TtActivity.objects
    if locked:
        # ``week_pattern`` and ``academic_term`` are nullable joins; lock only
        # the authoritative activity rows instead of PostgreSQL's outer-join
        # nullable side.
        queryset = queryset.select_for_update(of=("self",))
    return (
        queryset.select_related("academic_term", "week_pattern")
        .prefetch_related(
            "week",
            "week_pattern__week",
            "staff",
            "location",
            "student_set",
        )
        .order_by("id")
    )


def _latest_projected_states(
    activity_ids: tuple[int, ...],
) -> dict[int, dict[str, Any] | None]:
    aggregate_ids = [
        f"{settings.RESOURCE_BOOKING_INTEGRATION['DEPLOYMENT_ID']}:{activity_id}"
        for activity_id in activity_ids
    ]
    rows = IntegrationOutbox.objects.filter(
        source_scope=integration_source_scope(),
        aggregate_type="activity",
        aggregate_id__in=aggregate_ids,
        transaction_finalized=True,
    ).order_by("event_version", "id")
    latest: dict[str, IntegrationOutbox] = {}
    for row in rows:
        latest[str(row.aggregate_id)] = row
    result: dict[int, dict[str, Any] | None] = {}
    for activity_id, aggregate_id in zip(activity_ids, aggregate_ids, strict=True):
        row = latest.get(aggregate_id)
        replacement = ((row.payload if row else {}) or {}).get("committed_state", {}).get(
            "replacement", {}
        )
        result[activity_id] = replacement.get("current")
    return result


def _next_activity_versions(activity_ids: tuple[int, ...]) -> dict[int, int]:
    deployment_id = settings.RESOURCE_BOOKING_INTEGRATION["DEPLOYMENT_ID"]
    aggregate_ids = [f"{deployment_id}:{activity_id}" for activity_id in activity_ids]
    stored = {
        str(row.aggregate_id): int(row.version)
        for row in IntegrationAggregateVersion.objects.filter(
            aggregate_type="activity", aggregate_id__in=aggregate_ids
        )
    }
    versions: dict[int, int] = {}
    for activity_id, aggregate_id in zip(activity_ids, aggregate_ids, strict=True):
        latest_event_version = (
            IntegrationOutbox.objects.filter(
                aggregate_type="activity", aggregate_id=aggregate_id
            ).order_by("-event_version").values_list("event_version", flat=True).first()
            or 0
        )
        stored_version = stored.get(aggregate_id, 0)
        if int(latest_event_version) != stored_version:
            raise TermProjectionBackfillError(
                f"aggregate version cursor is inconsistent for activity {activity_id}"
            )
        versions[activity_id] = stored_version + 1
    return versions


def _committed_state(current: dict[str, Any], previous: dict[str, Any] | None) -> dict[str, Any]:
    old_scope = {
        "staff_ids": (previous or {}).get("staff_ids", []),
        "location_ids": (previous or {}).get("location_ids", []),
        "week_ids": (previous or {}).get("week_ids", []),
        "scheduled_start_slot": (previous or {}).get("scheduled_start_slot"),
        "duration_minutes": (previous or {}).get("duration_minutes"),
    }
    new_scope = {
        "staff_ids": current.get("staff_ids", []),
        "location_ids": current.get("location_ids", []),
        "week_ids": current.get("week_ids", []),
        "scheduled_start_slot": current.get("scheduled_start_slot"),
        "duration_minutes": current.get("duration_minutes"),
    }
    return {
        "replacement": {"previous": previous, "current": current},
        "affected_scope": {"previous": old_scope, "current": new_scope},
        "tombstone": False,
        "final_state_ref": source_identity("activity", current["id"]),
    }


def _preview_events(
    *,
    contract: TermProjectionBackfillContract,
    states: dict[int, dict[str, Any]],
    previous_states: dict[int, dict[str, Any] | None],
    versions: dict[int, int],
) -> list[SimpleNamespace]:
    source_scope = integration_source_scope()
    count = len(contract.activity_ids)
    rows: list[SimpleNamespace] = []
    for index, activity_id in enumerate(contract.activity_ids, start=1):
        committed_state = _committed_state(states[activity_id], previous_states[activity_id])
        member = {
            "event_id": str(
                uuid.uuid5(uuid.NAMESPACE_URL, f"{contract.operation_key}:preview:{activity_id}")
            ),
            "schema_version": settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
            "event_type": "timetabler.activity.snapshot",
            "aggregate": source_identity("activity", activity_id),
            "event_version": versions[activity_id],
            "ordering_key": source_scope,
            "change_set_id": str(contract.change_set_id),
            "transport": {
                "source_scope": source_scope,
                "sequence": contract.final_watermark,
                "transaction_id": str(contract.change_set_id),
                "transaction_index": index,
                "transaction_count": count,
                "member_only": True,
            },
            "origin": contract.operation_origin,
            # The execute path replaces this same-length sentinel with plan_sha256.
            "correlation_id": "0" * 64,
            "causation_id": None,
            "request_id": contract.operation_key,
            "actor_id": None,
            "committed_at": COMMITTED_AT_SIZE_SENTINEL,
            "committed_state_hash": canonical_hash(committed_state),
            "committed_state": committed_state,
        }
        rows.append(
            SimpleNamespace(
                transaction_index=index,
                source_scope=source_scope,
                transport_sequence=contract.final_watermark,
                change_set_id=contract.change_set_id,
                schema_version=settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
                payload=member,
            )
        )
    return rows


def _assert_term_and_states(
    *,
    contract: TermProjectionBackfillContract,
    locked: bool,
) -> tuple[list[TtActivity], dict[int, dict[str, Any]]]:
    term_query = TtAcademicTerm.objects
    if locked:
        term_query = term_query.select_for_update()
    term = term_query.filter(pk=contract.term_id).first()
    if term is None:
        raise TermProjectionBackfillError("academic term 26 is missing")
    if (
        term.start_date != contract.term_start
        or term.end_date != contract.term_end
        or int(term.status) != 1
        or settings.TIME_ZONE != contract.timezone
    ):
        raise TermProjectionBackfillError("academic term facts or timezone changed")

    activities = list(
        _activity_queryset(locked=locked).filter(academic_term_id=contract.term_id)
    )
    observed_ids = tuple(int(activity.id) for activity in activities)
    if (
        observed_ids != contract.activity_ids
        or _sha256(list(observed_ids)) != contract.activity_ids_sha256
    ):
        raise TermProjectionBackfillError("academic term activity identity set changed")
    if any(int(activity.status) != 1 for activity in activities):
        raise TermProjectionBackfillError("academic term contains an inactive activity")
    if any(not bool(activity.scheduled) for activity in activities):
        raise TermProjectionBackfillError("academic term scheduled state changed")

    states = {int(activity.id): capture_activity_state(activity) for activity in activities}
    missing_staff = tuple(
        activity_id for activity_id in contract.activity_ids if not states[activity_id]["staff_ids"]
    )
    missing_locations = tuple(
        activity_id for activity_id in contract.activity_ids if not states[activity_id]["location_ids"]
    )
    if (
        contract.expected_missing_staff_ids is not None
        and missing_staff != contract.expected_missing_staff_ids
    ):
        raise TermProjectionBackfillError("truthful missing-Staff activity set changed")
    if missing_locations:
        raise TermProjectionBackfillError("a scheduled source activity has no Location allocation")
    return activities, states


def build_plan(
    *,
    deployed_git_sha: str,
    contract: TermProjectionBackfillContract = TermProjectionBackfillContract(),
    locked: bool = False,
    check_readiness: bool = True,
) -> dict[str, Any]:
    if len(deployed_git_sha) != 40 or any(
        character not in "0123456789abcdef" for character in deployed_git_sha
    ):
        raise TermProjectionBackfillError("deployed Timetabler SHA is not exact")
    configuration = _configuration_attestation()
    _assert_configuration(configuration)
    database = _database_attestation()
    if database["vendor"] != "postgresql":
        raise TermProjectionBackfillError("backfill is restricted to attested PostgreSQL staging")
    if (
        contract.expected_database_configuration_sha256
        and database["configuration_sha256"] != contract.expected_database_configuration_sha256
    ):
        raise TermProjectionBackfillError("Timetabler database configuration identity changed")
    if (
        contract.expected_current_database_sha256
        and database["current_database_sha256"] != contract.expected_current_database_sha256
    ):
        raise TermProjectionBackfillError("Timetabler current database identity changed")
    readiness = (
        _assert_readiness(contract.entry_watermark)
        if check_readiness
        else {"transport_watermark": contract.entry_watermark}
    )
    _, states = _assert_term_and_states(contract=contract, locked=locked)
    previous_states = _latest_projected_states(contract.activity_ids)
    versions = _next_activity_versions(contract.activity_ids)
    preview_rows = _preview_events(
        contract=contract,
        states=states,
        previous_states=previous_states,
        versions=versions,
    )
    preview_envelope = build_transport_envelope(preview_rows)
    preview_bytes = len(_canonical_bytes(preview_envelope))
    occurrence_rows = [
        {"activity_id": activity_id, **occurrence}
        for activity_id in contract.activity_ids
        for occurrence in states[activity_id]["occurrences"]
    ]
    allocation_rows = [
        {
            "activity_id": activity_id,
            "staff_ids": states[activity_id]["staff_ids"],
            "location_ids": states[activity_id]["location_ids"],
        }
        for activity_id in contract.activity_ids
    ]
    missing_staff = [
        row["activity_id"] for row in allocation_rows if not row["staff_ids"]
    ]
    if (
        contract.expected_occurrence_count is not None
        and len(occurrence_rows) != contract.expected_occurrence_count
    ):
        raise TermProjectionBackfillError("absolute occurrence count changed")
    plan = {
        "schema_version": 1,
        "profile": "TERM26_EXISTING_SOURCE_STATE_BACKFILL",
        "operation_key": contract.operation_key,
        "operation_origin": contract.operation_origin,
        "deployed_git_sha": deployed_git_sha,
        "source_transaction_id": str(contract.change_set_id),
        "entry_source_watermark": contract.entry_watermark,
        "expected_final_source_watermark": contract.final_watermark,
        "academic_term": {
            "id": contract.term_id,
            "start_date": contract.term_start.isoformat(),
            "end_date": contract.term_end.isoformat(),
            "timezone": contract.timezone,
            "raw_activity_count": len(contract.activity_ids),
            "claimed_activity_count": 46,
            "count_mismatch_preserved": True,
            "activity_ids": list(contract.activity_ids),
            "activity_ids_sha256": contract.activity_ids_sha256,
        },
        "member_count": len(contract.activity_ids),
        "event_type": "timetabler.activity.snapshot",
        "absolute_occurrence_count": len(occurrence_rows),
        "absolute_occurrences_sha256": _sha256(occurrence_rows),
        "resource_allocations_sha256": _sha256(allocation_rows),
        "current_states_sha256": _sha256(
            [
                {"activity_id": activity_id, "state": states[activity_id]}
                for activity_id in contract.activity_ids
            ]
        ),
        "previous_states_sha256": _sha256(
            [
                {"activity_id": activity_id, "state": previous_states[activity_id]}
                for activity_id in contract.activity_ids
            ]
        ),
        "truthful_missing_staff_activity_ids": missing_staff,
        "truthful_missing_location_activity_ids": [],
        "receiver_contract": {
            "typed_resource_arrays_may_omit_one_resource_type": True,
            "missing_staff_is_not_fabricated": True,
            "resource_booking_commit_is_one_atomic_source_transaction": True,
        },
        "transport_preview": {
            "canonical_bytes": preview_bytes,
            "canonical_sha256": _sha256(preview_envelope),
            "events_hash": preview_envelope["events_hash"],
            "max_event_bytes": configuration["max_event_bytes"],
            "bytes_below_limit": configuration["max_event_bytes"] - preview_bytes,
            "within_limit": preview_bytes <= configuration["max_event_bytes"],
            "size_is_exact_for_same_length_runtime_ids_and_timestamps": True,
            "hash_is_preview_only": True,
        },
        "event_versions_sha256": _sha256(versions),
        "configuration": configuration,
        "configuration_sha256": _sha256(configuration),
        "database": database,
        "source_readiness": readiness,
        "guards": {
            "exact_term": True,
            "exact_45_activity_identity_set": True,
            "all_45_currently_scheduled": True,
            "no_domain_schedule_write": True,
            "no_staff_or_location_write": True,
            "no_engine_call": True,
            "no_resource_booking_access": True,
            "normal_outbox_capture_service_only": True,
            "one_finalized_change_set": True,
            "phase1_kafka_only": True,
            "reverse_delivery_enabled": False,
            "phase2_enabled": False,
            "dr_executed": False,
            "failed_load02_reclassified": False,
        },
        "mutation_executed": False,
    }
    if not plan["transport_preview"]["within_limit"]:
        raise TermProjectionBackfillError("planned canonical change set exceeds transport limit")
    plan["plan_sha256"] = _sha256(plan)
    return plan


def _existing_operation_rows(
    contract: TermProjectionBackfillContract,
) -> list[IntegrationOutbox]:
    return list(
        IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(), change_set_id=contract.change_set_id
        ).order_by("transaction_index", "id")
    )


def _actual_evidence(
    rows: list[IntegrationOutbox],
    *,
    plan_sha256: str,
    contract: TermProjectionBackfillContract,
    idempotent_retry: bool,
) -> dict[str, Any]:
    if len(rows) != len(contract.activity_ids):
        raise TermProjectionBackfillError("backfill change set has an incomplete member count")
    if not all(row.transaction_finalized for row in rows):
        raise TermProjectionBackfillError("backfill change set is not finalized")
    if [int(row.transaction_index) for row in rows] != list(
        range(1, len(contract.activity_ids) + 1)
    ) or {int(row.transaction_count) for row in rows} != {len(contract.activity_ids)}:
        raise TermProjectionBackfillError("backfill member boundaries are not exact")
    if {int(row.transport_sequence) for row in rows} != {contract.final_watermark}:
        raise TermProjectionBackfillError("backfill source sequence is not exact")
    if {str(row.event_type) for row in rows} != {"timetabler.activity.snapshot"}:
        raise TermProjectionBackfillError("backfill event types are not exact")
    if {str(row.aggregate_type) for row in rows} != {"activity"}:
        raise TermProjectionBackfillError("backfill aggregate types are not exact")
    if {str(row.origin) for row in rows} != {contract.operation_origin}:
        raise TermProjectionBackfillError("backfill origin is not exact")
    if {str(row.request_id) for row in rows} != {contract.operation_key}:
        raise TermProjectionBackfillError("backfill request identity is not exact")
    if {str(row.correlation_id) for row in rows} != {plan_sha256}:
        raise TermProjectionBackfillError("backfill plan correlation is not exact")
    if [str(row.aggregate_id).rsplit(":", 1)[-1] for row in rows] != [
        str(activity_id) for activity_id in contract.activity_ids
    ]:
        raise TermProjectionBackfillError("backfill aggregate order is not exact")
    for index, (row, activity_id) in enumerate(
        zip(rows, contract.activity_ids, strict=True), start=1
    ):
        payload = row.payload or {}
        committed_state = payload.get("committed_state") or {}
        current = (committed_state.get("replacement") or {}).get("current") or {}
        transport = payload.get("transport") or {}
        if (
            str(payload.get("change_set_id")) != str(contract.change_set_id)
            or int(payload.get("event_version") or 0) <= 0
            or int(current.get("id") or 0) != activity_id
            or not bool(current.get("scheduled"))
            or bool(committed_state.get("tombstone"))
            or str(payload.get("committed_state_hash")) != canonical_hash(committed_state)
            or int(transport.get("transaction_index") or 0) != index
            or int(transport.get("transaction_count") or 0) != len(contract.activity_ids)
            or int(transport.get("sequence") or 0) != contract.final_watermark
            or str(transport.get("transaction_id")) != str(contract.change_set_id)
        ):
            raise TermProjectionBackfillError(
                f"backfill member {index} content is not exact"
            )
    envelope = build_transport_envelope(rows)
    current_states = [
        {
            "activity_id": int(row.aggregate_id.rsplit(":", 1)[-1]),
            "state": row.payload["committed_state"]["replacement"]["current"],
        }
        for row in rows
    ]
    occurrence_rows = [
        {"activity_id": item["activity_id"], **occurrence}
        for item in current_states
        for occurrence in item["state"]["occurrences"]
    ]
    allocation_rows = [
        {
            "activity_id": item["activity_id"],
            "staff_ids": item["state"]["staff_ids"],
            "location_ids": item["state"]["location_ids"],
        }
        for item in current_states
    ]
    return {
        "source_transaction_id": str(contract.change_set_id),
        "source_sequence": contract.final_watermark,
        "entry_source_watermark": contract.entry_watermark,
        "member_count": len(rows),
        "event_type_counts": {"timetabler.activity.snapshot": len(rows)},
        "absolute_occurrence_count": len(occurrence_rows),
        "absolute_occurrences_sha256": _sha256(occurrence_rows),
        "resource_allocations_sha256": _sha256(allocation_rows),
        "current_states_sha256": _sha256(current_states),
        "truthful_missing_staff_activity_ids": [
            row["activity_id"] for row in allocation_rows if not row["staff_ids"]
        ],
        "canonical_bytes": len(_canonical_bytes(envelope)),
        "canonical_sha256": _sha256(envelope),
        "events_hash": envelope["events_hash"],
        "plan_sha256": plan_sha256,
        "statuses": sorted({str(row.status) for row in rows}),
        "published_at_min": min(
            (row.published_at.isoformat() for row in rows if row.published_at), default=None
        ),
        "published_at_max": max(
            (row.published_at.isoformat() for row in rows if row.published_at), default=None
        ),
        "idempotent_retry": idempotent_retry,
        "domain_mutation_executed": False,
        "engine_called": False,
        "resource_booking_accessed": False,
        "reverse_delivery_enabled": False,
        "phase2_enabled": False,
        "dr_executed": False,
        "failed_load02_reclassified": False,
    }


def execute_backfill(
    *,
    deployed_git_sha: str,
    expected_plan_sha256: str,
    contract: TermProjectionBackfillContract = TermProjectionBackfillContract(),
) -> dict[str, Any]:
    if len(expected_plan_sha256) != 64 or any(
        character not in "0123456789abcdef" for character in expected_plan_sha256
    ):
        raise TermProjectionBackfillError("exact plan SHA-256 is required")

    existing = _existing_operation_rows(contract)
    if existing:
        evidence = _actual_evidence(
            existing,
            plan_sha256=expected_plan_sha256,
            contract=contract,
            idempotent_retry=True,
        )
        activities, states = _assert_term_and_states(contract=contract, locked=False)
        del activities
        current_states_sha256 = _sha256(
            [
                {"activity_id": activity_id, "state": states[activity_id]}
                for activity_id in contract.activity_ids
            ]
        )
        if current_states_sha256 != evidence["current_states_sha256"]:
            raise TermProjectionBackfillError(
                "source domain state changed after the sealed backfill"
            )
        cursor = IntegrationTransportCursor.objects.filter(
            source_scope=integration_source_scope()
        ).first()
        if cursor is None or int(cursor.last_sequence) != contract.final_watermark:
            raise TermProjectionBackfillError(
                "source watermark advanced beyond the sealed backfill"
            )
        return evidence

    _assert_readiness(contract.entry_watermark)
    with transaction.atomic():
        _set_transaction_timeouts()
        locked_activities, _states = _assert_term_and_states(
            contract=contract, locked=True
        )
        cursor = IntegrationTransportCursor.objects.select_for_update().filter(
            source_scope=integration_source_scope()
        ).first()
        if cursor is None or int(cursor.last_sequence) != contract.entry_watermark:
            raise TermProjectionBackfillError("source watermark changed before atomic capture")
        if IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(), change_set_id=contract.change_set_id
        ).exists():
            raise TermProjectionBackfillError("backfill identity appeared during capture")

        plan = build_plan(
            deployed_git_sha=deployed_git_sha,
            contract=contract,
            locked=False,
            check_readiness=True,
        )
        if plan["plan_sha256"] != expected_plan_sha256:
            raise TermProjectionBackfillError("approved plan no longer matches source state")
        state_before = plan["current_states_sha256"]
        previous_states = _latest_projected_states(contract.activity_ids)
        activity_by_id = {int(activity.id): activity for activity in locked_activities}
        for activity_id in contract.activity_ids:
            activity = activity_by_id[activity_id]
            append_activity_event(
                activity=activity,
                activity_id=activity_id,
                action="snapshot",
                previous=previous_states[activity_id],
                change_set_id=contract.change_set_id,
                request_id=contract.operation_key,
                origin=contract.operation_origin,
                correlation_id=expected_plan_sha256,
                defer_change_set_finalization=True,
            )
        finalized = finalize_change_set(contract.change_set_id)
        if finalized != len(contract.activity_ids):
            raise IntegrationAtomicityError("backfill finalized an incomplete change set")
        state_after = _sha256(
            [
                {"activity_id": activity_id, "state": capture_activity_state(activity)}
                for activity_id, activity in zip(
                    contract.activity_ids, locked_activities, strict=True
                )
            ]
        )
        if state_after != state_before:
            raise TermProjectionBackfillError("source domain state changed during outbox capture")

    rows = _existing_operation_rows(contract)
    return _actual_evidence(
        rows,
        plan_sha256=expected_plan_sha256,
        contract=contract,
        idempotent_retry=False,
    )


def wait_for_published_backfill(
    *,
    expected_plan_sha256: str,
    contract: TermProjectionBackfillContract = TermProjectionBackfillContract(),
    timeout_seconds: int = 180,
    poll_seconds: float = 2.0,
) -> dict[str, Any]:
    deadline = time.monotonic() + timeout_seconds
    last_reason = "publisher did not reach the final fence"
    while time.monotonic() <= deadline:
        rows = _existing_operation_rows(contract)
        try:
            evidence = _actual_evidence(
                rows,
                plan_sha256=expected_plan_sha256,
                contract=contract,
                idempotent_retry=True,
            )
            readiness = _assert_readiness(contract.final_watermark)
            if {row.status for row in rows} != {IntegrationOutbox.Status.PUBLISHED}:
                raise TermProjectionBackfillError("backfill members are not all published")
            evidence["publisher_readiness"] = readiness
            evidence["statuses"] = [IntegrationOutbox.Status.PUBLISHED]
            evidence["verification_sha256"] = _sha256(evidence)
            return evidence
        except TermProjectionBackfillError as error:
            last_reason = str(error)
        time.sleep(poll_seconds)
    raise TermProjectionBackfillError(last_reason)
