#!/usr/bin/env python3
"""Guarded Timetabler-owned executor for Phase 1 LOAD-01..LOAD-08.

The runner is deliberately inert unless invoked by the separately guarded
workflow with an approved private manifest. Domain changes use authenticated
Timetabler APIs only. ORM access is read-only and exists solely for source
fences, immutable outbox evidence, fixture attestation, and local metrics.
Resource Booking credentials, APIs, databases, and processes are out of scope.
"""

from __future__ import annotations

import argparse
from collections import Counter
import copy
from concurrent.futures import Future, ThreadPoolExecutor, as_completed
from dataclasses import dataclass
import datetime
import hashlib
import json
import math
import os
from pathlib import Path
import re
import secrets
import subprocess
import sys
import threading
import time
from types import SimpleNamespace
from typing import Any
import stat
import uuid
from zoneinfo import ZoneInfo


SCRIPT_DIR = Path(__file__).resolve().parent
if str(SCRIPT_DIR) not in sys.path:
    sys.path.insert(0, str(SCRIPT_DIR))

import phase1_e2e_entry_gate as identity_gate  # noqa: E402


STAGES = (
    "entry",
    "load-01",
    "provision",
    "load-02",
    "load-03",
    "load-04",
    "load-05",
    "load-06",
    "load-07",
    "load-08",
    "cleanup",
    "final",
)
LOAD_EXECUTION_STAGES = {
    "load-02",
    "load-03",
    "load-04",
    "load-05",
    "load-06",
    "load-07",
    "load-08",
}
MUTATING_STAGES = LOAD_EXECUTION_STAGES | {"provision", "cleanup"}
EXECUTION_WORKLOAD_STAGES = LOAD_EXECUTION_STAGES | {"cleanup"}
OPERATIONAL_STAGES = {"load-07"}
ALLOWED_PATHS = {
    "/api/admin/staff/create",
    "/api/admin/staff/update",
    "/api/admin/staff/delete",
    "/api/admin/location/create",
    "/api/admin/location/update",
    "/api/admin/location/delete",
    "/api/admin/activity/create",
    "/api/admin/activity/delete",
    "/api/admin/resources/update-requirement",
    "/api/admin/preschedule-request",
    "/api/admin/schedule-request",
    "/api/admin/unschedule",
    "/api/admin/booking/create",
    "/api/admin/booking/schedule",
    "/api/admin/booking/unschedule",
    "/api/admin/booking/delete",
    "/api/admin/variant-create",
}
MODEL_MAP = {
    "activity": "TtActivity",
    "location": "TtLocation",
    "staff": "TtStaff",
}
RUN_ID_PATTERN = re.compile(r"^phase1-load-[0-9]{8}t[0-9]{6}z-[0-9a-f]{12,32}$")
SHA256_PATTERN = re.compile(r"^[0-9a-f]{64}$")
GIT_SHA_PATTERN = re.compile(r"^[0-9a-f]{40}$")
MAX_MANIFEST_BYTES = 1_048_576
MAX_EVIDENCE_BYTES = 64 * 1024 * 1024
MAX_SELECTORS = 256
MAX_SELECTOR_ROWS = 10_000
MAX_STEPS = 128
TRANSPORT_LIMIT_BYTES = 2_097_152
SUSTAINED_HARD_MEMORY_SECONDS = 300
EVIDENCE_EXCHANGE_ROOT = Path("/var/tmp/mayvins-timetabler-phase1-load")
TIMETABLER_EVIDENCE_OWNER_UID = 1073
RESOURCE_BOOKING_EVIDENCE_READER_UID = 1069
EVIDENCE_EXCHANGE_DIRECTORY_MODE = 0o755
EVIDENCE_EXCHANGE_FILE_MODE = 0o444
DEFAULT_THRESHOLDS = {
    "end_to_end_p95_ms": 5_000,
    "end_to_end_p99_ms": 15_000,
    "lag_drain_seconds": 300,
    "rb_apply_p95_ms": 2_000,
    "rb_apply_p99_ms": 5_000,
    "db_lock_wait_p95_ms": 500,
    "booking_api_p95_impact_percent": 20,
    "booking_api_error_increase_percentage_points": 0.1,
    "sustained_host_percent": 80,
    "hard_abort_memory_percent": 90,
}
DEFAULT_LOAD = {
    "p_tps": 2,
    "burst_tps": 10,
    "steady_seconds": 1_800,
    "burst_seconds": 300,
    "drain_seconds": 300,
    "restart_seconds": 300,
    "stability_seconds": 3_600,
}
REQUIRED_APPROVAL_ROLES = (
    "operator",
    "abort_authority",
    "timetabler_owner",
    "resource_booking_owner",
    "qa_owner",
    "product_owner",
    "operations_owner",
)
DEFAULT_PROFILE_TIMING = {
    "ENTRY": {"workload_seconds": 0, "drain_seconds": 0, "observation_seconds": 300},
    "LOAD-01": {"workload_seconds": 0, "drain_seconds": 0, "observation_seconds": 300},
    "PROVISION": {
        "workload_seconds": 5400,
        "drain_seconds": 300,
        "observation_seconds": 5700,
    },
    "LOAD-02": {
        "workload_seconds": 1800,
        "drain_seconds": 300,
        "observation_seconds": 2100,
    },
    "LOAD-03": {
        "workload_seconds": 300,
        "drain_seconds": 300,
        "observation_seconds": 600,
    },
    "LOAD-04": {
        "workload_seconds": 300,
        "drain_seconds": 300,
        "observation_seconds": 600,
    },
    "LOAD-05": {
        "workload_seconds": 300,
        "drain_seconds": 300,
        "observation_seconds": 600,
    },
    "LOAD-06": {
        "workload_seconds": 300,
        "drain_seconds": 300,
        "observation_seconds": 600,
    },
    "LOAD-07": {
        "workload_seconds": 300,
        "drain_seconds": 300,
        "observation_seconds": 600,
    },
    "LOAD-08": {
        "workload_seconds": 3600,
        "drain_seconds": 300,
        "observation_seconds": 3900,
    },
    "CLEANUP": {
        "workload_seconds": 600,
        "drain_seconds": 300,
        "observation_seconds": 900,
    },
    "FINAL": {"workload_seconds": 0, "drain_seconds": 0, "observation_seconds": 300},
}
LEGACY_TRAFFIC_CLASSIFICATIONS = (
    "staff_lifecycle",
    "location_lifecycle",
    "schedule",
    "reschedule",
    "terminal",
    "recurrence_schedule",
    "recurrence_terminal",
)
REQUIRED_TRAFFIC_CLASSIFICATIONS = ("schedule", "unschedule")
REQUIRED_TRAFFIC_STEP_CLASSIFICATIONS = (
    "schedule",
    "unschedule",
    "schedule",
    "unschedule",
)
ACADEMIC_TERM_LOAD_SCOPE = {
    "academic_term_id": 26,
    "academic_term_start_date": "2026-01-05",
    "academic_term_end_date": "2026-06-26",
    "transaction_types": ["schedule", "unschedule"],
    "minimum_recurring_occurrences_per_activity": 2,
    "recurring_resource_types": ["staff", "location"],
    "resource_booking_sync_required": True,
    "unschedule_releases_assignments": True,
    "timetabling_engine_contract_changed": False,
}
EXPECTED_STAGE_PREDECESSORS = {
    "entry": "E2E-FINAL",
    "load-01": "ENTRY",
    "provision": "LOAD-01",
    "load-02": "PROVISION",
    "load-03": "LOAD-02",
    "load-04": "LOAD-03",
    "load-05": "LOAD-04",
    "load-06": "LOAD-05",
    "load-07": "LOAD-06",
    "load-08": "LOAD-07",
    "cleanup": "LOAD-08",
    "final": "CLEANUP",
}
PROVISION_CORRECTION_REASON = "missing-timezone-offset"
PROVISION_CORRECTION_SOURCE_TIMEZONE = "Asia/Kuala_Lumpur"
LOAD02_CHECKPOINT_RUN_ID = "phase1-load-20260808t150644z-e538cd6b05ec7e218cd399e0"
LOAD02_CHECKPOINT_ARTIFACT_FILENAME = "load-02-checkpoint.json"
ABANDONED_LOAD02_CLEANUP_INITIAL_WATERMARK = 2640
LOAD02_ENGINE_AUDIT_ACTION = "audit-failed-load-02-engine-work"
ABANDONED_LOAD02_CLEANUP_ACTION = "cleanup-abandoned-load-02"
CURRENT_LOAD02_ENGINE_AUDIT_ACTION = "audit-current-load-02-engine-work"
CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION = "cleanup-current-abandoned-load-02"
CURRENT_LOAD02_RECOVERY_ACTIONS = {
    CURRENT_LOAD02_ENGINE_AUDIT_ACTION,
    CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION,
}
CURRENT_LOAD02_RECOVERY_KEYS = {
    "failed_run_id",
    "failed_github_run_id",
    "original_manifest_sha256",
    "original_timetabler_sha",
    "recovery_timetabler_sha",
    "abort_evidence_sha256",
    "settled_source_watermark",
    "approved_actions",
    "record_ref",
    "approved_at",
}
SAFE_REQUEST_SUFFIX_PATTERN = re.compile(r"[A-Za-z0-9][A-Za-z0-9._:-]{0,95}")
PROVISION_CORRECTION_ARTIFACT_FILENAMES = {
    "provision-corrected.json",
    "provision-corrected-v2.json",
    "provision-corrected-v3.json",
    "provision-corrected-v4.json",
    "provision-corrected-v5.json",
    "provision-corrected-v6.json",
    "provision-corrected-v7.json",
}
PROVISION_CORRECTION_KEYS = {
    "artifact_filename",
    "original_evidence_sha256",
    "original_timetabler_sha",
    "original_resource_booking_sha",
    "source_timezone",
    "reason",
}


class LoadHarnessError(RuntimeError):
    """A fail-closed load precondition or correctness invariant failed."""


@dataclass(frozen=True)
class LoadConfig:
    action: str
    run_id: str
    environment: str
    source_scope: str
    base_url: str
    email: str
    password: str
    manifest_json: str
    expected_sha: str
    github_run_id: str
    prior_evidence_sha256: str
    mutation_confirmed: bool
    operational_confirmed: bool
    restart_fence_sha256: str

    @classmethod
    def from_environment(cls) -> "LoadConfig":
        def required(name: str, *, strip: bool = True) -> str:
            value = os.environ.get(name, "")
            value = value.strip() if strip else value
            if not value:
                raise LoadHarnessError(f"required load setting is absent: {name}")
            return value

        action = required("TT_PHASE1_LOAD_ACTION")
        if action not in STAGES:
            raise LoadHarnessError("unsupported Phase 1 load action")
        return cls(
            action=action,
            run_id=required("TT_PHASE1_LOAD_RUN_ID"),
            environment=required("TT_PHASE1_LOAD_ENVIRONMENT"),
            source_scope=required("TT_PHASE1_LOAD_SOURCE_SCOPE"),
            base_url=required("TT_PHASE1_LOAD_BASE_URL"),
            email=required("TT_PHASE1_LOAD_EMAIL"),
            password=required("TT_PHASE1_LOAD_PASSWORD", strip=False),
            manifest_json=required("TT_PHASE1_LOAD_MANIFEST_JSON", strip=False),
            expected_sha=required("TT_PHASE1_LOAD_EXPECTED_SHA"),
            github_run_id=required("TT_PHASE1_LOAD_GITHUB_RUN_ID"),
            prior_evidence_sha256=required(
                "TT_PHASE1_LOAD_PRIOR_EVIDENCE_SHA256"
            ).lower(),
            mutation_confirmed=_env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION"),
            operational_confirmed=_env_true("TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"),
            restart_fence_sha256=os.environ.get(
                "TT_PHASE1_LOAD_RESTART_FENCE_SHA256", ""
            )
            .strip()
            .lower(),
        )


def _env_true(name: str) -> bool:
    return os.environ.get(name, "").strip().lower() in {"1", "true", "yes", "on"}


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


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


def _sha256_log_parts(value: str) -> list[str]:
    """Keep a non-secret attestation hash readable when GitHub masks the whole value."""

    if SHA256_PATTERN.fullmatch(value) is None:
        raise LoadHarnessError("attestation SHA-256 is invalid")
    return [value[:32], value[32:]]


def seal_evidence(payload: dict[str, Any]) -> dict[str, Any]:
    return {**payload, "evidence_sha256": _sha256(payload)}


def _timezone_qualified_iso(
    value: datetime.datetime | str, *, source_timezone: str | None = None
) -> str:
    """Serialize a timestamp with an offset, localizing only naive DB values."""
    if isinstance(value, str):
        try:
            parsed = datetime.datetime.fromisoformat(value.replace("Z", "+00:00"))
        except ValueError as error:
            raise LoadHarnessError("source transaction timestamp is invalid") from error
    elif isinstance(value, datetime.datetime):
        parsed = value
    else:
        raise LoadHarnessError("source transaction timestamp is invalid")
    if parsed.tzinfo is None or parsed.utcoffset() is None:
        if source_timezone is None:
            from django.conf import settings

            source_timezone = str(settings.TIME_ZONE)
        try:
            parsed = parsed.replace(tzinfo=ZoneInfo(source_timezone))
        except (KeyError, ValueError) as error:
            raise LoadHarnessError("source transaction timezone is invalid") from error
    return parsed.isoformat()


def evidence_exchange_attestation() -> dict[str, Any]:
    return {
        "path_sha256": hashlib.sha256(str(EVIDENCE_EXCHANGE_ROOT).encode()).hexdigest(),
        "timetabler_owner_uid": TIMETABLER_EVIDENCE_OWNER_UID,
        "resource_booking_reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
        "directory_mode": "0755",
        "file_mode": "0444",
        "credentials_forbidden": True,
    }


def _evidence_path(run_id: str, profile: str) -> tuple[Path, Path, Path]:
    validate_run_id(run_id)
    canonical_profiles = {
        *(f"LOAD-{index:02d}" for index in range(1, 9)),
        "PROVISION",
        "CLEANUP",
        "FINAL",
        "ENTRY",
    }
    if profile not in canonical_profiles:
        raise LoadHarnessError("evidence profile is not canonical")
    root = EVIDENCE_EXCHANGE_ROOT
    if not root.is_absolute() or str(root) != "/var/tmp/mayvins-timetabler-phase1-load":
        raise LoadHarnessError("evidence exchange root differs from the fixed contract")
    run_dir = root / run_id
    if run_dir.parent != root or run_dir.name != run_id:
        raise LoadHarnessError("evidence path failed traversal-safe validation")
    return root, run_dir, run_dir / f"{profile.lower()}.json"


def _ensure_evidence_exchange_directories(run_id: str) -> tuple[Path, Path]:
    root, run_dir, _ = _evidence_path(run_id, "ENTRY")
    if os.getuid() != TIMETABLER_EVIDENCE_OWNER_UID:
        raise LoadHarnessError("runtime uid differs from the evidence exchange owner")
    for path in (root, run_dir):
        try:
            path.mkdir(mode=EVIDENCE_EXCHANGE_DIRECTORY_MODE, exist_ok=False)
        except FileExistsError:
            pass
        info = path.stat(follow_symlinks=False)
        if (
            not stat.S_ISDIR(info.st_mode)
            or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        ):
            raise LoadHarnessError(
                "evidence exchange directory ownership or mode is unsafe"
            )
        os.chmod(path, EVIDENCE_EXCHANGE_DIRECTORY_MODE)
        if stat.S_IMODE(path.stat(follow_symlinks=False).st_mode) != (
            EVIDENCE_EXCHANGE_DIRECTORY_MODE
        ):
            raise LoadHarnessError("evidence exchange directory mode is unsafe")
    return root, run_dir


def _assert_exchange_payload_has_no_credentials(value: Any) -> None:
    if _contains_credential_key(value):
        raise LoadHarnessError(
            "evidence exchange payload contains credential-like data"
        )
    encoded = _canonical_bytes(value)
    for name, secret in os.environ.items():
        normalized = name.lower()
        if not any(
            marker in normalized
            for marker in ("password", "secret", "token", "authorization", "credential")
        ):
            continue
        if secret and len(secret) >= 8 and secret.encode() in encoded:
            raise LoadHarnessError(
                "evidence exchange payload contains a runtime credential"
            )


def _persist_immutable_exchange_file(
    *, target: Path, encoded: bytes, label: str
) -> os.stat_result:
    if target.exists():
        info = target.stat(follow_symlinks=False)
        existing = target.read_bytes() if stat.S_ISREG(info.st_mode) else b""
        if (
            existing != encoded
            or not stat.S_ISREG(info.st_mode)
            or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
            or stat.S_IMODE(info.st_mode) != EVIDENCE_EXCHANGE_FILE_MODE
        ):
            raise LoadHarnessError(
                f"immutable {label} path already contains different or unsafe data"
            )
        return info

    temporary = target.parent / f".{target.name}.{os.getpid()}.{uuid.uuid4().hex}.tmp"
    descriptor = os.open(
        temporary,
        os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0),
        0o600,
    )
    try:
        with os.fdopen(descriptor, "wb") as handle:
            handle.write(encoded)
            handle.flush()
            os.fsync(handle.fileno())
        os.chmod(temporary, EVIDENCE_EXCHANGE_FILE_MODE)
        try:
            os.link(temporary, target, follow_symlinks=False)
        except FileExistsError:
            info = target.stat(follow_symlinks=False)
            existing = target.read_bytes() if stat.S_ISREG(info.st_mode) else b""
            if (
                existing != encoded
                or not stat.S_ISREG(info.st_mode)
                or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
                or stat.S_IMODE(info.st_mode) != EVIDENCE_EXCHANGE_FILE_MODE
            ):
                raise LoadHarnessError(
                    f"immutable {label} path already contains different or unsafe data"
                )
    finally:
        if temporary.exists():
            temporary.unlink()
    info = target.stat(follow_symlinks=False)
    if (
        not stat.S_ISREG(info.st_mode)
        or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or stat.S_IMODE(info.st_mode) != EVIDENCE_EXCHANGE_FILE_MODE
    ):
        raise LoadHarnessError(f"persisted {label} ownership or mode is unsafe")
    return info


def persist_evidence(evidence: dict[str, Any]) -> dict[str, Any]:
    """Persist one immutable cross-user evidence artifact, mode 0444.

    A retry may reuse an identical file. It may never overwrite evidence with
    different bytes under the same run/profile identity.
    """

    run_id = str(evidence.get("run_id", ""))
    profile = str(evidence.get("profile", ""))
    _, _, target = _evidence_path(run_id, profile)
    _ensure_evidence_exchange_directories(run_id)
    if evidence.get("evidence_exchange") != evidence_exchange_attestation():
        raise LoadHarnessError(
            "evidence exchange attestation is absent or contradictory"
        )
    _assert_exchange_payload_has_no_credentials(evidence)
    encoded = _canonical_bytes(evidence) + b"\n"
    expected_sha = hashlib.sha256(
        _canonical_bytes(
            {key: value for key, value in evidence.items() if key != "evidence_sha256"}
        )
    ).hexdigest()
    if evidence.get("evidence_sha256") != expected_sha:
        raise LoadHarnessError("evidence hash is contradictory before persistence")
    info = _persist_immutable_exchange_file(
        target=target, encoded=encoded, label="evidence"
    )
    return {
        "evidence_path": str(target),
        "evidence_sha256": evidence["evidence_sha256"],
        "owner_uid": info.st_uid,
        "reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
        "mode": "0444",
        "bytes": info.st_size,
    }


def _read_immutable_exchange_json(
    *, root: Path, run_dir: Path, path: Path, label: str
) -> tuple[dict[str, Any], os.stat_result]:
    try:
        root_info = root.stat(follow_symlinks=False)
        run_info = run_dir.stat(follow_symlinks=False)
        info = path.stat(follow_symlinks=False)
    except FileNotFoundError as error:
        raise LoadHarnessError(f"required {label} file is missing") from error
    if (
        not stat.S_ISDIR(root_info.st_mode)
        or not stat.S_ISDIR(run_info.st_mode)
        or root_info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or run_info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or stat.S_IMODE(root_info.st_mode) != EVIDENCE_EXCHANGE_DIRECTORY_MODE
        or stat.S_IMODE(run_info.st_mode) != EVIDENCE_EXCHANGE_DIRECTORY_MODE
        or not stat.S_ISREG(info.st_mode)
        or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or stat.S_IMODE(info.st_mode) != EVIDENCE_EXCHANGE_FILE_MODE
        or info.st_size > MAX_EVIDENCE_BYTES
    ):
        raise LoadHarnessError(f"{label} ownership, mode, type, or size is unsafe")
    try:
        descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0))
        with os.fdopen(descriptor, "rb") as handle:
            encoded = handle.read(MAX_EVIDENCE_BYTES + 1)
    except OSError as error:
        raise LoadHarnessError(f"{label} could not be read safely") from error
    if len(encoded) > MAX_EVIDENCE_BYTES:
        raise LoadHarnessError(f"{label} exceeds its safety bound")
    try:
        payload = json.loads(encoded)
    except (UnicodeDecodeError, json.JSONDecodeError) as error:
        raise LoadHarnessError(f"{label} is not canonical JSON") from error
    if encoded != _canonical_bytes(payload) + b"\n":
        raise LoadHarnessError(f"{label} bytes are not canonical")
    _assert_exchange_payload_has_no_credentials(payload)
    embedded = payload.get("evidence_sha256")
    calculated = _sha256(
        {key: value for key, value in payload.items() if key != "evidence_sha256"}
    )
    if embedded != calculated:
        raise LoadHarnessError(f"{label} embedded digest is contradictory")
    return payload, info


def persist_abort_evidence(
    *, config: LoadConfig, runner: "LoadRunner", reason: str
) -> dict[str, Any]:
    """Preserve a bounded failure record without occupying the PASS profile path."""

    _, run_dir = _ensure_evidence_exchange_directories(config.run_id)
    target = run_dir / f"{config.action}-abort-{config.github_run_id}.json"
    payload = seal_evidence(
        {
            "schema_version": 1,
            "profile": config.action.upper(),
            "status": "ABORTED",
            "run_id": config.run_id,
            "git_sha": config.expected_sha,
            "github_run_id": config.github_run_id,
            "abort_reason": reason[:1000],
            "source_transactions_observed": runner.transactions,
            "api_measurements_observed": runner.api_measurements,
            "api_failures_observed": sorted(
                runner.api_failures,
                key=lambda item: (
                    int(item.get("iteration", 0)),
                    str(item.get("request_id") or ""),
                ),
            ),
            "source_metric_samples": runner.source_samples,
            "host_metric_samples": runner.host_samples,
            "load_executed": config.action in LOAD_EXECUTION_STAGES,
            "pass_claimed": False,
            "reverse_delivery_enabled": False,
            "phase2_enabled": False,
            "dr_executed": False,
            "evidence_exchange": evidence_exchange_attestation(),
        }
    )
    _assert_exchange_payload_has_no_credentials(payload)
    encoded = _canonical_bytes(payload) + b"\n"
    _persist_immutable_exchange_file(target=target, encoded=encoded, label="abort")
    return {
        "abort_evidence_path": str(target),
        "abort_evidence_sha256": payload["evidence_sha256"],
    }


def verify_predecessor_evidence(
    *,
    run_id: str,
    stage: str,
    manifest: dict[str, Any],
    supplied_sha256: str,
) -> dict[str, Any]:
    """Verify the immutable predecessor without putting future hashes in the manifest."""

    if not SHA256_PATTERN.fullmatch(supplied_sha256):
        raise LoadHarnessError("supplied predecessor evidence digest is invalid")
    predecessor = manifest["stage_predecessors"][stage]
    expected_profile = predecessor["profile"]
    if expected_profile == "E2E-FINAL":
        expected_root = predecessor["external_root_sha256"]
        if supplied_sha256 != expected_root:
            raise LoadHarnessError(
                "ENTRY predecessor differs from the accepted E2E final root"
            )
        return {
            "profile": expected_profile,
            "evidence_sha256": supplied_sha256,
            "external_accepted_root": True,
            "evidence_path": None,
        }

    root, run_dir, path = _evidence_path(run_id, expected_profile)
    declared_correction = manifest.get("evidence_corrections", {}).get("provision")
    correction = (
        declared_correction
        if stage == "load-02"
        and expected_profile == "PROVISION"
        and declared_correction
        else None
    )
    if correction:
        path = run_dir / correction["artifact_filename"]
    payload, info = _read_immutable_exchange_json(
        root=root, run_dir=run_dir, path=path, label="predecessor evidence"
    )
    calculated = payload["evidence_sha256"]
    if supplied_sha256 != calculated:
        raise LoadHarnessError("predecessor evidence digest is contradictory")
    if (
        payload.get("schema_version") != 1
        or payload.get("service") != "timetabler"
        or payload.get("run_id") != run_id
        or payload.get("profile") != expected_profile
        or payload.get("reverse_delivery_enabled") is not False
        or payload.get("phase2_enabled") is not False
        or payload.get("dr_executed") is not False
        or payload.get("evidence_exchange") != evidence_exchange_attestation()
    ):
        raise LoadHarnessError("predecessor evidence identity or safety flags differ")
    if correction and (
        payload.get("evidence_correction") != correction
        or payload.get("git_sha") != manifest["commits"]["timetabler"]
        or payload.get("resource_booking_git_sha_attestation")
        != manifest["commits"]["resource_booking"]
    ):
        raise LoadHarnessError(
            "corrected PROVISION predecessor differs from its manifest declaration"
        )
    if expected_profile == "PROVISION":
        _assert_provision_term_scope(payload, manifest)
    return {
        "profile": expected_profile,
        "evidence_sha256": calculated,
        "external_accepted_root": False,
        "evidence_path": str(path),
        "owner_uid": info.st_uid,
        "reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
        "mode": "0444",
    }


def _assert_provision_term_scope(
    provision_evidence: dict[str, Any], manifest: dict[str, Any]
) -> None:
    authoritative_term = manifest.get("accepted_baseline", {}).get(
        "authoritative_term"
    )
    if authoritative_term is None:
        return
    expected = {
        "academic_term_id": authoritative_term["academic_term_id"],
        "academic_term_start_date": authoritative_term["start_date"],
        "academic_term_end_date": authoritative_term["end_date"],
        "authoritative_term_scope_enforced": True,
        "term_template_donor_only": True,
        "template_candidate_defaults_authoritative": False,
        "candidate_defaults_replaced_before_scheduling": True,
        "run_owned_activity_final_allocations_empty_at_provision": True,
        "advisory_requirement_source_watermark_unchanged": True,
    }
    fixture_manifest = provision_evidence.get("fixture_manifest")
    stage_fixture_manifest = provision_evidence.get("stage_assertions", {}).get(
        "fixture_manifest"
    )
    if (
        not isinstance(fixture_manifest, dict)
        or fixture_manifest != stage_fixture_manifest
        or any(fixture_manifest.get(key) != value for key, value in expected.items())
        or fixture_manifest.get("activity_template_id")
        != fixture_manifest.get("donor_activity_template_id")
    ):
        raise LoadHarnessError(
            "PROVISION evidence is not bound to the authoritative academic-term scope"
        )


def _validate_original_provision_evidence(
    *,
    original: dict[str, Any],
    manifest: dict[str, Any],
    run_id: str,
    supplied_sha256: str,
) -> None:
    declaration = manifest["evidence_corrections"]["provision"]
    start = manifest["expected_stage_start_watermarks"]["provision"]
    final = manifest["expected_stage_final_watermarks"]["provision"]
    transactions = original.get("source_transactions")
    if (
        supplied_sha256 != declaration["original_evidence_sha256"]
        or original.get("evidence_sha256") != supplied_sha256
        or original.get("schema_version") != 1
        or original.get("profile") != "PROVISION"
        or original.get("service") != "timetabler"
        or original.get("run_id") != run_id
        or original.get("git_sha") != declaration["original_timetabler_sha"]
        or original.get("resource_booking_git_sha_attestation")
        != declaration["original_resource_booking_sha"]
        or original.get("entry_watermark") != start
        or original.get("final_watermark") != final
        or original.get("reverse_delivery_enabled") is not False
        or original.get("phase2_enabled") is not False
        or original.get("dr_executed") is not False
        or original.get("evidence_exchange") != evidence_exchange_attestation()
        or not isinstance(transactions, list)
        or len(transactions)
        != manifest["provision_plan"]["expected_source_transactions"]
        or any(not isinstance(item, dict) for item in transactions)
        or [item.get("source_sequence") for item in transactions]
        != list(range(start + 1, final + 1))
    ):
        raise LoadHarnessError(
            "original PROVISION evidence identity, fence, or transaction inventory differs"
        )
    fixture_manifest = original.get("fixture_manifest")
    stage_fixture_manifest = original.get("stage_assertions", {}).get(
        "fixture_manifest"
    )
    if (
        fixture_manifest != stage_fixture_manifest
        or not isinstance(fixture_manifest, dict)
        or fixture_manifest.get("counts") != manifest["provision_plan"]["counts"]
        or fixture_manifest.get("maximum_recurrence_week_count") != 12
        or fixture_manifest.get("complete_staff_location_presets") is not True
    ):
        raise LoadHarnessError(
            "original PROVISION fixture assertions differ from the approved inventory"
        )
    _assert_provision_term_scope(original, manifest)
    created = original.get("created_counts")
    reused = original.get("reused_counts")
    stage_assertions = original.get("stage_assertions", {})
    if (
        not isinstance(created, dict)
        or not isinstance(reused, dict)
        or stage_assertions.get("created_counts") != created
        or stage_assertions.get("reused_counts") != reused
        or any(
            created.get(key, 0) + reused.get(key, 0) != count
            for key, count in manifest["provision_plan"]["counts"].items()
        )
    ):
        raise LoadHarnessError(
            "original PROVISION created/reused assertions are incomplete"
        )


def _build_corrected_provision_evidence(
    *,
    original: dict[str, Any],
    manifest: dict[str, Any],
    github_run_id: str,
    inventory_counts: dict[str, int],
    source_progress_attestation: dict[str, Any],
) -> dict[str, Any]:
    declaration = copy.deepcopy(manifest["evidence_corrections"]["provision"])
    corrected = copy.deepcopy(
        {key: value for key, value in original.items() if key != "evidence_sha256"}
    )
    corrected["git_sha"] = manifest["commits"]["timetabler"]
    corrected["resource_booking_git_sha_attestation"] = manifest["commits"][
        "resource_booking"
    ]
    corrected["github_run_id"] = github_run_id
    for transaction in corrected["source_transactions"]:
        transaction["committed_at"] = _timezone_qualified_iso(
            transaction.get("committed_at"),
            source_timezone=declaration["source_timezone"],
        )
        transaction["published_at"] = _timezone_qualified_iso(
            transaction.get("published_at"),
            source_timezone=declaration["source_timezone"],
        )
    corrected["evidence_correction"] = declaration
    corrected["post_provision_source_progress"] = copy.deepcopy(
        source_progress_attestation
    )
    corrected["correction_action_assertions"] = {
        "api_authentication_performed": False,
        "domain_mutation": False,
        "database_write": False,
        "manual_outbox_edit": False,
        "process_restart": False,
        "load_executed": False,
        "dr_executed": False,
        "reverse_delivery_enabled": False,
        "original_evidence_preserved": True,
        "fixture_inventory_counts": inventory_counts,
        "fixture_inventory_sha256": _sha256(inventory_counts),
    }
    return seal_evidence(corrected)


