SSE Sidecar — stream AgentEvents to a browser¶
Many Orqest apps need a "side channel" that pushes observability data
(plan updates, memory recalls, trace spans) to the frontend without
cluttering the main chat stream. sse_sidecar subscribes to an
EventBus and produces SSE-formatted chunks you can hand straight to
any ASGI framework's streaming response.
The result: a single /events endpoint that streams everything a
long-running agent is doing, decoupled from the chat protocol and
immune to the Vercel AI data-stream shape.
Minimal FastAPI integration¶
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from orqest.observability import EventBus, sse_sidecar
app = FastAPI()
bus = EventBus() # shared with your agent's HookRunner
@app.get("/events")
async def events() -> StreamingResponse:
return StreamingResponse(
sse_sidecar(bus, heartbeat_s=15.0),
media_type="text/event-stream",
)
On the browser side, a plain EventSource('/events') receives one
message per AgentEvent. The SSE event: field holds
AgentEvent.event_type, the data: field holds a JSON payload.
Replay on reconnect¶
Pass replay=[...] to emit historical events before live streaming
starts — useful for refresh/reload scenarios where the client has lost
its in-memory view.
from collections import deque
buffer: deque[AgentEvent] = deque(maxlen=200)
bus.subscribe_all(buffer.append)
@app.get("/events")
async def events() -> StreamingResponse:
return StreamingResponse(
sse_sidecar(bus, replay=list(buffer)),
media_type="text/event-stream",
)
Heartbeats + backpressure¶
heartbeat_s(default 15.0) emits an SSE comment line every N seconds of idle time. Prevents proxies from closing the stream.queue_size(default 256) caps in-flight events per consumer. When a slow consumer fills the queue, the oldest event is dropped to make room for the newest — the stream never stalls the publisher.
Reference¶
sse_sidecar ¶
Stream :class:AgentEvent\s as Server-Sent Events.
sse_sidecar subscribes to an :class:EventBus, converts every event
it receives into an SSE-formatted string chunk, and yields those chunks
as an async iterator suitable for StreamingResponse in FastAPI or
any ASGI framework.
Event format follows the SSE spec:
event: <AgentEvent.event_type>
id: <iso-timestamp>-<counter>
data: <json-serialized AgentEvent payload>
A heartbeat comment (: keep-alive\n\n) is emitted every
heartbeat_s seconds to keep intermediaries from closing idle
connections. Optional replay of recent events lets a reconnecting
client catch up on what it missed.
sse_sidecar ¶
sse_sidecar(
bus: EventBus,
*,
replay: Iterable[AgentEvent] = (),
heartbeat_s: float = 15.0,
queue_size: int = 256,
) -> AsyncIterator[str]
Return an async iterator yielding SSE chunks for every event on bus.
Subscription to the bus happens eagerly (before the returned iterator is consumed) so no events published after this call are lost, even if the caller awaits something else before starting to iterate.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
bus
|
EventBus
|
Source :class: |
required |
replay
|
Iterable[AgentEvent]
|
Historical events to emit before live streaming starts. Useful for letting a reconnecting client re-hydrate its view. |
()
|
heartbeat_s
|
float
|
Interval in seconds for keep-alive comments. Set
large or |
15.0
|
queue_size
|
int
|
Bound on the in-flight event queue. The bus producer never blocks: when the queue is full the oldest event is evicted to make room for the newest, and if eviction races another producer the new event is dropped silently. This protects the producer at the cost of event loss under sustained backpressure — a bounded buffer with overflow drop, not a guaranteed-delivery ring buffer. Size up if your consumer is slow enough that loss matters, or wire a slower-but-reliable transport instead. |
256
|
Returns:
| Type | Description |
|---|---|
AsyncIterator[str]
|
SSE-formatted async iterator. Pass directly to an ASGI |
AsyncIterator[str]
|
|