Skip to content
src.fastware.sse
Edit
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).

More tools from this site

  • claudestream Drive Claude Code from Python: run it as a subprocess and read its output as typed events, with async and sync sessions, sandbox policies, and tools you define in Python
  • claudewheel A TUI Claude Code Launcher that lets you have more than one profile, manage sessions lifecycle, pick the exact CC version, model to use (even older unlisted ones), pick which GitHub account to use, etc.
  • dirstat Fast, single-binary directory statistics CLI: every file under a tree grouped by format, with counts, sizes, and lines of code, as a colored terminal table or as JSON
  • go-toml-edit Zero-dep TOML editing library for Go with comment preservation
  • howmuchleft The fastest Claude Code statusline: context window, 5-hour, and weekly limit usage as three customizable gradient bars, rendering in about 6 ms
  • orxtra
  • pgdesign
  • predraw Declarative rendering pipeline: describe a scene in JSON and get SVG, PNG and WebP out, with light and dark style tokens, reusable components and text converted to path outlines
  • reposummary Turn a git repository's history into a Markdown journal: pick a time window or revision range and get a readable digest of what changed, optionally narrated by an LLM
  • rlsbl Release orchestration and project scaffolding CLI that bumps versions, validates a structured JSONL changelog, tags only the commit CI verified, and publishes to npm, PyPI, Go and more
  • safegit git wrapper CLI that gives each commit its own temporary index and retries ref updates on conflict, so concurrent agents share one repository
  • saferm Command-line replacement for rm that archives every deletion with a mandatory reason and the context it ran in, so deleted files can be listed, inspected and restored
  • selfdoc Static Site Generator that builds a project's documentation site directly from its source code, so the docs can never drift from the code they describe, with SEO/AEO, first-class blog, search, and cross-project linking built in
  • strictcli
  • stricttest An always-on test-isolation floor: a pytest plugin and a Go env-hygiene module that make a test suite structurally unable to reach real credentials, the real HOME, the network, or the development repository.
  • wesktop A Python framework that turns an ASGI web app into a desktop application, serving it from a local Granian server and displaying it in a native OS window via pywebview
Search