def correct_provision_evidence_from_environment() -> dict[str, Any]:
    """Create a linked immutable timestamp correction without domain mutation."""
    if (
        os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip()
        != "correct-provision-evidence"
    ):
        raise LoadHarnessError(
            "correct-provision-evidence CLI action differs from guarded workflow action"
        )
    if _env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION") or _env_true(
        "TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"
    ):
        raise LoadHarnessError(
            "correct-provision-evidence requires false mutation and operational confirmations"
        )
    run_id = os.environ.get("TT_PHASE1_LOAD_RUN_ID", "").strip()
    validate_run_id(run_id)
    supplied_sha256 = (
        os.environ.get("TT_PHASE1_LOAD_PRIOR_EVIDENCE_SHA256", "").strip().lower()
    )
    if not SHA256_PATTERN.fullmatch(supplied_sha256):
        raise LoadHarnessError("original PROVISION evidence digest is invalid")
    github_run_id = os.environ.get("TT_PHASE1_LOAD_GITHUB_RUN_ID", "").strip()
    if not github_run_id.isdigit():
        raise LoadHarnessError("immutable correction GitHub run id is invalid")
    manifest = load_manifest(
        os.environ.get("TT_PHASE1_LOAD_MANIFEST_JSON", ""),
        run_id=run_id,
        environment=os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
        source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
        require_approval=True,
    )
    declaration = manifest.get("evidence_corrections", {}).get("provision")
    if not declaration:
        raise LoadHarnessError(
            "approved manifest lacks the PROVISION evidence correction declaration"
        )
    exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
    if exact_sha != manifest["commits"]["timetabler"]:
        raise LoadHarnessError(
            "deployed Timetabler SHA differs from the approved correction manifest"
        )
    assert_evidence_exchange_runtime(manifest, run_id=run_id)
    now = datetime.datetime.now(datetime.UTC)
    if not (
        _parse_iso8601(
            manifest["approved_window"]["starts_at"], "approved_window.starts_at"
        )
        <= now
        <= _parse_iso8601(
            manifest["approved_window"]["ends_at"], "approved_window.ends_at"
        )
    ):
        raise LoadHarnessError(
            "current time is outside the explicitly approved correction window"
        )
    expected_watermark = manifest["expected_stage_final_watermarks"]["provision"]
    metrics = _assert_source_safety(expected_watermark=None)
    source_progress = _load02_source_progress_attestation(
        manifest=manifest,
        run_id=run_id,
        current_watermark=metrics["transport_watermark"],
    )

    root, run_dir, original_path = _evidence_path(run_id, "PROVISION")
    original, _ = _read_immutable_exchange_json(
        root=root,
        run_dir=run_dir,
        path=original_path,
        label="original PROVISION evidence",
    )
    _validate_original_provision_evidence(
        original=original,
        manifest=manifest,
        run_id=run_id,
        supplied_sha256=supplied_sha256,
    )
    fixture_ids = _resolve_selectors(manifest)
    inventory_counts = {key: len(value) for key, value in sorted(fixture_ids.items())}
    if inventory_counts != {
        key: int(selector["expected_count"])
        for key, selector in sorted(manifest["fixture_selectors"].items())
    }:
        raise LoadHarnessError(
            "current run-owned fixture inventory differs from the sealed PROVISION assertions"
        )
    corrected = _build_corrected_provision_evidence(
        original=original,
        manifest=manifest,
        github_run_id=github_run_id,
        inventory_counts=inventory_counts,
        source_progress_attestation=source_progress,
    )
    final_metrics = _assert_source_safety(
        expected_watermark=source_progress["current_source_watermark"]
    )
    final_source_progress = _load02_source_progress_attestation(
        manifest=manifest,
        run_id=run_id,
        current_watermark=final_metrics["transport_watermark"],
    )
    if final_source_progress != source_progress:
        raise LoadHarnessError(
            "LOAD-02 durable prefix changed during PROVISION evidence correction"
        )
    target = run_dir / declaration["artifact_filename"]
    _assert_exchange_payload_has_no_credentials(corrected)
    info = _persist_immutable_exchange_file(
        target=target,
        encoded=_canonical_bytes(corrected) + b"\n",
        label="corrected PROVISION evidence",
    )
    original_after, _ = _read_immutable_exchange_json(
        root=root,
        run_dir=run_dir,
        path=original_path,
        label="original PROVISION evidence",
    )
    if original_after["evidence_sha256"] != supplied_sha256:
        raise LoadHarnessError("original PROVISION evidence changed during correction")
    return {
        "phase1_load": "provision_evidence_corrected",
        "profile": "PROVISION",
        "run_id": run_id,
        "entry_watermark": expected_watermark,
        "final_watermark": final_metrics["transport_watermark"],
        "publisher_last_sequence": final_metrics["publisher_last_sequence"],
        "load02_prefix_transaction_count": source_progress[
            "load02_prefix_transaction_count"
        ],
        "load02_prefix_source_transactions_sha256": source_progress[
            "load02_prefix_source_transactions_sha256"
        ],
        "original_evidence_sha256": supplied_sha256,
        "corrected_evidence_sha256": corrected["evidence_sha256"],
        "evidence_path": str(target),
        "owner_uid": info.st_uid,
        "reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
        "mode": "0444",
        "bytes": info.st_size,
        "api_authentication_performed": False,
        "mutation_executed": False,
        "operational_action_executed": False,
        "load_executed": False,
        "dr_executed": False,
        "reverse_delivery_enabled": False,
    }


def verify_restart_fence(
    *,
    run_id: str,
    manifest: dict[str, Any],
    supplied_sha256: str,
    now: datetime.datetime | None = None,
) -> dict[str, Any]:
    policy = manifest["workloads"]["load-07"]["restart_fence_policy"]
    if not SHA256_PATTERN.fullmatch(supplied_sha256):
        raise LoadHarnessError("LOAD-07 restart fence digest is absent or invalid")
    root, run_dir, _ = _evidence_path(run_id, "LOAD-07")
    path = run_dir / "load-07-fence.json"
    try:
        root_info = root.stat(follow_symlinks=False)
        run_info = run_dir.stat(follow_symlinks=False)
        info = path.stat(follow_symlinks=False)
    except FileNotFoundError as error:
        raise LoadHarnessError(
            "LOAD-07 short-lived fence artifact is missing"
        ) from error
    if (
        not stat.S_ISDIR(root_info.st_mode)
        or not stat.S_ISDIR(run_info.st_mode)
        or root_info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or run_info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or stat.S_IMODE(root_info.st_mode) != EVIDENCE_EXCHANGE_DIRECTORY_MODE
        or stat.S_IMODE(run_info.st_mode) != EVIDENCE_EXCHANGE_DIRECTORY_MODE
        or not stat.S_ISREG(info.st_mode)
        or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
        or stat.S_IMODE(info.st_mode) != EVIDENCE_EXCHANGE_FILE_MODE
        or info.st_size > 16_384
    ):
        raise LoadHarnessError(
            "LOAD-07 fence artifact ownership, mode, type, or size is unsafe"
        )
    descriptor = os.open(path, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0))
    with os.fdopen(descriptor, "rb") as handle:
        encoded = handle.read(16_385)
    try:
        artifact = json.loads(encoded)
    except (UnicodeDecodeError, json.JSONDecodeError) as error:
        raise LoadHarnessError(
            "LOAD-07 fence artifact is not canonical JSON"
        ) from error
    if encoded != _canonical_bytes(artifact) + b"\n":
        raise LoadHarnessError("LOAD-07 fence artifact bytes are not canonical")
    _assert_exchange_payload_has_no_credentials(artifact)
    unsigned = {
        key: value for key, value in artifact.items() if key != "coordination_sha256"
    }
    calculated = _sha256(unsigned)
    if (
        artifact.get("coordination_sha256") != calculated
        or supplied_sha256 != calculated
    ):
        raise LoadHarnessError("LOAD-07 fence artifact hash is contradictory")
    if (
        set(artifact)
        != {
            "schema_version",
            "run_id",
            "timetabler_sha",
            "resource_booking_sha",
            "execute_at_epoch",
            "nonce",
            "approval_ref",
            "created_at",
            "expires_at",
            "coordination_sha256",
        }
        or artifact["schema_version"] != 1
        or artifact["run_id"] != run_id
        or artifact["timetabler_sha"] != manifest["commits"]["timetabler"]
        or artifact["resource_booking_sha"] != manifest["commits"]["resource_booking"]
        or not isinstance(artifact["execute_at_epoch"], int)
        or not re.fullmatch(r"[0-9a-f]{64}", str(artifact["nonce"]))
        or not str(artifact["approval_ref"]).strip()
    ):
        raise LoadHarnessError("LOAD-07 fence artifact identity is invalid")
    created = _parse_iso8601(artifact["created_at"], "LOAD-07 fence created_at")
    expires = _parse_iso8601(artifact["expires_at"], "LOAD-07 fence expires_at")
    current = now or datetime.datetime.now(datetime.UTC)
    execute_at = datetime.datetime.fromtimestamp(
        artifact["execute_at_epoch"], tz=datetime.UTC
    )
    lead = (execute_at - current).total_seconds()
    if (
        not created <= current <= expires
        or expires < execute_at + datetime.timedelta(seconds=5)
        or not policy["minimum_lead_seconds"] <= lead <= policy["maximum_lead_seconds"]
    ):
        raise LoadHarnessError(
            "LOAD-07 fence artifact is expired or outside lead policy"
        )
    return {
        **artifact,
        "artifact_path": str(path),
        "mode": "0444",
        "owner_uid": info.st_uid,
        "reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
    }


def arm_restart_fence(
    *,
    run_id: str,
    manifest: dict[str, Any],
    supplied_predecessor_sha256: str,
    execute_at_epoch: int,
    approval_ref: str,
    now: datetime.datetime | None = None,
) -> dict[str, Any]:
    """Create or identically reuse the canonical LOAD-07 coordination fence."""

    current = now or datetime.datetime.now(datetime.UTC)
    if current.tzinfo is None or current.utcoffset() is None:
        raise LoadHarnessError("arm-fence current time must be timezone-aware")
    current = current.astimezone(datetime.UTC)
    normalized_approval = approval_ref.strip()
    if (
        not normalized_approval
        or "not-approved" in normalized_approval.lower()
        or "example" in normalized_approval.lower()
    ):
        raise LoadHarnessError("arm-fence requires a real recorded approval_ref")
    assert_evidence_exchange_runtime(manifest, run_id=run_id)
    predecessor = verify_predecessor_evidence(
        run_id=run_id,
        stage="load-07",
        manifest=manifest,
        supplied_sha256=supplied_predecessor_sha256,
    )
    source = _assert_source_safety(
        expected_watermark=manifest["expected_stage_start_watermarks"]["load-07"]
    )
    assert_approved_window(manifest, "LOAD-07", now=current)
    policy = manifest["workloads"]["load-07"]["restart_fence_policy"]
    lead = execute_at_epoch - current.timestamp()
    if not policy["minimum_lead_seconds"] <= lead <= policy["maximum_lead_seconds"]:
        raise LoadHarnessError(
            "arm-fence epoch is outside the 60–900 second lead policy"
        )

    _, run_dir, _ = _evidence_path(run_id, "LOAD-07")
    target = run_dir / policy["artifact_filename"]
    if target.exists():
        info = target.stat(follow_symlinks=False)
        if (
            not stat.S_ISREG(info.st_mode)
            or info.st_uid != TIMETABLER_EVIDENCE_OWNER_UID
            or stat.S_IMODE(info.st_mode) != EVIDENCE_EXCHANGE_FILE_MODE
        ):
            raise LoadHarnessError("existing LOAD-07 fence ownership or mode is unsafe")
        try:
            descriptor = os.open(target, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0))
            with os.fdopen(descriptor, "rb") as handle:
                existing = json.loads(handle.read(16_385))
        except (OSError, UnicodeDecodeError, json.JSONDecodeError) as error:
            raise LoadHarnessError("existing LOAD-07 fence is unreadable") from error
        digest = str(existing.get("coordination_sha256", ""))
        verified = verify_restart_fence(
            run_id=run_id,
            manifest=manifest,
            supplied_sha256=digest,
            now=current,
        )
        if (
            verified["execute_at_epoch"] != execute_at_epoch
            or verified["approval_ref"] != normalized_approval
        ):
            raise LoadHarnessError(
                "existing LOAD-07 fence differs from the requested immutable identity"
            )
        artifact = {
            key: value
            for key, value in verified.items()
            if key not in {"artifact_path", "mode", "owner_uid", "reader_uid"}
        }
    else:
        execute_at = datetime.datetime.fromtimestamp(execute_at_epoch, tz=datetime.UTC)
        unsigned = {
            "schema_version": 1,
            "run_id": run_id,
            "timetabler_sha": manifest["commits"]["timetabler"],
            "resource_booking_sha": manifest["commits"]["resource_booking"],
            "execute_at_epoch": execute_at_epoch,
            "nonce": secrets.token_hex(32),
            "approval_ref": normalized_approval,
            "created_at": current.isoformat().replace("+00:00", "Z"),
            "expires_at": (execute_at + datetime.timedelta(seconds=300))
            .isoformat()
            .replace("+00:00", "Z"),
        }
        artifact = {**unsigned, "coordination_sha256": _sha256(unsigned)}
        _assert_exchange_payload_has_no_credentials(artifact)
        _persist_immutable_exchange_file(
            target=target,
            encoded=_canonical_bytes(artifact) + b"\n",
            label="LOAD-07 fence",
        )
        verify_restart_fence(
            run_id=run_id,
            manifest=manifest,
            supplied_sha256=artifact["coordination_sha256"],
            now=current,
        )
    return {
        "schema_version": 1,
        "action": "ARM-FENCE",
        "run_id": run_id,
        "execute_at_epoch": execute_at_epoch,
        "coordination_sha256": artifact["coordination_sha256"],
        "prior_evidence_sha256": supplied_predecessor_sha256,
        "predecessor_profile": predecessor["profile"],
        "source_watermark": source["transport_watermark"],
        "evidence_exchange": evidence_exchange_attestation(),
        "owner_uid": TIMETABLER_EVIDENCE_OWNER_UID,
        "reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
        "directory_mode": "0755",
        "file_mode": "0444",
        "mutation_executed": False,
        "load_executed": False,
        "restart_executed": False,
        "database_mutation_executed": False,
        "domain_mutation_executed": False,
        "resource_booking_database_touched": False,
        "resource_booking_process_touched": False,
        "credentials_written": False,
        "reverse_delivery_enabled": False,
        "phase2_enabled": False,
        "dr_executed": False,
    }


def validate_run_id(run_id: str) -> None:
    if not RUN_ID_PATTERN.fullmatch(run_id):
        raise LoadHarnessError(
            "run id must match phase1-load-YYYYMMDDtHHMMSSz-<12..32 lowercase hex>"
        )


def _parse_iso8601(value: Any, label: str) -> datetime.datetime:
    if not isinstance(value, str) or not value.strip():
        raise LoadHarnessError(f"{label} must be an ISO-8601 timestamp")
    normalized = value.strip().replace("Z", "+00:00")
    try:
        parsed = datetime.datetime.fromisoformat(normalized)
    except ValueError as error:
        raise LoadHarnessError(f"{label} must be an ISO-8601 timestamp") from error
    if parsed.tzinfo is None or parsed.utcoffset() is None:
        raise LoadHarnessError(f"{label} must include an explicit UTC offset")
    return parsed.astimezone(datetime.UTC)


def _contains_credential_key(value: Any) -> bool:
    if isinstance(value, dict):
        for key, item in value.items():
            lowered = str(key).lower()
            if lowered == "credentials_forbidden" and item is True:
                continue
            if any(
                marker in lowered
                for marker in (
                    "password",
                    "secret",
                    "token",
                    "authorization",
                    "credential",
                    "bootstrap_servers",
                )
            ):
                return True
            if _contains_credential_key(item):
                return True
    elif isinstance(value, list):
        return any(_contains_credential_key(item) for item in value)
    return False


def _validate_selector(name: str, selector: Any, run_id: str) -> None:
    if not isinstance(selector, dict) or selector.get("model") not in MODEL_MAP:
        raise LoadHarnessError(f"selector {name} uses an unsupported model")
    filters = selector.get("filters")
    count = selector.get("expected_count")
    if not isinstance(filters, dict) or not filters:
        raise LoadHarnessError(f"selector {name} must use bounded filters")
    if not isinstance(count, int) or not 1 <= count <= MAX_SELECTOR_ROWS:
        raise LoadHarnessError(f"selector {name} expected_count is outside its bound")
    if selector.get("run_owned") is not True or run_id not in json.dumps(filters):
        raise LoadHarnessError(f"selector {name} is not explicitly run-owned")


def _validate_step(step: Any, *, allow_expected_rejection: bool = False) -> None:
    if not isinstance(step, dict) or step.get("kind") != "api":
        raise LoadHarnessError("load workload contains a non-API mutation")
    if step.get("target", "timetabler") != "timetabler":
        raise LoadHarnessError("load workload targets a non-Timetabler service")
    if step.get("path") not in ALLOWED_PATHS:
        raise LoadHarnessError("load workload uses an unsupported Timetabler path")
    if not isinstance(step.get("payload"), dict):
        raise LoadHarnessError("load workload API payload must be an object")
    expected = step.get("expected_source_transactions")
    allowed = {0, 1} if allow_expected_rejection else {1}
    if expected not in allowed:
        raise LoadHarnessError("load workload has an invalid source transaction count")
    if expected == 0 and step.get("expected_rejection") is not True:
        raise LoadHarnessError("a zero-transaction boundary step must expect rejection")


def _workload_steps(manifest: dict[str, Any], stage: str) -> list[dict[str, Any]]:
    raw = manifest["workloads"][stage].get("steps")
    if raw == "canonical_traffic_mix":
        raw = manifest.get("traffic_mix_steps")
    if not isinstance(raw, list):
        raise LoadHarnessError(f"{stage} workload needs bounded API steps")
    return raw


