#!/usr/bin/env python3
"""Provision one isolated, unscheduled E2E-10 rerun activity through TT APIs.

ORM access in this script is read-only. The only domain writes are authenticated
requests to existing Timetabler admin APIs. The source fence, original accepted
conflict allocation, transaction shape, advisory requirement delta, publisher,
and updated private manifest are all checked fail closed.
"""

from __future__ import annotations

import base64
from copy import deepcopy
import hashlib
import json
import os
from pathlib import Path
import re
import sys
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
import phase1_e2e_harness as harness  # noqa: E402
from phase1_e2e_provision import REQUIRED_PERMISSIONS, wait_for_publisher  # noqa: E402


EXPECTED_START_WATERMARK = 112
EXPECTED_END_WATERMARK = 113
ORIGINAL_ACTIVITY_ID = 1312
CONFLICT_STAFF_ID = 1758
CONFLICT_LOCATION_ID = 99
RERUN_ACTIVITY_SUFFIX = "conflict-rerun-02"
RERUN_SLOT = 21
CONFLICT_SENTINEL_EPOCH = 4_102_444_800


class ConflictRerunProvisionError(harness.HarnessError):
    pass


def _exact_sha_and_run() -> tuple[str, str]:
    git_sha = os.environ.get("TT_PHASE1_E2E_EXPECTED_SHA", "").strip()
    github_run_id = os.environ.get("TT_PHASE1_E2E_GITHUB_RUN_ID", "").strip()
    if not re.fullmatch(r"[0-9a-f]{40}", git_sha):
        raise ConflictRerunProvisionError("deployed git SHA attestation is unavailable")
    if not github_run_id.isdigit():
        raise ConflictRerunProvisionError("immutable GitHub workflow run id is unavailable")
    return git_sha, github_run_id


def _case(manifest: dict[str, Any], case_id: str) -> dict[str, Any]:
    return next(case for case in manifest["cases"] if case["id"] == case_id)


def build_rerun_manifest(
    manifest: dict[str, Any], *, activity_code: str, github_run_id: str
) -> dict[str, Any]:
    updated = deepcopy(manifest)
    updated["approval_ref"] = (
        f"github-actions:{github_run_id}:conflict-rerun-provision"
    )
    updated["conflict_armed"] = False
    updated["load_enabled"] = False
    updated["dr_enabled"] = False
    updated["variables"]["conflict_slot"] = RERUN_SLOT
    updated["reference_selectors"]["conflict_activity_id"]["filters"] = {
        "code": activity_code
    }
    updated["reference_selectors"]["all_fixture_activity_ids"][
        "expected_count"
    ] = 9
    fences = updated["expected_stage_start_watermarks"]
    fences.update(
        {
            "entry": EXPECTED_END_WATERMARK,
            "conflict": EXPECTED_END_WATERMARK,
            "restart": EXPECTED_END_WATERMARK + 1,
            "rollback": EXPECTED_END_WATERMARK + 1,
            "cleanup": EXPECTED_END_WATERMARK + 1,
            "final": EXPECTED_END_WATERMARK + 4,
        }
    )
    _case(updated, "E2E-10")["execute_at_epoch"] = CONFLICT_SENTINEL_EPOCH
    gate.load_and_validate_manifest(
        json.dumps(updated),
        environment="staging",
        source_scope=updated["source_scope"],
        require_approval=True,
    )
    harness.validate_execution_manifest(updated, updated["run_id"])
    return updated


