import datetime
import decimal
import hashlib
import json
import logging
import uuid
from collections.abc import Iterable
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError

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

from api.models import (
    IntegrationAggregateVersion,
    IntegrationOutbox,
    IntegrationTransportCursor,
    TtActivity,
    TtLocation,
    TtStaff,
)
from api.services.integration.context import current_context


logger = logging.getLogger("api.integration.outbox")


class IntegrationAtomicityError(RuntimeError):
    pass


def _canonical(value):
    if isinstance(value, (datetime.datetime, datetime.date, datetime.time)):
        return value.isoformat()
    if isinstance(value, (decimal.Decimal, uuid.UUID)):
        return str(value)
    if isinstance(value, dict):
        return {str(key): _canonical(item) for key, item in value.items()}
    if isinstance(value, (list, tuple, set)):
        return [_canonical(item) for item in value]
    return value


def _utc_now_iso() -> str:
    return datetime.datetime.now(datetime.UTC).isoformat().replace("+00:00", "Z")


def canonical_hash(value) -> str:
    encoded = json.dumps(
        _canonical(value),
        sort_keys=True,
        separators=(",", ":"),
        ensure_ascii=False,
        allow_nan=False,
    ).encode("utf-8")
    return hashlib.sha256(encoded).hexdigest()


def build_transport_envelope(
    events: Iterable[IntegrationOutbox], *, event_id: str | None = None
) -> dict:
    members = sorted(events, key=lambda event: event.transaction_index)
    if not members:
        raise IntegrationAtomicityError("A transport change set must contain at least one event")
    source_scopes = {event.source_scope for event in members}
    source_sequences = {event.transport_sequence for event in members}
    change_set_ids = {event.change_set_id for event in members}
    if len(source_scopes) != 1 or len(source_sequences) != 1 or len(change_set_ids) != 1:
        raise IntegrationAtomicityError("Transport change-set identity is inconsistent")
    maximum = settings.RESOURCE_BOOKING_INTEGRATION.get("MAX_EVENTS_PER_CHANGE_SET", 5000)
    if len(members) > maximum:
        raise IntegrationAtomicityError(
            f"Transport change set exceeds {maximum} canonical member events"
        )
    member_payloads = [event.payload for event in members]
    envelope = {
        "event_id": event_id or str(members[0].change_set_id),
        "schema_version": members[0].schema_version,
        "event_type": "timetabler.change_set.final_state",
        "source_system": "timetabler",
        "source_scope": members[0].source_scope,
        "source_sequence": members[0].transport_sequence,
        "ordering_key": members[0].source_scope,
        "source_transaction_id": str(members[0].change_set_id),
        "transaction_count": len(members),
        "complete": True,
        "committed_at": max(event.payload["committed_at"] for event in members),
        "events_hash": canonical_hash(member_payloads),
        "events": member_payloads,
    }
    encoded_size = len(
        json.dumps(envelope, sort_keys=True, separators=(",", ":"), ensure_ascii=False).encode(
            "utf-8"
        )
    )
    if encoded_size > settings.RESOURCE_BOOKING_INTEGRATION.get("MAX_EVENT_BYTES", 2097152):
        raise IntegrationAtomicityError(
            "Canonical integration change set exceeds the configured transport byte limit"
        )
    return envelope


def integration_enabled() -> bool:
    return bool(settings.RESOURCE_BOOKING_INTEGRATION["CAPTURE_ENABLED"])


def source_identity(aggregate_type: str, source_id) -> dict:
    deployment_id = settings.RESOURCE_BOOKING_INTEGRATION["DEPLOYMENT_ID"]
    return {
        "deployment_id": deployment_id,
        "aggregate_type": aggregate_type,
        "source_id": str(source_id),
    }


def aggregate_source_id(source_id) -> str:
    deployment_id = settings.RESOURCE_BOOKING_INTEGRATION["DEPLOYMENT_ID"]
    return f"{deployment_id}:{source_id}"


def integration_source_scope() -> str:
    config = settings.RESOURCE_BOOKING_INTEGRATION
    return config.get("SOURCE_SCOPE") or config.get("DEPLOYMENT_ID", "default")


