#!/usr/bin/env python3
"""Guarded Phase 1 E2E-01..13 executor for the deployed staging runtime.

All domain mutations go through authenticated application APIs. ORM access in
this module is read-only and is limited to fixture discovery and immutable
outbox evidence. E2E-09 deliberately republishes one existing canonical Kafka
record without changing its source row/version/sequence.
"""

from __future__ import annotations

import argparse
import hashlib
import json
import os
from pathlib import Path
import re
import sys
import time
from typing import Any


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 gate  # noqa: E402


CASE_STAGES = {
    "E2E-01": "functional",
    "E2E-02": "functional",
    "E2E-03": "functional",
    "E2E-04": "functional",
    "E2E-05": "functional",
    "E2E-06": "functional",
    "E2E-07": "functional",
    "E2E-08": "functional",
    "E2E-09": "redelivery",
    "E2E-10": "conflict",
    "E2E-11": "restart",
    "E2E-12": "rollback",
    "E2E-13": "final",
}
CASE_ALLOWED_PATHS = {
    "E2E-01": {
        "/api/admin/staff/create",
        "/api/admin/staff/update",
        "/api/admin/staff/delete",
        "/api/admin/location/create",
        "/api/admin/location/update",
        "/api/admin/location/delete",
    },
    "E2E-02": {"/api/admin/schedule-request"},
    "E2E-03": {"/api/admin/preschedule-request"},
    "E2E-04": {"/api/admin/schedule-request"},
    "E2E-05": {
        "/api/admin/booking/create",
        "/api/admin/booking/schedule",
    },
    "E2E-06": {"/api/admin/schedule-request"},
    "E2E-07": {
        "/api/admin/unschedule",
        "/api/admin/booking/unschedule",
        "/api/admin/booking/delete",
        "/api/admin/activity/delete",
        "/api/admin/variant-create",
    },
    "E2E-08": {"/api/admin/schedule-request"},
    "E2E-10": {"/api/admin/schedule-request"},
}
MODEL_MAP = {
    "academic_term": "TtAcademicTerm",
    "activity": "TtActivity",
    "activity_template": "TtActivityTemplate",
    "activity_type": "TtActivityType",
    "location": "TtLocation",
    "staff": "TtStaff",
    "week": "TtWeek",
    "week_pattern": "TtWeekPattern",
}
SELECTOR_FIELDS = {
    "id",
    "code",
    "name",
    "booking",
    "status",
    "academic_term_id",
    "activity_type_id",
    "activity_template_id",
    "is_booking",
    "code__startswith",
}
MAX_FIXTURE_IDS = 64
MIN_CONFLICT_LEAD_SECONDS = 30
MAX_CONFLICT_LEAD_SECONDS = 900
CLEANUP_PATHS = {
    "/api/admin/staff/delete",
    "/api/admin/location/delete",
    "/api/admin/unschedule",
    "/api/admin/booking/unschedule",
    "/api/admin/booking/delete",
    "/api/admin/activity/delete",
}
CLEANUP_FIXTURE_MODELS = {
    "all_fixture_activity_ids": "TtActivity",
    "all_fixture_staff_ids": "TtStaff",
    "all_fixture_location_ids": "TtLocation",
}


class HarnessError(gate.EntryGateError):
    pass


def _template_names(value: Any) -> set[str]:
    if isinstance(value, dict):
        return set().union(*(_template_names(item) for item in value.values()), set())
    if isinstance(value, list):
        return set().union(*(_template_names(item) for item in value), set())
    if isinstance(value, str):
        return set(re.findall(r"\{\{([a-zA-Z0-9_]+)}}", value))
    return set()


