"""Серверные события (SSE): шина публикации (ТЗ 3.14).

События грубозернистые: подписчики получают «kind + данные» и сами решают,
что перечитать (refetch вместо патчей состояния). Издатели — HTTP-хендлеры,
MCP-инструменты (потоки FastMCP) и фоновый detail_task: публикация всегда
идёт через loop.call_soon_threadsafe, поэтому работает из любого потока.
"""

import asyncio
import logging
from dataclasses import dataclass
from typing import Any

logger = logging.getLogger(__name__)

# Размер очереди подписчика: при переполнении события пропускаются —
# coarse-модель самовосстанавливается (клиент перечитывает данные целиком)
QUEUE_MAXSIZE = 64


@dataclass(frozen=True)
class Event:
    kind: str
    data: dict[str, Any] | None = None


class EventBus:
    """Реестр подписок по user_id (None при publish = broadcast всем)."""

    def __init__(self) -> None:
        # user_id -> очереди его соединений (несколько вкладок = несколько очередей)
        self._subscribers: dict[str, set[asyncio.Queue[Event]]] = {}
        # Loop берётся лениво при первой подписке: lifespan в тестах не стартует,
        # а после закрытия loop публикация молча сбрасывает его на новый
        self._loop: asyncio.AbstractEventLoop | None = None

    def subscribe(self, user_id: str) -> asyncio.Queue[Event]:
        queue: asyncio.Queue[Event] = asyncio.Queue(maxsize=QUEUE_MAXSIZE)
        self._subscribers.setdefault(user_id, set()).add(queue)
        try:
            self._loop = asyncio.get_running_loop()
        except RuntimeError:
            pass  # вне loop (не бывает для SSE-хендлера) — оставляем прежний
        return queue

    def unsubscribe(self, user_id: str, queue: asyncio.Queue[Event]) -> None:
        queues = self._subscribers.get(user_id)
        if queues is not None:
            queues.discard(queue)
            if not queues:
                del self._subscribers[user_id]

    def publish(self, user_id: str | None, kind: str, data: dict[str, Any] | None = None) -> None:
        """Положить событие всем подписчикам (None = всем). Потокобезопасно."""
        if not self._subscribers or self._loop is None:
            return  # подписчиков нет — клиент догонит refetch'ем при подключении
        event = Event(kind=kind, data=data)
        try:
            self._loop.call_soon_threadsafe(self._dispatch, user_id, event)
        except RuntimeError:
            # loop закрыт (перезапуск uvicorn --reload) — сбросим и ждём новую подписку
            self._loop = None

    def _dispatch(self, user_id: str | None, event: Event) -> None:
        if user_id is None:
            targets = list(self._subscribers.values())
        else:
            targets = [self._subscribers.get(user_id, set())]
        for queues in targets:
            for queue in queues:
                try:
                    queue.put_nowait(event)
                except asyncio.QueueFull:
                    logger.warning("SSE queue overflow for %r: event %r dropped", queue, event.kind)


bus = EventBus()


def publish(user_id: str | None, kind: str, data: dict[str, Any] | None = None) -> None:
    """Точка публикации для хендлеров/агентов (см. ТЗ 3.14)."""
    bus.publish(user_id, kind, data)