def _lock_aggregate(aggregate_type: str, aggregate_id: str) -> None:
    if connection.vendor == "postgresql":
        lock_key = f"{aggregate_type}:{aggregate_id}"
        with connection.cursor() as cursor:
            cursor.execute("SELECT pg_advisory_xact_lock(hashtextextended(%s, 0))", [lock_key])


def _next_version(aggregate_type: str, aggregate_id: str) -> int:
    _lock_aggregate(aggregate_type, aggregate_id)
    version_row, created = IntegrationAggregateVersion.objects.get_or_create(
        aggregate_type=aggregate_type,
        aggregate_id=aggregate_id,
        defaults={"version": 0},
    )
    if not created:
        version_row = IntegrationAggregateVersion.objects.select_for_update().get(pk=version_row.pk)
    version_row.version += 1
    version_row.save(update_fields=["version", "updated_at"])
    return version_row.version


def _next_transport_position(change_set_id: uuid.UUID) -> tuple[str, int, int]:
    """Allocate an event sequence while holding the scope cursor until commit.

    The cursor update participates in the caller's authoritative transaction. A
    concurrent transaction cannot allocate the next sequence until this one
    commits or rolls back, so sequence order is also database commit order.
    """

    source_scope = integration_source_scope()
    finalized = IntegrationOutbox.objects.filter(
        source_scope=source_scope,
        change_set_id=change_set_id,
        transaction_finalized=True,
    ).exists()
    if finalized:
        raise IntegrationAtomicityError(f"Change set {change_set_id} is already finalized")

    existing = list(
        IntegrationOutbox.objects.select_for_update()
        .filter(source_scope=source_scope, change_set_id=change_set_id)
        .order_by("transaction_index")
    )
    if existing:
        transport_sequence = existing[0].transport_sequence
    else:
        cursor, created = IntegrationTransportCursor.objects.get_or_create(
            source_scope=source_scope,
            defaults={"last_sequence": 0},
        )
        if not created:
            cursor = IntegrationTransportCursor.objects.select_for_update().get(pk=cursor.pk)
        cursor.last_sequence += 1
        cursor.save(update_fields=["last_sequence", "updated_at"])
        transport_sequence = cursor.last_sequence
    return source_scope, transport_sequence, len(existing) + 1


def finalize_change_set(change_set_id: uuid.UUID) -> int:
    if not connection.in_atomic_block:
        raise IntegrationAtomicityError("Integration change sets must finalize inside transaction.atomic()")
    source_scope = integration_source_scope()
    events = list(
        IntegrationOutbox.objects.select_for_update()
        .filter(source_scope=source_scope, change_set_id=change_set_id)
        .order_by("transport_sequence", "transaction_index", "id")
    )
    count = len(events)
    for index, event in enumerate(events, start=1):
        event.transaction_index = index
        event.transaction_count = count
        event.transaction_finalized = True
        event.payload["transport"] = {
            "source_scope": source_scope,
            "sequence": event.transport_sequence,
            "transaction_id": str(change_set_id),
            "transaction_index": index,
            "transaction_count": count,
            "member_only": True,
        }
    if events:
        # Validate the exact single-record broker envelope before the source
        # transaction can commit; fragments are never independently complete.
        build_transport_envelope(events)
        IntegrationOutbox.objects.bulk_update(
            events,
            ["transaction_index", "transaction_count", "transaction_finalized", "payload"],
        )
    return count