def validate_execution_manifest(manifest: dict[str, Any], run_id: str) -> None:
    if manifest.get("run_id") != run_id:
        raise HarnessError("private manifest run_id does not match the workflow run")
    if manifest.get("dedicated_fixture_prefix") != run_id:
        raise HarnessError(
            "private manifest must use the unique run id as its fixture prefix"
        )
    if (
        manifest.get("load_enabled") is not False
        or manifest.get("dr_enabled") is not False
    ):
        raise HarnessError("private manifest must explicitly disable LOAD and DR")
    if not isinstance(manifest.get("conflict_armed"), bool):
        raise HarnessError("private manifest must explicitly declare conflict_armed")
    if [case["id"] for case in manifest["cases"]] != list(CASE_STAGES):
        raise HarnessError("private manifest cases must use canonical E2E-01..13 order")
    variables = manifest.get("variables")
    if not isinstance(variables, dict) or len(variables) > MAX_FIXTURE_IDS:
        raise HarnessError("private manifest variables must be a bounded object")
    for name in variables:
        if not re.fullmatch(r"[a-zA-Z][a-zA-Z0-9_]{0,63}", name):
            raise HarnessError(f"private manifest variable name is invalid: {name}")
        if any(
            forbidden in name.lower()
            for forbidden in (
                "password",
                "token",
                "secret",
                "signature",
                "authorization",
            )
        ):
            raise HarnessError(f"private manifest variable is credential-like: {name}")

    cases = {case["id"]: case for case in manifest["cases"]}
    for case_id, stage in CASE_STAGES.items():
        case = cases[case_id]
        if case.get("stage") != stage:
            raise HarnessError(f"{case_id} must be assigned to stage {stage}")
        if stage in {"functional", "conflict"}:
            steps = case.get("steps")
            if not isinstance(steps, list) or not steps:
                raise HarnessError(f"{case_id} must define executable API steps")
            allowed_paths = CASE_ALLOWED_PATHS[case_id]
            seen_paths = set()
            for step in steps:
                if not isinstance(step, dict) or step.get("kind") != "api":
                    raise HarnessError(f"{case_id} contains a non-API domain mutation")
                path = step.get("path")
                target = step.get("target", "timetabler")
                if target != "timetabler":
                    raise HarnessError(
                        f"{case_id} uses forbidden non-Timetabler target {target}"
                    )
                if path not in allowed_paths:
                    raise HarnessError(
                        f"{case_id} uses unsupported Timetabler path {path}"
                    )
                seen_paths.add(path)
                if not isinstance(step.get("payload"), dict):
                    raise HarnessError(f"{case_id} API step payload must be an object")
                expected = step.get("expected_source_transactions")
                if expected not in (0, 1):
                    raise HarnessError(
                        f"{case_id} must declare zero or one logical source transaction per API step"
                    )
                capture = step.get("capture", {})
                if not isinstance(capture, dict):
                    raise HarnessError(f"{case_id} step capture must be an object")
                for name, selector in capture.items():
                    validate_selector(
                        name,
                        selector,
                        run_id=run_id,
                        require_run_tag=True,
                    )
            required = allowed_paths if case_id == "E2E-01" else set()
            if required and seen_paths != required:
                raise HarnessError(
                    f"{case_id} does not cover its complete lifecycle path set"
                )
        cleanup_steps = case.get("cleanup_steps", [])
        if not isinstance(cleanup_steps, list):
            raise HarnessError(f"{case_id} cleanup_steps must be an array")
        for step in cleanup_steps:
            if (
                not isinstance(step, dict)
                or step.get("kind") != "api"
                or step.get("target", "timetabler") != "timetabler"
                or step.get("path") not in CLEANUP_PATHS
                or not isinstance(step.get("payload"), dict)
                or step.get("expected_source_transactions") not in (0, 1)
            ):
                raise HarnessError(f"{case_id} contains an unsupported cleanup step")
        if case_id == "E2E-10":
            if len(case["steps"]) != 1:
                raise HarnessError("E2E-10 must contain exactly one Timetabler step")
            if not isinstance(case.get("execute_at_epoch"), int):
                raise HarnessError("E2E-10 requires an integer execute_at_epoch")
    selectors = manifest.get("reference_selectors")
    if not isinstance(selectors, dict) or not selectors:
        raise HarnessError("private manifest must define dedicated reference selectors")
    for name, selector in selectors.items():
        validate_selector(name, selector, run_id=run_id, require_run_tag=True)

    base_variables = {"run_id", *variables, *selectors}
    functional_variables = set(base_variables)
    for case in manifest["cases"]:
        if case["stage"] != "functional":
            continue
        for step in case["steps"]:
            missing = _template_names(step["payload"]) - functional_variables
            if missing:
                raise HarnessError(
                    f"{case['id']} has unresolved pre-mutation variables: {sorted(missing)}"
                )
            capture = step.get("capture", {})
            for selector in capture.values():
                missing = _template_names(selector) - functional_variables
                if missing:
                    raise HarnessError(
                        f"{case['id']} capture has unresolved variables: {sorted(missing)}"
                    )
            functional_variables.update(capture)

    conflict = cases["E2E-10"]
    missing = _template_names(conflict["steps"][0]["payload"]) - base_variables
    if missing:
        raise HarnessError(
            f"E2E-10 has unresolved pre-mutation variables: {sorted(missing)}"
        )
    for case in reversed(manifest["cases"]):
        for step in case.get("cleanup_steps", []):
            missing = _template_names(step["payload"]) - base_variables
            if missing:
                raise HarnessError(
                    f"{case['id']} cleanup has unresolved variables: {sorted(missing)}"
                )


