On this page
SSE Broadcaster with typed event registration, per-client async queues, automatic disconnect pruning, optional heartbeat, and the sse_route helper.
#src.fastware.sse
#src.fastware.sse
SSE (Server-Sent Events) broadcaster with typed event registration, per-client async queues, automatic disconnect pruning, and strict mode enforcement.
#Broadcaster
Manages SSE client connections and broadcasts typed events.
Event types must be registered via register_event before they can be broadcast. In strict mode (the default), broadcasting an unregistered event raises ValueError. Pass strict=False to skip validation.
#register_event
def register_event(self, name: str) -> NoneDeclare an allowed event type.
#event_types
def event_types(self) -> frozenset[str]Currently registered event types.
#_format_sse
def _format_sse(self, event: str, data: dict[str, Any] | str) -> strFormat a payload as an SSE wire message.
Dict payloads are serialized with msgspec (project convention). A multi-line payload is emitted as one data: line per line, per the SSE spec, so a stray newline in the payload cannot terminate the event early or inject additional SSE fields.
#broadcast
def broadcast(self, event: str, data: dict[str, Any] | str) -> NoneSend an event to all connected clients.
Prunes clients whose queues are full (they fell behind and are presumed disconnected or stuck).
Raises ValueError if event was not previously registered and the broadcaster is in strict mode.
#_event_generator
async def _event_generator(self, queue: asyncio.Queue[str], initial: list[tuple[str, dict[str, Any] | str]] | None=None) -> AsyncGenerator[str, None]Yield SSE messages from a per-client queue.
The queue is registered as a client only once iteration begins, and the finally block guarantees it is unregistered when the generator is closed (e.g. client disconnect). Registering here — rather than in stream() — ensures a StreamResponse whose body is never consumed does not leak a queue into self._clients.
initial events are formatted and yielded once, before the queue loop, so a connection can be primed with current state (e.g. the current build id on the update channel). They bypass the strict registration check -- the caller controls them, not a broadcast.
When heartbeat_interval is set, yields SSE comment heartbeats (": heartbeat\n\n") if no real message arrives within the interval.
#stream
async def stream(self, request: Request, initial: list[tuple[str, dict[str, Any] | str]] | None=None) -> StreamResponseReturn a StreamResponse for an SSE endpoint.
Creates a per-client queue and wraps the async generator in the framework's streaming response type. The queue is registered as a client by _event_generator when iteration starts, not here, so an unconsumed response never leaks a queue.
initial is an optional list of (event, data) pairs sent to this connection before any broadcast, so a client can be primed with the current state on connect.
#client_count
def client_count(self) -> intNumber of currently connected SSE clients.
#sse_route
def sse_route(broadcaster: Broadcaster)Return an async handler suitable for router.add_route("GET", "/events", handler).