def _assert_initial_manifest(manifest: dict[str, Any], run_id: str) -> None:
    harness.validate_execution_manifest(manifest, run_id)
    if manifest["conflict_armed"] is not False:
        raise ConflictRerunProvisionError(
            "conflict rerun preparation requires conflict_armed=false"
        )
    if manifest["load_enabled"] is not False or manifest["dr_enabled"] is not False:
        raise ConflictRerunProvisionError("LOAD and DR must remain disabled")
    if manifest["variables"].get("conflict_slot") != 20:
        raise ConflictRerunProvisionError(
            "initial manifest must retain the accepted slot-20 evidence"
        )
    conflict_selector = manifest["reference_selectors"].get("conflict_activity_id")
    if (
        not isinstance(conflict_selector, dict)
        or conflict_selector.get("filters") != {"code": f"{run_id}-conflict"}
        or conflict_selector.get("expected_count") != 1
    ):
        raise ConflictRerunProvisionError(
            "initial manifest must retain the original conflict activity selector"
        )
    activity_selector = manifest["reference_selectors"].get(
        "all_fixture_activity_ids"
    )
    if (
        not isinstance(activity_selector, dict)
        or activity_selector.get("filters")
        != {"code__startswith": f"{run_id}-"}
        or activity_selector.get("expected_count") != 8
    ):
        raise ConflictRerunProvisionError(
            "initial manifest must resolve the exact eight pre-rerun activities"
        )
    fences = manifest["expected_stage_start_watermarks"]
    if fences["entry"] != EXPECTED_START_WATERMARK:
        raise ConflictRerunProvisionError("entry fence must equal 112 before preparation")
    if fences["conflict"] != EXPECTED_START_WATERMARK:
        raise ConflictRerunProvisionError(
            "conflict fence must equal 112 before preparation"
        )


def _source_weeks(activity) -> list[int]:
    weeks = (
        activity.week_pattern.week.all()
        if activity.week_pattern_id
        else activity.week.all()
    )
    values = list(weeks.order_by("id").values_list("id", flat=True))
    if len(values) != 2:
        raise ConflictRerunProvisionError(
            "original conflict activity must contain exactly two occurrence weeks"
        )
    return [int(value) for value in values]


def _assert_original_allocation(run_id: str):
    from api.models import TtActivity, TtLocation, TtSetting, TtStaff

    activity_code = f"{run_id}-conflict"
    activity = TtActivity.objects.filter(code=activity_code).first()
    staff = TtStaff.objects.filter(
        pk=CONFLICT_STAFF_ID,
        code=f"{run_id}-fixture-staff-conflict",
        status=1,
    ).first()
    location = TtLocation.objects.filter(
        pk=CONFLICT_LOCATION_ID,
        code=f"{run_id}-fixture-location-conflict",
        status=1,
    ).first()
    if activity is None or activity.pk != ORIGINAL_ACTIVITY_ID:
        raise ConflictRerunProvisionError(
            "original conflict activity identity does not match sequence-112 evidence"
        )
    if staff is None or location is None:
        raise ConflictRerunProvisionError(
            "dedicated conflict Staff/Location identity is unavailable"
        )
    if not activity.scheduled:
        raise ConflictRerunProvisionError(
            "original conflict activity is no longer allocated"
        )
    if set(activity.staff.values_list("id", flat=True)) != {CONFLICT_STAFF_ID}:
        raise ConflictRerunProvisionError("original conflict Staff allocation changed")
    if set(activity.location.values_list("id", flat=True)) != {CONFLICT_LOCATION_ID}:
        raise ConflictRerunProvisionError("original conflict Location allocation changed")
    slot_per_week = int(
        TtSetting.get_multiple_setting({"slot_per_week"})["slot_per_week"]
    )
    if int(activity.scheduled_start_slot) % slot_per_week != 20:
        raise ConflictRerunProvisionError("original conflict slot is no longer 20")
    return activity


def _resolve_new_activity(code: str, actor_id: int):
    from api.models import TtActivity

    rows = list(TtActivity.objects.filter(code=code).order_by("id")[:2])
    if len(rows) != 1:
        raise ConflictRerunProvisionError(
            "rerun activity did not resolve to exactly one row"
        )
    activity = rows[0]
    if activity.created_by != actor_id or activity.status != 1 or activity.scheduled:
        raise ConflictRerunProvisionError(
            "rerun activity is not active, unscheduled, and E2E-owned"
        )
    return activity


def _attest_created_event(activity_id: int, transaction: dict[str, Any]):
    from api.models import IntegrationOutbox
    from api.services.integration.outbox import (
        aggregate_source_id,
        integration_source_scope,
    )

    rows = list(
        IntegrationOutbox.objects.filter(
            source_scope=integration_source_scope(),
            transport_sequence=EXPECTED_END_WATERMARK,
        ).order_by("transaction_index")[:2]
    )
    if len(rows) != 1:
        raise ConflictRerunProvisionError(
            "sequence 113 must contain exactly one finalized outbox member"
        )
    event = rows[0]
    if (
        event.aggregate_type != "activity"
        or event.aggregate_id != aggregate_source_id(activity_id)
        or event.event_type != "timetabler.activity.created"
        or str(event.change_set_id) != transaction["source_transaction_id"]
        or event.transaction_index != 1
        or event.transaction_count != 1
        or not event.transaction_finalized
    ):
        raise ConflictRerunProvisionError(
            "sequence 113 is not the exact new activity-created transaction"
        )
    return event