def validate_manifest(
    manifest: Any,
    *,
    run_id: str,
    environment: str = "staging",
    source_scope: str = "default",
    require_approval: bool,
    require_authoritative_term_baseline: bool = False,
) -> dict[str, Any]:
    if not isinstance(manifest, dict) or manifest.get("schema_version") != 1:
        raise LoadHarnessError("load manifest schema_version must equal 1")
    if manifest.get("run_id") != run_id:
        raise LoadHarnessError("load manifest run_id differs from workflow input")
    validate_run_id(run_id)
    if manifest.get("environment") != environment or environment != "staging":
        raise LoadHarnessError(
            "load manifest must target the attested staging environment"
        )
    if manifest.get("source_scope") != source_scope:
        raise LoadHarnessError("load manifest source_scope differs from runtime")
    if _contains_credential_key(manifest):
        raise LoadHarnessError("load manifest contains credential-like configuration")
    if manifest.get("production") is not False:
        raise LoadHarnessError("load manifest must explicitly prohibit production")
    if manifest.get("reverse_delivery_enabled") is not False:
        raise LoadHarnessError(
            "load manifest must explicitly keep reverse delivery false"
        )
    if manifest.get("phase2_enabled") is not False:
        raise LoadHarnessError("load manifest must explicitly keep Phase 2 false")
    if manifest.get("dr_executed") is not False:
        raise LoadHarnessError("load manifest must explicitly exclude DR")
    if manifest.get("load_acceptance_scope") != ACADEMIC_TERM_LOAD_SCOPE:
        raise LoadHarnessError(
            "load acceptance must be term-26 schedule/unschedule recurring assignments only"
        )
    if manifest.get("evidence_exchange") != evidence_exchange_attestation():
        raise LoadHarnessError(
            "load manifest evidence exchange attestation differs from the fixed contract"
        )
    if require_approval:
        if manifest.get("approved") is not True:
            raise LoadHarnessError("load manifest is not approved")
        if not str(manifest.get("approval_ref") or "").strip():
            raise LoadHarnessError("load manifest approval_ref is absent")
        serialized = json.dumps(manifest, sort_keys=True).lower()
        if "example only" in serialized or "not-approved" in serialized:
            raise LoadHarnessError(
                "approved load manifest retains example placeholders"
            )

    commits = manifest.get("commits")
    if not isinstance(commits, dict) or any(
        not GIT_SHA_PATTERN.fullmatch(str(commits.get(key, "")))
        for key in ("timetabler", "resource_booking")
    ):
        raise LoadHarnessError("load manifest requires exact TT and RB commit SHAs")
    if require_approval and any(value == "0" * 40 for value in commits.values()):
        raise LoadHarnessError(
            "approved load manifest retains zero commit placeholders"
        )
    corrections = manifest.get("evidence_corrections")
    if corrections is not None:
        if not isinstance(corrections, dict) or set(corrections) != {"provision"}:
            raise LoadHarnessError(
                "evidence_corrections may declare only the provision correction"
            )
        provision_correction = corrections.get("provision")
        if (
            not isinstance(provision_correction, dict)
            or set(provision_correction) != PROVISION_CORRECTION_KEYS
            or provision_correction.get("artifact_filename")
            not in PROVISION_CORRECTION_ARTIFACT_FILENAMES
            or not SHA256_PATTERN.fullmatch(
                str(provision_correction.get("original_evidence_sha256", ""))
            )
            or not GIT_SHA_PATTERN.fullmatch(
                str(provision_correction.get("original_timetabler_sha", ""))
            )
            or not GIT_SHA_PATTERN.fullmatch(
                str(provision_correction.get("original_resource_booking_sha", ""))
            )
            or provision_correction.get("source_timezone")
            != PROVISION_CORRECTION_SOURCE_TIMEZONE
            or provision_correction.get("reason") != PROVISION_CORRECTION_REASON
        ):
            raise LoadHarnessError(
                "provision evidence correction declaration is incomplete"
            )
        try:
            ZoneInfo(provision_correction["source_timezone"])
        except (KeyError, ValueError) as error:
            raise LoadHarnessError(
                "provision evidence correction timezone is invalid"
            ) from error
        if require_approval and (
            provision_correction["original_evidence_sha256"] == "0" * 64
            or provision_correction["original_timetabler_sha"] == "0" * 40
            or provision_correction["original_resource_booking_sha"] == "0" * 40
        ):
            raise LoadHarnessError(
                "approved provision evidence correction retains placeholders"
            )
    config = manifest.get("configuration_snapshot")
    if (
        not isinstance(config, dict)
        or not SHA256_PATTERN.fullmatch(str(config.get("sha256", "")))
        or not isinstance(config.get("captured_at"), str)
        or config.get("reverse_delivery_enabled") is not False
        or config.get("phase") != "phase1"
        or config.get("transport") != "kafka"
        or not SHA256_PATTERN.fullmatch(
            str(config.get("resource_booking_database_sha256", ""))
        )
        or not SHA256_PATTERN.fullmatch(
            str(config.get("timetabler_database_sha256", ""))
        )
    ):
        raise LoadHarnessError(
            "versioned non-secret configuration snapshot is incomplete"
        )
    if require_approval and any(
        config[key] == "0" * 64
        for key in (
            "sha256",
            "resource_booking_database_sha256",
            "timetabler_database_sha256",
        )
    ):
        raise LoadHarnessError(
            "approved load manifest retains zero configuration hashes"
        )

    approval = manifest.get("threshold_approval")
    thresholds = manifest.get("thresholds")
    if not isinstance(approval, dict) or not str(approval.get("ref") or "").strip():
        raise LoadHarnessError("threshold approval record is incomplete")
    if approval.get("uses_defaults") is True:
        if thresholds != DEFAULT_THRESHOLDS:
            raise LoadHarnessError(
                "default threshold values differ from the approved plan"
            )
    elif (
        not isinstance(thresholds, dict)
        or set(thresholds) != set(DEFAULT_THRESHOLDS)
        or not SHA256_PATTERN.fullmatch(str(approval.get("decision_sha256", "")))
    ):
        raise LoadHarnessError("replacement thresholds lack a versioned joint decision")

    operators = manifest.get("operators")
    if not isinstance(operators, dict):
        raise LoadHarnessError("named operator and abort authority are required")
    if set(operators) != set(REQUIRED_APPROVAL_ROLES):
        raise LoadHarnessError(
            "canonical cross-service approval role inventory is incomplete"
        )
    for key in REQUIRED_APPROVAL_ROLES:
        item = operators.get(key)
        if not isinstance(item, dict) or not all(
            str(item.get(field) or "").strip()
            for field in ("name", "role", "record_ref", "approved_at")
        ):
            raise LoadHarnessError(f"named load role is incomplete: {key}")
        _parse_iso8601(item["approved_at"], f"operators.{key}.approved_at")

    window = manifest.get("approved_window")
    if not isinstance(window, dict) or set(window) != {"starts_at", "ends_at"}:
        raise LoadHarnessError("explicit approved load window is incomplete")
    window_start = _parse_iso8601(window["starts_at"], "approved_window.starts_at")
    window_end = _parse_iso8601(window["ends_at"], "approved_window.ends_at")
    if window_end <= window_start:
        raise LoadHarnessError("approved load window ends before it starts")
    if any(
        _parse_iso8601(operators[key]["approved_at"], f"operators.{key}.approved_at")
        > window_start
        for key in REQUIRED_APPROVAL_ROLES
    ):
        raise LoadHarnessError(
            "every named approval must precede the approved load window"
        )

    profile_timing = manifest.get("profile_timing")
    if profile_timing != DEFAULT_PROFILE_TIMING:
        raise LoadHarnessError(
            "canonical profile workload/observation timing differs from the plan"
        )

    load = manifest.get("load")
    if not isinstance(load, dict):
        raise LoadHarnessError("load parameters are absent")
    effective = {key: load.get(key) for key in DEFAULT_LOAD}
    if effective["p_tps"] != 2 or effective["burst_tps"] != 10:
        if not (
            str(load.get("rate_approval_ref") or "").strip()
            and SHA256_PATTERN.fullmatch(str(load.get("telemetry_sha256", "")))
        ):
            raise LoadHarnessError("non-default P/5P lacks approved telemetry evidence")
    for key, expected in DEFAULT_LOAD.items():
        if key in {"p_tps", "burst_tps"}:
            continue
        if effective[key] != expected:
            raise LoadHarnessError(f"load duration differs from plan: {key}")
    if load.get("burst_tps") != 5 * load.get("p_tps"):
        raise LoadHarnessError("burst rate must equal 5P")
    rate_fidelity = manifest.get("rate_fidelity")
    if rate_fidelity != {
        "maximum_submission_skew_ms": 1000,
        "minimum_achieved_tps_ratio": 0.95,
    }:
        raise LoadHarnessError("approved rate-fidelity tolerance differs from the plan")

    starts = manifest.get("expected_stage_start_watermarks")
    finals = manifest.get("expected_stage_final_watermarks")
    predecessors = manifest.get("stage_predecessors")
    for value, label in ((starts, "start"), (finals, "final")):
        if not isinstance(value, dict) or set(value) != set(STAGES):
            raise LoadHarnessError(f"manifest must declare every stage {label} fence")
    if any(
        not isinstance(value, int) or value < 1
        for value in (*starts.values(), *finals.values())
    ):
        raise LoadHarnessError("stage watermarks must be positive integers")
    if any(finals[stage] < starts[stage] for stage in STAGES):
        raise LoadHarnessError("a stage final watermark precedes its start")
    if not isinstance(predecessors, dict) or set(predecessors) != set(STAGES):
        raise LoadHarnessError(
            "manifest must declare the immutable stage predecessor graph"
        )
    for stage, expected_profile in EXPECTED_STAGE_PREDECESSORS.items():
        predecessor = predecessors.get(stage)
        if (
            not isinstance(predecessor, dict)
            or predecessor.get("profile") != expected_profile
        ):
            raise LoadHarnessError("manifest stage predecessor chain is invalid")
        allowed_keys = (
            {"profile", "external_root_sha256"} if stage == "entry" else {"profile"}
        )
        if set(predecessor) != allowed_keys:
            raise LoadHarnessError(
                "manifest predecessor contains mutable/future hash data"
            )
        if stage == "entry" and not SHA256_PATTERN.fullmatch(
            str(predecessor.get("external_root_sha256", ""))
        ):
            raise LoadHarnessError("ENTRY predecessor lacks its accepted E2E root hash")
    for previous, current in zip(STAGES, STAGES[1:]):
        if finals[previous] != starts[current]:
            raise LoadHarnessError("manifest stage watermark chain is discontinuous")

    selectors = manifest.get("fixture_selectors")
    if not isinstance(selectors, dict) or not 1 <= len(selectors) <= MAX_SELECTORS:
        raise LoadHarnessError("bounded run-owned fixture selectors are required")
    for name, selector in selectors.items():
        _validate_selector(name, selector, run_id)

    workloads = manifest.get("workloads")
    if not isinstance(workloads, dict) or set(workloads) != set(
        EXECUTION_WORKLOAD_STAGES
    ):
        raise LoadHarnessError(
            "manifest must declare each mutating load/cleanup workload"
        )
    for stage, workload in workloads.items():
        if not isinstance(workload, dict):
            raise LoadHarnessError(f"{stage} workload must be an object")
        steps = _workload_steps(manifest, stage)
        if not isinstance(steps, list) or not 1 <= len(steps) <= MAX_STEPS:
            raise LoadHarnessError(f"{stage} workload needs bounded API steps")
        for step in steps:
            _validate_step(step, allow_expected_rejection=stage == "load-06")
    for stage in ("load-02", "load-03", "load-07", "load-08"):
        workload = workloads[stage]
        steps = _workload_steps(manifest, stage)
        classifications = [step.get("workload_classification") for step in steps]
        if tuple(classifications) != REQUIRED_TRAFFIC_STEP_CLASSIFICATIONS:
            raise LoadHarnessError(
                f"{stage} must alternate only term-26 schedule and unschedule traffic"
            )
        for step in steps:
            if step.get("fixture_selector") not in selectors or not isinstance(
                step.get("fixture_variable"), str
            ):
                raise LoadHarnessError(
                    f"{stage} traffic step lacks a run-owned fixture shard"
                )
        expected_requests = {
            "load-02": load["p_tps"] * load["steady_seconds"],
            "load-03": load["burst_tps"] * load["burst_seconds"],
            "load-07": load["p_tps"] * load["restart_seconds"],
            "load-08": load["p_tps"] * load["stability_seconds"],
        }[stage]
        computed_counts = {
            classification: sum(
                1
                for index in range(expected_requests)
                if classifications[index % len(classifications)] == classification
            )
            for classification in REQUIRED_TRAFFIC_CLASSIFICATIONS
        }
        if workload.get("expected_class_counts") != computed_counts:
            raise LoadHarnessError(
                f"{stage} deterministic per-class counts are inconsistent"
            )
    if workloads["load-05"].get("maximum_recurrence_defined") is not True:
        raise LoadHarnessError("LOAD-05 maximum recurrence must be canonically defined")

    boundary = workloads["load-06"]
    if (
        boundary.get("limit_bytes") != TRANSPORT_LIMIT_BYTES
        or len(boundary["steps"]) != 2
        or boundary["steps"][0].get("expected_source_transactions") != 1
        or boundary["steps"][1].get("expected_source_transactions") != 0
        or boundary["steps"][1].get("expected_rejection") is not True
        or not isinstance(boundary.get("rollback_selector"), str)
        or boundary["rollback_selector"] not in selectors
        or not isinstance(boundary.get("below_boundary_tolerance_bytes"), int)
        or not 1 <= boundary["below_boundary_tolerance_bytes"] <= 4096
        or not str(boundary.get("boundary_approval_ref") or "").strip()
        or boundary.get("calibration_read_only") is not True
        or boundary.get("slot_variable") not in manifest.get("variables", {})
        or boundary.get("over_slot_variable") not in manifest.get("variables", {})
        or boundary.get("slot_variable") == boundary.get("over_slot_variable")
    ):
        raise LoadHarnessError("LOAD-06 boundary contract is incomplete")

    cleanup = workloads["cleanup"]
    activity_batch_max = cleanup.get("activity_batch_max_members")
    activity_selector = cleanup.get("activity_selector")
    if (
        not isinstance(activity_batch_max, int)
        or not 1 <= activity_batch_max <= 128
        or activity_selector not in selectors
        or not isinstance(cleanup.get("transport_margin_bytes"), int)
        or not 65_536 <= cleanup["transport_margin_bytes"] <= 524_288
        or cleanup.get("expected_activity_batch_count")
        != math.ceil(
            selectors[activity_selector]["expected_count"] / activity_batch_max
        )
        or cleanup.get("expected_source_transactions")
        != cleanup["expected_activity_batch_count"] + 2
        or len(cleanup["steps"]) != 3
        or cleanup["steps"][0]["path"] != "/api/admin/activity/delete"
        or any(
            step.get("fixture_selector") not in selectors
            for step in cleanup["steps"][1:]
        )
    ):
        raise LoadHarnessError("cleanup bounded batch contract is incomplete")

    restart = workloads["load-07"].get("restart_fence_policy")
    if restart != {
        "armed_by_separate_artifact": True,
        "arm_action": "arm-fence",
        "requires_predecessor_profile": "LOAD-06",
        "artifact_filename": "load-07-fence.json",
        "minimum_lead_seconds": 60,
        "maximum_lead_seconds": 900,
        "maximum_execution_skew_seconds": 5,
        "resource_booking_controlled_by_timetabler": False,
        "rb_observer_must_start_before_epoch": True,
        "domain_mutation": False,
        "database_mutation": False,
        "process_restart": False,
    }:
        raise LoadHarnessError("LOAD-07 immutable restart-fence policy is invalid")

    provision = manifest.get("provision_plan")
    # Scheduling and unscheduling a term-26 activity can each take up to a
    # minute. The rate pools therefore cover that complete two-step
    # asynchronous lifecycle, not merely the engine's 30-second reservation
    # TTL.  This remains a load-harness safety boundary; it does not alter the
    # Timetabler--Scheduling Engine request or response contract.
    required_counts = {
        "lifecycle_staff": 32,
        "lifecycle_location": 32,
        "allocation_staff": 640,
        "allocation_location": 640,
        "traffic_activity": 320,
        "rate_recurrence_activity": 320,
        "recurrence_activity": 32,
        "bulk_activity": 100,
        "boundary_activity": 1200,
    }
    if not isinstance(provision, dict) or provision.get("counts") != required_counts:
        raise LoadHarnessError(
            "provision plan does not match the approved realistic fixture inventory"
        )
    if (
        provision.get("supported_api_only") is not True
        or provision.get("idempotent_retry") is not True
        or provision.get("direct_database_mutation") is not False
        or provision.get("expected_source_transactions")
        != sum(required_counts.values())
        or not isinstance(provision.get("maximum_recurrence_week_count"), int)
        or provision["maximum_recurrence_week_count"] < 1
        or not str(provision.get("maximum_recurrence_approval_ref") or "").strip()
        or provision.get("engine_reservation_ttl_seconds") != 30
        or provision.get("engine_operation_settlement_seconds") != 60
        or provision.get("engine_lifecycle_safety_margin_seconds") != 4
        or provision.get("disjoint_rate_resource_shards") is not True
    ):
        raise LoadHarnessError("provision safety/count contract is incomplete")
    if (
        not str(workloads["load-05"].get("recurrence_contract_ref") or "").strip()
        or workloads["load-05"].get("expected_occurrence_count")
        != required_counts["recurrence_activity"]
        * provision["maximum_recurrence_week_count"]
        or workloads["load-05"].get("expected_occurrence_count") != 384
        or workloads["load-05"]["steps"][0].get("payload", {}).get("activity_ids")
        != "{{maximum_recurrence_activity_ids}}"
    ):
        raise LoadHarnessError(
            "LOAD-05 approved maximum recurrence count is inconsistent"
        )
    traffic_selector_names = {
        step["fixture_selector"] for step in manifest["traffic_mix_steps"]
    }
    traffic_selector_order = tuple(
        step["fixture_selector"] for step in manifest["traffic_mix_steps"]
    )
    burst_reuse_seconds = (
        required_counts["traffic_activity"]
        * len(manifest["traffic_mix_steps"])
        / load["burst_tps"]
    )
    minimum_lifecycle_seconds = (
        2 * provision["engine_operation_settlement_seconds"]
        + provision["engine_lifecycle_safety_margin_seconds"]
    )
    if (
        len(traffic_selector_names) != 2
        or traffic_selector_order
        != (
            "traffic_activity_ids",
            "traffic_activity_ids",
            "traffic_recurrence_activity_ids",
            "traffic_recurrence_activity_ids",
        )
        or selectors["lifecycle_staff_ids"]["expected_count"] != 32
        or selectors["lifecycle_location_ids"]["expected_count"] != 32
        or selectors["traffic_activity_ids"]["expected_count"]
        != required_counts["traffic_activity"]
        or selectors["traffic_recurrence_activity_ids"]["expected_count"]
        != required_counts["rate_recurrence_activity"]
        or selectors["maximum_recurrence_activity_ids"]["expected_count"] != 32
        or burst_reuse_seconds <= minimum_lifecycle_seconds
        or required_counts["allocation_staff"]
        < required_counts["traffic_activity"]
        + required_counts["rate_recurrence_activity"]
        or required_counts["allocation_location"]
        < required_counts["traffic_activity"]
        + required_counts["rate_recurrence_activity"]
    ):
        raise LoadHarnessError(
            "realistic traffic requires disjoint lifecycle-safe fixture sharding"
        )

    bulk = workloads["load-04"]
    if (
        not isinstance(bulk.get("approved_real_maximum_activity_count"), int)
        or bulk["approved_real_maximum_activity_count"] < 1
        or not isinstance(bulk.get("approved_maximum_member_count"), int)
        or bulk["approved_maximum_member_count"] < 1
        or bulk["approved_real_maximum_activity_count"]
        != bulk["approved_maximum_member_count"]
        or not str(bulk.get("maximum_approval_ref") or "").strip()
        or bulk.get("fixture_selector") not in selectors
        or selectors[bulk["fixture_selector"]]["expected_count"]
        != bulk["approved_real_maximum_activity_count"]
    ):
        raise LoadHarnessError(
            "LOAD-04 approved real maximum/member reference is incomplete"
        )

    baseline = manifest.get("accepted_baseline")
    if not isinstance(baseline, dict) or any(
        isinstance(baseline.get(key), bool)
        or not isinstance(baseline.get(key), int)
        or baseline[key] < 1
        for key in ("resources", "activities", "source_watermark")
    ):
        raise LoadHarnessError(
            "accepted Phase 1 baseline counts and source fence must be positive integers"
        )
    if any(
        not SHA256_PATTERN.fullmatch(str(baseline.get(key, "")))
        for key in ("timetabler_evidence_sha256", "resource_booking_evidence_sha256")
    ):
        raise LoadHarnessError("accepted baseline evidence hashes are absent")
    authoritative_term = baseline.get("authoritative_term")
    if require_authoritative_term_baseline or baseline["source_watermark"] >= 2640:
        expected_term = {
            "academic_term_id": 26,
            "start_date": "2026-01-05",
            "end_date": "2026-06-26",
            "activity_count": 45,
            "first_activity_id": 672,
            "last_activity_id": 716,
            "activity_ids_sha256": "1b599dcd8709fcdb29b292b110ef1e5cb011929ba371011ef1ae055c9e408453",
            "source_transaction_id": "fc63fa1f-09e6-5af6-b4a9-699cebfd6693",
            "source_watermark": 2640,
        }
        if authoritative_term != expected_term:
            raise LoadHarnessError(
                "accepted baseline does not bind authoritative term 26 activities 672-716"
            )
        if baseline["source_watermark"] < authoritative_term["source_watermark"]:
            raise LoadHarnessError(
                "accepted baseline source fence predates the authoritative term projection"
            )
    historical = manifest.get("historical_projection_milestone")
    if not isinstance(historical, dict) or {
        "resources": historical.get("resources"),
        "activities": historical.get("activities"),
        "projection_revision_transactions": historical.get(
            "projection_revision_transactions"
        ),
    } != {"resources": 1794, "activities": 46, "projection_revision_transactions": 19}:
        raise LoadHarnessError(
            "historical fixed-watermark milestone provenance differs"
        )
    historical_record = {
        key: value for key, value in historical.items() if key != "record_sha256"
    }
    if (
        set(historical_record)
        != {
            "resources",
            "activities",
            "projection_revision_transactions",
            "timetabler_run_id",
            "timetabler_run_url",
            "resource_booking_run_id",
            "resource_booking_run_url",
        }
        or not str(historical["timetabler_run_id"]).isdigit()
        or not str(historical["resource_booking_run_id"]).isdigit()
        or historical["timetabler_run_url"]
        != (
            "https://github.com/Mayvins/timetabler-be/actions/runs/"
            f"{historical['timetabler_run_id']}"
        )
        or historical["resource_booking_run_url"]
        != (
            "https://github.com/Mayvins/resource-booking-be/actions/runs/"
            f"{historical['resource_booking_run_id']}"
        )
        or historical["record_sha256"] == "0" * 64
        or historical["record_sha256"] != _sha256(historical_record)
    ):
        raise LoadHarnessError(
            "historical milestone provenance record is contradictory"
        )
    if predecessors["entry"]["external_root_sha256"] != manifest.get(
        "e2e_final_evidence_sha256"
    ):
        raise LoadHarnessError(
            "ENTRY predecessor does not match accepted E2E final evidence"
        )
    receiver = manifest.get("load01_receiver_attestation")
    if (
        not isinstance(receiver, dict)
        or receiver.get("unfiltered_reconciliation") is not True
        or receiver.get("source_watermark") != starts["load-01"]
        or not SHA256_PATTERN.fullmatch(str(receiver.get("evidence_sha256", "")))
        or not isinstance(receiver.get("completed_at"), str)
        or receiver.get("counters")
        != {
            "drift": 0,
            "repaired": 0,
            "unresolved": 0,
            "quarantined": 0,
            "dead_letter": 0,
            "gap": 0,
            "mirror_only_violations": 0,
            "backlog": 0,
        }
        or receiver.get("resources") != baseline["resources"]
        or receiver.get("activities") != baseline["activities"]
    ):
        raise LoadHarnessError(
            "LOAD-01 current unfiltered receiver attestation is incomplete"
        )
    if require_approval and receiver["evidence_sha256"] == "0" * 64:
        raise LoadHarnessError(
            "approved LOAD-01 receiver evidence hash is a placeholder"
        )
    _parse_iso8601(receiver["completed_at"], "load01_receiver_attestation.completed_at")

    expected_deltas = {
        "entry": 0,
        "load-01": 0,
        "provision": provision["expected_source_transactions"],
        "load-02": load["p_tps"] * load["steady_seconds"],
        "load-03": load["burst_tps"] * load["burst_seconds"],
        "load-04": 1,
        "load-05": 1,
        "load-06": 1,
        "load-07": load["p_tps"] * load["restart_seconds"],
        "load-08": load["p_tps"] * load["stability_seconds"],
        "cleanup": workloads["cleanup"]["expected_source_transactions"],
        "final": 0,
    }
    if any(finals[stage] - starts[stage] != expected_deltas[stage] for stage in STAGES):
        raise LoadHarnessError(
            "manifest stage watermark delta differs from its exact workload"
        )
    return manifest


def assert_approved_window(manifest: dict[str, Any], profile: str, *, now=None) -> None:
    now = now or datetime.datetime.now(datetime.UTC)
    start = _parse_iso8601(
        manifest["approved_window"]["starts_at"], "approved_window.starts_at"
    )
    end = _parse_iso8601(
        manifest["approved_window"]["ends_at"], "approved_window.ends_at"
    )
    if not start <= now <= end:
        raise LoadHarnessError(
            "current time is outside the explicitly approved load window"
        )
    observation_seconds = manifest["profile_timing"][profile]["observation_seconds"]
    if now + datetime.timedelta(seconds=observation_seconds) > end:
        raise LoadHarnessError(
            "approved load window cannot contain the complete profile observation"
        )


def load_manifest(raw: str, **kwargs: Any) -> dict[str, Any]:
    if len(raw.encode()) > MAX_MANIFEST_BYTES:
        raise LoadHarnessError("load manifest exceeds its safety bound")
    try:
        parsed = json.loads(raw)
    except json.JSONDecodeError as error:
        raise LoadHarnessError("load manifest is not valid JSON") from error
    return validate_manifest(parsed, **kwargs)


def assert_evidence_exchange_runtime(manifest: dict[str, Any], *, run_id: str) -> None:
    if manifest.get("evidence_exchange") != evidence_exchange_attestation():
        raise LoadHarnessError("runtime evidence exchange contract is contradictory")
    _ensure_evidence_exchange_directories(run_id)


def read_only_configuration_attestation() -> dict[str, Any]:
    """Emit only hashed/non-secret facts needed to build the immutable manifest."""

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

    exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
    if not GIT_SHA_PATTERN.fullmatch(exact_sha):
        raise LoadHarnessError("read-only attestation requires the exact deployed SHA")
    integration = settings.RESOURCE_BOOKING_INTEGRATION
    topic = str(integration.get("KAFKA_TOPIC") or "")
    database = connection.settings_dict
    database_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 ""),
    }
    timetabler_database_sha256 = _sha256(database_fingerprint)
    configuration = {
        "capture_enabled": bool(integration["CAPTURE_ENABLED"]),
        "publish_enabled": bool(integration["PUBLISH_ENABLED"]),
        "activation_approved": bool(integration["ACTIVATION_APPROVED"]),
        "phase": str(integration["PHASE"]),
        "reverse_delivery_enabled": bool(integration["REVERSE_DELIVERY_ENABLED"]),
        "transport": str(integration["TRANSPORT"]),
        "source_scope": str(integration["SOURCE_SCOPE"]),
        "schema_version": int(integration["SCHEMA_VERSION"]),
        "topic_sha256": hashlib.sha256(topic.encode()).hexdigest(),
        "max_event_bytes": int(integration["MAX_EVENT_BYTES"]),
        "timetabler_database_sha256": timetabler_database_sha256,
    }
    configuration_sha256 = _sha256(configuration)
    metrics = _assert_source_safety()
    if (
        configuration["phase"] != "phase1"
        or configuration["reverse_delivery_enabled"]
        or configuration["transport"] != "kafka"
    ):
        raise LoadHarnessError(
            "read-only configuration attestation is not Phase 1 safe"
        )
    return {
        "schema_version": 1,
        "profile": "ATTEST",
        "service": "timetabler",
        "repository": "Mayvins/timetabler-be",
        "git_sha": exact_sha,
        "source": metrics,
        "configuration": configuration,
        "configuration_sha256": configuration_sha256,
        "configuration_sha256_log_parts": _sha256_log_parts(configuration_sha256),
        "timetabler_database_sha256_log_parts": _sha256_log_parts(
            timetabler_database_sha256
        ),
        "runtime_uid": os.getuid(),
        "evidence_exchange": evidence_exchange_attestation(),
        "evidence_exchange_owner_matches_runtime": os.getuid()
        == TIMETABLER_EVIDENCE_OWNER_UID,
        "production": False,
        "mutation_executed": False,
        "load_executed": False,
        "reverse_delivery_enabled": False,
        "phase2_enabled": False,
        "dr_executed": False,
    }


def _render(value: Any, variables: dict[str, Any]) -> Any:
    if isinstance(value, dict):
        return {key: _render(item, variables) for key, item in value.items()}
    if isinstance(value, list):
        return [_render(item, variables) for item in value]
    if not isinstance(value, str):
        return value
    whole = re.fullmatch(r"\{\{([a-zA-Z0-9_]+)}}", value)
    if whole:
        name = whole.group(1)
        if name not in variables:
            raise LoadHarnessError(f"manifest variable is unresolved: {name}")
        return variables[name]
    result = value
    for name in re.findall(r"\{\{([a-zA-Z0-9_]+)}}", value):
        if name not in variables:
            raise LoadHarnessError(f"manifest variable is unresolved: {name}")
        result = result.replace(f"{{{{{name}}}}}", str(variables[name]))
    return result


def _resolve_selectors(
    manifest: dict[str, Any], *, include_expected_absent: bool = True
) -> dict[str, list[int]]:
    from api import models

    resolved: dict[str, list[int]] = {}
    for name, selector in manifest["fixture_selectors"].items():
        if (
            not include_expected_absent
            and selector.get("expected_absent_after_cleanup") is True
        ):
            continue
        model = getattr(models, MODEL_MAP[selector["model"]])
        expected = selector["expected_count"]
        rows = list(
            model.objects.filter(**selector["filters"])
            .order_by("id")
            .values_list("id", flat=True)[: expected + 1]
        )
        if len(rows) != expected:
            raise LoadHarnessError(
                f"run-owned selector {name} resolved {len(rows)} rows, expected {expected}"
            )
        values = [int(row) for row in rows]
        resolved[name] = values
    return resolved


def _resolve_cleanup_selectors(
    manifest: dict[str, Any], run_id: str
) -> dict[str, list[int]]:
    """Recover the immutable provision inventory, including already-deleted retry batches."""

    from api import models
    from api.models import IntegrationOutbox

    provision_rows = list(
        IntegrationOutbox.objects.filter(
            request_id__startswith=f"{run_id}-provision-",
            transaction_finalized=True,
            event_type__in=(
                "timetabler.staff.created",
                "timetabler.location.created",
                "timetabler.activity.created",
            ),
        ).values_list("aggregate_type", "payload")
    )
    inventory: list[tuple[str, int, str]] = []
    for aggregate_type, payload in provision_rows:
        current = (
            payload.get("committed_state", {}).get("replacement", {}).get("current")
        )
        if not isinstance(current, dict):
            raise LoadHarnessError(
                "provision inventory contains a non-current create event"
            )
        source_id = current.get("id")
        code = current.get("code")
        if (
            not isinstance(source_id, int)
            or not isinstance(code, str)
            or run_id not in code
        ):
            raise LoadHarnessError("provision inventory identity is incomplete")
        inventory.append((aggregate_type, source_id, code))

    resolved: dict[str, list[int]] = {}
    for name, selector in manifest["fixture_selectors"].items():
        filters = selector["filters"]
        prefix = filters.get("code__startswith")
        if set(filters) != {"code__startswith"} or not isinstance(prefix, str):
            raise LoadHarnessError(
                "cleanup retry supports only exact run-code selectors"
            )
        aggregate_type = selector["model"]
        ids = sorted(
            {
                source_id
                for row_type, source_id, code in inventory
                if row_type == aggregate_type and code.startswith(prefix)
            }
        )
        if len(ids) != selector["expected_count"]:
            raise LoadHarnessError(
                "immutable provision inventory is incomplete for cleanup"
            )
        model = getattr(models, MODEL_MAP[aggregate_type])
        present = list(
            model.objects.filter(id__in=ids).order_by("id").values_list("id", "code")
        )
        if any(not code.startswith(prefix) for _, code in present):
            raise LoadHarnessError("cleanup selector collides with a non-run-owned row")
        resolved[name] = ids
    return resolved


def _abandoned_cleanup_retry_state(
    *,
    run_id: str,
    initial_watermark: int,
    current_watermark: int,
    expected_current_watermark: int,
    cleanup_transaction_count: int,
    completed_rows: list[tuple[str, int]],
) -> dict[str, int]:
    """Prove an abandoned cleanup retry is one exact contiguous prefix."""

    prefix = f"{run_id}-cleanup-"
    completed: dict[int, int] = {}
    for request_id, sequence in completed_rows:
        suffix = str(request_id).removeprefix(prefix)
        if not str(request_id).startswith(prefix) or not re.fullmatch(r"\d{6}", suffix):
            raise LoadHarnessError("abandoned cleanup has an invalid request id")
        index = int(suffix)
        if not 1 <= index <= cleanup_transaction_count:
            raise LoadHarnessError("abandoned cleanup request id exceeds its fixed bound")
        if index in completed and completed[index] != int(sequence):
            raise LoadHarnessError("abandoned cleanup request maps to multiple sequences")
        completed[index] = int(sequence)
    completed_count = len(completed)
    if set(completed) != set(range(1, completed_count + 1)):
        raise LoadHarnessError("abandoned cleanup transactions are not a contiguous prefix")
    if any(
        completed[index]
        != initial_watermark + index
        for index in completed
    ):
        raise LoadHarnessError("abandoned cleanup transaction sequence is outside its fence")
    approved_current = initial_watermark + completed_count
    if (
        expected_current_watermark != approved_current
        or current_watermark != approved_current
    ):
        raise LoadHarnessError("abandoned cleanup source watermark differs from its exact fence")
    return {
        "initial_watermark": initial_watermark,
        "current_watermark": approved_current,
        "completed_transaction_count": completed_count,
        "remaining_transaction_count": cleanup_transaction_count - completed_count,
        "final_watermark": initial_watermark + cleanup_transaction_count,
    }


def _current_watermark() -> int:
    from django.db.models import Max

    from api.models import IntegrationOutbox
    from api.services.integration.outbox import integration_source_scope

    return int(
        IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(), transaction_finalized=True
        ).aggregate(value=Max("transport_sequence"))["value"]
        or 0
    )


def _source_metrics() -> dict[str, Any]:
    from django.db import connection
    from django.db.models import Max, Min
    from django.utils import timezone

    from api.models import (
        EngineResponseQuarantine,
        IntegrationOutbox,
        PostCommitDelivery,
    )
    from api.services.integration.outbox import integration_source_scope
    from api.services.integration.readiness import phase1_readiness

    scope = integration_source_scope()
    rows = IntegrationOutbox.objects.filter(source_scope=scope)
    backlog = rows.filter(status__in=("pending", "publishing", "retry"))
    oldest = backlog.aggregate(value=Min("created_at"))["value"]
    sequence = rows.aggregate(
        minimum=Min("transport_sequence"), maximum=Max("transport_sequence")
    )
    readiness = phase1_readiness(require_enabled=True)
    now = timezone.now()
    return {
        "captured_at": now.isoformat(),
        "transport_watermark": int(sequence["maximum"] or 0),
        "minimum_sequence": int(sequence["minimum"] or 0),
        "publisher_last_sequence": int(readiness["publisher_last_sequence"]),
        "publisher_live": bool(readiness["publisher_live"]),
        "publisher_status": readiness["publisher_status"],
        "outbox_depth": backlog.count(),
        "oldest_outbox_age_seconds": (
            max(0.0, (now - oldest).total_seconds()) if oldest else 0.0
        ),
        "outbox_dead_letter_count": rows.filter(status="dead_letter").count(),
        "post_commit_dead_letter_count": PostCommitDelivery.objects.filter(
            status="dead_letter"
        ).count(),
        "engine_response_quarantine_count": EngineResponseQuarantine.objects.count(),
        "reverse_delivery_enabled": bool(readiness["reverse_delivery_enabled"]),
        "database_connection": {
            "vendor": connection.vendor,
            "usable": bool(connection.is_usable()),
            "in_atomic_block": bool(connection.in_atomic_block),
        },
    }


def _assert_source_safety(*, expected_watermark: int | None = None) -> dict[str, Any]:
    metrics = _source_metrics()
    if metrics["reverse_delivery_enabled"]:
        raise LoadHarnessError("reverse delivery is enabled")
    if not metrics["publisher_live"]:
        raise LoadHarnessError("ordered publisher is not live")
    if metrics["publisher_last_sequence"] != metrics["transport_watermark"]:
        raise LoadHarnessError("ordered publisher is not caught up at the source fence")
    for key in (
        "outbox_dead_letter_count",
        "post_commit_dead_letter_count",
        "engine_response_quarantine_count",
        "outbox_depth",
    ):
        if metrics[key] != 0:
            raise LoadHarnessError(f"source correctness counter is non-zero: {key}")
    if (
        expected_watermark is not None
        and metrics["transport_watermark"] != expected_watermark
    ):
        raise LoadHarnessError("source watermark differs from the exact approved fence")
    return metrics


def _read_proc_host_sample() -> dict[str, Any]:
    memory: dict[str, int] = {}
    with open("/proc/meminfo", encoding="utf-8") as handle:
        for line in handle:
            key, raw = line.split(":", 1)
            value = raw.strip().split()[0]
            if value.isdigit():
                memory[key] = int(value)
    total = memory.get("MemTotal", 0)
    available = memory.get("MemAvailable", 0)
    if total <= 0 or not 0 <= available <= total:
        raise LoadHarnessError("host memory telemetry is unavailable")
    with open("/proc/stat", encoding="utf-8") as handle:
        cpu = handle.readline().split()
    if not cpu or cpu[0] != "cpu" or len(cpu) < 8:
        raise LoadHarnessError("host CPU telemetry is unavailable")
    ticks = [int(value) for value in cpu[1:]]
    idle = ticks[3] + (ticks[4] if len(ticks) > 4 else 0)
    return {
        "captured_at_epoch": time.time(),
        "memory_percent": round((total - available) * 100 / total, 3),
        "cpu_total_ticks": sum(ticks),
        "cpu_idle_ticks": idle,
    }


class HostSampler:
    def __init__(
        self,
        *,
        interval_seconds: int,
        hard_memory_percent: float,
        require_processes: bool = False,
        publisher_restart_window: tuple[float, float] | None = None,
    ):
        self.interval_seconds = interval_seconds
        self.hard_memory_percent = hard_memory_percent
        self.require_processes = require_processes
        self.publisher_restart_window = publisher_restart_window
        self.stop_event = threading.Event()
        self.abort_event = threading.Event()
        self.samples: list[dict[str, Any]] = []
        self.error: str | None = None
        self._thread: threading.Thread | None = None
        self._memory_breach_started_at: float | None = None

    def start(self) -> None:
        self._thread = threading.Thread(target=self._run, daemon=True)
        self._thread.start()

    def stop(self) -> list[dict[str, Any]]:
        self.stop_event.set()
        if self._thread:
            self._thread.join(timeout=self.interval_seconds + 2)
        if self.error:
            raise LoadHarnessError(self.error)
        return self.samples

    def _run(self) -> None:
        previous: dict[str, Any] | None = None
        try:
            while not self.stop_event.is_set():
                current = _read_proc_host_sample()
                if previous:
                    total_delta = (
                        current["cpu_total_ticks"] - previous["cpu_total_ticks"]
                    )
                    idle_delta = current["cpu_idle_ticks"] - previous["cpu_idle_ticks"]
                    current["cpu_percent"] = (
                        round((total_delta - idle_delta) * 100 / total_delta, 3)
                        if total_delta > 0
                        else 0.0
                    )
                else:
                    current["cpu_percent"] = None
                current.pop("cpu_total_ticks")
                current.pop("cpu_idle_ticks")
                current["process_state"] = _process_counts()
                self.samples.append(current)
                previous = _read_proc_host_sample()
                if self.require_processes and (
                    current["process_state"]["application_processes"] < 1
                    or (
                        current["process_state"]["publisher_processes"] != 1
                        and not (
                            self.publisher_restart_window
                            and self.publisher_restart_window[0]
                            <= time.time()
                            <= self.publisher_restart_window[1]
                        )
                    )
                ):
                    self.error = (
                        "required Timetabler application/publisher process failed"
                    )
                    self.abort_event.set()
                    return
                if self._memory_breach_is_sustained(
                    current["memory_percent"], now_monotonic=time.monotonic()
                ):
                    self.error = (
                        "host memory exceeded the 90 percent hard-abort threshold "
                        "continuously for 300 seconds"
                    )
                    self.abort_event.set()
                    return
                self.stop_event.wait(self.interval_seconds)
        except (OSError, ValueError, LoadHarnessError) as error:
            self.error = f"host telemetry failed: {type(error).__name__}"
            self.abort_event.set()

    def _memory_breach_is_sustained(
        self, memory_percent: float, *, now_monotonic: float
    ) -> bool:
        if memory_percent <= self.hard_memory_percent:
            self._memory_breach_started_at = None
            return False
        if self._memory_breach_started_at is None:
            self._memory_breach_started_at = now_monotonic
            return False
        return (
            now_monotonic - self._memory_breach_started_at
            >= SUSTAINED_HARD_MEMORY_SECONDS
        )


def _process_counts() -> dict[str, int]:
    app_dir = str(Path(__file__).resolve().parents[1])
    publisher = 0
    application = 0
    proc = Path("/proc")
    for child in proc.iterdir():
        if not child.name.isdigit():
            continue
        try:
            command = (
                (child / "cmdline")
                .read_bytes()
                .replace(b"\0", b" ")
                .decode("utf-8", errors="replace")
            )
            process_cwd = str((child / "cwd").resolve())
        except (OSError, PermissionError):
            continue
        if app_dir not in command and process_cwd != app_dir:
            continue
        if "publish_resource_booking_outbox" in command:
            publisher += 1
        elif any(
            marker in command
            for marker in ("gunicorn", "manage.py runserver", "uvicorn")
        ):
            application += 1
    return {"publisher_processes": publisher, "application_processes": application}