def append_outbox_event(
    *,
    event_type: str,
    aggregate_type: str,
    source_id,
    payload: dict,
    change_set_id: uuid.UUID | None = None,
    origin: str | None = None,
    correlation_id: str | None = None,
    causation_id: str | None = None,
    request_id: str | None = None,
    actor_id: int | None = None,
    replay_of: IntegrationOutbox | None = None,
    replay_requested_by: int | None = None,
    replay_reason: str | None = None,
    defer_change_set_finalization: bool = False,
) -> IntegrationOutbox | None:
    if not integration_enabled():
        return None
    if not connection.in_atomic_block:
        raise IntegrationAtomicityError(
            f"{event_type} for {aggregate_type}:{source_id} must be appended inside transaction.atomic()"
        )

    context = current_context()
    aggregate_id = aggregate_source_id(source_id)
    event_id = uuid.uuid4()
    resolved_change_set_id = change_set_id or uuid.uuid4()
    source_scope, transport_sequence, transaction_index = _next_transport_position(
        resolved_change_set_id
    )
    ordering_key = source_scope
    version = _next_version(aggregate_type, aggregate_id)
    committed_state = _canonical(payload)
    envelope = {
        "event_id": str(event_id),
        "schema_version": settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
        "event_type": event_type,
        "aggregate": source_identity(aggregate_type, source_id),
        "event_version": version,
        "ordering_key": ordering_key,
        "change_set_id": str(resolved_change_set_id),
        "transport": {
            "source_scope": source_scope,
            "sequence": transport_sequence,
            "transaction_id": str(resolved_change_set_id),
            "transaction_index": transaction_index,
            "transaction_count": 0,
            "member_only": True,
        },
        "origin": origin or context.origin,
        "correlation_id": correlation_id or context.correlation_id,
        "causation_id": causation_id or context.causation_id,
        "request_id": request_id or context.request_id,
        "actor_id": actor_id if actor_id is not None else context.actor_id,
        "committed_at": _utc_now_iso(),
        "committed_state_hash": canonical_hash(committed_state),
        "committed_state": committed_state,
    }
    encoded_size = len(
        json.dumps(envelope, sort_keys=True, separators=(",", ":"), ensure_ascii=False).encode(
            "utf-8"
        )
    )
    if encoded_size > settings.RESOURCE_BOOKING_INTEGRATION["MAX_EVENT_BYTES"]:
        raise IntegrationAtomicityError(
            f"Canonical integration event exceeds {settings.RESOURCE_BOOKING_INTEGRATION['MAX_EVENT_BYTES']} bytes"
        )
    event = IntegrationOutbox.objects.create(
        event_id=event_id,
        schema_version=settings.RESOURCE_BOOKING_INTEGRATION["SCHEMA_VERSION"],
        event_type=event_type,
        aggregate_type=aggregate_type,
        aggregate_id=aggregate_id,
        event_version=version,
        source_scope=source_scope,
        transport_sequence=transport_sequence,
        transaction_index=transaction_index,
        transaction_count=0,
        transaction_finalized=False,
        ordering_key=ordering_key,
        change_set_id=resolved_change_set_id,
        origin=origin or context.origin,
        correlation_id=correlation_id or context.correlation_id,
        causation_id=causation_id or context.causation_id,
        request_id=request_id or context.request_id,
        actor_id=actor_id if actor_id is not None else context.actor_id,
        payload=envelope,
        replay_of=replay_of,
        replay_requested_by=replay_requested_by,
        replay_reason=replay_reason,
    )
    if not defer_change_set_finalization:
        finalize_change_set(resolved_change_set_id)
        event.refresh_from_db(
            fields=["transaction_index", "transaction_count", "transaction_finalized", "payload"]
        )
    logger.info(
        "integration_outbox_appended",
        extra={
            "event_id": str(event.event_id),
            "event_type": event.event_type,
            "aggregate_type": event.aggregate_type,
            "aggregate_id": event.aggregate_id,
            "event_version": event.event_version,
            "change_set_id": str(event.change_set_id),
            "correlation_id": event.correlation_id,
        },
    )
    return event


def capture_resource_state(resource) -> dict:
    if isinstance(resource, TtStaff):
        aggregate_type = "staff"
        fields = (
            "code",
            "name",
            "email",
            "desc",
            "department_id",
            "zone_id",
            "is_part_time",
            "maximum_period",
            "contract_period",
            "shared_with_all_department",
            "availability_id",
            "availability_pattern",
            "start_preference_id",
            "start_preference_pattern",
            "usage_preference_id",
            "usage_preference_pattern",
            "status",
            "created_by",
            "updated_by",
            "created_at",
            "updated_at",
        )
    elif isinstance(resource, TtLocation):
        aggregate_type = "location"
        fields = (
            "code",
            "name",
            "desc",
            "department_id",
            "zone_id",
            "capacity",
            "area",
            "maximum_period",
            "contract_period",
            "shared_with_all_department",
            "availability_id",
            "availability_pattern",
            "start_preference_id",
            "start_preference_pattern",
            "usage_preference_id",
            "usage_preference_pattern",
            "status",
            "created_by",
            "updated_by",
            "created_at",
            "updated_at",
        )
    else:
        raise TypeError(f"Unsupported integration resource: {type(resource)!r}")

    state = {field: _canonical(getattr(resource, field, None)) for field in fields}
    state.update(
        {
            "identity": source_identity(aggregate_type, resource.pk),
            "id": resource.pk,
            "record_state": "active" if resource.status == resource.STATUS_TO_CODE["active"] else "archived",
        }
    )
    return state


