fastware v0.6.0 /src.fastware.sse
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

python
def register_event(self, name: str) -> None

Declare an allowed event type.

#event_types

python
def event_types(self) -> frozenset[str]

Currently registered event types.

#_format_sse

python
def _format_sse(self, event: str, data: dict[str, Any] | str) -> str

Format 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

python
def broadcast(self, event: str, data: dict[str, Any] | str) -> None

Send 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

python
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

python
async def stream(self, request: Request, initial: list[tuple[str, dict[str, Any] | str]] | None=None) -> StreamResponse

Return 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

python
def client_count(self) -> int

Number of currently connected SSE clients.

#sse_route

python
def sse_route(broadcaster: Broadcaster)

Return an async handler suitable for router.add_route("GET", "/events", handler).

Search