def _percentiles(values: list[float]) -> dict[str, Any]:
    if not values:
        return {
            "count": 0,
            "method": "nearest-rank",
            "p50": None,
            "p95": None,
            "p99": None,
        }
    ordered = sorted(values)

    def nearest(percent: float) -> float:
        return round(ordered[max(0, math.ceil(percent * len(ordered)) - 1)], 3)

    return {
        "count": len(ordered),
        "method": "nearest-rank",
        "p50": nearest(0.50),
        "p95": nearest(0.95),
        "p99": nearest(0.99),
    }


def _shutdown_executors(
    executors: list[ThreadPoolExecutor], *, abort_pending: bool
) -> None:
    for executor in executors:
        executor.shutdown(wait=True, cancel_futures=abort_pending)


def _rate_fidelity_metrics(
    measurements: list[dict[str, Any]], *, rate: int, policy: dict[str, Any]
) -> dict[str, float]:
    if not measurements or rate <= 0:
        raise LoadHarnessError("rate-fidelity measurement set is empty")
    started_epochs = sorted(item["submitted_at_epoch"] for item in measurements)
    actual_span = max(1 / rate, started_epochs[-1] - started_epochs[0] + 1 / rate)
    achieved_tps = len(measurements) / actual_span
    maximum_skew_ms = max(
        item["submission_scheduling_skew_ms"] for item in measurements
    )
    if (
        achieved_tps < rate * policy["minimum_achieved_tps_ratio"]
        or maximum_skew_ms > policy["maximum_submission_skew_ms"]
    ):
        raise LoadHarnessError(
            "accepted request starts did not sustain the approved rate"
        )
    return {
        "achieved_request_start_tps": round(achieved_tps, 6),
        "maximum_submission_scheduling_skew_ms": maximum_skew_ms,
    }


def _has_sustained_host_breach(
    samples: list[dict[str, Any]], *, metric: str, threshold: float, seconds: int = 300
) -> bool:
    breach_started: float | None = None
    for sample in sorted(samples, key=lambda item: item["captured_at_epoch"]):
        value = sample.get(metric)
        if value is None or value <= threshold:
            breach_started = None
            continue
        if breach_started is None:
            breach_started = sample["captured_at_epoch"]
        elif sample["captured_at_epoch"] - breach_started >= seconds:
            return True
    return False


def _assert_complete_transaction_members(members: list[Any]) -> None:
    member_count = len(members)
    if (
        member_count < 1
        or any(not member.transaction_finalized for member in members)
        or any(member.transaction_count != member_count for member in members)
        or [member.transaction_index for member in members]
        != list(range(1, member_count + 1))
    ):
        raise LoadHarnessError(
            "source transaction is not finalized with complete member boundaries"
        )


def _transaction_evidence(
    start: int, end: int, request_ids: set[str]
) -> list[dict[str, Any]]:
    from api.models import IntegrationOutbox
    from api.services.integration.outbox import (
        build_transport_envelope,
        canonical_hash,
        integration_source_scope,
    )

    rows = list(
        IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(),
            transport_sequence__gt=start,
            transport_sequence__lte=end,
        ).order_by("transport_sequence", "transaction_index")
    )
    grouped: dict[int, list[Any]] = {}
    for row in rows:
        grouped.setdefault(int(row.transport_sequence), []).append(row)
    expected_sequences = list(range(start + 1, end + 1))
    if list(grouped) != expected_sequences:
        raise LoadHarnessError(
            "source sequence gap or unattributed sequence range detected"
        )

    evidence = []
    seen_requests: set[str] = set()
    for sequence, members in grouped.items():
        _assert_complete_transaction_members(members)
        if any(member.request_id not in request_ids for member in members):
            raise LoadHarnessError(
                "source transaction is not attributable to this load run"
            )
        request_set = {str(member.request_id) for member in members}
        if len(request_set) != 1:
            raise LoadHarnessError(
                "one source transaction has contradictory request attribution"
            )
        request_id = next(iter(request_set))
        seen_requests.add(request_id)
        if any(
            member.payload.get("committed_state_hash")
            != canonical_hash(member.payload.get("committed_state"))
            for member in members
        ):
            raise LoadHarnessError(
                "canonical member committed-state hash is contradictory"
            )
        envelope = build_transport_envelope(members)
        if envelope["events_hash"] != canonical_hash(envelope["events"]):
            raise LoadHarnessError("canonical transaction events hash is contradictory")
        statuses = {member.status for member in members}
        if statuses != {"published"}:
            raise LoadHarnessError("source transaction is not completely published")
        created = min(member.created_at for member in members)
        published = max(member.published_at for member in members)
        if published is None:
            raise LoadHarnessError(
                "published source transaction lacks publication time"
            )
        evidence.append(
            {
                "source_sequence": sequence,
                "source_transaction_id": str(members[0].change_set_id),
                "request_id": request_id,
                "transaction_count": len(members),
                "member_count": len(members),
                "event_types": sorted({member.event_type for member in members}),
                "events_hash": envelope["events_hash"],
                "canonical_bytes": len(_canonical_bytes(envelope)),
                "outbox_to_publish_ms": round(
                    (published - created).total_seconds() * 1000, 3
                ),
                "committed_at": _timezone_qualified_iso(created),
                "published_at": _timezone_qualified_iso(published),
            }
        )
    if seen_requests != request_ids:
        missing = sorted(request_ids - seen_requests)[:10]
        raise LoadHarnessError(
            f"accepted requests lack exact source transactions: {missing}"
        )
    return evidence


def _rate_request_id(run_id: str, stage: str, iteration: int) -> str:
    return f"{run_id}-{stage}-{iteration:06d}"


def _classify_rate_transaction_range(
    *,
    transactions: list[dict[str, Any]],
    steps: list[dict[str, Any]],
    run_id: str,
    stage: str,
    approved_start: int,
    first_iteration: int,
) -> None:
    for offset, transaction in enumerate(transactions):
        iteration = first_iteration + offset
        expected_request_id = _rate_request_id(run_id, stage, iteration)
        if (
            transaction.get("source_sequence") != approved_start + iteration
            or transaction.get("request_id") != expected_request_id
        ):
            raise LoadHarnessError(
                "rate source evidence is not the exact deterministic contiguous sequence"
            )
        transaction["workload_classification"] = steps[(iteration - 1) % len(steps)][
            "workload_classification"
        ]


def _load02_resume_prefix_evidence(
    *,
    manifest: dict[str, Any],
    run_id: str,
    current_watermark: int,
) -> list[dict[str, Any]]:
    approved_start = manifest["expected_stage_start_watermarks"]["load-02"]
    approved_final = manifest["expected_stage_final_watermarks"]["load-02"]
    if current_watermark == approved_start:
        return []
    if not approved_start < current_watermark < approved_final:
        raise LoadHarnessError(
            "LOAD-02 watermark is outside its strict partial-prefix resume fence"
        )
    prefix_count = current_watermark - approved_start
    request_ids = {
        _rate_request_id(run_id, "load-02", iteration)
        for iteration in range(1, prefix_count + 1)
    }
    transactions = _transaction_evidence(approved_start, current_watermark, request_ids)
    if len(transactions) != prefix_count:
        raise LoadHarnessError("LOAD-02 durable prefix transaction count differs")
    _classify_rate_transaction_range(
        transactions=transactions,
        steps=_workload_steps(manifest, "load-02"),
        run_id=run_id,
        stage="load-02",
        approved_start=approved_start,
        first_iteration=1,
    )
    return transactions


def _load02_source_progress_attestation(
    *,
    manifest: dict[str, Any],
    run_id: str,
    current_watermark: int,
) -> dict[str, Any]:
    prefix_transactions = _load02_resume_prefix_evidence(
        manifest=manifest,
        run_id=run_id,
        current_watermark=current_watermark,
    )
    return {
        "historical_provision_final_watermark": manifest[
            "expected_stage_final_watermarks"
        ]["provision"],
        "load02_approved_final_watermark": manifest["expected_stage_final_watermarks"][
            "load-02"
        ],
        "current_source_watermark": current_watermark,
        "load02_prefix_transaction_count": len(prefix_transactions),
        "load02_prefix_source_transactions_sha256": _sha256(prefix_transactions),
        "load02_prefix_verified": True,
        "source_transactions_in_corrected_provision_evidence": manifest[
            "provision_plan"
        ]["expected_source_transactions"],
        "prefix_source_transactions_in_corrected_provision_evidence": 0,
    }


def _bounded_request_attribution(
    *, actual_request_id: Any, expected_request_id: str, run_id: str
) -> dict[str, Any]:
    actual = "" if actual_request_id is None else str(actual_request_id)
    run_prefix = f"{run_id}-"
    safe_suffix: str | None = None
    if actual.startswith(run_prefix):
        candidate = actual[len(run_prefix) :]
        if SAFE_REQUEST_SUFFIX_PATTERN.fullmatch(candidate):
            safe_suffix = candidate
    if actual == expected_request_id:
        classification = "exact_expected_request_id"
    elif actual.startswith(f"{run_id}-load-02-"):
        classification = "same_run_load02_variant"
    elif actual.startswith(run_prefix):
        classification = "same_run_other_stage_or_shape"
    elif not actual:
        classification = "missing_request_id"
    else:
        classification = "foreign_or_unclassified_request_id"
    return {
        "classification": classification,
        "request_id_present": bool(actual),
        "request_id_length": len(actual),
        "request_id_sha256": hashlib.sha256(actual.encode("utf-8")).hexdigest(),
        "safe_same_run_suffix": safe_suffix,
    }


def _integer_ranges(values: list[int]) -> list[dict[str, int]]:
    if not values:
        return []
    ranges: list[dict[str, int]] = []
    start = previous = values[0]
    for value in values[1:]:
        if value == previous + 1:
            previous = value
            continue
        ranges.append({"start": start, "end": previous})
        start = previous = value
    ranges.append({"start": start, "end": previous})
    return ranges


def _diagnose_load02_prefix_rows(
    *, rows: list[Any], manifest: dict[str, Any], run_id: str, current_watermark: int
) -> dict[str, Any]:
    approved_start = manifest["expected_stage_start_watermarks"]["load-02"]
    approved_final = manifest["expected_stage_final_watermarks"]["load-02"]
    if not approved_start < current_watermark < approved_final:
        raise LoadHarnessError(
            "LOAD-02 diagnostic watermark is outside its strict open recovery fence"
        )
    grouped: dict[int, list[Any]] = {}
    for row in rows:
        grouped.setdefault(int(row.transport_sequence), []).append(row)
    expected_sequences = list(range(approved_start + 1, current_watermark + 1))
    missing_sequences = [
        sequence for sequence in expected_sequences if sequence not in grouped
    ]
    unexpected_count = 0
    classification_counts: dict[str, int] = {}
    first_unexpected: dict[str, Any] | None = None
    iteration_counts: dict[int, int] = {}
    sequence_iteration_pairs: list[dict[str, int]] = []
    invalid_shape_transaction_count = 0
    incomplete_transaction_count = 0
    unpublished_transaction_count = 0
    request_pattern = re.compile(rf"{re.escape(run_id)}-load-02-(\d{{6}})")
    for sequence in expected_sequences:
        members = grouped.get(sequence)
        if not members:
            continue
        iteration = sequence - approved_start
        expected_request_id = _rate_request_id(run_id, "load-02", iteration)
        request_values = sorted(
            {
                "" if member.request_id is None else str(member.request_id)
                for member in members
            }
        )
        member_count = len(members)
        complete = (
            all(member.transaction_finalized for member in members)
            and all(member.transaction_count == member_count for member in members)
            and [member.transaction_index for member in members]
            == list(range(1, member_count + 1))
        )
        published = all(member.status == "published" for member in members)
        exact_request = request_values == [expected_request_id]
        if not complete:
            incomplete_transaction_count += 1
        if not published:
            unpublished_transaction_count += 1
        parsed_iterations = []
        for request_value in request_values:
            match = request_pattern.fullmatch(request_value)
            if match:
                iteration = int(match.group(1))
                parsed_iterations.append(iteration)
                iteration_counts[iteration] = iteration_counts.get(iteration, 0) + 1
        if len(request_values) == 1 and len(parsed_iterations) == 1:
            sequence_iteration_pairs.append(
                {
                    "source_sequence": sequence,
                    "request_iteration": parsed_iterations[0],
                }
            )
        else:
            invalid_shape_transaction_count += 1
        if complete and published and exact_request:
            continue
        unexpected_count += 1
        attributions = [
            _bounded_request_attribution(
                actual_request_id=value,
                expected_request_id=expected_request_id,
                run_id=run_id,
            )
            for value in request_values
        ]
        for attribution in attributions:
            classification = attribution["classification"]
            classification_counts[classification] = (
                classification_counts.get(classification, 0) + 1
            )
        if first_unexpected is None:
            first_unexpected = {
                "source_sequence": sequence,
                "expected_iteration": iteration,
                "expected_safe_suffix": f"load-02-{iteration:06d}",
                "member_count": member_count,
                "transaction_complete": complete,
                "transaction_published": published,
                "request_attributions": attributions,
                "event_types": sorted({str(member.event_type) for member in members}),
            }
    approved_iteration_count = approved_final - approved_start
    expected_current_count = current_watermark - approved_start
    present_in_range = sorted(
        iteration
        for iteration in iteration_counts
        if 1 <= iteration <= approved_iteration_count
    )
    missing_current_prefix = sorted(
        set(range(1, expected_current_count + 1)) - set(present_in_range)
    )
    missing_approved = sorted(
        set(range(1, approved_iteration_count + 1)) - set(present_in_range)
    )
    out_of_range = sorted(
        iteration
        for iteration in iteration_counts
        if not 1 <= iteration <= approved_iteration_count
    )
    duplicates = [
        {"iteration": iteration, "transaction_count": count}
        for iteration, count in sorted(iteration_counts.items())
        if count != 1
    ]
    exact_current_prefix_bijection = (
        not missing_sequences
        and len(grouped) == expected_current_count
        and incomplete_transaction_count == 0
        and unpublished_transaction_count == 0
        and invalid_shape_transaction_count == 0
        and not out_of_range
        and not duplicates
        and present_in_range == list(range(1, expected_current_count + 1))
    )
    return {
        "approved_entry_watermark": approved_start,
        "current_watermark": current_watermark,
        "approved_final_watermark": approved_final,
        "expected_prefix_transaction_count": current_watermark - approved_start,
        "observed_sequence_count": len(grouped),
        "missing_sequence_count": len(missing_sequences),
        "first_missing_source_sequence": missing_sequences[0]
        if missing_sequences
        else None,
        "unexpected_transaction_count": unexpected_count,
        "unexpected_classification_counts": classification_counts,
        "first_unexpected": first_unexpected,
        "request_iteration_inventory": {
            "approved_iteration_count": approved_iteration_count,
            "unique_present_iteration_count": len(present_in_range),
            "present_iteration_ranges": _integer_ranges(present_in_range),
            "present_iterations_sha256": _sha256(present_in_range),
            "missing_current_prefix_iteration_count": len(missing_current_prefix),
            "missing_current_prefix_iteration_ranges": _integer_ranges(
                missing_current_prefix
            ),
            "missing_current_prefix_iterations_sha256": _sha256(missing_current_prefix),
            "missing_approved_iteration_count": len(missing_approved),
            "missing_approved_iteration_ranges": _integer_ranges(missing_approved),
            "missing_approved_iterations_sha256": _sha256(missing_approved),
            "duplicate_iteration_count": len(duplicates),
            "duplicate_iterations": duplicates,
            "out_of_range_iteration_count": len(out_of_range),
            "out_of_range_iterations": out_of_range,
            "invalid_request_shape_transaction_count": invalid_shape_transaction_count,
            "sequence_iteration_mapping_sha256": _sha256(sequence_iteration_pairs),
            "exact_current_prefix_bijection": exact_current_prefix_bijection,
        },
        "incomplete_transaction_count": incomplete_transaction_count,
        "unpublished_transaction_count": unpublished_transaction_count,
        "strict_prefix_attributable": not missing_sequences and unexpected_count == 0,
        "resume_or_correction_permitted": False,
        "diagnostic_read_only": True,
    }


def _read_load02_prefix_rows(
    *, source_scope: str, start: int, current_watermark: int
) -> list[Any]:
    from api.models import IntegrationOutbox

    return list(
        IntegrationOutbox.objects.filter(
            source_scope=source_scope,
            transport_sequence__gt=start,
            transport_sequence__lte=current_watermark,
        ).order_by("transport_sequence", "transaction_index")
    )


def diagnose_load02_prefix_from_environment() -> dict[str, Any]:
    """Report bounded attribution facts without writing evidence or changing state."""

    if os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip() != "diagnose-load-02-prefix":
        raise LoadHarnessError(
            "diagnose-load-02-prefix CLI action differs from guarded workflow action"
        )
    if _env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION") or _env_true(
        "TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"
    ):
        raise LoadHarnessError(
            "diagnose-load-02-prefix requires false mutation and operational confirmations"
        )
    run_id = os.environ.get("TT_PHASE1_LOAD_RUN_ID", "").strip()
    if run_id != LOAD02_CHECKPOINT_RUN_ID:
        raise LoadHarnessError(
            "diagnose-load-02-prefix run id differs from the fixed recovery"
        )
    validate_run_id(run_id)
    github_run_id = os.environ.get("TT_PHASE1_LOAD_GITHUB_RUN_ID", "").strip()
    if not github_run_id.isdigit():
        raise LoadHarnessError("LOAD-02 diagnostic GitHub run id is invalid")
    manifest = load_manifest(
        os.environ.get("TT_PHASE1_LOAD_MANIFEST_JSON", ""),
        run_id=run_id,
        environment=os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
        source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
        require_approval=True,
    )
    exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
    if exact_sha != manifest["commits"]["timetabler"]:
        raise LoadHarnessError(
            "deployed Timetabler SHA differs from the diagnostic manifest"
        )
    metrics = _assert_source_safety(expected_watermark=None)
    current_watermark = metrics["transport_watermark"]
    start = manifest["expected_stage_start_watermarks"]["load-02"]
    final = manifest["expected_stage_final_watermarks"]["load-02"]
    if not start < current_watermark < final:
        raise LoadHarnessError(
            "LOAD-02 diagnostic watermark is outside its strict open recovery fence"
        )
    rows = _read_load02_prefix_rows(
        source_scope=manifest["source_scope"],
        start=start,
        current_watermark=current_watermark,
    )
    diagnosis = _diagnose_load02_prefix_rows(
        rows=rows,
        manifest=manifest,
        run_id=run_id,
        current_watermark=current_watermark,
    )
    final_metrics = _assert_source_safety(expected_watermark=current_watermark)
    final_rows = _read_load02_prefix_rows(
        source_scope=manifest["source_scope"],
        start=start,
        current_watermark=current_watermark,
    )
    final_diagnosis = _diagnose_load02_prefix_rows(
        rows=final_rows,
        manifest=manifest,
        run_id=run_id,
        current_watermark=current_watermark,
    )
    if final_diagnosis != diagnosis:
        raise LoadHarnessError("LOAD-02 request inventory changed during diagnosis")
    result = {
        "phase1_load": "load02_prefix_diagnosed",
        "run_id": run_id,
        "github_run_id": github_run_id,
        "git_sha": exact_sha,
        "resource_booking_git_sha_attestation": manifest["commits"]["resource_booking"],
        "diagnosis": diagnosis,
        "two_read_inventory_stable": True,
        "final_source_metrics": final_metrics,
        "evidence_written": False,
        "api_authentication_performed": False,
        "mutation_executed": False,
        "load_executed": False,
        "restart_executed": False,
        "database_write_executed": False,
        "manual_outbox_edit": False,
        "resource_booking_database_touched": False,
        "resource_booking_process_touched": False,
        "phase2_enabled": False,
        "reverse_delivery_enabled": False,
        "dr_executed": False,
    }
    _assert_exchange_payload_has_no_credentials(result)
    return result


def _failed_load02_expected_operations(
    *, manifest: dict[str, Any], run_id: str, fixtures: dict[str, list[int]]
) -> list[dict[str, Any]]:
    """Build the immutable logical-operation inventory without calling an API."""

    steps = _workload_steps(manifest, "load-02")
    total = int(manifest["load"]["p_tps"]) * int(manifest["load"]["steady_seconds"])
    operations: list[dict[str, Any]] = []
    for iteration in range(1, total + 1):
        step = steps[(iteration - 1) % len(steps)]
        cycle = (iteration - 1) // len(steps)
        selector_name = step["fixture_selector"]
        pool = fixtures[selector_name]
        rate_shards = max(len(fixtures[item["fixture_selector"]]) for item in steps)
        target_id = int(pool[(cycle % rate_shards) % len(pool)])
        operations.append(
            {
                "iteration": iteration,
                "request_id": _rate_request_id(run_id, "load-02", iteration),
                "workload_classification": step["workload_classification"],
                "path": step["path"],
                "target_id": target_id,
                "slot": step["payload"].get("slot"),
            }
        )
    return operations


def _safe_int_list(value: Any) -> list[int] | None:
    if not isinstance(value, list):
        return None
    try:
        return [int(item) for item in value]
    except (TypeError, ValueError):
        return None


def _classify_failed_load02_api_log(
    *,
    url: Any,
    params: Any,
    fixture_sets: dict[str, set[int]],
    slot_a: int,
    slot_b: int,
) -> str | None:
    """Classify only an exact run-owned fixture request; never return its payload."""

    if not isinstance(url, str) or not isinstance(params, dict):
        return None
    path = next(
        (
            candidate
            for candidate in (
                "/api/admin/staff/update",
                "/api/admin/location/update",
                "/api/admin/schedule-request",
                "/api/admin/unschedule",
            )
            if url.endswith(candidate)
        ),
        None,
    )
    if path is None:
        return None
    key = (
        "activity_ids"
        if path
        in {
            "/api/admin/schedule-request",
            "/api/admin/unschedule",
        }
        else "id"
    )
    identifiers = _safe_int_list(params.get(key))
    if identifiers is None or len(identifiers) != 1:
        return None
    target_id = identifiers[0]
    if path == "/api/admin/staff/update":
        return (
            "staff_lifecycle"
            if target_id in fixture_sets["lifecycle_staff_ids"]
            else None
        )
    if path == "/api/admin/location/update":
        return (
            "location_lifecycle"
            if target_id in fixture_sets["lifecycle_location_ids"]
            else None
        )
    if path == "/api/admin/unschedule":
        if target_id in fixture_sets["traffic_activity_ids"]:
            return "terminal"
        if target_id in fixture_sets["traffic_recurrence_activity_ids"]:
            return "recurrence_terminal"
        return None
    try:
        slot = int(params.get("slot"))
    except (TypeError, ValueError):
        return None
    if target_id in fixture_sets["traffic_activity_ids"]:
        if slot == slot_a:
            return "schedule"
        if slot == slot_b:
            return "reschedule"
    if target_id in fixture_sets["traffic_recurrence_activity_ids"] and slot == slot_a:
        return "recurrence_schedule"
    return None


def _parse_engine_response_for_audit(
    raw: Any, *, delivery_id: str, expected_activity_id: int
) -> dict[str, Any]:
    """Parse the bounded final-response facts needed for terminal classification."""

    if raw in (None, ""):
        return {
            "present": False,
            "protocol_valid": False,
            "status": None,
            "assignment_count": 0,
            "successful_assignment": False,
            "response_sha256": None,
            "canonical_response_hash": None,
        }
    response_sha256 = hashlib.sha256(str(raw).encode("utf-8")).hexdigest()
    try:
        response = json.loads(raw) if isinstance(raw, str) else raw
        if not isinstance(response, dict):
            raise ValueError("response is not an object")
        response_id = str(response.get("request_id") or "")
        schedule = response.get("schedule", [])
        schedule = json.loads(schedule) if isinstance(schedule, str) else schedule
        if not isinstance(schedule, list):
            raise ValueError("schedule is not an array")
        activity_ids: list[int] = []
        successful = False
        for item in schedule:
            if not isinstance(item, dict):
                raise ValueError("assignment is not an object")
            activity_id = int(item["activity"])
            activity_ids.append(activity_id)
            if activity_id == expected_activity_id and item.get("start_slot") not in (
                None,
                "",
            ):
                int(item["start_slot"])
                successful = True
            for key in ("teaching_staff", "location"):
                resource_ids = item.get(key) or []
                resource_ids = (
                    json.loads(resource_ids)
                    if isinstance(resource_ids, str)
                    else resource_ids
                )
                if not isinstance(resource_ids, list):
                    raise ValueError("assignment resources are not an array")
                [int(resource_id) for resource_id in resource_ids]
        protocol_valid = (
            response_id == delivery_id
            and activity_ids == [expected_activity_id]
            and isinstance(response.get("status"), str)
        )
        canonical_response_hash = hashlib.sha256(
            json.dumps(
                response, sort_keys=True, separators=(",", ":"), default=str
            ).encode("utf-8")
        ).hexdigest()
        return {
            "present": True,
            "protocol_valid": protocol_valid,
            "status": str(response.get("status") or "").lower(),
            "assignment_count": len(schedule),
            "successful_assignment": successful,
            "response_sha256": response_sha256,
            "canonical_response_hash": canonical_response_hash,
        }
    except (KeyError, TypeError, ValueError, json.JSONDecodeError):
        return {
            "present": True,
            "protocol_valid": False,
            "status": None,
            "assignment_count": 0,
            "successful_assignment": False,
            "response_sha256": response_sha256,
            "canonical_response_hash": None,
        }


def _summarize_audit_identifiers(values: list[Any]) -> dict[str, Any]:
    normalized = sorted(str(value) for value in values)
    return {
        "count": len(normalized),
        "sha256": _sha256(normalized),
    }


def _summarize_numeric_identifiers(values: list[Any]) -> dict[str, Any]:
    normalized = sorted(int(value) for value in values)
    return {
        "count": len(normalized),
        "minimum": normalized[0] if normalized else None,
        "maximum": normalized[-1] if normalized else None,
        "sha256": _sha256(normalized),
    }


def _audit_outbox_transactions(rows: list[Any]) -> list[dict[str, Any]]:
    grouped: dict[tuple[int, str], list[Any]] = {}
    for row in rows:
        grouped.setdefault(
            (int(row.transport_sequence), str(row.change_set_id)), []
        ).append(row)
    transactions: list[dict[str, Any]] = []
    for (sequence, change_set_id), members in sorted(grouped.items()):
        members = sorted(members, key=lambda member: member.transaction_index)
        request_ids = sorted({str(member.request_id or "") for member in members})
        complete = (
            bool(members)
            and all(member.transaction_finalized for member in members)
            and all(member.transaction_count == len(members) for member in members)
            and [member.transaction_index for member in members]
            == list(range(1, len(members) + 1))
        )
        transactions.append(
            {
                "source_sequence": sequence,
                "change_set_id": change_set_id,
                "request_ids": request_ids,
                "member_count": len(members),
                "complete": complete,
                "published": all(member.status == "published" for member in members),
                "event_types": sorted({str(member.event_type) for member in members}),
                "aggregate_identity_sha256": _sha256(
                    sorted(
                        f"{member.aggregate_type}:{member.aggregate_id}"
                        for member in members
                    )
                ),
            }
        )
    return transactions


def _audit_run_owned_activity_states(
    *,
    manifest: dict[str, Any],
    fixtures: dict[str, list[int]],
    allow_absent: bool = False,
) -> dict[str, Any]:
    from api.models import TtActivity

    activity_selector_names = (
        "traffic_activity_ids",
        "traffic_recurrence_activity_ids",
        "maximum_recurrence_activity_ids",
        "bulk_activity_ids",
        "boundary_activity_ids",
    )
    all_ids = sorted(
        {
            activity_id
            for name in activity_selector_names
            for activity_id in fixtures[name]
        }
    )
    activities = list(
        TtActivity.objects.filter(id__in=all_ids)
        .prefetch_related("staff", "location", "week")
        .order_by("id")
    )
    present_ids = [int(activity.id) for activity in activities]
    if present_ids != all_ids and not allow_absent:
        raise LoadHarnessError("run-owned activity audit inventory is incomplete")
    states: dict[int, dict[str, Any]] = {}
    for activity in activities:
        states[int(activity.id)] = {
            "id": int(activity.id),
            "status": int(activity.status),
            "scheduled": int(activity.scheduled or 0),
            "scheduled_start_slot": activity.scheduled_start_slot,
            "staff_ids": sorted(
                int(value) for value in activity.staff.values_list("id", flat=True)
            ),
            "location_ids": sorted(
                int(value) for value in activity.location.values_list("id", flat=True)
            ),
            "week_ids": sorted(
                int(value) for value in activity.week.values_list("id", flat=True)
            ),
        }
    result: dict[str, Any] = {}
    for name in activity_selector_names:
        selector_states = [
            states[activity_id]
            for activity_id in fixtures[name]
            if activity_id in states
        ]
        absent_ids = [
            activity_id for activity_id in fixtures[name] if activity_id not in states
        ]
        slot_counts: dict[str, int] = {}
        for state in selector_states:
            slot = str(state["scheduled_start_slot"])
            slot_counts[slot] = slot_counts.get(slot, 0) + 1
        result[name] = {
            "activity_count": len(fixtures[name]),
            "present_activity_count": len(selector_states),
            "absent_activity_count": len(absent_ids),
            "absent_activity_ids_sha256": _sha256(absent_ids),
            "scheduled_count": sum(
                state["scheduled"] == 1 for state in selector_states
            ),
            "staff_relation_count": sum(
                len(state["staff_ids"]) for state in selector_states
            ),
            "location_relation_count": sum(
                len(state["location_ids"]) for state in selector_states
            ),
            "slot_counts": dict(sorted(slot_counts.items())),
            "canonical_state_sha256": _sha256(selector_states),
        }
    return result