def append_resource_event(
    resource,
    *,
    action: str,
    previous: dict | None = None,
    tombstone: bool = False,
    change_set_id: uuid.UUID | None = None,
    actor_id: int | None = None,
    origin: str | None = None,
    correlation_id: str | None = None,
    request_id: str | None = None,
    defer_change_set_finalization: bool = False,
) -> IntegrationOutbox | None:
    aggregate_type = "staff" if isinstance(resource, TtStaff) else "location"
    current = None if tombstone else capture_resource_state(resource)
    event_type = f"timetabler.{aggregate_type}.{action}"
    return append_outbox_event(
        event_type=event_type,
        aggregate_type=aggregate_type,
        source_id=resource.pk,
        payload={
            "replacement": {"previous": previous, "current": current},
            "tombstone": tombstone,
        },
        change_set_id=change_set_id,
        actor_id=actor_id,
        origin=origin,
        correlation_id=correlation_id,
        request_id=request_id,
        defer_change_set_finalization=defer_change_set_finalization,
    )


def _activity_weeks(activity: TtActivity):
    if activity.week_pattern_id:
        return list(activity.week_pattern.week.all().order_by("start_date", "id"))
    return list(activity.week.all().order_by("start_date", "id"))


def _activity_occurrences(activity: TtActivity, weeks) -> list[dict]:
    if (
        not activity.scheduled
        or activity.scheduled_day is None
        or activity.scheduled_start_time is None
        or activity.duration is None
    ):
        return []
    try:
        source_timezone = ZoneInfo(settings.TIME_ZONE)
    except ZoneInfoNotFoundError as error:
        raise IntegrationAtomicityError(
            f"TIME_ZONE {settings.TIME_ZONE!r} cannot produce deterministic occurrences"
        ) from error

    occurrences = []
    for week in weeks:
        local_date = week.start_date + datetime.timedelta(days=activity.scheduled_day)
        naive_start = datetime.datetime.combine(local_date, activity.scheduled_start_time)
        fold = settings.RESOURCE_BOOKING_INTEGRATION["OCCURRENCE_DST_FOLD"]
        local_start = naive_start.replace(tzinfo=source_timezone, fold=fold)
        utc_start = local_start.astimezone(datetime.UTC)
        if utc_start.astimezone(source_timezone).replace(tzinfo=None) != naive_start:
            raise IntegrationAtomicityError(
                f"Activity {activity.pk} resolves to a nonexistent local time in {settings.TIME_ZONE}"
            )
        utc_end = utc_start + datetime.timedelta(minutes=activity.duration)
        local_end = utc_end.astimezone(source_timezone)
        occurrences.append(
            {
                "occurrence_ref": (
                    f"{settings.RESOURCE_BOOKING_INTEGRATION['DEPLOYMENT_ID']}:"
                    f"activity:{activity.pk}:week:{week.pk}"
                ),
                "week_ref": source_identity("week", week.pk),
                "week_number": week.week,
                "week_start_date": week.start_date.isoformat(),
                "local_date": local_date.isoformat(),
                "local_start": local_start.isoformat(),
                "local_end": local_end.isoformat(),
                "utc_start": utc_start.isoformat().replace("+00:00", "Z"),
                "utc_end": utc_end.isoformat().replace("+00:00", "Z"),
                "timezone": settings.TIME_ZONE,
                "utc_offset": local_start.isoformat()[-6:],
                "dst_fold": fold,
            }
        )
    return occurrences