def validate_selector(
    name: str,
    selector: Any,
    *,
    run_id: str,
    require_run_tag: bool,
) -> None:
    if not isinstance(selector, dict) or selector.get("model") not in MODEL_MAP:
        raise HarnessError(f"fixture selector {name} uses an unsupported model")
    filters = selector.get("filters")
    if (
        not isinstance(filters, dict)
        or not filters
        or not set(filters) <= SELECTOR_FIELDS
    ):
        raise HarnessError(f"fixture selector {name} uses unsupported filters")
    expected_count = selector.get("expected_count")
    if (
        not isinstance(expected_count, int)
        or expected_count < 1
        or expected_count > MAX_FIXTURE_IDS
    ):
        raise HarnessError(
            f"fixture selector {name} expected_count must be between 1 and {MAX_FIXTURE_IDS}"
        )
    if require_run_tag and not selector.get("dedicated_e2e"):
        raise HarnessError(f"fixture selector {name} is not marked dedicated_e2e")
    if require_run_tag:
        encoded = json.dumps(filters, sort_keys=True)
        if run_id not in encoded and "PHASE1-E2E" not in encoded.upper():
            raise HarnessError(f"fixture selector {name} lacks a dedicated E2E marker")


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 HarnessError(f"manifest variable is unresolved: {name}")
        return variables[name]
    rendered = value
    for name in re.findall(r"\{\{([a-zA-Z0-9_]+)}}", value):
        if name not in variables:
            raise HarnessError(f"manifest variable is unresolved: {name}")
        rendered = rendered.replace(f"{{{{{name}}}}}", str(variables[name]))
    return rendered


def resolve_selector(name: str, selector: dict[str, Any]) -> int | list[int]:
    from api import models

    model = getattr(models, MODEL_MAP[selector["model"]])
    expected_count = selector["expected_count"]
    rows = list(
        model.objects.filter(**selector["filters"])
        .order_by("id")
        .values_list("id", flat=True)[: expected_count + 1]
    )
    if len(rows) != expected_count:
        raise HarnessError(
            f"fixture selector {name} resolved {len(rows)} rows instead of {expected_count}"
        )
    values = [int(value) for value in rows]
    return values[0] if expected_count == 1 else values


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 _assert_local_safety() -> int:
    from api.models import (
        EngineResponseQuarantine,
        IntegrationOutbox,
        PostCommitDelivery,
    )
    from api.services.integration.readiness import phase1_readiness

    readiness = phase1_readiness(require_enabled=True)
    if not readiness["ready"] or readiness["reverse_delivery_enabled"]:
        raise HarnessError(
            "Timetabler Phase 1 readiness or reverse-delivery guard failed"
        )
    if IntegrationOutbox.objects.filter(status="dead_letter").exists():
        raise HarnessError("Timetabler outbox dead letter blocks E2E execution")
    if PostCommitDelivery.objects.filter(status="dead_letter").exists():
        raise HarnessError(
            "Timetabler post-commit delivery dead letter blocks E2E execution"
        )
    if EngineResponseQuarantine.objects.exists():
        raise HarnessError("Timetabler engine-response quarantine blocks E2E execution")
    return int(readiness["transport_watermark"])