def _read_failed_load02_engine_audit_snapshot(
    *,
    manifest: dict[str, Any],
    run_id: str,
    fixtures: dict[str, list[int]] | None = None,
    allow_absent_activities: bool = False,
) -> dict[str, Any]:
    """Perform bounded ORM reads for the one fixed, run-owned failed workload."""

    from django.db.models import Q
    from django.conf import settings

    from api.models import (
        AppliedEngineResponse,
        EngineResponseQuarantine,
        IncomingApiAdmin,
        IntegrationOutbox,
        KafkaLog,
        PostCommitDelivery,
    )

    fixtures = fixtures or _resolve_selectors(manifest)
    operations = _failed_load02_expected_operations(
        manifest=manifest, run_id=run_id, fixtures=fixtures
    )
    all_request_ids = [operation["request_id"] for operation in operations]
    schedule_operations = [
        operation
        for operation in operations
        if operation["workload_classification"]
        in {"schedule", "reschedule", "recurrence_schedule"}
    ]
    schedule_request_ids = [
        operation["request_id"] for operation in schedule_operations
    ]
    deliveries = list(
        PostCommitDelivery.objects.filter(correlation_id__in=all_request_ids).order_by(
            "correlation_id", "created_at", "id"
        )[: len(all_request_ids) * 2 + 1]
    )
    if len(deliveries) > len(all_request_ids) * 2:
        raise LoadHarnessError("run-correlated delivery inventory exceeds its bound")
    delivery_ids = [str(delivery.delivery_id) for delivery in deliveries]
    kafka_logs = list(
        KafkaLog.objects.filter(request_id__in=delivery_ids).order_by("request_id")
    )
    receipts = list(
        AppliedEngineResponse.objects.filter(request_id__in=delivery_ids).order_by(
            "request_id"
        )
    )
    quarantines = list(
        EngineResponseQuarantine.objects.filter(request_id__in=delivery_ids).order_by(
            "request_id", "fingerprint"
        )
    )
    receipt_change_sets = [str(receipt.change_set_id) for receipt in receipts]
    outbox_rows = list(
        IntegrationOutbox.objects.filter(
            Q(request_id__in=all_request_ids)
            | Q(request_id__in=delivery_ids)
            | Q(change_set_id__in=receipt_change_sets)
        ).order_by("transport_sequence", "transaction_index")
    )

    fixture_sets = {name: set(values) for name, values in fixtures.items()}
    window_start = _parse_iso8601(
        manifest["approved_window"]["starts_at"], "approved_window.starts_at"
    )
    window_end = _parse_iso8601(
        manifest["approved_window"]["ends_at"], "approved_window.ends_at"
    )
    if not settings.USE_TZ:
        source_timezone = ZoneInfo(settings.TIME_ZONE)
        window_start = window_start.astimezone(source_timezone).replace(tzinfo=None)
        window_end = window_end.astimezone(source_timezone).replace(tzinfo=None)
    audit_time_anchors = [
        row.created_at for row in [*deliveries, *outbox_rows] if row.created_at
    ]
    if not audit_time_anchors:
        raise LoadHarnessError("failed LOAD-02 audit has no run-correlated time anchor")
    window_start = max(
        window_start,
        min(audit_time_anchors) - datetime.timedelta(minutes=5),
    )
    window_end = min(
        window_end,
        max(audit_time_anchors) + datetime.timedelta(minutes=5),
    )
    if window_start >= window_end:
        raise LoadHarnessError("failed LOAD-02 audit time window is contradictory")
    api_rows = list(
        IncomingApiAdmin.objects.filter(
            Q(url__endswith="/api/admin/staff/update")
            | Q(url__endswith="/api/admin/location/update")
            | Q(url__endswith="/api/admin/schedule-request")
            | Q(url__endswith="/api/admin/unschedule"),
            created_at__gte=window_start,
            created_at__lte=window_end,
        ).order_by("id")[: len(all_request_ids) * 2 + 1]
    )
    if len(api_rows) > len(all_request_ids) * 2:
        raise LoadHarnessError("candidate API audit log inventory exceeds its bound")
    classified_api_logs = []
    for row in api_rows:
        classification = _classify_failed_load02_api_log(
            url=row.url,
            params=row.incoming_params,
            fixture_sets=fixture_sets,
            slot_a=int(manifest["variables"]["slot_a"]),
            slot_b=int(manifest["variables"]["slot_b"]),
        )
        if classification:
            classified_api_logs.append(
                {
                    "id": int(row.id),
                    "workload_classification": classification,
                    "scode": row.scode,
                    "created_at": _timezone_qualified_iso(row.created_at),
                }
            )

    delivery_records = []
    for delivery in deliveries:
        activities = (
            _safe_int_list(delivery.payload.get("activities"))
            if isinstance(delivery.payload, dict)
            else None
        )
        delivery_records.append(
            {
                "delivery_id": str(delivery.delivery_id),
                "correlation_id": str(delivery.correlation_id or ""),
                "channel": str(delivery.channel),
                "method": str(delivery.method),
                "status": str(delivery.status),
                "payload_activity_ids": activities,
                "created_at": _timezone_qualified_iso(delivery.created_at),
                "published_at": (
                    _timezone_qualified_iso(delivery.published_at)
                    if delivery.published_at
                    else None
                ),
            }
        )
    kafka_records = []
    expected_by_request = {
        operation["request_id"]: operation for operation in schedule_operations
    }
    correlation_by_delivery = {
        record["delivery_id"]: record["correlation_id"] for record in delivery_records
    }
    for row in kafka_logs:
        delivery_id = str(row.request_id)
        expected = expected_by_request.get(correlation_by_delivery.get(delivery_id, ""))
        response = _parse_engine_response_for_audit(
            row.response_data,
            delivery_id=delivery_id,
            expected_activity_id=int(expected["target_id"]) if expected else -1,
        )
        request_payload_valid = False
        try:
            request_payload = (
                json.loads(row.request_data)
                if isinstance(row.request_data, str)
                else row.request_data
            )
            request_payload_valid = (
                isinstance(request_payload, dict)
                and expected is not None
                and str(request_payload.get("request_id")) == delivery_id
                and _safe_int_list(request_payload.get("activities"))
                == [int(expected["target_id"])]
            )
        except (TypeError, ValueError, json.JSONDecodeError):
            request_payload_valid = False
        kafka_records.append(
            {
                "request_id": delivery_id,
                "request_present": bool(row.request_at and row.request_data),
                "request_payload_valid": request_payload_valid,
                "response_at": (
                    _timezone_qualified_iso(row.response_at)
                    if row.response_at
                    else None
                ),
                "response": response,
            }
        )
    receipt_records = [
        {
            "request_id": str(receipt.request_id),
            "response_hash": str(receipt.response_hash),
            "change_set_id": str(receipt.change_set_id),
            "outcome": str(receipt.outcome),
            "applied_at": _timezone_qualified_iso(receipt.applied_at),
        }
        for receipt in receipts
    ]
    quarantine_records = [
        {
            "request_id": str(row.request_id or ""),
            "fingerprint": str(row.fingerprint),
            "response_hash": str(row.response_hash),
            "reason": str(row.reason),
            "occurrences": int(row.occurrences),
        }
        for row in quarantines
    ]
    return {
        "operations": operations,
        "schedule_request_ids": schedule_request_ids,
        "api_logs": classified_api_logs,
        "deliveries": delivery_records,
        "kafka_logs": kafka_records,
        "receipts": receipt_records,
        "quarantines": quarantine_records,
        "outbox_transactions": _audit_outbox_transactions(outbox_rows),
        "activity_states": _audit_run_owned_activity_states(
            manifest=manifest,
            fixtures=fixtures,
            allow_absent=allow_absent_activities,
        ),
    }


def _classify_failed_load02_engine_audit(
    *,
    snapshot: dict[str, Any],
    manifest: dict[str, Any],
    allow_partial_prefix: bool = False,
) -> dict[str, Any]:
    class_names = list(LEGACY_TRAFFIC_CLASSIFICATIONS)
    operations = snapshot["operations"]
    submitted_operation_count = len(snapshot["api_logs"])
    partial_prefix_attributable = False
    partial_submission_set_attributable = False
    if allow_partial_prefix:
        if not 0 < submitted_operation_count <= len(operations):
            raise LoadHarnessError(
                "partial LOAD-02 recovery has no bounded submitted operation set"
            )
        operations_by_request = {
            operation["request_id"]: operation for operation in operations
        }
        attributed_request_ids = {
            delivery["correlation_id"]
            for delivery in snapshot["deliveries"]
            if delivery["correlation_id"] in operations_by_request
        }
        attributed_request_ids.update(
            request_id
            for transaction in snapshot["outbox_transactions"]
            for request_id in transaction["request_ids"]
            if request_id in operations_by_request
        )
        submitted_operations = sorted(
            (
                operations_by_request[request_id]
                for request_id in attributed_request_ids
            ),
            key=lambda operation: operation["iteration"],
        )
        observed_class_counts = Counter(
            row["workload_classification"] for row in snapshot["api_logs"]
        )
        accepted_class_counts = Counter(
            row["workload_classification"]
            for row in snapshot["api_logs"]
            if isinstance(row["scode"], int) and 200 <= row["scode"] < 300
        )
        attributed_class_counts = Counter(
            operation["workload_classification"]
            for operation in submitted_operations
        )
        if accepted_class_counts != attributed_class_counts:
            class_count_deltas = {
                name: accepted_class_counts[name] - attributed_class_counts[name]
                for name in class_names
                if accepted_class_counts[name] != attributed_class_counts[name]
            }
            raise LoadHarnessError(
                "partial LOAD-02 accepted API inventory does not match the exact "
                "durably attributed operation set "
                f"(submitted={submitted_operation_count}, "
                "class_count_deltas="
                f"{json.dumps(class_count_deltas, sort_keys=True, separators=(',', ':'))})"
            )
        if any(
            not isinstance(row["scode"], int)
            or not (200 <= row["scode"] < 300 or row["scode"] >= 400)
            for row in snapshot["api_logs"]
        ):
            raise LoadHarnessError(
                "partial LOAD-02 API inventory contains an ambiguous HTTP outcome"
            )
        expected_counts = {
            name: observed_class_counts[name] for name in class_names
        }
        operations = submitted_operations
        partial_prefix_attributable = (
            [operation["iteration"] for operation in operations]
            == list(range(1, len(operations) + 1))
        )
        partial_submission_set_attributable = True
    else:
        expected_counts = manifest["workloads"]["load-02"][
            "expected_class_counts"
        ]
    class_results = {
        name: {
            "expected": int(expected_counts[name]),
            "accepted": 0,
            "applied": 0,
            "failed": 0,
            "no_op": 0,
            "ambiguous": 0,
            "api_log_identifiers": {},
        }
        for name in class_names
    }
    api_logs_by_class = {name: [] for name in class_names}
    for row in snapshot["api_logs"]:
        api_logs_by_class[row["workload_classification"]].append(row)
    for name, rows in api_logs_by_class.items():
        accepted = [
            row
            for row in rows
            if isinstance(row["scode"], int) and 200 <= row["scode"] < 300
        ]
        failed = [
            row for row in rows if isinstance(row["scode"], int) and row["scode"] >= 400
        ]
        unknown = len(rows) - len(accepted) - len(failed)
        class_results[name]["accepted"] = len(accepted)
        class_results[name]["failed"] = len(failed)
        class_results[name]["ambiguous"] += unknown + abs(
            class_results[name]["expected"] - len(rows)
        )
        class_results[name]["api_log_identifiers"] = _summarize_numeric_identifiers(
            [row["id"] for row in rows]
        )

    deliveries_by_correlation: dict[str, list[dict[str, Any]]] = {}
    for delivery in snapshot["deliveries"]:
        deliveries_by_correlation.setdefault(delivery["correlation_id"], []).append(
            delivery
        )
    kafka_by_request = {row["request_id"]: row for row in snapshot["kafka_logs"]}
    receipts_by_request: dict[str, list[dict[str, Any]]] = {}
    for receipt in snapshot["receipts"]:
        receipts_by_request.setdefault(receipt["request_id"], []).append(receipt)
    quarantines_by_request: dict[str, list[dict[str, Any]]] = {}
    for quarantine in snapshot["quarantines"]:
        quarantines_by_request.setdefault(quarantine["request_id"], []).append(
            quarantine
        )
    outbox_by_request: dict[str, list[dict[str, Any]]] = {}
    outbox_by_change_set: dict[str, list[dict[str, Any]]] = {}
    for transaction in snapshot["outbox_transactions"]:
        for request_id in transaction["request_ids"]:
            outbox_by_request.setdefault(request_id, []).append(transaction)
        outbox_by_change_set.setdefault(transaction["change_set_id"], []).append(
            transaction
        )

    engine_terminal_status_counts: dict[str, int] = {}
    delivery_uuid_values: list[str] = []
    correlation_values: list[str] = []
    receipt_values: list[str] = []
    change_set_values: list[str] = []
    sequence_values: list[int] = []
    nonterminal_iterations: list[int] = []
    ambiguous_iterations: list[int] = []
    quarantine_reason_counts: dict[str, int] = {}
    engine_expected_count = 0
    engine_terminal_count = 0
    quarantined_operation_count = 0
    engine_candidate_delivery_ids: set[str] = set()

    for operation in operations:
        name = operation["workload_classification"]
        request_id = operation["request_id"]
        direct_transactions = outbox_by_request.get(request_id, [])
        if name not in {"schedule", "reschedule", "recurrence_schedule"}:
            valid = [
                transaction
                for transaction in direct_transactions
                if transaction["complete"] and transaction["published"]
            ]
            if len(valid) == 1 and len(direct_transactions) == 1:
                class_results[name]["applied"] += 1
                sequence_values.append(valid[0]["source_sequence"])
                change_set_values.append(valid[0]["change_set_id"])
            elif direct_transactions:
                class_results[name]["ambiguous"] += 1
                ambiguous_iterations.append(operation["iteration"])
            continue

        engine_expected_count += 1
        deliveries = [
            delivery
            for delivery in deliveries_by_correlation.get(request_id, [])
            if _is_engine_request_delivery(
                delivery, expected_activity_id=int(operation["target_id"])
            )
        ]
        engine_candidate_delivery_ids.update(
            str(delivery["delivery_id"]) for delivery in deliveries
        )
        if len(deliveries) != 1:
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            nonterminal_iterations.append(operation["iteration"])
            engine_terminal_status_counts["missing_or_duplicate_delivery"] = (
                engine_terminal_status_counts.get("missing_or_duplicate_delivery", 0)
                + 1
            )
            continue
        delivery = deliveries[0]
        delivery_id = delivery["delivery_id"]
        delivery_uuid_values.append(delivery_id)
        correlation_values.append(request_id)
        if (
            delivery["channel"] != "kafka"
            or delivery["method"] != "schedule"
            or delivery["payload_activity_ids"] != [operation["target_id"]]
        ):
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            nonterminal_iterations.append(operation["iteration"])
            engine_terminal_status_counts["invalid_delivery_contract"] = (
                engine_terminal_status_counts.get("invalid_delivery_contract", 0) + 1
            )
            continue
        if delivery["status"] == "dead_letter":
            class_results[name]["failed"] += 1
            engine_terminal_count += 1
            engine_terminal_status_counts["delivery_dead_letter"] = (
                engine_terminal_status_counts.get("delivery_dead_letter", 0) + 1
            )
            continue
        if delivery["status"] != "published":
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            nonterminal_iterations.append(operation["iteration"])
            key = f"delivery_{delivery['status']}"
            engine_terminal_status_counts[key] = (
                engine_terminal_status_counts.get(key, 0) + 1
            )
            continue
        kafka = kafka_by_request.get(delivery_id)
        if (
            not kafka
            or not kafka["request_present"]
            or not kafka["request_payload_valid"]
        ):
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            nonterminal_iterations.append(operation["iteration"])
            engine_terminal_status_counts["missing_or_invalid_kafka_request"] = (
                engine_terminal_status_counts.get("missing_or_invalid_kafka_request", 0)
                + 1
            )
            continue
        if not kafka["response_at"]:
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            nonterminal_iterations.append(operation["iteration"])
            engine_terminal_status_counts["awaiting_engine_response"] = (
                engine_terminal_status_counts.get("awaiting_engine_response", 0) + 1
            )
            continue
        receipt_rows = receipts_by_request.get(delivery_id, [])
        quarantine_rows = quarantines_by_request.get(delivery_id, [])
        if receipt_rows and quarantine_rows:
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            engine_terminal_status_counts["receipt_quarantine_conflict"] = (
                engine_terminal_status_counts.get("receipt_quarantine_conflict", 0) + 1
            )
            continue
        if quarantine_rows:
            class_results[name]["failed"] += 1
            engine_terminal_count += 1
            quarantined_operation_count += 1
            for row in quarantine_rows:
                quarantine_reason_counts[row["reason"]] = (
                    quarantine_reason_counts.get(row["reason"], 0) + 1
                )
            engine_terminal_status_counts["quarantined"] = (
                engine_terminal_status_counts.get("quarantined", 0) + 1
            )
            continue
        if len(receipt_rows) != 1:
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            nonterminal_iterations.append(operation["iteration"])
            engine_terminal_status_counts["response_without_terminal_receipt"] = (
                engine_terminal_status_counts.get(
                    "response_without_terminal_receipt", 0
                )
                + 1
            )
            continue
        receipt = receipt_rows[0]
        receipt_values.append(delivery_id)
        response = kafka["response"]
        if (
            not response["protocol_valid"]
            or receipt["response_hash"] != response["canonical_response_hash"]
        ):
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            engine_terminal_status_counts["receipt_response_mismatch"] = (
                engine_terminal_status_counts.get("receipt_response_mismatch", 0) + 1
            )
            continue
        receipt_transactions = outbox_by_change_set.get(receipt["change_set_id"], [])
        if (
            any(
                not transaction["complete"] or not transaction["published"]
                for transaction in receipt_transactions
            )
            or len(receipt_transactions) > 1
        ):
            class_results[name]["ambiguous"] += 1
            ambiguous_iterations.append(operation["iteration"])
            engine_terminal_status_counts["invalid_outbox_change_set"] = (
                engine_terminal_status_counts.get("invalid_outbox_change_set", 0) + 1
            )
            continue
        engine_terminal_count += 1
        if receipt_transactions:
            if response["status"] != "success" or not response["successful_assignment"]:
                class_results[name]["ambiguous"] += 1
                ambiguous_iterations.append(operation["iteration"])
                engine_terminal_status_counts["outbox_response_contradiction"] = (
                    engine_terminal_status_counts.get(
                        "outbox_response_contradiction", 0
                    )
                    + 1
                )
                continue
            transaction = receipt_transactions[0]
            class_results[name]["applied"] += 1
            change_set_values.append(transaction["change_set_id"])
            sequence_values.append(transaction["source_sequence"])
            engine_terminal_status_counts["applied"] = (
                engine_terminal_status_counts.get("applied", 0) + 1
            )
        elif response["status"] != "success" or not response["successful_assignment"]:
            class_results[name]["failed"] += 1
            engine_terminal_status_counts["engine_failed_or_no_slot"] = (
                engine_terminal_status_counts.get("engine_failed_or_no_slot", 0) + 1
            )
        else:
            class_results[name]["no_op"] += 1
            engine_terminal_status_counts["applied_receipt_no_state_change"] = (
                engine_terminal_status_counts.get("applied_receipt_no_state_change", 0)
                + 1
            )

    for name in (
        "staff_lifecycle",
        "location_lifecycle",
        "terminal",
        "recurrence_terminal",
        "unschedule",
    ):
        if name not in class_results:
            continue
        result = class_results[name]
        if result["ambiguous"] == 0 and result["accepted"] >= result["applied"]:
            result["no_op"] = result["accepted"] - result["applied"]
        else:
            result["ambiguous"] += max(0, result["applied"] - result["accepted"])

    auxiliary_deliveries = [
        delivery
        for delivery in snapshot["deliveries"]
        if str(delivery["delivery_id"]) not in engine_candidate_delivery_ids
    ]
    auxiliary_nonterminal_count = sum(
        delivery["status"] not in {"published", "dead_letter"}
        for delivery in auxiliary_deliveries
    )
    all_engine_terminal = engine_terminal_count == engine_expected_count
    operator_replay_could_mutate = quarantined_operation_count > 0
    late_mutation_possible = (
        not all_engine_terminal
        or bool(nonterminal_iterations)
        or operator_replay_could_mutate
        or auxiliary_nonterminal_count > 0
    )
    total_ambiguous = sum(result["ambiguous"] for result in class_results.values())
    return {
        "workload_class_counts": class_results,
        "http_api_log_identifiers": _summarize_numeric_identifiers(
            [row["id"] for row in snapshot["api_logs"]]
        ),
        "post_commit_delivery": {
            "all_run_correlation_ids": _summarize_audit_identifiers(
                [delivery["correlation_id"] for delivery in snapshot["deliveries"]]
            ),
            "verified_engine_correlation_ids": _summarize_audit_identifiers(
                correlation_values
            ),
            "generated_engine_delivery_uuids": _summarize_audit_identifiers(
                delivery_uuid_values
            ),
            "status_counts": dict(
                sorted(
                    {
                        status: sum(
                            delivery["status"] == status
                            for delivery in snapshot["deliveries"]
                        )
                        for status in sorted(
                            {delivery["status"] for delivery in snapshot["deliveries"]}
                        )
                    }.items()
                )
            ),
            "auxiliary_delivery_count": len(auxiliary_deliveries),
            "auxiliary_nonterminal_count": auxiliary_nonterminal_count,
        },
        "engine_delivery": {
            "expected_operation_count": engine_expected_count,
            "terminal_operation_count": engine_terminal_count,
            "terminal_status_counts": dict(
                sorted(engine_terminal_status_counts.items())
            ),
            "nonterminal_iteration_ranges": _integer_ranges(
                sorted(set(nonterminal_iterations))
            ),
            "nonterminal_iterations_sha256": _sha256(
                sorted(set(nonterminal_iterations))
            ),
            "ambiguous_iteration_ranges": _integer_ranges(
                sorted(set(ambiguous_iterations))
            ),
            "ambiguous_iterations_sha256": _sha256(sorted(set(ambiguous_iterations))),
        },
        "kafka_logs": {
            "request_identifiers": _summarize_audit_identifiers(
                [row["request_id"] for row in snapshot["kafka_logs"]]
            ),
            "response_count": sum(
                bool(row["response_at"]) for row in snapshot["kafka_logs"]
            ),
            "valid_request_count": sum(
                bool(row["request_payload_valid"]) for row in snapshot["kafka_logs"]
            ),
        },
        "applied_engine_response_receipts": {
            "all_request_identifiers": _summarize_audit_identifiers(
                [receipt["request_id"] for receipt in snapshot["receipts"]]
            ),
            "verified_request_identifiers": _summarize_audit_identifiers(
                receipt_values
            ),
            "change_set_identifiers": _summarize_audit_identifiers(
                [receipt["change_set_id"] for receipt in snapshot["receipts"]]
            ),
        },
        "engine_quarantine": {
            "operation_count": quarantined_operation_count,
            "record_count": len(snapshot["quarantines"]),
            "reason_counts": dict(sorted(quarantine_reason_counts.items())),
            "fingerprint_identifiers": _summarize_audit_identifiers(
                [row["fingerprint"] for row in snapshot["quarantines"]]
            ),
        },
        "source_change_sets": {
            "transaction_count": len(snapshot["outbox_transactions"]),
            "all_source_sequence_ranges": _integer_ranges(
                sorted(
                    {
                        transaction["source_sequence"]
                        for transaction in snapshot["outbox_transactions"]
                    }
                )
            ),
            "all_source_sequences_sha256": _sha256(
                sorted(
                    {
                        transaction["source_sequence"]
                        for transaction in snapshot["outbox_transactions"]
                    }
                )
            ),
            "verified_attributable_source_sequence_ranges": _integer_ranges(
                sorted(set(sequence_values))
            ),
            "verified_attributable_source_sequences_sha256": _sha256(
                sorted(set(sequence_values))
            ),
            "all_change_set_identifiers": _summarize_audit_identifiers(
                [
                    transaction["change_set_id"]
                    for transaction in snapshot["outbox_transactions"]
                ]
            ),
            "verified_attributable_change_set_identifiers": _summarize_audit_identifiers(
                change_set_values
            ),
            "all_complete_and_published": all(
                transaction["complete"] and transaction["published"]
                for transaction in snapshot["outbox_transactions"]
            ),
        },
        "run_owned_activity_final_states": snapshot["activity_states"],
        "submitted_operation_count": submitted_operation_count,
        "partial_prefix_attributable": partial_prefix_attributable,
        "partial_submission_set_attributable": partial_submission_set_attributable,
        "all_engine_work_terminal": all_engine_terminal,
        "automatic_late_mutation_possible": not all_engine_terminal,
        "operator_replay_could_mutate": operator_replay_could_mutate,
        "late_mutation_possible": late_mutation_possible,
        "ambiguous_operation_count": total_ambiguous,
        "schema_limitations": {
            "incoming_api_admin_retains_request_id": False,
            "http_results_are_exact_per_class_not_per_iteration": True,
        },
        "load02_acceptance_passed": False,
        "cleanup_permitted": False,
        "containment_passed": (
            all_engine_terminal
            and total_ambiguous == 0
            and auxiliary_nonterminal_count == 0
            and not operator_replay_could_mutate
        ),
    }


def _provision_boundary_cleanup_audit(
    *,
    manifest: dict[str, Any],
    expected_watermark: int,
    correlated_delivery_count: int,
    correlated_outbox_count: int,
    activity_states: dict[str, Any],
) -> dict[str, Any] | None:
    """Prove that a failed cross-service gate never entered LOAD-02.

    Resource Booking can fail closed after Timetabler provisioning but before
    its local ENTRY/LOAD-01 evidence exists.  At that exact boundary there is
    deliberately no LOAD-02 API inventory to classify.  Recovery is safe only
    when the source fence is still the declared LOAD-02 start, no durable
    LOAD-02 correlation exists, and every provisioned activity remains wholly
    unallocated.
    """

    load02_start = manifest["expected_stage_start_watermarks"]["load-02"]
    if expected_watermark != load02_start:
        return None
    if correlated_delivery_count or correlated_outbox_count:
        raise LoadHarnessError(
            "provision-boundary cleanup found durable LOAD-02 work"
        )
    activity_selector_names = (
        "traffic_activity_ids",
        "traffic_recurrence_activity_ids",
        "maximum_recurrence_activity_ids",
        "bulk_activity_ids",
        "boundary_activity_ids",
    )
    for name in activity_selector_names:
        state = activity_states.get(name)
        expected_count = manifest["fixture_selectors"][name]["expected_count"]
        if not isinstance(state, dict) or any(
            (
                state.get("activity_count") != expected_count,
                state.get("present_activity_count") != expected_count,
                state.get("absent_activity_count") != 0,
                state.get("scheduled_count") != 0,
                state.get("staff_relation_count") != 0,
                state.get("location_relation_count") != 0,
            )
        ):
            raise LoadHarnessError(
                "provision-boundary activity inventory is not safely unallocated"
            )
    return {
        "containment_passed": True,
        "all_engine_work_terminal": True,
        "late_mutation_possible": False,
        "provision_boundary_recovery": True,
        "load02_durable_delivery_count": correlated_delivery_count,
        "load02_durable_outbox_count": correlated_outbox_count,
        "activity_states": activity_states,
    }


def _provision_slot_attestation(
    provision_evidence: dict[str, Any], manifest: dict[str, Any], run_id: str
) -> dict[str, Any]:
    """Expose only the sealed slot mismatch needed to build a corrected manifest."""

    fixture = provision_evidence.get("fixture_manifest")
    stage_fixture = provision_evidence.get("stage_assertions", {}).get(
        "fixture_manifest"
    )
    if (
        provision_evidence.get("profile") != "PROVISION"
        or provision_evidence.get("service") != "timetabler"
        or provision_evidence.get("run_id") != run_id
        or not isinstance(fixture, dict)
        or fixture != stage_fixture
    ):
        raise LoadHarnessError("sealed PROVISION fixture manifest is contradictory")
    try:
        actual = {
            "slot_a": int(fixture["slot_a"]),
            "slot_b": int(fixture["slot_b"]),
        }
        declared = {
            "slot_a": int(manifest["variables"]["slot_a"]),
            "slot_b": int(manifest["variables"]["slot_b"]),
        }
    except (KeyError, TypeError, ValueError) as error:
        raise LoadHarnessError("sealed PROVISION slots are invalid") from error
    if actual["slot_a"] < 0 or actual["slot_b"] < 0 or actual["slot_a"] == actual["slot_b"]:
        raise LoadHarnessError("sealed PROVISION slots are not two distinct positions")
    return {
        "actual": actual,
        "declared": declared,
        "matches_manifest": actual == declared,
        "source": "immutable PROVISION fixture manifest",
    }


def _assert_current_load02_recovery_control(
    *,
    manifest: dict[str, Any],
    run_id: str,
    action: str,
    exact_sha: str,
) -> dict[str, Any]:
    """Bind a recovery-only deployment to the immutable failed run.

    The normal acceptance chain remains bound to ``commits.timetabler``.  A
    later harness may execute only the two current-run containment actions
    when the replacement manifest embeds the canonical digest of the original
    manifest and the exact failed-run/abort/fence identity.
    """

    control = manifest.get("recovery_control")
    if (
        action not in CURRENT_LOAD02_RECOVERY_ACTIONS
        or not isinstance(control, dict)
        or set(control) != CURRENT_LOAD02_RECOVERY_KEYS
        or control.get("failed_run_id") != run_id
        or not str(control.get("failed_github_run_id") or "").isdigit()
        or not SHA256_PATTERN.fullmatch(
            str(control.get("original_manifest_sha256") or "")
        )
        or not GIT_SHA_PATTERN.fullmatch(
            str(control.get("original_timetabler_sha") or "")
        )
        or not GIT_SHA_PATTERN.fullmatch(
            str(control.get("recovery_timetabler_sha") or "")
        )
        or not SHA256_PATTERN.fullmatch(
            str(control.get("abort_evidence_sha256") or "")
        )
        or isinstance(control.get("settled_source_watermark"), bool)
        or not isinstance(control.get("settled_source_watermark"), int)
        or control["settled_source_watermark"] < 1
        or control.get("approved_actions")
        != sorted(CURRENT_LOAD02_RECOVERY_ACTIONS)
        or not str(control.get("record_ref") or "").strip()
        or not isinstance(control.get("approved_at"), str)
    ):
        raise LoadHarnessError("current LOAD-02 recovery control is incomplete")
    original_manifest = {
        key: value for key, value in manifest.items() if key != "recovery_control"
    }
    if (
        _sha256(original_manifest) != control["original_manifest_sha256"]
        or manifest["commits"]["timetabler"]
        != control["original_timetabler_sha"]
        or exact_sha != control["recovery_timetabler_sha"]
        or action not in control["approved_actions"]
    ):
        raise LoadHarnessError(
            "current LOAD-02 recovery deployment differs from its sealed control"
        )
    approved_at = _parse_iso8601(
        control["approved_at"], "recovery_control.approved_at"
    )
    window_end = _parse_iso8601(
        manifest["approved_window"]["ends_at"], "approved_window.ends_at"
    )
    if approved_at > window_end:
        raise LoadHarnessError(
            "current LOAD-02 recovery approval falls outside the load window"
        )
    return control


