"""Event registry, dispatcher, and subscriber decorator.

Synchronous by default. When `EVENTS_ASYNC_DISPATCH=True` and a Celery app
is configured, dispatch will enqueue (left as a TODO seam — Phase 8).
"""

from __future__ import annotations

from collections.abc import Callable, Iterable
from dataclasses import dataclass, field
from threading import Lock
from typing import Any

import structlog

from simorgh.core.audit import record_event

_log = structlog.get_logger("simorgh.events")


class EventError(Exception):
    """Raised for registry / dispatch misuse."""


@dataclass(frozen=True)
class EventSpec:
    name: str
    description: str
    payload_keys: tuple[str, ...] = ()


Handler = Callable[[dict[str, Any]], None]
CatchAllHandler = Callable[[str, dict[str, Any]], None]


@dataclass
class _Registry:
    events: dict[str, EventSpec] = field(default_factory=dict)
    subscribers: dict[str, list[Handler]] = field(default_factory=dict)
    catchall_handlers: list[CatchAllHandler] = field(default_factory=list)
    lock: Lock = field(default_factory=Lock)


_REGISTRY = _Registry()


def register_event(
    name: str,
    *,
    description: str = "",
    payload_keys: Iterable[str] = (),
) -> EventSpec:
    """Declare an event name. Idempotent (re-registration returns same spec)."""
    if "." not in name:
        raise EventError(f"event name must be dotted ('module.action'), got {name!r}")
    spec = EventSpec(name=name, description=description, payload_keys=tuple(payload_keys))
    with _REGISTRY.lock:
        existing = _REGISTRY.events.get(name)
        if existing and existing != spec:
            raise EventError(f"event {name!r} already registered with different spec")
        _REGISTRY.events[name] = spec
    return spec


def registered_events() -> dict[str, EventSpec]:
    return dict(_REGISTRY.events)


def subscribe(name: str) -> Callable[[Handler], Handler]:
    """Register `func` as a handler for `name`. Multiple handlers allowed."""

    def decorator(func: Handler) -> Handler:
        with _REGISTRY.lock:
            _REGISTRY.subscribers.setdefault(name, []).append(func)
        return func

    return decorator


def subscribe_all(func: CatchAllHandler) -> CatchAllHandler:
    """Register *func* as a catch-all handler invoked for **every** event.

    Handler signature: ``(event_name: str, payload: dict) -> None``.
    Catch-all handlers are called after per-event subscribers.  Used by
    the Automation Engine to route any dispatched event to matching
    ``AutomationRule`` rows without needing per-event subscriptions.
    """
    with _REGISTRY.lock:
        _REGISTRY.catchall_handlers.append(func)
    return func


def clear_subscribers(name: str | None = None) -> None:
    """Test helper — wipes subscribers (and optionally registry entries)."""
    with _REGISTRY.lock:
        if name is None:
            _REGISTRY.subscribers.clear()
            _REGISTRY.catchall_handlers.clear()
        else:
            _REGISTRY.subscribers.pop(name, None)


def dispatch(name: str, payload: dict[str, Any] | None = None, *, audit: bool = True) -> None:
    """Synchronously fan out an event to all subscribers.

    Each handler runs in isolation; exceptions are logged but do not abort
    other handlers (so a flaky listener can't take down a write path).
    """
    _dispatch_sync(name, payload, audit=audit)


def _dispatch_sync(
    name: str, payload: dict[str, Any] | None = None, *, audit: bool = True
) -> None:
    payload = payload or {}
    spec = _REGISTRY.events.get(name)
    if spec is None:
        raise EventError(f"unknown event {name!r} — call register_event() first")

    missing = [k for k in spec.payload_keys if k not in payload]
    if missing:
        raise EventError(f"event {name!r} missing payload keys: {missing}")

    if audit:
        record_event(
            "events.dispatched",
            resource_type="event",
            resource_id=name,
            after=payload,
        )

    for handler in list(_REGISTRY.subscribers.get(name, ())):
        try:
            handler(payload)
        except Exception as exc:
            _log.warning(
                "events.handler_failed",
                event_name=name,
                handler=f"{handler.__module__}.{handler.__qualname__}",
                error=str(exc),
            )

    # Invoke catch-all handlers (registered via subscribe_all).
    for handler in list(_REGISTRY.catchall_handlers):
        try:
            handler(name, payload)
        except Exception as exc:
            _log.warning(
                "events.catchall_handler_failed",
                event_name=name,
                handler=f"{handler.__module__}.{handler.__qualname__}",
                error=str(exc),
            )


def dispatch_async(
    name: str,
    payload: dict[str, Any] | None = None,
    *,
    delay_seconds: int = 0,
    max_attempts: int = 5,
    tenant_id: int | None = None,
) -> None:
    """Persist an outbox row inside the current transaction.

    The row is delivered after ``transaction.on_commit`` fires; if the
    transaction rolls back the event is silently dropped — that's the
    whole point of the transactional outbox pattern.

    If ``EVENTS_ASYNC_DISPATCH`` is False (default in tests) the event is
    delivered synchronously, exactly like :func:`dispatch`.
    """
    from django.conf import settings
    from django.db import transaction
    from django.utils import timezone

    payload = payload or {}
    spec = _REGISTRY.events.get(name)
    if spec is None:
        raise EventError(f"unknown event {name!r} — call register_event() first")
    missing = [k for k in spec.payload_keys if k not in payload]
    if missing:
        raise EventError(f"event {name!r} missing payload keys: {missing}")

    if not getattr(settings, "EVENTS_ASYNC_DISPATCH", False):
        _dispatch_sync(name, payload)
        return

    from datetime import timedelta

    from simorgh.apps.events.models import EventOutboxEntry

    next_retry_at = timezone.now() + timedelta(seconds=delay_seconds)
    entry = EventOutboxEntry.objects.create(
        name=name,
        payload=payload,
        next_retry_at=next_retry_at,
        max_attempts=max_attempts,
        tenant_id_hint=tenant_id,
    )

    def _enqueue() -> None:
        from simorgh.apps.events.tasks import flush_outbox

        flush_outbox.delay()

    transaction.on_commit(_enqueue)
    _log.info("events.outbox.enqueued", event_name=name, entry_id=str(entry.pk))