def _wait_for_publisher_readiness(
    *, expected_watermark: int, timeout_seconds: int = 120
) -> dict[str, Any]:
    from api.models import (
        EngineResponseQuarantine,
        IntegrationOutbox,
        PostCommitDelivery,
    )
    from api.services.integration.readiness import phase1_readiness

    deadline = time.monotonic() + timeout_seconds
    while True:
        readiness = phase1_readiness(require_enabled=True)
        blocked = (
            not readiness["ready"]
            or readiness["reverse_delivery_enabled"]
            or readiness["dead_letter_count"] != 0
            or IntegrationOutbox.objects.filter(status="dead_letter").exists()
            or PostCommitDelivery.objects.filter(status="dead_letter").exists()
            or EngineResponseQuarantine.objects.exists()
        )
        if blocked:
            raise HarnessError(
                "Timetabler publisher readiness, dead-letter, quarantine, or "
                "reverse-delivery guard failed"
            )
        if readiness["transport_watermark"] != expected_watermark:
            raise HarnessError(
                "cleanup source watermark differs from its exact approved fence"
            )
        if (
            readiness["publisher_last_sequence"] == expected_watermark
            and readiness["publisher_live"]
        ):
            return {
                "ready": True,
                "publisher_live": True,
                "publisher_status": readiness["publisher_status"],
                "publisher_last_sequence": expected_watermark,
                "transport_watermark": expected_watermark,
                "dead_letter_count": 0,
                "post_commit_dead_letter_count": 0,
                "engine_response_quarantine_count": 0,
                "reverse_delivery_enabled": False,
            }
        if time.monotonic() >= deadline:
            raise HarnessError(
                "timed out waiting for publisher to reach cleanup watermark"
            )
        time.sleep(1)


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

    deliveries = {
        str(value)
        for value in PostCommitDelivery.objects.filter(
            correlation_id=request_id
        ).values_list("delivery_id", flat=True)
    }
    rows = list(
        IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(),
            transport_sequence__gt=start,
            transport_sequence__lte=end,
        ).order_by("transport_sequence", "transaction_index")
    )
    evidence = []
    for sequence in range(start + 1, end + 1):
        members = [row for row in rows if row.transport_sequence == sequence]
        if not members:
            raise HarnessError("source sequence gap appeared during E2E execution")
        count = members[0].transaction_count
        if (
            len(members) != count
            or [row.transaction_index for row in members] != list(range(1, count + 1))
            or not all(row.transaction_finalized for row in members)
            or len({row.change_set_id for row in members}) != 1
        ):
            raise HarnessError(
                "partial source transaction appeared during E2E execution"
            )
        attributable = all(
            row.correlation_id == request_id
            or row.request_id == request_id
            or row.request_id in deliveries
            for row in members
        )
        if not attributable:
            raise HarnessError(
                "unrelated source transaction overlapped the isolated E2E step"
            )
        if any(row.status == "dead_letter" for row in members):
            raise HarnessError(
                "source transaction entered dead letter during E2E execution"
            )
        evidence.append(
            {
                "source_sequence": sequence,
                "source_transaction_id": str(members[0].change_set_id),
                "transaction_count": count,
                "event_types": sorted({row.event_type for row in members}),
            }
        )
    return evidence


def _wait_for_transaction_delta(
    *, start: int, expected: int, request_id: str, timeout_seconds: int
) -> tuple[int, list[dict[str, Any]]]:
    deadline = time.monotonic() + timeout_seconds
    settle_deadline = time.monotonic() + min(10, timeout_seconds)
    while True:
        end = _current_watermark()
        delta = end - start
        if delta > expected:
            raise HarnessError(
                "unexpected source transaction count; isolated run aborted"
            )
        if delta == expected and (expected > 0 or time.monotonic() >= settle_deadline):
            return end, _transaction_evidence(start, end, request_id)
        if time.monotonic() >= deadline:
            raise HarnessError("timed out waiting for exact source transaction delta")
        time.sleep(1)


