"""
Simple in-process pub/sub for Server-Sent Events (SSE).

Each SSE client calls subscribe() to get a Queue, then polls it.
publish() fans out events to all active queues.
Slow/full queues are silently dropped to avoid blocking the publisher.
"""

import asyncio
from typing import Any

_subscribers: list[asyncio.Queue] = []


def subscribe() -> asyncio.Queue:
    """Register a new SSE client. Returns a queue to read events from."""
    q: asyncio.Queue = asyncio.Queue(maxsize=20)
    _subscribers.append(q)
    return q


def unsubscribe(q: asyncio.Queue) -> None:
    """Remove a client queue (called on disconnect)."""
    try:
        _subscribers.remove(q)
    except ValueError:
        pass


def publish(event: dict[str, Any]) -> None:
    """Broadcast an event to all connected SSE clients.
    Clients whose queue is full are removed (assumed stale/slow)."""
    dead = []
    for q in _subscribers:
        try:
            q.put_nowait(event)
        except asyncio.QueueFull:
            dead.append(q)
    for q in dead:
        unsubscribe(q)