def capture_activity_state(activity: TtActivity) -> dict:
    weeks = _activity_weeks(activity)
    week_ids = [week.id for week in weeks]
    academic_term = activity.academic_term
    state = {
        "identity": source_identity("activity", activity.pk),
        "id": activity.pk,
        "code": activity.code,
        "name": activity.name,
        "academic_term_id": activity.academic_term_id,
        "module_id": activity.module_id,
        "activity_template_id": activity.activity_template_id,
        "activity_type_id": activity.activity_type_id,
        "is_booking": bool(activity.is_booking),
        "is_jta": bool(activity.is_jta),
        "jta_parent_id": activity.jta_parent_id,
        "is_variant": bool(activity.is_variant),
        "variant_parent_id": activity.variant_parent_id,
        "status": activity.status,
        "scheduled": bool(activity.scheduled),
        "scheduled_day": activity.scheduled_day,
        "scheduled_start_time": _canonical(activity.scheduled_start_time),
        "scheduled_start_slot": activity.scheduled_start_slot,
        "duration_minutes": activity.duration,
        "slot_required": activity.slot_required,
        "week_pattern_id": activity.week_pattern_id,
        "week_ids": week_ids,
        "calendar": {
            "academic_term_ref": (
                source_identity("academic_term", activity.academic_term_id)
                if activity.academic_term_id
                else None
            ),
            "academic_term_code": getattr(academic_term, "code", None),
            "academic_term_start_date": _canonical(getattr(academic_term, "start_date", None)),
            "academic_term_end_date": _canonical(getattr(academic_term, "end_date", None)),
            "timezone": settings.TIME_ZONE,
            "weeks": [
                {
                    "week_ref": source_identity("week", week.pk),
                    "week_number": week.week,
                    "week_start_date": week.start_date.isoformat(),
                }
                for week in weeks
            ],
        },
        "occurrences": _activity_occurrences(activity, weeks),
        "staff_ids": sorted(activity.staff.values_list("id", flat=True)),
        "location_ids": sorted(activity.location.values_list("id", flat=True)),
        "student_set_ids": sorted(activity.student_set.values_list("id", flat=True)),
        "timezone": settings.TIME_ZONE,
    }
    return _canonical(state)


def capture_activity_states(activities: Iterable[TtActivity] | Iterable[int]) -> dict[int, dict]:
    items = list(activities)
    if not items:
        return {}
    if isinstance(items[0], int):
        queryset = (
            TtActivity.objects.filter(id__in=items)
            .select_related("week_pattern")
            .select_related("academic_term")
            .prefetch_related("week", "week_pattern__week", "staff", "location", "student_set")
            .order_by("id")
        )
    else:
        queryset = items
    return {activity.id: capture_activity_state(activity) for activity in queryset}


def append_activity_event(
    *,
    activity: TtActivity | None,
    activity_id: int,
    action: str,
    previous: dict | None,
    tombstone: bool = False,
    change_set_id: uuid.UUID | None = None,
    request_id: str | None = None,
    correlation_id: str | None = None,
    actor_id: int | None = None,
    origin: str | None = None,
    defer_change_set_finalization: bool = False,
) -> IntegrationOutbox | None:
    current = None if tombstone or activity is None else capture_activity_state(activity)
    old_scope = {
        "staff_ids": (previous or {}).get("staff_ids", []),
        "location_ids": (previous or {}).get("location_ids", []),
        "week_ids": (previous or {}).get("week_ids", []),
        "scheduled_start_slot": (previous or {}).get("scheduled_start_slot"),
        "duration_minutes": (previous or {}).get("duration_minutes"),
    }
    new_scope = {
        "staff_ids": (current or {}).get("staff_ids", []),
        "location_ids": (current or {}).get("location_ids", []),
        "week_ids": (current or {}).get("week_ids", []),
        "scheduled_start_slot": (current or {}).get("scheduled_start_slot"),
        "duration_minutes": (current or {}).get("duration_minutes"),
    }
    return append_outbox_event(
        event_type=f"timetabler.activity.{action}",
        aggregate_type="activity",
        source_id=activity_id,
        payload={
            "replacement": {"previous": previous, "current": current},
            "affected_scope": {"previous": old_scope, "current": new_scope},
            "tombstone": tombstone,
            "final_state_ref": source_identity("activity", activity_id),
        },
        change_set_id=change_set_id,
        request_id=request_id,
        correlation_id=correlation_id,
        actor_id=actor_id,
        origin=origin,
        defer_change_set_finalization=defer_change_set_finalization,
    )