class ScenarioRunner:
    def __init__(
        self,
        *,
        config: gate.EntryConfig,
        manifest: dict[str, Any],
        session: gate.TimetablerAdminSession,
        resolve_fixtures: bool = True,
    ):
        self.config = config
        self.manifest = manifest
        self.session = session
        self.variables: dict[str, Any] = {
            "run_id": config.run_id,
            **manifest["variables"],
        }
        self.fixture_ids: dict[str, int | list[int]] = {}
        self.evidence: list[dict[str, Any]] = []
        self.cleanup_assertions: dict[str, Any] | None = None
        self.publisher_readiness: dict[str, Any] | None = None
        self.final_assertions: dict[str, Any] | None = None
        if resolve_fixtures:
            for name, selector in manifest["reference_selectors"].items():
                value = resolve_selector(name, selector)
                self.variables[name] = value
                self.fixture_ids[name] = value
        fixture_id_count = sum(
            len(value) if isinstance(value, list) else 1
            for value in self.fixture_ids.values()
        )
        if fixture_id_count > MAX_FIXTURE_IDS:
            raise HarnessError("private manifest exceeds the bounded fixture ID limit")

    def _capture(self, capture: dict[str, Any]) -> None:
        for name, raw_selector in capture.items():
            selector = _render(raw_selector, self.variables)
            validate_selector(
                name,
                selector,
                run_id=self.config.run_id,
                require_run_tag=True,
            )
            value = resolve_selector(name, selector)
            self.variables[name] = value
            self.fixture_ids[name] = value

    def _run_tt_step(self, case_id: str, index: int, step: dict[str, Any]) -> None:
        request_id = f"{self.config.run_id}-{case_id.lower()}-{index:02d}"
        payload = _render(step["payload"], self.variables)
        before = _assert_local_safety()
        self.session.post(step["path"], payload, request_id=request_id)
        if step.get("capture"):
            self._capture(step["capture"])
        after, transactions = _wait_for_transaction_delta(
            start=before,
            expected=step["expected_source_transactions"],
            request_id=request_id,
            timeout_seconds=min(120, int(step.get("timeout_seconds", 60))),
        )
        self.evidence.append(
            {
                "case": case_id,
                "step": step.get("name", str(index)),
                "before_watermark": before,
                "after_watermark": after,
                "transactions": transactions,
            }
        )

    def run_functional(self) -> None:
        for case in self.manifest["cases"]:
            if case["stage"] != "functional":
                continue
            for index, step in enumerate(case["steps"], start=1):
                self._run_tt_step(case["id"], index, step)

    def run_cleanup(self) -> None:
        from api import models
        from api.models import IntegrationOutbox
        from api.services.integration.outbox import integration_source_scope

        start_watermark = _current_watermark()
        expected_final_watermark = self.manifest["expected_stage_start_watermarks"][
            "final"
        ]
        expected_delta = sum(
            step["expected_source_transactions"]
            for case in self.manifest["cases"]
            for step in case.get("cleanup_steps", [])
        )
        if expected_final_watermark != start_watermark + expected_delta:
            raise HarnessError(
                "cleanup transaction count does not reach the approved final fence"
            )
        source_scope = integration_source_scope()
        durable_members_before = IntegrationOutbox.objects.filter(
            source_scope=source_scope,
            transport_sequence__lte=start_watermark,
        ).count()
        for case in reversed(self.manifest["cases"]):
            for index, step in enumerate(case.get("cleanup_steps", []), start=1):
                self._run_tt_step(case["id"], index, step)

        final_watermark = _current_watermark()
        if final_watermark != expected_final_watermark:
            raise HarnessError("cleanup did not reach the exact approved final fence")
        remaining = {}
        deleted_counts = {}
        for fixture_name, model_name in CLEANUP_FIXTURE_MODELS.items():
            raw_ids = self.fixture_ids[fixture_name]
            fixture_ids = raw_ids if isinstance(raw_ids, list) else [raw_ids]
            model = getattr(models, model_name)
            remaining[fixture_name] = list(
                model.objects.filter(id__in=fixture_ids)
                .order_by("id")
                .values_list("id", flat=True)
            )
            deleted_counts[fixture_name] = len(fixture_ids)
        if any(remaining.values()):
            raise HarnessError("run-owned fixture rows remain after supported cleanup")

        durable_members_after = IntegrationOutbox.objects.filter(
            source_scope=source_scope,
            transport_sequence__lte=start_watermark,
        ).count()
        if durable_members_after != durable_members_before:
            raise HarnessError("cleanup changed pre-existing durable outbox evidence")

        self.publisher_readiness = _wait_for_publisher_readiness(
            expected_watermark=expected_final_watermark
        )
        self.cleanup_assertions = {
            "supported_application_paths_only": True,
            "direct_database_mutation": False,
            "deleted_fixture_counts": deleted_counts,
            "remaining_fixture_ids": remaining,
            "pre_cleanup_durable_outbox_members": durable_members_before,
            "post_cleanup_preexisting_outbox_members": durable_members_after,
            "durable_state_deleted": False,
            "expected_source_transaction_count": expected_delta,
            "actual_source_transaction_count": final_watermark - start_watermark,
            "expected_final_watermark": expected_final_watermark,
        }

    def run_conflict(self) -> None:
        case = next(case for case in self.manifest["cases"] if case["id"] == "E2E-10")
        if self.manifest["conflict_armed"] is not True:
            raise HarnessError("E2E-10 manifest is not armed for coordinated execution")
        execute_at = case["execute_at_epoch"]
        now = time.time()
        lead_seconds = execute_at - now
        if not MIN_CONFLICT_LEAD_SECONDS <= lead_seconds <= MAX_CONFLICT_LEAD_SECONDS:
            raise HarnessError(
                "E2E-10 execute_at_epoch must be 30 to 900 seconds in the future"
            )
        step = case["steps"][0]
        request_id = f"{self.config.run_id}-e2e-10-tt"
        payload = _render(step["payload"], self.variables)
        while execute_at - time.time() > 1:
            time.sleep(min(1, execute_at - time.time() - 1))
        before = _assert_local_safety()
        remaining = execute_at - time.time()
        if remaining > 0:
            time.sleep(remaining)
        submitted_at = time.time()
        self.session.post(step["path"], payload, request_id=request_id)
        after, transactions = _wait_for_transaction_delta(
            start=before,
            expected=step["expected_source_transactions"],
            request_id=request_id,
            timeout_seconds=min(120, int(step.get("timeout_seconds", 120))),
        )
        self.evidence.append(
            {
                "case": "E2E-10",
                "step": step.get("name", "tt-contender"),
                "contender": "timetabler",
                "execute_at_epoch": execute_at,
                "submitted_at_epoch": submitted_at,
                "submission_skew_ms": round((submitted_at - execute_at) * 1000),
                "before_watermark": before,
                "after_watermark": after,
                "transactions": transactions,
                "aggregate_result_required": True,
            }
        )

    def run_redelivery(self) -> None:
        from api.models import IntegrationOutbox
        from api.services.integration.publisher import (
            KafkaOutboxTransport,
            TransportChangeSet,
        )

        case = next(case for case in self.manifest["cases"] if case["id"] == "E2E-09")
        correlation_prefix = case.get("correlation_prefix", self.config.run_id)
        source_sequence = (
            IntegrationOutbox.objects.filter(
                correlation_id__startswith=correlation_prefix,
                transaction_finalized=True,
                status="published",
            )
            .order_by("-transport_sequence")
            .values_list("transport_sequence", flat=True)
            .first()
        )
        if source_sequence is None:
            raise HarnessError(
                "E2E-09 could not resolve an immutable run-owned transaction"
            )
        rows = tuple(
            IntegrationOutbox.objects.filter(
                transport_sequence=source_sequence
            ).order_by("transaction_index")
        )
        before = _assert_local_safety()
        if before < source_sequence:
            raise HarnessError(
                "E2E-09 source transaction is outside the committed watermark"
            )
        KafkaOutboxTransport().publish(TransportChangeSet(rows))
        delay = min(30, max(1, int(case.get("delay_seconds", 5))))
        time.sleep(delay)
        after = _assert_local_safety()
        if after != before:
            raise HarnessError("E2E-09 redelivery invented a new source transaction")
        self.evidence.append(
            {
                "case": "E2E-09",
                "source_sequence": int(source_sequence),
                "source_transaction_id": str(rows[0].change_set_id),
                "before_watermark": before,
                "after_watermark": after,
                "new_source_transactions": 0,
            }
        )

    def run_final(self) -> None:
        from api.models import IntegrationOutbox
        from api.services.integration.outbox import integration_source_scope

        before = _assert_local_safety()
        expected_watermark = self.manifest["expected_stage_start_watermarks"][
            "final"
        ]
        if before != expected_watermark:
            raise HarnessError("final watermark differs from its exact approved fence")
        fixture_absence = gate.final_fixture_absence(
            self.manifest, self.config.run_id
        )
        self.publisher_readiness = _wait_for_publisher_readiness(
            expected_watermark=expected_watermark
        )

        cleanup_start = self.manifest["expected_stage_start_watermarks"]["cleanup"]
        cleanup_steps = [
            step
            for case in self.manifest["cases"]
            for step in case.get("cleanup_steps", [])
        ]
        expected_transactions = (
            (
                "/api/admin/activity/delete",
                "all_fixture_activity_ids",
                "timetabler.activity.deleted",
            ),
            (
                "/api/admin/staff/delete",
                "all_fixture_staff_ids",
                "timetabler.staff.deleted",
            ),
            (
                "/api/admin/location/delete",
                "all_fixture_location_ids",
                "timetabler.location.deleted",
            ),
        )
        if [step["path"] for step in cleanup_steps] != [
            item[0] for item in expected_transactions
        ]:
            raise HarnessError("final gate cleanup contract differs from the approved plan")

        durable_transactions = []
        for offset, (_, selector_name, event_type) in enumerate(
            expected_transactions, start=1
        ):
            sequence = cleanup_start + offset
            request_id = f"{self.config.run_id}-e2e-13-{offset:02d}"
            transaction = _transaction_evidence(
                sequence - 1, sequence, request_id
            )[0]
            expected_count = self.manifest["reference_selectors"][selector_name][
                "expected_count"
            ]
            if (
                transaction["transaction_count"] != expected_count
                or transaction["event_types"] != [event_type]
            ):
                raise HarnessError(
                    "cleanup durable transaction differs from final evidence contract"
                )
            statuses = set(
                IntegrationOutbox.objects.filter(
                    source_scope=integration_source_scope(),
                    transport_sequence=sequence,
                ).values_list("status", flat=True)
            )
            if statuses != {"published"}:
                raise HarnessError("cleanup durable transaction is not fully published")
            durable_transactions.append(transaction)

        durable_outbox_members = IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(),
            transport_sequence__lte=expected_watermark,
        ).count()
        if durable_outbox_members < sum(
            transaction["transaction_count"] for transaction in durable_transactions
        ):
            raise HarnessError("durable outbox evidence is incomplete at final gate")
        self.final_assertions = {
            **fixture_absence,
            "durable_state_deleted": False,
            "durable_outbox_member_count": durable_outbox_members,
            "cleanup_transactions": durable_transactions,
            "source_watermark_unchanged": _current_watermark() == before,
            "mutation_executed": False,
            "resource_booking_process_touched": False,
        }
        if not self.final_assertions["source_watermark_unchanged"]:
            raise HarnessError("final read-only gate changed the source watermark")
        self.evidence.append(
            {
                "case": "E2E-13",
                "final_watermark": before,
                "timetabler_source_evidence_complete": True,
                "e2e13_completion_claim": False,
                "requires_immutable_resource_booking_reconciliation_evidence": True,
                "mutation_executed": False,
            }
        )