def execute() -> dict[str, Any]:
    config = gate.EntryConfig.from_environment()
    gate.validate_run_id(config.run_id)
    git_sha, github_run_id = _exact_sha_and_run()
    gate.bootstrap_django_runtime()
    manifest = gate.load_and_validate_manifest(
        config.fixture_json,
        environment=config.environment,
        source_scope=config.source_scope,
        require_approval=True,
    )
    _assert_initial_manifest(manifest, config.run_id)

    before = harness._assert_local_safety()
    readiness = wait_for_publisher()
    if (
        before != EXPECTED_START_WATERMARK
        or readiness["transport_watermark"] != EXPECTED_START_WATERMARK
        or readiness["publisher_last_sequence"] != EXPECTED_START_WATERMARK
    ):
        raise ConflictRerunProvisionError(
            "conflict rerun preparation requires publisher-complete fence 112"
        )

    original = _assert_original_allocation(config.run_id)
    rerun_code = f"{config.run_id}-{RERUN_ACTIVITY_SUFFIX}"
    from api.models import TtActivity

    if TtActivity.objects.filter(code=rerun_code).exists():
        raise ConflictRerunProvisionError(
            "rerun activity already exists; refusing a second preparation mutation"
        )
    week_ids = _source_weeks(original)

    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()
        if not REQUIRED_PERMISSIONS <= session.permissions or session.user_id is None:
            raise ConflictRerunProvisionError(
                "dedicated Timetabler E2E identity lacks preparation permissions"
            )
        create_request_id = f"{config.run_id}-{RERUN_ACTIVITY_SUFFIX}-create"
        session.post(
            "/api/admin/activity/create",
            {
                "code": rerun_code,
                "name": f"{config.run_id} conflict rerun 02",
                "duration": int(original.duration),
                "planned_size": 1,
                "academic_term_id": int(original.academic_term_id),
                "activity_template_id": int(original.activity_template_id),
                "week_pattern": week_ids,
            },
            request_id=create_request_id,
        )
        after_create, transactions = harness._wait_for_transaction_delta(
            start=before,
            expected=1,
            request_id=create_request_id,
            timeout_seconds=90,
        )
        expected_transaction = {
            "event_types": ["timetabler.activity.created"],
            "source_sequence": EXPECTED_END_WATERMARK,
            "transaction_count": 1,
        }
        if (
            after_create != EXPECTED_END_WATERMARK
            or len(transactions) != 1
            or any(
                transactions[0][key] != value
                for key, value in expected_transaction.items()
            )
        ):
            raise ConflictRerunProvisionError(
                "rerun activity creation did not produce the exact sequence-113 transaction"
            )
        activity = _resolve_new_activity(rerun_code, int(session.user_id))
        source_event = _attest_created_event(int(activity.pk), transactions[0])
        if activity.pk == ORIGINAL_ACTIVITY_ID:
            raise ConflictRerunProvisionError("rerun activity reused the original identity")
        if activity.staff.exists() or activity.location.exists():
            raise ConflictRerunProvisionError(
                "rerun activity exposed an allocation before scheduling"
            )

        requirement_request_id = (
            f"{config.run_id}-{RERUN_ACTIVITY_SUFFIX}-requirements"
        )
        requirement_before = harness._assert_local_safety()
        if requirement_before != EXPECTED_END_WATERMARK:
            raise ConflictRerunProvisionError(
                "source fence changed before advisory requirement configuration"
            )
        session.post(
            "/api/admin/resources/update-requirement",
            {
                "activity_id": int(activity.pk),
                "staff_requirement": 1,
                "staff_preset": [CONFLICT_STAFF_ID],
                "staff_suitability": [],
                "location_requirement": 1,
                "location_preset": [CONFLICT_LOCATION_ID],
                "location_suitability": [],
            },
            request_id=requirement_request_id,
        )
        requirement_after, requirement_transactions = (
            harness._wait_for_transaction_delta(
                start=requirement_before,
                expected=0,
                request_id=requirement_request_id,
                timeout_seconds=30,
            )
        )
        activity = _resolve_new_activity(rerun_code, int(session.user_id))
        if requirement_after != EXPECTED_END_WATERMARK or requirement_transactions:
            raise ConflictRerunProvisionError(
                "advisory requirements unexpectedly created a source transaction"
            )
        if set(activity.staff_preset.values_list("id", flat=True)) != {
            CONFLICT_STAFF_ID
        }:
            raise ConflictRerunProvisionError("rerun Staff preset is incomplete")
        if set(activity.location_preset.values_list("id", flat=True)) != {
            CONFLICT_LOCATION_ID
        }:
            raise ConflictRerunProvisionError("rerun Location preset is incomplete")
        if activity.staff.exists() or activity.location.exists() or activity.scheduled:
            raise ConflictRerunProvisionError(
                "advisory requirements exposed a premature allocation"
            )

        published = wait_for_publisher()
        if (
            published["transport_watermark"] != EXPECTED_END_WATERMARK
            or published["publisher_last_sequence"] != EXPECTED_END_WATERMARK
            or published["dead_letter_count"] != 0
            or published["reverse_delivery_enabled"]
        ):
            raise ConflictRerunProvisionError(
                "publisher did not complete the exact sequence-113 fence"
            )
        source_event.refresh_from_db()
        if source_event.status != "published":
            raise ConflictRerunProvisionError(
                "sequence-113 activity-created event is not published"
            )

        updated_manifest = build_rerun_manifest(
            manifest, activity_code=rerun_code, github_run_id=github_run_id
        )
        if len(harness.resolve_selector(
            "all_fixture_activity_ids",
            updated_manifest["reference_selectors"]["all_fixture_activity_ids"],
        )) != 9:
            raise ConflictRerunProvisionError(
                "updated manifest did not resolve exactly nine run-owned activities"
            )
        conflict_basis = gate.conflict_contender_basis(updated_manifest)
    finally:
        session.logout()

    canonical = json.dumps(
        updated_manifest, sort_keys=True, separators=(",", ":"), ensure_ascii=True
    )
    manifest_sha256 = hashlib.sha256(canonical.encode()).hexdigest()
    transaction = transactions[0]
    summary = {
        "conflict_rerun_fixture_provisioning": "passed",
        "run_id": config.run_id,
        "git_sha": git_sha,
        "github_run_id": github_run_id,
        "before_watermark": before,
        "after_watermark": EXPECTED_END_WATERMARK,
        "publisher_last_sequence": published["publisher_last_sequence"],
        "activity_id": int(activity.pk),
        "activity_code": rerun_code,
        "original_activity_id": ORIGINAL_ACTIVITY_ID,
        "conflict_staff_id": CONFLICT_STAFF_ID,
        "conflict_location_id": CONFLICT_LOCATION_ID,
        "source_transaction": transaction,
        "source_event": {
            "aggregate_type": source_event.aggregate_type,
            "aggregate_id": source_event.aggregate_id,
            "event_type": source_event.event_type,
            "event_version": source_event.event_version,
            "status": source_event.status,
        },
        "requirement_source_transaction_count": 0,
        "conflict_contender_basis": conflict_basis,
        "manifest_sha256": manifest_sha256,
        "conflict_armed": False,
        "reverse_delivery_enabled": False,
        "load_executed": False,
        "dr_executed": False,
    }
    print(json.dumps(summary, sort_keys=True))
    print(
        "TT_PHASE1_E2E_MANIFEST_BASE64="
        + base64.b64encode(canonical.encode()).decode()
    )
    return summary


def main() -> int:
    try:
        execute()
        return 0
    except (ConflictRerunProvisionError, gate.EntryGateError) as error:
        print(
            json.dumps(
                {"conflict_rerun_fixture_provisioning": "blocked", "reason": str(error)},
                sort_keys=True,
            ),
            file=sys.stderr,
        )
        return 3


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