def audit_failed_load02_engine_work_from_environment() -> dict[str, Any]:
    """Correlate the failed asynchronous load without writes or credentials."""

    action = os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip()
    if action not in {
        LOAD02_ENGINE_AUDIT_ACTION,
        CURRENT_LOAD02_ENGINE_AUDIT_ACTION,
    }:
        raise LoadHarnessError(
            "failed LOAD-02 audit CLI action differs from guarded workflow action"
        )
    if _env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION") or _env_true(
        "TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"
    ):
        raise LoadHarnessError(
            "failed LOAD-02 audit requires false mutation and operational confirmations"
        )
    run_id = os.environ.get("TT_PHASE1_LOAD_RUN_ID", "").strip()
    if action == LOAD02_ENGINE_AUDIT_ACTION and run_id != LOAD02_CHECKPOINT_RUN_ID:
        raise LoadHarnessError(
            "failed LOAD-02 audit run id differs from the fixed recovery"
        )
    validate_run_id(run_id)
    github_run_id = os.environ.get("TT_PHASE1_LOAD_GITHUB_RUN_ID", "").strip()
    if not github_run_id.isdigit():
        raise LoadHarnessError("failed LOAD-02 audit GitHub run id is invalid")
    manifest = load_manifest(
        os.environ.get("TT_PHASE1_LOAD_MANIFEST_JSON", ""),
        run_id=run_id,
        environment=os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
        source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
        require_approval=True,
    )
    exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
    recovery_control = None
    if action == CURRENT_LOAD02_ENGINE_AUDIT_ACTION:
        recovery_control = _assert_current_load02_recovery_control(
            manifest=manifest,
            run_id=run_id,
            action=action,
            exact_sha=exact_sha,
        )
    elif exact_sha != manifest["commits"]["timetabler"]:
        raise LoadHarnessError(
            "deployed Timetabler SHA differs from the audit manifest"
        )
    if os.environ.get("RB_INTEGRATION_PHASE", "phase1").strip().lower() != "phase1":
        raise LoadHarnessError("failed LOAD-02 audit requires Phase 1 mode")
    before_metrics = _source_metrics()
    current_watermark = before_metrics["transport_watermark"]
    start = manifest["expected_stage_start_watermarks"]["load-02"]
    final = manifest["expected_stage_final_watermarks"]["load-02"]
    if not start < current_watermark < final:
        raise LoadHarnessError(
            "failed LOAD-02 audit watermark is outside its recovery fence"
        )
    if (
        recovery_control is not None
        and current_watermark != recovery_control["settled_source_watermark"]
    ):
        raise LoadHarnessError(
            "current LOAD-02 source fence differs from its recovery control"
        )
    if before_metrics["reverse_delivery_enabled"]:
        raise LoadHarnessError("failed LOAD-02 audit found reverse delivery enabled")
    first_snapshot = _read_failed_load02_engine_audit_snapshot(
        manifest=manifest, run_id=run_id
    )
    second_snapshot = _read_failed_load02_engine_audit_snapshot(
        manifest=manifest, run_id=run_id
    )
    first_digest = _sha256(first_snapshot)
    second_digest = _sha256(second_snapshot)
    if first_digest != second_digest:
        raise LoadHarnessError(
            "failed LOAD-02 engine inventory changed across two reads"
        )
    after_metrics = _source_metrics()
    stable_source_keys = (
        "transport_watermark",
        "publisher_last_sequence",
        "outbox_depth",
        "outbox_dead_letter_count",
        "post_commit_dead_letter_count",
        "engine_response_quarantine_count",
        "reverse_delivery_enabled",
    )
    if any(before_metrics[key] != after_metrics[key] for key in stable_source_keys):
        raise LoadHarnessError("failed LOAD-02 source metrics changed across the audit")
    audit = _classify_failed_load02_engine_audit(
        snapshot=second_snapshot,
        manifest=manifest,
        allow_partial_prefix=action == CURRENT_LOAD02_ENGINE_AUDIT_ACTION,
    )
    source_caught_up = (
        after_metrics["publisher_last_sequence"] == current_watermark
        and after_metrics["outbox_depth"] == 0
        and after_metrics["outbox_dead_letter_count"] == 0
    )
    source_global_correctness_clear = (
        source_caught_up
        and after_metrics["post_commit_dead_letter_count"] == 0
        and after_metrics["engine_response_quarantine_count"] == 0
    )
    audit["source_caught_up"] = source_caught_up
    audit["source_global_correctness_clear"] = source_global_correctness_clear
    audit["containment_passed"] = bool(
        audit["containment_passed"] and source_global_correctness_clear
    )
    result = {
        "phase1_load": "failed_load02_engine_work_audited",
        "run_id": run_id,
        "github_run_id": github_run_id,
        "git_sha": exact_sha,
        "resource_booking_git_sha_attestation": manifest["commits"]["resource_booking"],
        "source_entry_watermark": start,
        "source_current_watermark": current_watermark,
        "source_final_fence": final,
        "two_read_inventory_sha256": second_digest,
        "two_read_inventory_stable": True,
        "audit": audit,
        "final_source_metrics": after_metrics,
        "evidence_written": False,
        "api_authentication_performed": False,
        "mutation_executed": False,
        "load_executed": False,
        "restart_executed": False,
        "database_write_executed": False,
        "manual_outbox_edit": False,
        "kafka_publish_executed": False,
        "cleanup_executed": False,
        "resource_booking_database_touched": False,
        "resource_booking_process_touched": False,
        "phase2_enabled": False,
        "reverse_delivery_enabled": False,
        "dr_executed": False,
    }
    _assert_exchange_payload_has_no_credentials(result)
    return result


def cleanup_abandoned_load02_from_environment() -> dict[str, Any]:
    """Remove only the fixed failed LOAD-02 fixtures through supported APIs.

    The operation is deliberately separate from the acceptance-chain CLEANUP
    stage.  It first proves every asynchronous engine request is terminal and
    that no automatic or operator replay can still mutate the fixture set.
    Only the immutable provision inventory for the fixed failed run can be
    selected; the authoritative academic-term activities are therefore out of
    scope by construction.
    """

    action = os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip()
    if action not in {
        ABANDONED_LOAD02_CLEANUP_ACTION,
        CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION,
    }:
        raise LoadHarnessError(
            "abandoned LOAD-02 cleanup CLI action differs from guarded workflow action"
        )
    if not _env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION") or _env_true(
        "TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"
    ):
        raise LoadHarnessError(
            "abandoned LOAD-02 cleanup requires mutation confirmation only"
        )
    run_id = os.environ.get("TT_PHASE1_LOAD_RUN_ID", "").strip()
    if action == ABANDONED_LOAD02_CLEANUP_ACTION and run_id != LOAD02_CHECKPOINT_RUN_ID:
        raise LoadHarnessError("abandoned cleanup run id differs from the fixed failed run")
    validate_run_id(run_id)
    expected_watermark_raw = os.environ.get(
        "TT_PHASE1_RECOVERY_EXPECTED_SOURCE_WATERMARK", ""
    ).strip()
    if not expected_watermark_raw.isdigit() or int(expected_watermark_raw) < 1:
        raise LoadHarnessError("abandoned cleanup requires an exact positive source fence")
    expected_watermark = int(expected_watermark_raw)
    github_run_id = os.environ.get("TT_PHASE1_LOAD_GITHUB_RUN_ID", "").strip()
    exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
    if not github_run_id.isdigit() or not GIT_SHA_PATTERN.fullmatch(exact_sha):
        raise LoadHarnessError("abandoned cleanup deployment identity is invalid")
    manifest_json = os.environ.get("TT_PHASE1_LOAD_MANIFEST_JSON", "")
    manifest = load_manifest(
        manifest_json,
        run_id=run_id,
        environment=os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
        source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
        require_approval=True,
        require_authoritative_term_baseline=(
            action == CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION
        ),
    )
    recovery_control = None
    if action == CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION:
        recovery_control = _assert_current_load02_recovery_control(
            manifest=manifest,
            run_id=run_id,
            action=action,
            exact_sha=exact_sha,
        )
        if expected_watermark != recovery_control["settled_source_watermark"]:
            raise LoadHarnessError(
                "current abandoned cleanup fence differs from recovery control"
            )
    elif exact_sha != manifest["commits"]["timetabler"]:
        raise LoadHarnessError(
            "deployed Timetabler SHA differs from the cleanup manifest"
        )
    if os.environ.get("RB_INTEGRATION_PHASE", "phase1").strip().lower() != "phase1":
        raise LoadHarnessError("abandoned cleanup requires Phase 1 mode")

    root, run_dir, provision_path = _evidence_path(run_id, "PROVISION")
    provision_evidence, _ = _read_immutable_exchange_json(
        root=root,
        run_dir=run_dir,
        path=provision_path,
        label="PROVISION evidence for cleanup slot attestation",
    )
    provision_slots = _provision_slot_attestation(
        provision_evidence, manifest, run_id
    )
    if (
        recovery_control is not None
        and provision_evidence.get("git_sha")
        != recovery_control["original_timetabler_sha"]
    ):
        raise LoadHarnessError(
            "PROVISION evidence differs from the recovery-controlled deployment"
        )

    from api.models import IntegrationOutbox, PostCommitDelivery

    cleanup_transactions = manifest["workloads"]["cleanup"][
        "expected_source_transactions"
    ]
    completed_cleanup_rows = list(
        IntegrationOutbox.objects.filter(
            request_id__startswith=f"{run_id}-cleanup-",
            transaction_finalized=True,
        )
        .values_list("request_id", "transport_sequence")
        .distinct()
        .order_by("transport_sequence", "request_id")
    )
    before_metrics = _source_metrics()
    if action == ABANDONED_LOAD02_CLEANUP_ACTION:
        initial_cleanup_watermark = ABANDONED_LOAD02_CLEANUP_INITIAL_WATERMARK
    elif completed_cleanup_rows:
        initial_cleanup_watermark = min(
            int(sequence) for _, sequence in completed_cleanup_rows
        ) - 1
    else:
        initial_cleanup_watermark = expected_watermark
    retry_state = _abandoned_cleanup_retry_state(
        run_id=run_id,
        initial_watermark=initial_cleanup_watermark,
        current_watermark=before_metrics["transport_watermark"],
        expected_current_watermark=expected_watermark,
        cleanup_transaction_count=cleanup_transactions,
        completed_rows=[
            (str(request_id), int(sequence))
            for request_id, sequence in completed_cleanup_rows
        ],
    )
    recovery_fixtures = _resolve_cleanup_selectors(manifest, run_id)
    load02_request_ids = [
        operation["request_id"]
        for operation in _failed_load02_expected_operations(
            manifest=manifest,
            run_id=run_id,
            fixtures=recovery_fixtures,
        )
    ]
    activity_states = _audit_run_owned_activity_states(
        manifest=manifest,
        fixtures=recovery_fixtures,
        allow_absent=False,
    )
    audit = _provision_boundary_cleanup_audit(
        manifest=manifest,
        expected_watermark=expected_watermark,
        correlated_delivery_count=PostCommitDelivery.objects.filter(
            correlation_id__in=load02_request_ids
        ).count(),
        correlated_outbox_count=IntegrationOutbox.objects.filter(
            request_id__in=load02_request_ids
        ).count(),
        activity_states=activity_states,
    )
    if audit is not None:
        first_snapshot = {
            "provision_boundary_recovery": True,
            "activity_states": activity_states,
        }
        second_snapshot = copy.deepcopy(first_snapshot)
    else:
        first_snapshot = _read_failed_load02_engine_audit_snapshot(
            manifest=manifest,
            run_id=run_id,
            fixtures=recovery_fixtures,
            allow_absent_activities=True,
        )
        second_snapshot = _read_failed_load02_engine_audit_snapshot(
            manifest=manifest,
            run_id=run_id,
            fixtures=recovery_fixtures,
            allow_absent_activities=True,
        )
        if _sha256(first_snapshot) != _sha256(second_snapshot):
            raise LoadHarnessError(
                "failed LOAD-02 engine inventory changed across two reads"
            )
        audit = _classify_failed_load02_engine_audit(
            snapshot=second_snapshot,
            manifest=manifest,
            allow_partial_prefix=(
                action == CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION
            ),
        )
    stable_metrics = _source_metrics()
    stable_keys = (
        "transport_watermark",
        "publisher_last_sequence",
        "outbox_depth",
        "outbox_dead_letter_count",
        "post_commit_dead_letter_count",
        "engine_response_quarantine_count",
        "reverse_delivery_enabled",
    )
    if any(before_metrics[key] != stable_metrics[key] for key in stable_keys):
        raise LoadHarnessError("failed LOAD-02 state changed during containment audit")
    source_clear = (
        stable_metrics["publisher_last_sequence"] == expected_watermark
        and stable_metrics["outbox_depth"] == 0
        and stable_metrics["outbox_dead_letter_count"] == 0
        and stable_metrics["post_commit_dead_letter_count"] == 0
        and stable_metrics["engine_response_quarantine_count"] == 0
        and stable_metrics["reverse_delivery_enabled"] is False
    )
    if not (
        audit["containment_passed"]
        and audit["all_engine_work_terminal"]
        and audit["late_mutation_possible"] is False
        and source_clear
    ):
        raise LoadHarnessError(
            "failed LOAD-02 engine work is not safely contained; cleanup remains blocked"
        )

    cleanup_manifest = copy.deepcopy(manifest)
    cleanup_manifest["expected_stage_start_watermarks"]["cleanup"] = retry_state[
        "initial_watermark"
    ]
    cleanup_manifest["expected_stage_final_watermarks"]["cleanup"] = (
        retry_state["final_watermark"]
    )
    config = LoadConfig(
        action="cleanup",
        run_id=run_id,
        environment="staging",
        source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
        base_url=os.environ.get("TT_PHASE1_LOAD_BASE_URL", "").strip(),
        email=os.environ.get("TT_PHASE1_LOAD_EMAIL", "").strip(),
        password=os.environ.get("TT_PHASE1_LOAD_PASSWORD", ""),
        manifest_json=manifest_json,
        expected_sha=exact_sha,
        github_run_id=github_run_id,
        prior_evidence_sha256="",
        mutation_confirmed=True,
        operational_confirmed=False,
        restart_fence_sha256="",
    )
    session = identity_gate.TimetablerAdminSession(
        base_url=identity_gate.validate_base_url(config.base_url, "Timetabler base URL"),
        email=config.email,
        password=config.password,
        run_id=config.run_id,
    )
    session.login_and_attest()
    runner = LoadRunner(config=config, manifest=cleanup_manifest, session=session)
    runner.run_cleanup()
    final_metrics = _assert_source_safety(
        expected_watermark=retry_state["final_watermark"]
    )
    evidence = seal_evidence(
        {
            "schema_version": 1,
            "action": action,
            "repository": "Mayvins/timetabler-be",
            "git_sha": exact_sha,
            "github_run_id": github_run_id,
            "run_id": run_id,
            "source_entry_watermark": retry_state["initial_watermark"],
            "source_resume_watermark": retry_state["current_watermark"],
            "source_final_watermark": final_metrics["transport_watermark"],
            "source_transaction_count": cleanup_transactions,
            "completed_transaction_count_before_retry": retry_state[
                "completed_transaction_count"
            ],
            "remaining_transaction_count_at_retry": retry_state[
                "remaining_transaction_count"
            ],
            "containment_inventory_sha256": _sha256(second_snapshot),
            "containment_audit_sha256": _sha256(audit),
            "containment_passed": True,
            "all_engine_work_terminal": True,
            "late_mutation_possible": False,
            "cleanup_assertions": runner.stage_assertions,
            "cleanup_transactions_sha256": _sha256(runner.transactions),
            "provision_slot_attestation": provision_slots,
            "supported_timetabler_apis_only": True,
            "authoritative_term_activity_ids_touched": False,
            "direct_database_mutation": False,
            "resource_booking_database_touched": False,
            "reverse_delivery_enabled": False,
            "phase2_enabled": False,
            "dr_executed": False,
        }
    )
    _assert_exchange_payload_has_no_credentials(evidence)
    return evidence


def checkpoint_load02_from_environment() -> dict[str, Any]:
    """Seal one read-only durable LOAD-02 prefix checkpoint for the fixed recovery."""

    if os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip() != "checkpoint-load-02":
        raise LoadHarnessError(
            "checkpoint-load-02 CLI action differs from guarded workflow action"
        )
    if _env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION") or _env_true(
        "TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"
    ):
        raise LoadHarnessError(
            "checkpoint-load-02 requires false mutation and operational confirmations"
        )
    run_id = os.environ.get("TT_PHASE1_LOAD_RUN_ID", "").strip()
    if run_id != LOAD02_CHECKPOINT_RUN_ID:
        raise LoadHarnessError(
            "checkpoint-load-02 run id differs from the fixed recovery"
        )
    validate_run_id(run_id)
    github_run_id = os.environ.get("TT_PHASE1_LOAD_GITHUB_RUN_ID", "").strip()
    if not github_run_id.isdigit():
        raise LoadHarnessError("checkpoint-load-02 GitHub run id is invalid")
    supplied_sha256 = (
        os.environ.get("TT_PHASE1_LOAD_PRIOR_EVIDENCE_SHA256", "").strip().lower()
    )
    manifest = load_manifest(
        os.environ.get("TT_PHASE1_LOAD_MANIFEST_JSON", ""),
        run_id=run_id,
        environment=os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
        source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
        require_approval=True,
    )
    assert_evidence_exchange_runtime(manifest, run_id=run_id)
    exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
    if (
        not GIT_SHA_PATTERN.fullmatch(exact_sha)
        or exact_sha != manifest["commits"]["timetabler"]
    ):
        raise LoadHarnessError(
            "deployed Timetabler SHA differs from the checkpoint manifest"
        )
    predecessor = verify_predecessor_evidence(
        run_id=run_id,
        stage="load-02",
        manifest=manifest,
        supplied_sha256=supplied_sha256,
    )
    assert_approved_window(manifest, "LOAD-02")
    metrics = _assert_source_safety(expected_watermark=None)
    current_watermark = metrics["transport_watermark"]
    approved_start = manifest["expected_stage_start_watermarks"]["load-02"]
    approved_final = manifest["expected_stage_final_watermarks"]["load-02"]
    if not approved_start < current_watermark < approved_final:
        raise LoadHarnessError(
            "checkpoint-load-02 watermark is outside its strict open recovery fence"
        )
    prefix_transactions = _load02_resume_prefix_evidence(
        manifest=manifest,
        run_id=run_id,
        current_watermark=current_watermark,
    )
    final_metrics = _assert_source_safety(expected_watermark=current_watermark)
    final_prefix_transactions = _load02_resume_prefix_evidence(
        manifest=manifest,
        run_id=run_id,
        current_watermark=final_metrics["transport_watermark"],
    )
    if final_prefix_transactions != prefix_transactions:
        raise LoadHarnessError("LOAD-02 durable prefix changed during checkpointing")
    prefix_sha256 = _sha256(prefix_transactions)
    artifact = seal_evidence(
        {
            "schema_version": 1,
            "artifact_type": "timetabler.phase1.load02.checkpoint",
            "profile": "LOAD-02",
            "stage": "checkpoint-load-02",
            "service": "timetabler",
            "repository": "Mayvins/timetabler-be",
            "run_id": run_id,
            "environment": os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
            "production": False,
            "git_sha": exact_sha,
            "resource_booking_git_sha_attestation": manifest["commits"][
                "resource_booking"
            ],
            "github_run_id": github_run_id,
            "source_scope": manifest["source_scope"],
            "entry_watermark": manifest["expected_stage_start_watermarks"]["load-02"],
            "current_watermark": current_watermark,
            "approved_final_watermark": manifest["expected_stage_final_watermarks"][
                "load-02"
            ],
            "prefix_transaction_count": len(prefix_transactions),
            "source_transactions_sha256": prefix_sha256,
            "source_transactions": prefix_transactions,
            "prior_evidence_sha256": supplied_sha256,
            "predecessor_evidence": predecessor,
            "configuration_snapshot_sha256": manifest["configuration_snapshot"][
                "sha256"
            ],
            "evidence_exchange": manifest["evidence_exchange"],
            "final_source_metrics": final_metrics,
            "checkpoint_assertions": {
                "exact_contiguous_deterministic_prefix_verified": True,
                "all_source_transactions_finalized": True,
                "all_source_transactions_published": True,
                "unrelated_source_transactions": 0,
                "mutation_executed": False,
                "load_executed": False,
                "process_restart_executed": False,
                "database_write_executed": False,
                "manual_outbox_edit": False,
                "resource_booking_database_touched": False,
                "resource_booking_process_touched": False,
                "phase2_enabled": False,
                "reverse_delivery_enabled": False,
                "dr_executed": False,
            },
            "mutation_executed": False,
            "load_executed": False,
            "restart_executed": False,
            "database_write_executed": False,
            "manual_outbox_edit": False,
            "resource_booking_database_touched": False,
            "resource_booking_process_touched": False,
            "phase2_enabled": False,
            "reverse_delivery_enabled": False,
            "dr_executed": False,
        }
    )
    _assert_exchange_payload_has_no_credentials(artifact)
    _, run_dir = _ensure_evidence_exchange_directories(run_id)
    target = run_dir / LOAD02_CHECKPOINT_ARTIFACT_FILENAME
    info = _persist_immutable_exchange_file(
        target=target,
        encoded=_canonical_bytes(artifact) + b"\n",
        label="LOAD-02 checkpoint",
    )
    return {
        "phase1_load": "load02_checkpoint_persisted",
        "profile": "LOAD-02",
        "run_id": run_id,
        "entry_watermark": manifest["expected_stage_start_watermarks"]["load-02"],
        "current_watermark": current_watermark,
        "approved_final_watermark": manifest["expected_stage_final_watermarks"][
            "load-02"
        ],
        "prefix_transaction_count": len(prefix_transactions),
        "source_transactions_sha256": prefix_sha256,
        "evidence_sha256": artifact["evidence_sha256"],
        "evidence_path": str(target),
        "owner_uid": info.st_uid,
        "reader_uid": RESOURCE_BOOKING_EVIDENCE_READER_UID,
        "mode": "0444",
        "bytes": info.st_size,
        "mutation_executed": False,
        "load_executed": False,
        "restart_executed": False,
        "database_write_executed": False,
        "resource_booking_database_touched": False,
        "resource_booking_process_touched": False,
        "phase2_enabled": False,
        "reverse_delivery_enabled": False,
        "dr_executed": False,
    }


def _rate_execution_plan(
    *,
    manifest: dict[str, Any],
    run_id: str,
    stage: str,
    rate: int,
    duration_seconds: int,
    current_watermark: int,
) -> dict[str, Any]:
    total_approved_requests = rate * duration_seconds
    approved_start = manifest["expected_stage_start_watermarks"][stage]
    approved_final = manifest["expected_stage_final_watermarks"][stage]
    if approved_final != approved_start + total_approved_requests:
        raise LoadHarnessError(
            "rate workload count does not reach its approved final fence"
        )
    if current_watermark != approved_start:
        raise LoadHarnessError(
            "engine-aware rate workload differs from its exact entry fence and is not resumable"
        )
    prefix_transactions = []
    prefix_count = len(prefix_transactions)
    active_request_count = total_approved_requests - prefix_count
    if active_request_count <= 0:
        raise LoadHarnessError("LOAD-02 resume has no strictly bounded active segment")
    return {
        "approved_start": approved_start,
        "approved_final": approved_final,
        "resume_watermark": current_watermark,
        "total_approved_requests": total_approved_requests,
        "prefix_transactions": prefix_transactions,
        "prefix_count": prefix_count,
        "active_request_count": active_request_count,
        "active_duration_seconds": active_request_count / rate,
    }


def _wait_for_stage_publication(
    *, start: int, expected_delta: int, request_ids: set[str], timeout_seconds: int
) -> tuple[dict[str, Any], list[dict[str, Any]]]:
    deadline = time.monotonic() + timeout_seconds
    expected_end = start + expected_delta
    while time.monotonic() < deadline:
        current = _current_watermark()
        if current > expected_end:
            raise LoadHarnessError("source watermark exceeded the approved stage fence")
        if current == expected_end:
            try:
                metrics = _assert_source_safety(expected_watermark=expected_end)
                return metrics, _transaction_evidence(start, expected_end, request_ids)
            except LoadHarnessError as error:
                if "caught up" not in str(error) and "outbox_depth" not in str(error):
                    raise
        time.sleep(1)
    raise LoadHarnessError(
        "source publication did not reach the exact stage fence in time"
    )


def _is_engine_request_delivery(
    delivery: Any, *, expected_activity_id: int
) -> bool:
    """Distinguish the engine request from correlated downstream deliveries."""

    if isinstance(delivery, dict):
        channel = delivery.get("channel")
        method = delivery.get("method")
        activity_ids = delivery.get("payload_activity_ids")
    else:
        channel = delivery.channel
        method = delivery.method
        payload = delivery.payload if isinstance(delivery.payload, dict) else {}
        activity_ids = _safe_int_list(payload.get("activities"))
    return (
        channel == "kafka"
        and method == "schedule"
        and activity_ids == [expected_activity_id]
    )


def _wait_for_engine_applied_transaction(
    *, correlation_id: str, expected_activity_id: int, timeout_seconds: int
) -> dict[str, Any]:
    """Wait for one exact engine response, atomic receipt, and published change set."""

    from api.models import (
        AppliedEngineResponse,
        EngineResponseQuarantine,
        IntegrationOutbox,
        KafkaLog,
        PostCommitDelivery,
    )

    deadline = time.monotonic() + timeout_seconds
    while time.monotonic() < deadline:
        correlated_deliveries = list(
            PostCommitDelivery.objects.filter(correlation_id=correlation_id)
            .order_by("id")[:3]
        )
        if len(correlated_deliveries) > 2:
            raise LoadHarnessError(
                "engine-aware request exceeded its bounded correlated delivery inventory"
            )
        deliveries = [
            delivery
            for delivery in correlated_deliveries
            if _is_engine_request_delivery(
                delivery, expected_activity_id=expected_activity_id
            )
        ]
        if len(deliveries) > 1:
            raise LoadHarnessError(
                "engine-aware request created more than one matching engine delivery"
            )
        if not deliveries:
            if correlated_deliveries:
                raise LoadHarnessError("engine-aware delivery contract is invalid")
            time.sleep(0.2)
            continue
        delivery = deliveries[0]
        delivery_id = str(delivery.delivery_id)
        if (
            delivery.channel != "kafka"
            or delivery.method != "schedule"
            or _safe_int_list(delivery.payload.get("activities"))
            != [expected_activity_id]
        ):
            raise LoadHarnessError("engine-aware delivery contract is invalid")
        if delivery.status == PostCommitDelivery.Status.DEAD_LETTER:
            raise LoadHarnessError("engine-aware delivery reached dead letter")
        if delivery.status != PostCommitDelivery.Status.PUBLISHED:
            time.sleep(0.2)
            continue
        if EngineResponseQuarantine.objects.filter(request_id=delivery_id).exists():
            raise LoadHarnessError("engine-aware response was quarantined")
        kafka_rows = list(KafkaLog.objects.filter(request_id=delivery_id)[:2])
        if len(kafka_rows) > 1:
            raise LoadHarnessError("engine-aware Kafka correlation is not unique")
        if not kafka_rows or not kafka_rows[0].response_at:
            time.sleep(0.2)
            continue
        parsed = _parse_engine_response_for_audit(
            kafka_rows[0].response_data,
            delivery_id=delivery_id,
            expected_activity_id=expected_activity_id,
        )
        if (
            not parsed["protocol_valid"]
            or parsed["status"] != "success"
            or not parsed["successful_assignment"]
        ):
            raise LoadHarnessError(
                "Scheduling Engine did not return one successful assignment"
            )
        receipts = list(
            AppliedEngineResponse.objects.filter(request_id=delivery_id).order_by("id")[:2]
        )
        if len(receipts) != 1:
            time.sleep(0.2)
            continue
        receipt = receipts[0]
        if (
            str(receipt.response_hash) != parsed["canonical_response_hash"]
            or str(receipt.outcome) != AppliedEngineResponse.Outcome.APPLIED
        ):
            raise LoadHarnessError("engine-aware applied receipt is contradictory")
        rows = list(
            IntegrationOutbox.objects.filter(change_set_id=receipt.change_set_id)
            .order_by("transport_sequence", "transaction_index")
        )
        sequences = sorted({int(row.transport_sequence) for row in rows})
        if len(sequences) != 1:
            raise LoadHarnessError(
                "successful engine response did not emit exactly one source transaction"
            )
        try:
            transaction = _transaction_evidence(
                sequences[0] - 1,
                sequences[0],
                {delivery_id},
            )[0]
        except LoadHarnessError as error:
            if str(error) != "source transaction is not completely published":
                raise
            time.sleep(0.2)
            continue
        return {
            "engine_delivery_id": delivery_id,
            "engine_response_sha256": parsed["canonical_response_hash"],
            "engine_change_set_id": str(receipt.change_set_id),
            "engine_source_sequence": sequences[0],
            "source_request_id": delivery_id,
            "source_transaction": transaction,
        }
    raise LoadHarnessError("Scheduling Engine completion exceeded its bounded timeout")


def _recurring_assignment_state(
    *,
    manifest: dict[str, Any],
    activity_id: int,
    expect_scheduled: bool,
) -> dict[str, Any]:
    """Prove one rate activity's term, recurrence, and resource transition."""

    from api.models import TtActivity
    from api.services.integration.outbox import capture_activity_state

    activity = (
        TtActivity.objects.filter(id=activity_id)
        .select_related("academic_term", "week_pattern")
        .prefetch_related(
            "week",
            "week_pattern__week",
            "staff",
            "location",
            "staff_preset",
            "location_preset",
            "student_set",
        )
        .first()
    )
    if activity is None:
        raise LoadHarnessError("rate activity disappeared during assignment attestation")
    state = capture_activity_state(activity)
    scope = manifest["load_acceptance_scope"]
    actual_scope = {
        "academic_term_id": state.get("academic_term_id"),
        "academic_term_start_date": state.get("calendar", {}).get(
            "academic_term_start_date"
        ),
        "academic_term_end_date": state.get("calendar", {}).get(
            "academic_term_end_date"
        ),
    }
    expected_scope = {
        key: scope[key]
        for key in (
            "academic_term_id",
            "academic_term_start_date",
            "academic_term_end_date",
        )
    }
    recurrence_count = len(state.get("week_ids") or [])
    occurrence_count = len(state.get("occurrences") or [])
    identity = state.get("identity")
    staff_ids = state.get("staff_ids") or []
    location_ids = state.get("location_ids") or []
    preset_staff_ids = list(activity.staff_preset.values_list("id", flat=True))
    preset_location_ids = list(activity.location_preset.values_list("id", flat=True))
    expected_assignment_count = 1 if expect_scheduled else 0
    if (
        actual_scope != expected_scope
        or not isinstance(identity, dict)
        or identity.get("aggregate_type") != "activity"
        or not str(identity.get("deployment_id") or "")
        or not str(identity.get("source_id") or "")
        or recurrence_count
        < scope["minimum_recurring_occurrences_per_activity"]
        or bool(state.get("scheduled")) is not expect_scheduled
        or len(staff_ids) != expected_assignment_count
        or len(location_ids) != expected_assignment_count
        or len(preset_staff_ids) != 1
        or len(preset_location_ids) != 1
        or (expect_scheduled and staff_ids != preset_staff_ids)
        or (expect_scheduled and location_ids != preset_location_ids)
        or occurrence_count != (recurrence_count if expect_scheduled else 0)
    ):
        raise LoadHarnessError(
            "term-26 recurring staff/location assignment transition is incomplete"
        )
    return {
        **actual_scope,
        "activity_source_ref_sha256": _sha256(
            f"{identity['deployment_id']}:{identity['source_id']}"
        ),
        "recurrence_occurrence_count": recurrence_count,
        "staff_assignment_count": len(staff_ids),
        "location_assignment_count": len(location_ids),
        "staff_preset_count": len(preset_staff_ids),
        "location_preset_count": len(preset_location_ids),
        "assignment_transition": "assigned" if expect_scheduled else "released",
    }