def seal_evidence(payload: dict[str, Any]) -> dict[str, Any]:
    canonical = json.dumps(
        payload, sort_keys=True, separators=(",", ":"), ensure_ascii=True
    ).encode()
    return {**payload, "evidence_sha256": hashlib.sha256(canonical).hexdigest()}


def load_runtime() -> tuple[gate.EntryConfig, dict[str, Any]]:
    config = gate.EntryConfig.from_environment()
    gate.validate_run_id(config.run_id)
    manifest = gate.load_and_validate_manifest(
        config.fixture_json,
        environment=config.environment,
        source_scope=config.source_scope,
        require_approval=True,
    )
    validate_execution_manifest(manifest, config.run_id)
    if not re.fullmatch(
        r"[0-9a-f]{40}", os.environ.get("TT_PHASE1_E2E_EXPECTED_SHA", "")
    ):
        raise HarnessError("deployed git SHA attestation is unavailable")
    if not os.environ.get("TT_PHASE1_E2E_GITHUB_RUN_ID", "").isdigit():
        raise HarnessError("immutable GitHub workflow run id is unavailable")
    return config, manifest


def execute(stage: str) -> dict[str, Any]:
    config, manifest = load_runtime()
    entry = gate.execute_entry(config)
    if (
        entry["before_outbox_watermark"]
        != manifest["expected_stage_start_watermarks"][stage]
    ):
        raise HarnessError("entry watermark does not match the approved manifest")
    gate.bootstrap_django_runtime()
    session = gate.TimetablerAdminSession(
        base_url=gate.validate_base_url(
            config.timetabler_base_url, "Timetabler base URL"
        ),
        email=config.timetabler_email,
        password=config.timetabler_password,
        run_id=config.run_id,
    )
    try:
        session.login_and_attest()
        runner = ScenarioRunner(
            config=config,
            manifest=manifest,
            session=session,
            resolve_fixtures=stage != "final",
        )
        if stage == "functional":
            runner.run_functional()
        elif stage == "redelivery":
            runner.run_redelivery()
        elif stage == "conflict":
            runner.run_conflict()
        elif stage == "final":
            runner.run_final()
        elif stage == "cleanup":
            runner.run_cleanup()
        else:
            raise HarnessError(f"unsupported harness stage: {stage}")
    finally:
        session.logout()
    payload = {
        "schema_version": 1,
        "service": "timetabler",
        "repository": "Mayvins/timetabler-be",
        "git_sha": os.environ.get("TT_PHASE1_E2E_EXPECTED_SHA", ""),
        "github_run_id": os.environ.get("TT_PHASE1_E2E_GITHUB_RUN_ID", ""),
        "run_id": config.run_id,
        "source_scope": config.source_scope,
        "stage": stage,
        "entry_watermark": entry["before_outbox_watermark"],
        "final_watermark": _current_watermark(),
        "fixture_ids": dict(sorted(runner.fixture_ids.items())),
        "evidence": runner.evidence,
        "reverse_delivery_enabled": False,
        "load_executed": False,
        "dr_executed": False,
    }
    if runner.cleanup_assertions is not None:
        payload["cleanup_assertions"] = runner.cleanup_assertions
    if runner.publisher_readiness is not None:
        payload["publisher_readiness"] = runner.publisher_readiness
    if runner.final_assertions is not None:
        payload["final_assertions"] = runner.final_assertions
    return seal_evidence(payload)


def main() -> int:
    parser = argparse.ArgumentParser()
    parser.add_argument(
        "stage",
        choices=(
            "validate",
            "functional",
            "redelivery",
            "conflict",
            "cleanup",
            "final",
        ),
    )
    parser.add_argument("--manifest")
    parser.add_argument("--run-id", default="phase1-e2e-example")
    args = parser.parse_args()
    try:
        if args.stage == "validate":
            if not args.manifest:
                raise HarnessError("--manifest is required for local validation")
            manifest = json.loads(Path(args.manifest).read_text())
            gate.load_and_validate_manifest(
                json.dumps(manifest),
                environment="staging",
                source_scope="default",
                require_approval=False,
            )
            validate_execution_manifest(manifest, args.run_id)
            print(json.dumps({"manifest": "valid", "case_count": 13}, sort_keys=True))
            return 0
        print(json.dumps(execute(args.stage), sort_keys=True))
        return 0
    except (HarnessError, gate.EntryGateError) as error:
        print(
            json.dumps({"e2e_stage": "blocked", "reason": str(error)}, sort_keys=True),
            file=sys.stderr,
        )
        return 3


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