def _selector_state_hash(manifest: dict[str, Any], selector_name: str) -> str:
    from api import models

    selector = manifest["fixture_selectors"][selector_name]
    model = getattr(models, MODEL_MAP[selector["model"]])
    rows = list(
        model.objects.filter(**selector["filters"])
        .order_by("id")
        .values()[: selector["expected_count"] + 1]
    )
    if len(rows) != selector["expected_count"]:
        raise LoadHarnessError("rollback selector changed cardinality")
    return _sha256(rows)


def _calibrate_activity_allocation(
    *,
    activity_ids: list[int],
    slot: int,
    request_id: str,
    actor_id: int,
    sequence_offset: int = 1,
    prior_scheduled_activity_ids: set[int] | None = None,
    prior_slot: int | None = None,
) -> dict[str, Any]:
    """Read-only canonical sizing for one proposed allocation change set."""

    from django.conf import settings
    from django.test.utils import override_settings

    from api.models import (
        IntegrationAggregateVersion,
        IntegrationTransportCursor,
        TtActivity,
        TtSetting,
    )
    from api.services.integration.outbox import (
        aggregate_source_id,
        canonical_hash,
        capture_activity_state,
        integration_source_scope,
        source_identity,
        build_transport_envelope,
    )

    activities = list(
        TtActivity.objects.filter(id__in=activity_ids)
        .select_related("academic_term", "week_pattern")
        .prefetch_related(
            "week",
            "week_pattern__week",
            "staff",
            "location",
            "student_set",
            "staff_preset",
            "location_preset",
        )
        .order_by("id")
    )
    if len(activities) != len(activity_ids):
        raise LoadHarnessError("boundary calibration candidate set is incomplete")
    values = TtSetting.get_multiple_setting({"minute_per_slot", "slot_per_day"})
    minute_per_slot = int(values["minute_per_slot"])
    slot_per_day = int(values["slot_per_day"])
    timezone_name = settings.TIME_ZONE
    source_timezone = ZoneInfo(timezone_name)
    scope = integration_source_scope()
    cursor = IntegrationTransportCursor.objects.filter(source_scope=scope).first()
    sequence = int(cursor.last_sequence if cursor else 0) + sequence_offset
    change_set_id = uuid.UUID("00000000-0000-4000-8000-000000000001")
    committed_at = "2026-08-08T00:00:00.000000Z"
    version_rows = {
        row.aggregate_id: int(row.version)
        for row in IntegrationAggregateVersion.objects.filter(
            aggregate_type="activity",
            aggregate_id__in=[
                aggregate_source_id(activity.id) for activity in activities
            ],
        )
    }
    prior_scheduled_activity_ids = prior_scheduled_activity_ids or set()

    def projected_state(
        activity, state: dict[str, Any], projected_slot: int
    ) -> dict[str, Any]:
        projected = copy.deepcopy(state)
        projected_day = projected_slot // slot_per_day
        projected_minute = (projected_slot % slot_per_day) * minute_per_slot
        projected_time = datetime.time(projected_minute // 60, projected_minute % 60)
        projected.update(
            {
                "scheduled": True,
                "scheduled_day": projected_day,
                "scheduled_start_time": projected_time.isoformat(),
                "scheduled_start_slot": projected_slot,
                "staff_ids": sorted(activity.staff_preset.values_list("id", flat=True)),
                "location_ids": sorted(
                    activity.location_preset.values_list("id", flat=True)
                ),
            }
        )
        occurrences = []
        for week in projected["calendar"]["weeks"]:
            local_date = datetime.date.fromisoformat(
                week["week_start_date"]
            ) + datetime.timedelta(days=projected_day)
            local_start = datetime.datetime.combine(
                local_date, projected_time, tzinfo=source_timezone
            )
            utc_start = local_start.astimezone(datetime.UTC)
            utc_end = utc_start + datetime.timedelta(
                minutes=projected["duration_minutes"]
            )
            occurrences.append(
                {
                    "occurrence_ref": (
                        f"{settings.RESOURCE_BOOKING_INTEGRATION['DEPLOYMENT_ID']}:"
                        f"activity:{activity.id}:week:{week['week_ref']['source_id']}"
                    ),
                    "week_ref": week["week_ref"],
                    "week_number": week["week_number"],
                    "week_start_date": week["week_start_date"],
                    "local_date": local_date.isoformat(),
                    "local_start": local_start.isoformat(),
                    "local_end": utc_end.astimezone(source_timezone).isoformat(),
                    "utc_start": utc_start.isoformat().replace("+00:00", "Z"),
                    "utc_end": utc_end.isoformat().replace("+00:00", "Z"),
                    "timezone": timezone_name,
                    "utc_offset": local_start.isoformat()[-6:],
                    "dst_fold": settings.RESOURCE_BOOKING_INTEGRATION[
                        "OCCURRENCE_DST_FOLD"
                    ],
                }
            )
        projected["occurrences"] = occurrences
        return projected

    events = []
    for index, activity in enumerate(activities, start=1):
        previous = capture_activity_state(activity)
        prior_applied = activity.id in prior_scheduled_activity_ids
        if prior_applied:
            if prior_slot is None:
                raise LoadHarnessError("boundary calibration prior slot is absent")
            previous = projected_state(activity, previous, prior_slot)
        current = projected_state(activity, previous, slot)
        committed_state = {
            "replacement": {"previous": previous, "current": current},
            "affected_scope": {
                "previous": {
                    "staff_ids": previous["staff_ids"],
                    "location_ids": previous["location_ids"],
                    "week_ids": previous["week_ids"],
                    "scheduled_start_slot": previous["scheduled_start_slot"],
                    "duration_minutes": previous["duration_minutes"],
                },
                "current": {
                    "staff_ids": current["staff_ids"],
                    "location_ids": current["location_ids"],
                    "week_ids": current["week_ids"],
                    "scheduled_start_slot": current["scheduled_start_slot"],
                    "duration_minutes": current["duration_minutes"],
                },
            },
            "tombstone": False,
            "final_state_ref": source_identity("activity", activity.id),
        }
        aggregate_id = aggregate_source_id(activity.id)
        event_id = uuid.UUID(int=index)
        payload = {
            "event_id": str(event_id),
            "schema_version": settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
            "event_type": "timetabler.activity.allocation_replaced",
            "aggregate": source_identity("activity", activity.id),
            "event_version": version_rows.get(aggregate_id, 0) + 1 + int(prior_applied),
            "ordering_key": scope,
            "change_set_id": str(change_set_id),
            "transport": {
                "source_scope": scope,
                "sequence": sequence,
                "transaction_id": str(change_set_id),
                "transaction_index": index,
                "transaction_count": len(activities),
                "member_only": True,
            },
            "origin": "timetabler",
            "correlation_id": request_id,
            "causation_id": None,
            "request_id": request_id,
            "actor_id": actor_id,
            "committed_at": committed_at,
            "committed_state_hash": canonical_hash(committed_state),
            "committed_state": committed_state,
        }
        events.append(
            SimpleNamespace(
                payload=payload,
                transaction_index=index,
                source_scope=scope,
                transport_sequence=sequence,
                change_set_id=change_set_id,
                schema_version=settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
                correlation_id=request_id,
            )
        )
    expanded = {
        **settings.RESOURCE_BOOKING_INTEGRATION,
        "MAX_EVENT_BYTES": 16 * 1024 * 1024,
    }
    with override_settings(RESOURCE_BOOKING_INTEGRATION=expanded):
        envelope = build_transport_envelope(events)
    return {
        "activity_ids": [activity.id for activity in activities],
        "member_count": len(events),
        "canonical_bytes": len(_canonical_bytes(envelope)),
        "events_hash": envelope["events_hash"],
        "read_only": True,
    }


def _find_boundary_calibration(
    *,
    candidate_ids: list[int],
    slot: int,
    request_id: str,
    actor_id: int,
    tolerance: int,
    sequence_offset: int = 1,
    prior_scheduled_activity_ids: set[int] | None = None,
    prior_slot: int | None = None,
) -> tuple[dict[str, Any], dict[str, Any]]:
    low = 1
    high = len(candidate_ids)
    best: dict[str, Any] | None = None
    while low <= high:
        middle = (low + high) // 2
        candidate = _calibrate_activity_allocation(
            activity_ids=candidate_ids[:middle],
            slot=slot,
            request_id=request_id,
            actor_id=actor_id,
            sequence_offset=sequence_offset,
            prior_scheduled_activity_ids=prior_scheduled_activity_ids,
            prior_slot=prior_slot,
        )
        if candidate["canonical_bytes"] < TRANSPORT_LIMIT_BYTES:
            best = candidate
            low = middle + 1
        else:
            high = middle - 1
    if best is None or TRANSPORT_LIMIT_BYTES - best["canonical_bytes"] > tolerance:
        raise LoadHarnessError(
            "read-only calibration found no admissible case near 2 MiB"
        )
    over_count = best["member_count"] + 1
    if over_count > len(candidate_ids):
        raise LoadHarnessError("read-only calibration has no over-boundary candidate")
    over = _calibrate_activity_allocation(
        activity_ids=candidate_ids[:over_count],
        slot=slot,
        request_id=request_id,
        actor_id=actor_id,
        sequence_offset=sequence_offset,
        prior_scheduled_activity_ids=prior_scheduled_activity_ids,
        prior_slot=prior_slot,
    )
    if over["canonical_bytes"] <= TRANSPORT_LIMIT_BYTES:
        raise LoadHarnessError(
            "read-only calibration over-boundary case is not over 2 MiB"
        )
    return best, over


def _calibrate_activity_deletion(
    *,
    activity_ids: list[int],
    request_id: str,
    actor_id: int,
    sequence_offset: int,
) -> dict[str, Any]:
    """Read-only canonical sizing for one proposed activity-delete batch."""

    from django.conf import settings
    from django.test.utils import override_settings

    from api.models import (
        IntegrationAggregateVersion,
        IntegrationTransportCursor,
        TtActivity,
    )
    from api.services.integration.outbox import (
        aggregate_source_id,
        build_transport_envelope,
        canonical_hash,
        capture_activity_state,
        integration_source_scope,
        source_identity,
    )

    activities = list(
        TtActivity.objects.filter(id__in=activity_ids)
        .select_related("academic_term", "week_pattern")
        .prefetch_related(
            "week",
            "week_pattern__week",
            "staff",
            "location",
            "student_set",
        )
        .order_by("id")
    )
    if len(activities) != len(activity_ids):
        raise LoadHarnessError("cleanup calibration candidate set is incomplete")
    scope = integration_source_scope()
    cursor = IntegrationTransportCursor.objects.filter(source_scope=scope).first()
    sequence = int(cursor.last_sequence if cursor else 0) + sequence_offset
    change_set_id = uuid.UUID("00000000-0000-4000-8000-000000000002")
    committed_at = "2026-08-08T00:00:00.000000Z"
    version_rows = {
        row.aggregate_id: int(row.version)
        for row in IntegrationAggregateVersion.objects.filter(
            aggregate_type="activity",
            aggregate_id__in=[
                aggregate_source_id(activity.id) for activity in activities
            ],
        )
    }
    events = []
    for index, activity in enumerate(activities, start=1):
        previous = capture_activity_state(activity)
        committed_state = {
            "replacement": {"previous": previous, "current": None},
            "affected_scope": {
                "previous": {
                    "staff_ids": previous["staff_ids"],
                    "location_ids": previous["location_ids"],
                    "week_ids": previous["week_ids"],
                    "scheduled_start_slot": previous["scheduled_start_slot"],
                    "duration_minutes": previous["duration_minutes"],
                },
                "current": {
                    "staff_ids": [],
                    "location_ids": [],
                    "week_ids": [],
                    "scheduled_start_slot": None,
                    "duration_minutes": None,
                },
            },
            "tombstone": True,
            "final_state_ref": source_identity("activity", activity.id),
        }
        aggregate_id = aggregate_source_id(activity.id)
        payload = {
            "event_id": str(uuid.UUID(int=index)),
            "schema_version": settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
            "event_type": "timetabler.activity.deleted",
            "aggregate": source_identity("activity", activity.id),
            "event_version": version_rows.get(aggregate_id, 0) + 1,
            "ordering_key": scope,
            "change_set_id": str(change_set_id),
            "transport": {
                "source_scope": scope,
                "sequence": sequence,
                "transaction_id": str(change_set_id),
                "transaction_index": index,
                "transaction_count": len(activities),
                "member_only": True,
            },
            "origin": "timetabler",
            "correlation_id": request_id,
            "causation_id": None,
            "request_id": request_id,
            "actor_id": actor_id,
            "committed_at": committed_at,
            "committed_state_hash": canonical_hash(committed_state),
            "committed_state": committed_state,
        }
        events.append(
            SimpleNamespace(
                payload=payload,
                transaction_index=index,
                source_scope=scope,
                transport_sequence=sequence,
                change_set_id=change_set_id,
                schema_version=settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
                correlation_id=request_id,
            )
        )
    expanded = {
        **settings.RESOURCE_BOOKING_INTEGRATION,
        "MAX_EVENT_BYTES": 16 * 1024 * 1024,
    }
    with override_settings(RESOURCE_BOOKING_INTEGRATION=expanded):
        envelope = build_transport_envelope(events)
    return {
        "member_count": len(events),
        "canonical_bytes": len(_canonical_bytes(envelope)),
        "events_hash": envelope["events_hash"],
        "read_only": True,
    }


def _consume_completed_futures(
    futures: list[Future[dict[str, Any]]],
    consumed: set[Future[dict[str, Any]]],
) -> list[dict[str, Any]]:
    """Return newly completed work immediately, propagating its first failure."""

    completed: list[dict[str, Any]] = []
    for future in futures:
        if future in consumed or not future.done():
            continue
        completed.append(future.result())
        consumed.add(future)
    return completed


class LoadRunner:
    def __init__(
        self,
        *,
        config: LoadConfig,
        manifest: dict[str, Any],
        session: identity_gate.TimetablerAdminSession,
    ):
        self.config = config
        self.manifest = manifest
        self.session = session
        self.fixtures = (
            _resolve_cleanup_selectors(manifest, config.run_id)
            if config.action == "cleanup"
            else _resolve_selectors(manifest)
        )
        self.variables: dict[str, Any] = {
            "run_id": config.run_id,
            **manifest.get("variables", {}),
            **self.fixtures,
        }
        self.source_samples: list[dict[str, Any]] = []
        self.host_samples: list[dict[str, Any]] = []
        self.transactions: list[dict[str, Any]] = []
        self.api_measurements: list[dict[str, Any]] = []
        self.api_failures: list[dict[str, Any]] = []
        self._api_failures_lock = threading.Lock()
        self.stage_assertions: dict[str, Any] = {}
        self.abort_reason: str | None = None
        self.restart_fence: dict[str, Any] | None = None

    def _submit_step(
        self,
        stage: str,
        iteration: int,
        step: dict[str, Any],
        *,
        shard_index: int | None = None,
        cycle: int | None = None,
        scheduled_at_epoch: float | None = None,
    ) -> dict[str, Any]:
        request_id = _rate_request_id(self.config.run_id, stage, iteration)
        values = {
            **self.variables,
            "iteration": iteration,
            "cycle": cycle if cycle is not None else iteration - 1,
            "lifecycle_status": 2 if (cycle or 0) % 2 else 1,
            "request_id": request_id,
        }
        selector_name = step.get("fixture_selector")
        fixture_variable = step.get("fixture_variable")
        if selector_name and fixture_variable:
            pool = self.fixtures[selector_name]
            if not isinstance(pool, list) or shard_index is None:
                raise LoadHarnessError(
                    "realistic traffic fixture pool is not a bounded array"
                )
            values[fixture_variable] = pool[shard_index % len(pool)]
        payload = _render(step["payload"], values)
        classification = step.get("workload_classification", stage)
        started_wall = time.time()
        started = time.monotonic()
        try:
            self.session.post(step["path"], payload, request_id=request_id)
        except identity_gate.EntryGateError as error:
            attributed = error.with_safe_details(
                stage=stage,
                iteration=iteration,
                request_id=request_id,
                path=step["path"],
                workload_classification=classification,
                fixture_selector=selector_name,
            )
            with self._api_failures_lock:
                self.api_failures.append(attributed.safe_details)
            raise attributed from error
        elapsed_ms = (time.monotonic() - started) * 1000
        engine_result: dict[str, Any] | None = None
        if step["path"] == "/api/admin/schedule-request":
            activity_ids = _safe_int_list(payload.get("activity_ids"))
            if activity_ids is None or len(activity_ids) != 1:
                raise LoadHarnessError(
                    "engine-aware rate request must target exactly one activity"
                )
            engine_result = _wait_for_engine_applied_transaction(
                correlation_id=request_id,
                expected_activity_id=activity_ids[0],
                timeout_seconds=int(step.get("timeout_seconds", 120)),
            )
        recurring_assignment: dict[str, Any] | None = None
        if classification in REQUIRED_TRAFFIC_CLASSIFICATIONS:
            activity_ids = _safe_int_list(payload.get("activity_ids"))
            if activity_ids is None or len(activity_ids) != 1:
                raise LoadHarnessError(
                    "recurring assignment rate step must target exactly one activity"
                )
            recurring_assignment = _recurring_assignment_state(
                manifest=self.manifest,
                activity_id=activity_ids[0],
                expect_scheduled=classification == "schedule",
            )
        completed_ms = (time.monotonic() - started) * 1000
        return {
            "request_id": request_id,
            "source_request_id": (
                engine_result["source_request_id"] if engine_result else request_id
            ),
            "path": step["path"],
            "submitted_at_epoch": started_wall,
            "scheduled_at_epoch": scheduled_at_epoch,
            "submission_scheduling_skew_ms": round(
                max(0.0, started_wall - (scheduled_at_epoch or started_wall)) * 1000,
                3,
            ),
            "api_latency_ms": round(elapsed_ms, 3),
            # The full request/commit latency is a conservative upper bound on
            # database lock wait. A pass at <=500 ms therefore proves the
            # narrower lock wait cannot exceed the threshold.
            "database_lock_wait_upper_bound_ms": round(elapsed_ms, 3),
            "engine_aware": engine_result is not None,
            "engine_completion_ms": (
                round(completed_ms, 3) if engine_result else None
            ),
            "engine_delivery_id": (
                engine_result["engine_delivery_id"] if engine_result else None
            ),
            "engine_change_set_id": (
                engine_result["engine_change_set_id"] if engine_result else None
            ),
            "engine_source_sequence": (
                engine_result["engine_source_sequence"] if engine_result else None
            ),
            "expected_source_transactions": step["expected_source_transactions"],
            "workload_classification": classification,
            "recurring_assignment": recurring_assignment,
            "fixture_shard": shard_index,
            "traffic_cycle": cycle,
        }

    def run_rate(self, stage: str, *, rate: int, duration_seconds: int) -> None:
        workload = self.manifest["workloads"][stage]
        steps = _workload_steps(self.manifest, stage)
        plan = _rate_execution_plan(
            manifest=self.manifest,
            run_id=self.config.run_id,
            stage=stage,
            rate=rate,
            duration_seconds=duration_seconds,
            current_watermark=_current_watermark(),
        )
        total_approved_requests = plan["total_approved_requests"]
        expected_final = plan["approved_final"]
        resume_watermark = plan["resume_watermark"]
        prefix_transactions = plan["prefix_transactions"]
        prefix_count = plan["prefix_count"]
        active_request_count = plan["active_request_count"]
        active_duration_seconds = plan["active_duration_seconds"]
        if expected_final != resume_watermark + active_request_count:
            raise LoadHarnessError(
                "active rate segment does not terminate at the exact final fence"
            )
        sampler = HostSampler(
            interval_seconds=int(workload.get("sample_interval_seconds", 5)),
            hard_memory_percent=float(
                self.manifest["thresholds"]["hard_abort_memory_percent"]
            ),
            require_processes=True,
            publisher_restart_window=(
                (
                    self.restart_fence["execute_at_epoch"] - 1,
                    self.restart_fence["execute_at_epoch"] + 30,
                )
                if stage == "load-07" and self.restart_fence
                else None
            ),
        )
        sampler.start()
        futures: list[Future[dict[str, Any]]] = []
        consumed_futures: set[Future[dict[str, Any]]] = set()
        pool_sizes = {len(self.fixtures[step["fixture_selector"]]) for step in steps}
        lifecycle_pool_size = len(self.fixtures["traffic_activity_ids"])
        recurrence_pool_size = len(self.fixtures["traffic_recurrence_activity_ids"])
        provision = self.manifest["provision_plan"]
        lifecycle_minimum_seconds = (
            2 * int(provision["engine_operation_settlement_seconds"])
            + int(provision["engine_lifecycle_safety_margin_seconds"])
        )
        if (
            pool_sizes != {lifecycle_pool_size}
            or lifecycle_pool_size != recurrence_pool_size
            or lifecycle_pool_size * len(steps) / rate
            <= lifecycle_minimum_seconds
        ):
            raise LoadHarnessError(
                "realistic traffic does not have lifecycle-safe fixture sharding"
            )
        rate_shards = lifecycle_pool_size
        shard_executors = [
            ThreadPoolExecutor(max_workers=1) for _ in range(rate_shards)
        ]
        started = time.monotonic()
        started_wall = time.time()
        next_source_sample = started
        abort_pending = False
        try:
            for active_index in range(active_request_count):
                if sampler.abort_event.is_set():
                    raise LoadHarnessError(
                        "host hard-abort signal stopped load submission"
                    )
                iteration = prefix_count + active_index + 1
                due = started + active_index / rate
                due_epoch = started_wall + active_index / rate
                remaining = due - time.monotonic()
                if remaining > 0:
                    time.sleep(remaining)
                self.api_measurements.extend(
                    _consume_completed_futures(futures, consumed_futures)
                )
                step = steps[(iteration - 1) % len(steps)]
                cycle = (iteration - 1) // len(steps)
                shard_index = cycle % rate_shards
                futures.append(
                    shard_executors[shard_index].submit(
                        self._submit_step,
                        stage,
                        iteration,
                        step,
                        shard_index=shard_index,
                        cycle=cycle,
                        scheduled_at_epoch=due_epoch,
                    )
                )
                if time.monotonic() >= next_source_sample:
                    self.source_samples.append(_source_metrics())
                    next_source_sample = time.monotonic() + int(
                        workload.get("sample_interval_seconds", 5)
                    )
            for future in as_completed(
                [future for future in futures if future not in consumed_futures]
            ):
                self.api_measurements.append(future.result())
        except BaseException as error:
            abort_pending = True
            self.abort_reason = f"{type(error).__name__}: {error}"
            for future in futures:
                future.cancel()
            raise
        finally:
            _shutdown_executors(shard_executors, abort_pending=abort_pending)
            try:
                self.host_samples.extend(sampler.stop())
            except LoadHarnessError:
                if not abort_pending:
                    raise
        request_ids = {item["request_id"] for item in self.api_measurements}
        source_request_ids = {
            item["source_request_id"] for item in self.api_measurements
        }
        if (
            len(request_ids) != active_request_count
            or len(source_request_ids) != active_request_count
        ):
            raise LoadHarnessError("load API request attribution is incomplete")
        metrics, active_transactions = _wait_for_stage_publication(
            start=resume_watermark,
            expected_delta=active_request_count,
            request_ids=source_request_ids,
            timeout_seconds=self.manifest["load"]["drain_seconds"],
        )
        self.source_samples.append(metrics)
        active_measurements = {
            item["source_request_id"]: item
            for item in self.api_measurements
        }
        for transaction in active_transactions:
            measurement = active_measurements[transaction["request_id"]]
            transaction["workload_classification"] = measurement[
                "workload_classification"
            ]
            transaction["recurring_assignment"] = copy.deepcopy(
                measurement["recurring_assignment"]
            )
        transactions = [*prefix_transactions, *active_transactions]
        if len(transactions) != total_approved_requests:
            raise LoadHarnessError(
                "rate source evidence does not contain the full approved transaction count"
            )
        self.transactions.extend(transactions)
        latencies = [item["api_latency_ms"] for item in self.api_measurements]
        lock_bounds = [
            item["database_lock_wait_upper_bound_ms"] for item in self.api_measurements
        ]
        fidelity = self.manifest["rate_fidelity"]
        rate_metrics = _rate_fidelity_metrics(
            self.api_measurements, rate=rate, policy=fidelity
        )
        lock_percentiles = _percentiles(lock_bounds)
        if lock_percentiles["p95"] > self.manifest["thresholds"]["db_lock_wait_p95_ms"]:
            raise LoadHarnessError(
                "database lock-wait upper-bound p95 exceeds threshold"
            )
        sustained_threshold = self.manifest["thresholds"]["sustained_host_percent"]
        if any(
            _has_sustained_host_breach(
                self.host_samples, metric=metric, threshold=sustained_threshold
            )
            for metric in ("cpu_percent", "memory_percent")
        ):
            raise LoadHarnessError(
                "host CPU or memory exceeded 80 percent for 300 seconds"
            )
        publish_percentiles = _percentiles(
            [transaction["outbox_to_publish_ms"] for transaction in transactions]
        )
        if (
            publish_percentiles["p95"]
            > self.manifest["thresholds"]["end_to_end_p95_ms"]
            or publish_percentiles["p99"]
            > self.manifest["thresholds"]["end_to_end_p99_ms"]
        ):
            raise LoadHarnessError(
                "source outbox-to-publish latency exceeds end-to-end budget"
            )
        full_class_counts = {
            classification: sum(
                1
                for transaction in transactions
                if transaction["workload_classification"] == classification
            )
            for classification in REQUIRED_TRAFFIC_CLASSIFICATIONS
        }
        if full_class_counts != workload["expected_class_counts"]:
            raise LoadHarnessError(
                "full rate transaction classifications differ from the manifest"
            )
        active_class_counts = {
            classification: sum(
                1
                for item in self.api_measurements
                if item["workload_classification"] == classification
            )
            for classification in REQUIRED_TRAFFIC_CLASSIFICATIONS
        }
        recurring_states = [
            transaction.get("recurring_assignment") for transaction in transactions
        ]
        scope = self.manifest["load_acceptance_scope"]
        expected_term = {
            key: scope[key]
            for key in (
                "academic_term_id",
                "academic_term_start_date",
                "academic_term_end_date",
            )
        }
        for transaction, state in zip(transactions, recurring_states, strict=True):
            classification = transaction["workload_classification"]
            expected_transition = (
                "assigned" if classification == "schedule" else "released"
            )
            expected_assignments = 1 if classification == "schedule" else 0
            if (
                not isinstance(state, dict)
                or any(state.get(key) != value for key, value in expected_term.items())
                or state.get("recurrence_occurrence_count", 0)
                < scope["minimum_recurring_occurrences_per_activity"]
                or not SHA256_PATTERN.fullmatch(
                    str(state.get("activity_source_ref_sha256") or "")
                )
                or state.get("staff_assignment_count") != expected_assignments
                or state.get("location_assignment_count") != expected_assignments
                or state.get("staff_preset_count") != 1
                or state.get("location_preset_count") != 1
                or state.get("assignment_transition") != expected_transition
            ):
                raise LoadHarnessError(
                    "rate source evidence lacks a complete recurring assignment transition"
                )
        engine_measurements = [
            item for item in self.api_measurements if item["engine_aware"]
        ]
        expected_engine_count = sum(
            count
            for classification, count in workload["expected_class_counts"].items()
            if classification == "schedule"
        )
        if len(engine_measurements) != expected_engine_count:
            raise LoadHarnessError(
                "engine-aware terminal application count differs from the traffic mix"
            )
        self.stage_assertions.update(
            {
                "workload_rate_tps": rate,
                "workload_duration_seconds": active_duration_seconds,
                "approved_workload_duration_seconds": duration_seconds,
                "accepted_request_count": active_request_count,
                "total_approved_request_count": total_approved_requests,
                "resumed_prefix_transaction_count": prefix_count,
                "active_request_count": active_request_count,
                "exact_resume_watermark": resume_watermark,
                "prefix_verification": "exact contiguous finalized published durable outbox",
                **rate_metrics,
                "rate_fidelity_policy": fidelity,
                "per_class_counts": full_class_counts,
                "active_per_class_counts": active_class_counts,
                "measurement_scope": {
                    "api_latency": "active resumed segment only",
                    "database_lock_wait": "active resumed segment only",
                    "rate_fidelity": "active resumed segment only",
                    "host_metrics": "active resumed segment only",
                    "source_publish": "full approved profile including durable prefix",
                    "active_request_count": active_request_count,
                    "source_transaction_count": total_approved_requests,
                },
                "complete_staff_location_allocations": True,
                "deterministic_recurrence_in_mix": True,
                "load_acceptance_scope": copy.deepcopy(
                    self.manifest["load_acceptance_scope"]
                ),
                "all_transactions_term_26_schedule_or_unschedule": True,
                "schedule_assigns_recurring_staff_and_location": True,
                "unschedule_releases_recurring_staff_and_location": True,
                "per_fixture_dependency_order": "single-worker shard executors",
                "engine_completion_required_before_same_shard_dependency": True,
                "engine_terminal_applied_count": len(engine_measurements),
                "engine_completion_ms": _percentiles(
                    [item["engine_completion_ms"] for item in engine_measurements]
                ),
                "source_attribution": (
                    "client request id for synchronous APIs; post-commit delivery id "
                    "for atomic Scheduling Engine responses"
                ),
                "api_latency_ms": _percentiles(latencies),
                "outbox_to_publish_latency_ms": publish_percentiles,
                "database_lock_wait_upper_bound_ms": lock_percentiles,
                "database_lock_wait_measurement": "full supported-API commit latency; conservative upper bound",
                "drain_completed_within_seconds": self.manifest["load"][
                    "drain_seconds"
                ],
            }
        )

    def run_one(self, stage: str, step: dict[str, Any], *, iteration: int = 1) -> None:
        start = _current_watermark()
        measurement = self._submit_step(stage, iteration, step)
        self.api_measurements.append(measurement)
        expected = int(step["expected_source_transactions"])
        metrics, transactions = _wait_for_stage_publication(
            start=start,
            expected_delta=expected,
            request_ids={measurement["request_id"]} if expected else set(),
            timeout_seconds=int(step.get("timeout_seconds", 300)),
        )
        self.source_samples.append(metrics)
        self.transactions.extend(transactions)

    def run_load_04(self) -> None:
        workload = self.manifest["workloads"]["load-04"]
        selected = self.fixtures[workload["fixture_selector"]]
        if len(selected) != workload["approved_real_maximum_activity_count"]:
            raise LoadHarnessError(
                "LOAD-04 approved real maximum fixture count differs"
            )
        self.run_one("load-04", workload["steps"][0])
        if len(self.transactions) != 1:
            raise LoadHarnessError(
                "LOAD-04 did not produce one complete bulk transaction"
            )
        maximum = workload.get(
            "maximum_approved_serialized_bytes", TRANSPORT_LIMIT_BYTES - 1
        )
        actual = self.transactions[0]["canonical_bytes"]
        actual_members = self.transactions[0]["member_count"]
        if (
            actual > maximum
            or actual >= TRANSPORT_LIMIT_BYTES
            or actual_members != workload["approved_maximum_member_count"]
        ):
            raise LoadHarnessError(
                "LOAD-04 canonical bulk transaction exceeds its size budget"
            )
        self.stage_assertions = {
            "largest_real_approved_bulk": True,
            "single_atomic_transaction": True,
            "approved_maximum_member_count": workload["approved_maximum_member_count"],
            "approved_real_maximum_activity_count": workload[
                "approved_real_maximum_activity_count"
            ],
            "actual_member_count": actual_members,
            "maximum_approval_ref": workload["maximum_approval_ref"],
            "canonical_serialized_bytes": actual,
            "transport_limit_bytes": TRANSPORT_LIMIT_BYTES,
        }

    def run_load_05(self) -> None:
        workload = self.manifest["workloads"]["load-05"]
        if workload.get("maximum_recurrence_defined") is not True:
            raise LoadHarnessError("LOAD-05 maximum recurrence is not jointly defined")
        self.run_one("load-05", workload["steps"][0])
        occurrences = 0
        for transaction in self.transactions:
            from api.models import IntegrationOutbox

            members = IntegrationOutbox.objects.filter(
                transport_sequence=transaction["source_sequence"]
            ).values_list("payload", flat=True)
            for payload in members:
                state = (
                    payload.get("committed_state", {})
                    .get("replacement", {})
                    .get("current")
                )
                if isinstance(state, dict):
                    occurrences += len(state.get("occurrences", []))
        expected = workload.get("expected_occurrence_count")
        if occurrences != expected:
            raise LoadHarnessError("LOAD-05 canonical occurrence count is incomplete")
        self.stage_assertions = {
            "maximum_recurrence_defined": True,
            "recurrence_contract_ref": workload.get("recurrence_contract_ref"),
            "canonical_occurrence_count": occurrences,
        }

    def run_load_06(self) -> None:
        workload = self.manifest["workloads"]["load-06"]
        below, over = workload["steps"]
        candidates = self.fixtures[workload["rollback_selector"]]
        if not isinstance(candidates, list):
            raise LoadHarnessError(
                "LOAD-06 boundary selector must be a candidate array"
            )
        slot = int(self.variables[workload["slot_variable"]])
        over_slot = int(self.variables[workload["over_slot_variable"]])
        under_request = f"{self.config.run_id}-load-06-000001"
        under_calibration, _ = _find_boundary_calibration(
            candidate_ids=candidates,
            slot=slot,
            request_id=under_request,
            actor_id=int(self.session.user_id or 0),
            tolerance=workload["below_boundary_tolerance_bytes"],
        )
        request_id = f"{self.config.run_id}-load-06-000002"
        _, over_calibration = _find_boundary_calibration(
            candidate_ids=candidates,
            slot=over_slot,
            request_id=request_id,
            actor_id=int(self.session.user_id or 0),
            tolerance=workload["below_boundary_tolerance_bytes"],
            sequence_offset=2,
            prior_scheduled_activity_ids=set(under_calibration["activity_ids"]),
            prior_slot=slot,
        )
        calibrated_at = datetime.datetime.now(datetime.UTC).isoformat()
        below = copy.deepcopy(below)
        below["payload"]["activity_ids"] = under_calibration["activity_ids"]
        self.run_one("load-06", below, iteration=1)
        if len(self.transactions) != 1:
            raise LoadHarnessError(
                "LOAD-06 admissible case did not commit exactly once"
            )
        below_size = self.transactions[0]["canonical_bytes"]
        if (
            below_size != under_calibration["canonical_bytes"]
            or below_size >= TRANSPORT_LIMIT_BYTES
            or TRANSPORT_LIMIT_BYTES - below_size
            > workload["below_boundary_tolerance_bytes"]
        ):
            raise LoadHarnessError(
                "LOAD-06 admissible case does not match near-boundary calibration"
            )

        before = _assert_source_safety()["transport_watermark"]
        state_before = _selector_state_hash(
            self.manifest, workload["rollback_selector"]
        )
        over = copy.deepcopy(over)
        over["payload"]["activity_ids"] = over_calibration["activity_ids"]
        over["payload"]["slot"] = over_slot
        payload = _render(
            over["payload"],
            {**self.variables, "iteration": 2, "request_id": request_id},
        )
        rejected = False
        rejection = ""
        try:
            self.session.post(over["path"], payload, request_id=request_id)
        except identity_gate.EntryGateError as error:
            rejected = True
            rejection = str(error)
        if not rejected:
            raise LoadHarnessError("LOAD-06 over-boundary case was accepted")
        time.sleep(int(over.get("settle_seconds", 10)))
        after = _assert_source_safety(expected_watermark=before)["transport_watermark"]
        state_after = _selector_state_hash(self.manifest, workload["rollback_selector"])
        if state_after != state_before:
            raise LoadHarnessError(
                "LOAD-06 over-boundary rejection left partial domain state"
            )
        self.stage_assertions = {
            "transport_limit_bytes": TRANSPORT_LIMIT_BYTES,
            "boundary_approval_ref": workload["boundary_approval_ref"],
            "below_boundary_tolerance_bytes": workload[
                "below_boundary_tolerance_bytes"
            ],
            "read_only_calibration": {
                "completed_before_any_mutation_at": calibrated_at,
                "under": under_calibration,
                "over": over_calibration,
            },
            "admissible_serialized_bytes": below_size,
            "admissible_applied_once": True,
            "over_boundary_rejected": True,
            "over_boundary_rejection_class": rejection[:160],
            "over_boundary_source_watermark_unchanged": after == before,
            "over_boundary_domain_state_unchanged": state_after == state_before,
            "partial_publication": False,
        }

    def run_load_07(self) -> None:
        if self.restart_fence is None:
            raise LoadHarnessError("LOAD-07 verified short-lived fence is absent")
        fence = self.restart_fence
        lead = fence["execute_at_epoch"] - time.time()
        if not 60 <= lead <= 900:
            raise LoadHarnessError(
                "LOAD-07 execute_at_epoch must be 60 to 900 seconds ahead"
            )

        operation_result: dict[str, Any] = {}
        operation_error: list[Exception] = []

        def restart_at_fence() -> None:
            try:
                remaining = fence["execute_at_epoch"] - time.time()
                if remaining > 0:
                    time.sleep(remaining)
                env = {
                    **os.environ,
                    "TT_PHASE1_LOAD_RUN_ID": self.config.run_id,
                    "TT_PHASE1_LOAD_COORDINATION_SHA256": fence["coordination_sha256"],
                    "TT_PHASE1_LOAD_EXECUTE_AT_EPOCH": str(fence["execute_at_epoch"]),
                }
                result = subprocess.run(
                    ["bash", str(SCRIPT_DIR / "phase1_load_operations.sh")],
                    check=True,
                    capture_output=True,
                    text=True,
                    timeout=180,
                    env=env,
                )
                parsed = json.loads(result.stdout.strip().splitlines()[-1])
                operation_result.update(parsed)
            except (
                Exception
            ) as error:  # preserve exact fail-closed outcome for the main thread
                operation_error.append(error)

        thread = threading.Thread(target=restart_at_fence, daemon=True)
        thread.start()
        self.run_rate(
            "load-07",
            rate=int(self.manifest["load"]["p_tps"]),
            duration_seconds=int(self.manifest["load"]["restart_seconds"]),
        )
        thread.join(timeout=240)
        if thread.is_alive() or operation_error:
            raise LoadHarnessError(
                "LOAD-07 Timetabler publisher restart did not complete safely"
            )
        operation_valid = (
            operation_result.get("operation") == "publisher_restart"
            and operation_result.get("run_id") == self.config.run_id
            and operation_result.get("coordination_sha256")
            == fence["coordination_sha256"]
            and operation_result.get("execute_at_epoch") == fence["execute_at_epoch"]
            and isinstance(
                operation_result.get("actual_execute_at_epoch"), (int, float)
            )
            and isinstance(
                operation_result.get("execute_at_skew_seconds"), (int, float)
            )
            and 0 <= operation_result["execute_at_skew_seconds"] <= 5
            and operation_result.get("publisher_count") == 1
            and operation_result.get("publisher_restart_count") == 1
            and operation_result.get("publisher_liveness") is True
            and operation_result.get("durable_state_deleted") is False
            and operation_result.get("reverse_delivery_enabled") is False
            and operation_result.get("resource_booking_process_touched") is False
        )
        if not operation_valid:
            raise LoadHarnessError(
                "LOAD-07 publisher restart evidence is contradictory"
            )
        self.stage_assertions["publisher_restart"] = operation_result
        self.stage_assertions["coordination_fence"] = {
            "execute_at_epoch": fence["execute_at_epoch"],
            "coordination_sha256": fence["coordination_sha256"],
            "resource_booking_consumer_restart_owned_externally": True,
        }

    def run_cleanup(self) -> None:
        workload = self.manifest["workloads"]["cleanup"]
        from api.models import IntegrationOutbox, TtActivity

        approved_start = self.manifest["expected_stage_start_watermarks"]["cleanup"]
        activity_ids = self.fixtures[workload["activity_selector"]]
        batch_max = workload["activity_batch_max_members"]
        batches = [
            activity_ids[index : index + batch_max]
            for index in range(0, len(activity_ids), batch_max)
        ]
        if len(batches) != workload["expected_activity_batch_count"]:
            raise LoadHarnessError(
                "cleanup activity batch count differs from the manifest"
            )

        def completed_request(request_id: str) -> bool:
            sequences = list(
                IntegrationOutbox.objects.filter(request_id=request_id)
                .values_list("transport_sequence", flat=True)
                .distinct()[:2]
            )
            if len(sequences) > 1:
                raise LoadHarnessError(
                    "cleanup request id maps to multiple source transactions"
                )
            return len(sequences) == 1

        calibrations = []
        pending: list[tuple[int, list[int]]] = []
        for index, batch in enumerate(batches, start=1):
            request_id = f"{self.config.run_id}-cleanup-{index:06d}"
            present = set(
                TtActivity.objects.filter(id__in=batch).values_list("id", flat=True)
            )
            if completed_request(request_id):
                if present:
                    raise LoadHarnessError(
                        "completed cleanup batch still has domain rows"
                    )
                prior = _transaction_evidence(
                    approved_start + index - 1,
                    approved_start + index,
                    {request_id},
                )[0]
                calibration = {
                    "member_count": prior["member_count"],
                    "canonical_bytes": prior["canonical_bytes"],
                    "read_only": True,
                    "retry_reconstructed_from_immutable_transaction": True,
                }
            else:
                if present != set(batch):
                    raise LoadHarnessError(
                        "cleanup batch is partially absent without its transaction"
                    )
                calibration = _calibrate_activity_deletion(
                    activity_ids=batch,
                    request_id=request_id,
                    actor_id=int(self.session.user_id or 0),
                    sequence_offset=index,
                )
                pending.append((index, batch))
            if (
                calibration["member_count"] > batch_max
                or calibration["canonical_bytes"]
                >= TRANSPORT_LIMIT_BYTES - workload["transport_margin_bytes"]
            ):
                raise LoadHarnessError(
                    "cleanup batch exceeds its approved transport margin"
                )
            calibrations.append({"batch_index": index, **calibration})

        activity_template = copy.deepcopy(workload["steps"][0])
        for index, batch in pending:
            step = copy.deepcopy(activity_template)
            step["payload"]["id"] = batch
            self.run_one("cleanup", step, iteration=index)

        resource_steps = workload["steps"][1:]
        for offset, step in enumerate(resource_steps, start=len(batches) + 1):
            request_id = f"{self.config.run_id}-cleanup-{offset:06d}"
            selector_name = step["fixture_selector"]
            from api import models

            selector = self.manifest["fixture_selectors"][selector_name]
            model = getattr(models, MODEL_MAP[selector["model"]])
            present = set(
                model.objects.filter(id__in=self.fixtures[selector_name]).values_list(
                    "id", flat=True
                )
            )
            if completed_request(request_id):
                if present:
                    raise LoadHarnessError(
                        "completed resource cleanup still has domain rows"
                    )
                continue
            if present != set(self.fixtures[selector_name]):
                raise LoadHarnessError(
                    "resource cleanup is partially absent without its transaction"
                )
            self.run_one("cleanup", step, iteration=offset)

        all_request_ids = {
            f"{self.config.run_id}-cleanup-{index:06d}"
            for index in range(1, workload["expected_source_transactions"] + 1)
        }
        final_metrics, complete_transactions = _wait_for_stage_publication(
            start=approved_start,
            expected_delta=workload["expected_source_transactions"],
            request_ids=all_request_ids,
            timeout_seconds=self.manifest["profile_timing"]["CLEANUP"]["drain_seconds"],
        )
        for index, transaction in enumerate(complete_transactions, start=1):
            transaction["workload_classification"] = (
                "activity_cleanup" if index <= len(batches) else "resource_cleanup"
            )
        self.transactions = complete_transactions
        self.source_samples.append(final_metrics)
        remaining = {}
        from api import models

        for name in workload.get("absence_selectors", []):
            selector = self.manifest["fixture_selectors"][name]
            model = getattr(models, MODEL_MAP[selector["model"]])
            remaining[name] = list(
                model.objects.filter(**selector["filters"])
                .order_by("id")
                .values_list("id", flat=True)[: MAX_SELECTOR_ROWS + 1]
            )
        if any(remaining.values()):
            raise LoadHarnessError("run-owned fixtures remain after supported cleanup")
        self.stage_assertions = {
            "supported_timetabler_apis_only": True,
            "direct_database_mutation": False,
            "retry_safe_stable_request_ids": True,
            "activity_batch_max_members": batch_max,
            "activity_batch_count": len(batches),
            "expected_source_transactions": workload["expected_source_transactions"],
            "transport_margin_bytes": workload["transport_margin_bytes"],
            "read_only_pre_mutation_calibration": calibrations,
            "resource_booking_touched": False,
            "run_owned_fixtures_absent": True,
            "preexisting_term_template_preserved": True,
            "remaining_fixture_ids": remaining,
            "preexisting_durable_outbox_preserved": True,
        }


def _entry_evidence(
    config: LoadConfig,
    manifest: dict[str, Any],
    predecessor: dict[str, Any],
) -> dict[str, Any]:
    expected = manifest["expected_stage_start_watermarks"][config.action]
    metrics = _assert_source_safety(
        expected_watermark=(
            None if config.action == "cleanup" else expected
        )
    )
    if config.action == "cleanup" and not (
        expected
        <= metrics["transport_watermark"]
        <= manifest["expected_stage_final_watermarks"]["cleanup"]
    ):
        raise LoadHarnessError("cleanup retry watermark is outside its approved fences")
    entry = {
        "source": metrics,
        "configuration_snapshot_sha256": manifest["configuration_snapshot"]["sha256"],
        "commits": manifest["commits"],
        "threshold_approval": manifest["threshold_approval"],
        "operators": manifest["operators"],
        "process_state": _process_counts(),
        "predecessor_evidence": predecessor,
    }
    if config.action == "load-02":
        entry["load02_resume"] = {
            "approved_entry_watermark": expected,
            "resume_watermark": metrics["transport_watermark"],
            "resumed_prefix_transaction_count": 0,
            "prefix_verified": False,
            "engine_aware_non_resumable": True,
        }
    return entry


def _baseline_evidence(manifest: dict[str, Any]) -> dict[str, Any]:
    baseline = manifest["accepted_baseline"]
    return {
        **baseline,
        "baseline_attested_only": True,
        "current_source_domain_rows_not_substituted_for_fixed_baseline": True,
        "historical_projection_milestone": manifest["historical_projection_milestone"],
        "e2e_final_evidence_sha256": manifest.get("e2e_final_evidence_sha256"),
        "current_unfiltered_receiver_attestation": manifest[
            "load01_receiver_attestation"
        ],
    }


def run_provision(config: LoadConfig) -> dict[str, Any]:
    """Dispatch the guarded API-only provision stage through its focused runner."""
    import phase1_load_provision

    try:
        return phase1_load_provision.execute(config)
    except RuntimeError as error:
        raise LoadHarnessError(str(error)) from error


def execute(config: LoadConfig) -> dict[str, Any]:
    if config.action == "provision":
        return run_provision(config)
    identity_gate.bootstrap_django_runtime()
    manifest = load_manifest(
        config.manifest_json,
        run_id=config.run_id,
        environment=config.environment,
        source_scope=config.source_scope,
        require_approval=True,
    )
    assert_evidence_exchange_runtime(manifest, run_id=config.run_id)
    if config.expected_sha != manifest["commits"]["timetabler"]:
        raise LoadHarnessError("deployed Timetabler SHA differs from approved manifest")
    if not GIT_SHA_PATTERN.fullmatch(config.expected_sha):
        raise LoadHarnessError("exact deployed Timetabler SHA is invalid")
    if not config.github_run_id.isdigit():
        raise LoadHarnessError("immutable GitHub workflow run id is invalid")
    predecessor = verify_predecessor_evidence(
        run_id=config.run_id,
        stage=config.action,
        manifest=manifest,
        supplied_sha256=config.prior_evidence_sha256,
    )
    restart_fence = None
    if config.action == "load-07":
        restart_fence = verify_restart_fence(
            run_id=config.run_id,
            manifest=manifest,
            supplied_sha256=config.restart_fence_sha256,
        )
    elif config.restart_fence_sha256:
        raise LoadHarnessError("non-LOAD-07 stage received a restart fence digest")
    assert_approved_window(manifest, config.action.upper())
    if config.action in MUTATING_STAGES and not config.mutation_confirmed:
        raise LoadHarnessError(
            "mutating load stage lacks explicit workflow confirmation"
        )
    if config.action not in MUTATING_STAGES and config.mutation_confirmed:
        raise LoadHarnessError("read-only load stage received a mutation confirmation")
    if config.action in OPERATIONAL_STAGES and not config.operational_confirmed:
        raise LoadHarnessError("LOAD-07 lacks explicit operational confirmation")
    if config.action not in OPERATIONAL_STAGES and config.operational_confirmed:
        raise LoadHarnessError(
            "non-operational stage received operational confirmation"
        )

    entry = _entry_evidence(config, manifest, predecessor)
    session = identity_gate.TimetablerAdminSession(
        base_url=identity_gate.validate_base_url(
            config.base_url, "Timetabler base URL"
        ),
        email=config.email,
        password=config.password,
        run_id=config.run_id,
    )
    session.login_and_attest()
    runner: LoadRunner | None = None
    try:
        if config.action in {"entry", "load-01", "final"}:
            fixtures = (
                _resolve_selectors(manifest, include_expected_absent=False)
                if config.action == "final"
                else {}
            )
            stage_assertions: dict[str, Any] = {
                "read_only": True,
                "fixture_count": sum(
                    len(value) if isinstance(value, list) else 1
                    for value in fixtures.values()
                ),
            }
            if config.action == "load-01":
                stage_assertions["accepted_baseline"] = _baseline_evidence(manifest)
            if config.action == "final":
                cleanup = manifest["workloads"]["cleanup"]
                from api import models

                remaining = {}
                for name in cleanup.get("absence_selectors", []):
                    selector = manifest["fixture_selectors"][name]
                    model = getattr(models, MODEL_MAP[selector["model"]])
                    remaining[name] = list(
                        model.objects.filter(**selector["filters"])
                        .order_by("id")
                        .values_list("id", flat=True)[: MAX_SELECTOR_ROWS + 1]
                    )
                if any(remaining.values()):
                    raise LoadHarnessError(
                        "final read-only gate found run-owned fixtures"
                    )
                stage_assertions.update(
                    {
                        "run_owned_fixtures_absent": True,
                        "unfiltered_cross_service_reconciliation_required": True,
                        "resource_booking_reconciliation_claimed_by_timetabler": False,
                        "durable_outbox_preserved": True,
                    }
                )
            transactions: list[dict[str, Any]] = []
            api_measurements: list[dict[str, Any]] = []
            source_samples = [entry["source"]]
            host_samples: list[dict[str, Any]] = []
        else:
            runner = LoadRunner(config=config, manifest=manifest, session=session)
            runner.restart_fence = restart_fence
            try:
                if config.action == "load-02":
                    runner.run_rate(
                        "load-02",
                        rate=int(manifest["load"]["p_tps"]),
                        duration_seconds=int(manifest["load"]["steady_seconds"]),
                    )
                elif config.action == "load-03":
                    runner.run_rate(
                        "load-03",
                        rate=int(manifest["load"]["burst_tps"]),
                        duration_seconds=int(manifest["load"]["burst_seconds"]),
                    )
                elif config.action == "load-04":
                    runner.run_load_04()
                elif config.action == "load-05":
                    runner.run_load_05()
                elif config.action == "load-06":
                    runner.run_load_06()
                elif config.action == "load-07":
                    runner.run_load_07()
                elif config.action == "load-08":
                    runner.run_rate(
                        "load-08",
                        rate=int(manifest["load"]["p_tps"]),
                        duration_seconds=int(manifest["load"]["stability_seconds"]),
                    )
                elif config.action == "cleanup":
                    runner.run_cleanup()
                else:
                    raise LoadHarnessError("unsupported mutating load stage")
            except (LoadHarnessError, identity_gate.EntryGateError) as error:
                abort = persist_abort_evidence(
                    config=config, runner=runner, reason=str(error)
                )
                raise LoadHarnessError(
                    f"{error}; abort evidence sha256={abort['abort_evidence_sha256']}"
                ) from error
            stage_assertions = runner.stage_assertions
            transactions = runner.transactions
            api_measurements = runner.api_measurements
            source_samples = runner.source_samples
            host_samples = runner.host_samples
    finally:
        session.logout()

    final_metrics = _assert_source_safety(
        expected_watermark=manifest["expected_stage_final_watermarks"][config.action]
    )
    for transaction in transactions:
        transaction.setdefault("workload_classification", config.action.upper())
    return seal_evidence(
        {
            "schema_version": 1,
            "profile": config.action.upper(),
            "service": "timetabler",
            "repository": "Mayvins/timetabler-be",
            "stage": config.action,
            "run_id": config.run_id,
            "environment": config.environment,
            "production": False,
            "git_sha": config.expected_sha,
            "resource_booking_git_sha_attestation": manifest["commits"][
                "resource_booking"
            ],
            "github_run_id": config.github_run_id,
            "source_scope": config.source_scope,
            "entry_watermark": manifest["expected_stage_start_watermarks"][
                config.action
            ],
            "resume_watermark": entry["source"]["transport_watermark"],
            "final_watermark": final_metrics["transport_watermark"],
            "prior_evidence_sha256": config.prior_evidence_sha256,
            "predecessor_evidence": predecessor,
            "configuration_snapshot_sha256": manifest["configuration_snapshot"][
                "sha256"
            ],
            "evidence_exchange": manifest["evidence_exchange"],
            "workload_parameters": manifest["load"],
            "profile_timing": manifest["profile_timing"][config.action.upper()],
            "approved_window": manifest["approved_window"],
            "thresholds": manifest["thresholds"],
            "entry_evidence": entry,
            "stage_assertions": stage_assertions,
            "source_transactions": transactions,
            "api_measurements": api_measurements,
            "source_metric_samples": source_samples,
            "host_metric_samples": host_samples,
            "final_source_metrics": final_metrics,
            "load_executed": config.action in LOAD_EXECUTION_STAGES,
            "mutation_executed": config.action in MUTATING_STAGES,
            "restart_executed": config.action == "load-07",
            "reverse_delivery_enabled": False,
            "phase2_enabled": False,
            "dr_executed": False,
            "resource_booking_database_touched": False,
            "resource_booking_process_touched": False,
            "direct_database_mutation": False,
            "manual_outbox_edit": False,
            "unfiltered_resource_booking_reconciliation_required": config.action
            in {"cleanup", "final"},
        }
    )


def main() -> int:
    parser = argparse.ArgumentParser()
    parser.add_argument(
        "action",
        choices=(
            "validate",
            "attest",
            "correct-provision-evidence",
            "diagnose-load-02-prefix",
            LOAD02_ENGINE_AUDIT_ACTION,
            CURRENT_LOAD02_ENGINE_AUDIT_ACTION,
            ABANDONED_LOAD02_CLEANUP_ACTION,
            CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION,
            "checkpoint-load-02",
            "arm-fence",
            *STAGES,
        ),
    )
    parser.add_argument("--manifest")
    parser.add_argument("--run-id")
    args = parser.parse_args()
    try:
        if args.action == "attest":
            identity_gate.bootstrap_django_runtime()
            print(json.dumps(read_only_configuration_attestation(), sort_keys=True))
            return 0
        if args.action == "correct-provision-evidence":
            if (
                os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip()
                != "correct-provision-evidence"
            ):
                raise LoadHarnessError(
                    "correct-provision-evidence CLI action differs from guarded workflow action"
                )
            identity_gate.bootstrap_django_runtime()
            print(
                json.dumps(
                    correct_provision_evidence_from_environment(), sort_keys=True
                )
            )
            return 0
        if args.action == "diagnose-load-02-prefix":
            identity_gate.bootstrap_django_runtime()
            print(json.dumps(diagnose_load02_prefix_from_environment(), sort_keys=True))
            return 0
        if args.action in {
            LOAD02_ENGINE_AUDIT_ACTION,
            CURRENT_LOAD02_ENGINE_AUDIT_ACTION,
        }:
            identity_gate.bootstrap_django_runtime()
            print(
                json.dumps(
                    audit_failed_load02_engine_work_from_environment(),
                    sort_keys=True,
                )
            )
            return 0
        if args.action in {
            ABANDONED_LOAD02_CLEANUP_ACTION,
            CURRENT_ABANDONED_LOAD02_CLEANUP_ACTION,
        }:
            identity_gate.bootstrap_django_runtime()
            print(
                json.dumps(
                    cleanup_abandoned_load02_from_environment(),
                    sort_keys=True,
                )
            )
            return 0
        if args.action == "checkpoint-load-02":
            identity_gate.bootstrap_django_runtime()
            print(json.dumps(checkpoint_load02_from_environment(), sort_keys=True))
            return 0
        if args.action == "arm-fence":
            if os.environ.get("TT_PHASE1_LOAD_ACTION", "").strip() != "arm-fence":
                raise LoadHarnessError(
                    "arm-fence CLI action differs from guarded workflow action"
                )
            if _env_true("TT_PHASE1_LOAD_CONFIRM_MUTATION") or _env_true(
                "TT_PHASE1_LOAD_CONFIRM_OPERATIONAL"
            ):
                raise LoadHarnessError(
                    "arm-fence must not receive mutation or operational confirmation"
                )
            run_id = os.environ.get("TT_PHASE1_LOAD_RUN_ID", "").strip()
            validate_run_id(run_id)
            raw = os.environ.get("TT_PHASE1_LOAD_MANIFEST_JSON", "")
            identity_gate.bootstrap_django_runtime()
            manifest = load_manifest(
                raw,
                run_id=run_id,
                environment=os.environ.get("TT_PHASE1_LOAD_ENVIRONMENT", "").strip(),
                source_scope=os.environ.get("TT_PHASE1_LOAD_SOURCE_SCOPE", "").strip(),
                require_approval=True,
            )
            exact_sha = os.environ.get("TT_PHASE1_LOAD_EXPECTED_SHA", "").strip()
            if exact_sha != manifest["commits"]["timetabler"]:
                raise LoadHarnessError(
                    "deployed Timetabler SHA differs from the approved manifest"
                )
            raw_epoch = os.environ.get("TT_PHASE1_LOAD_EXECUTE_AT_EPOCH", "").strip()
            if not raw_epoch.isdigit():
                raise LoadHarnessError("arm-fence execute_at_epoch is invalid")
            result = arm_restart_fence(
                run_id=run_id,
                manifest=manifest,
                supplied_predecessor_sha256=os.environ.get(
                    "TT_PHASE1_LOAD_PRIOR_EVIDENCE_SHA256", ""
                )
                .strip()
                .lower(),
                execute_at_epoch=int(raw_epoch),
                approval_ref=os.environ.get("TT_PHASE1_LOAD_FENCE_APPROVAL_REF", ""),
            )
            print(json.dumps(result, sort_keys=True))
            return 0
        if args.action == "validate":
            if not args.manifest or not args.run_id:
                raise LoadHarnessError(
                    "local validation requires --manifest and --run-id"
                )
            raw = Path(args.manifest).read_text()
            manifest = load_manifest(
                raw,
                run_id=args.run_id,
                environment="staging",
                source_scope="default",
                require_approval=False,
            )
            print(
                json.dumps(
                    {
                        "load_manifest": "valid",
                        "stage_count": len(STAGES),
                        "workload_count": len(manifest["workloads"]),
                        "load_dispatched": False,
                    },
                    sort_keys=True,
                )
            )
            return 0
        config = LoadConfig.from_environment()
        if config.action != args.action:
            raise LoadHarnessError("CLI action differs from guarded workflow action")
        evidence = execute(config)
        persisted = persist_evidence(evidence)
        print(
            json.dumps(
                {
                    "phase1_load": "evidence_persisted",
                    "profile": evidence["profile"],
                    "run_id": evidence["run_id"],
                    "entry_watermark": evidence["entry_watermark"],
                    "final_watermark": evidence["final_watermark"],
                    "load_executed": evidence["load_executed"],
                    "dr_executed": False,
                    "reverse_delivery_enabled": False,
                    **persisted,
                },
                sort_keys=True,
            )
        )
        return 0
    except (LoadHarnessError, identity_gate.EntryGateError) as error:
        print(
            json.dumps(
                {"phase1_load": "blocked", "reason": str(error)}, sort_keys=True
            ),
            file=sys.stderr,
        )
        return 3


if __name__ == "__main__":
    raise SystemExit(main())
