How SSE Event Streaming Works for Real-Time Agent Communication in Pi Web
Pi Web uses a three-layer Server-Sent Events (SSE) pipeline—comprising an API route, a stream factory, and a client-side connection manager—to stream real-time agent events from the server to the browser.
The agegr/pi-web repository implements a robust real-time communication system between its browser UI and backend agent sessions. This article breaks down the complete SSE implementation, from the HTTP endpoint that initiates the stream to the React hooks that consume it.
The Three-Component SSE Architecture
Pi Web's SSE pipeline consists of three tightly-coupled pieces working in sequence:
- API route (
app/api/agent/[id]/events/route.ts) – opens the SSE endpoint and hands aReadableStreamto the HTTP response - Stream factory (
lib/agent-event-stream.ts) – manages heartbeats, session readiness, and event buffering - Client connection manager (
lib/agent-event-connection.ts+hooks/useAgentSession.ts) – handlesEventSourcelifecycle, reconnection, and UI state updates
Step 1: Opening the SSE Endpoint on the Server
The SSE connection begins when the browser requests /api/agent/<sessionId>/events. The route handler in app/api/agent/[id]/events/route.ts manages session lookup and stream initialization:
// app/api/agent/[id]/events/route.ts
export async function GET(req: Request, { params }: { params: Promise<{ id: string }> }) {
const { id } = await params;
// Fast-path: reuse an already-alive AgentSessionWrapper
const session = getRpcSession(id);
const sessionPromise = session?.isAlive()
? Promise.resolve(session)
: startRpcSession(id, await resolveSessionPath(id), undefined).then(r => r.session);
// Create the SSE stream that will push events to the client
const stream = createAgentEventStream(req, id, sessionPromise);
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
},
});
}
Key behaviors:
- Fast path – if
AgentSessionWrapperalready exists viagetRpcSession(id), the route reuses it - Cold start – otherwise
startRpcSessionloads the.jsonlfile and creates a new wrapper - SSE headers –
X-Accel-Buffering: nodisables nginx buffering;Cache-Controlprevents proxy caching
Step 2: Creating the Event Stream with Heartbeats and Buffering
The createAgentEventStream function in lib/agent-event-stream.ts constructs a ReadableStream<Uint8Array> that follows the SSE wire format (data: …\n\n). Its lifecycle manages connection health and event ordering:
// lib/agent-event-stream.ts
export function createAgentEventStream(
req: Request,
sessionId: string,
sessionPromise: Promise<AgentEventStreamSession>,
): ReadableStream<Uint8Array> {
// ...setup, heartbeat, cleanup...
const publishSession = async () => {
const session = await sessionPromise; // ← wait for the AgentSessionWrapper
const bufferedEvents: AgentEventLike[] = [];
let snapshotPublished = false;
const handleEvent = (event: AgentEventLike) => {
if (!snapshotPublished) bufferedEvents.push(event); // buffer until snapshot
else forwardEvent(event, snapshot);
};
const stopListening = session.onEvent(handleEvent);
// …store unsubscribe for later cleanup…
// Send the *connected* envelope first
encode({ type: "connected", sessionId, isStreaming: session.isStreaming });
// Flush any events that arrived before snapshot was ready
for (const ev of bufferedEvents) forwardEvent(ev, snapshot);
if (snapshot !== undefined && snapshot !== null) {
encode({ type: "message_start", message: snapshot });
}
snapshotPublished = true;
};
// Kick off the async publish logic
void publishSession();
// Heartbeat every 30 s to keep the HTTP connection alive
heartbeat = setInterval(() => enqueueText(":\n\n"), HEARTBEAT_INTERVAL_MS);
// ...cancel, abort handling, etc.
}
Critical implementation details:
- Buffering guarantees – events arriving before session readiness are stored in
bufferedEventsand replayed after the snapshot, ensuring no data loss - Heartbeat mechanism – a colon line (
:\n\n) sent every 30 seconds prevents proxy timeouts on idle connections - Startup error reporting – failures to initialize the session emit a
startup_errorevent that the client surfaces to users
Step 3: Managing the Client-Side Connection
The AgentEventConnection class in lib/agent-event-connection.ts wraps the native EventSource with robust retry logic:
// lib/agent-event-connection.ts
export class AgentEventConnection {
// ...
async ensureConnected(sessionId: string): Promise<void> {
while (true) {
// Open a new EventSource if none exists or the id changed
if (!this.current || this.current.sessionId !== sessionId) this.current = this.open(sessionId);
// Wait for the connection's "ready" promise (the `connected` event)
await this.current.attempt.promise;
if (this.current.source.readyState === EVENT_SOURCE_OPEN) return;
// If the source stays CONNECTING we discard it and retry
this.discard(this.current, new AgentEventConnectionError("closed"));
}
}
private open(sessionId: string): Connection {
// Build the EventSource pointing at the SSE endpoint
const source = this.options.createSource(sessionId);
// Attach message and error handlers
source.onmessage = (msg) => {
const event = JSON.parse(msg.data) as AgentEventLike;
if (event.type === "connected") { /* mark ready */ }
else if (event.type === "startup_error") { /* fail */ }
this.options.onEvent(event);
};
source.onerror = () => this.fail(connection, new AgentEventConnectionError("closed"));
// ...
}
}
Connection management features:
- Readiness handshake – the first
connectedenvelope resolves the internalreadypromise, signaling that the agent is live - Exponential backoff reconnection – on any error, the class schedules retry after
reconnectDelayMs - Graceful shutdown – when
agent_endoragent_settledarrives, a grace timer (EVENT_STREAM_IDLE_GRACE_MS) closes the connection after inactivity
Step 4: Integrating with React via useAgentSession
The useAgentSession.ts hook instantiates a single AgentEventConnection and orchestrates its lifecycle with UI state:
// hooks/useAgentSession.ts (excerpt)
eventConnectionRef.current = new AgentEventConnection({
createSource: (sid) => new EventSource(`/api/agent/${encodeURIComponent(sid)}/events`),
onEvent: (event) => handleAgentEventRef.current?.(event as AgentEvent),
shouldMaintain: (sid) => (
sessionHookMountedRef.current &&
sessionIdRef.current === sid &&
(agentRunningRef.current ||
eventStreamGraceActiveRef.current ||
(sessionPropIdRef.current === sid && sessionRunningRef.current))
),
readinessTimeoutMs: EVENT_STREAM_READY_TIMEOUT_MS,
reconnectDelayMs: EVENT_STREAM_RECONNECT_DELAY_MS,
onUnexpectedError: (error) => console.error("Failed to maintain the agent event stream:", error),
});
Hook responsibilities:
shouldMaintainpredicate – mirrors the server-side "running-state" poll to keep streams alive only during active sessionshandleAgentEventdispatch – translates raw events (agent_start,message_start,message_update, etc.) into React state updates- State reconciliation – background polls every
AGENT_STATE_RECONCILE_MSplus visibility/online listeners recover from missed events
Example: Handling a New Assistant Message in the UI
When the SSE stream delivers a message_start event, the UI updates through the streaming reducer:
// inside handleAgentEvent()
case "message_start":
const msg = event.message as AssistantMessage;
dispatch({ type: "snapshot", message: msg }); // feed streaming reducer
setAgentPhase(null); // hide "waiting model" spinner
break;
Complete Event Flow
| Phase | Action | Component |
|---|---|---|
| 1 | User opens chat | useAgentSession ensures AgentEventConnection exists |
| 2 | Browser creates EventSource |
Points to /api/agent/:id/events |
| 3 | Server returns ReadableStream |
Immediate heartbeat, then waits for AgentSessionWrapper |
| 4 | Session emits events | AgentSessionWrapper broadcasts via onEvent listeners |
| 5 | Stream factory forwards events | createAgentEventStream encodes SSE frames, flushes buffer |
| 6 | Client updates UI | React state changes drive real-time chat rendering |
| 7 | Idle detection | Grace period expires, connection closes |
Summary
- SSE endpoint (
app/api/agent/[id]/events/route.ts) – reuses or creates agent sessions, returns properly headersReadableStream - Stream factory (
lib/agent-event-stream.ts) – implements heartbeats, event buffering, and ordered delivery guarantees - Client manager (
lib/agent-event-connection.ts) – wrapsEventSourcewith handshake, retry, and graceful shutdown logic - React integration (
hooks/useAgentSession.ts) – maintains connection lifecycle, reconciles state, and maps events to UI updates
Frequently Asked Questions
What prevents SSE connection drops in Pi Web?
The createAgentEventStream function sends a colon heartbeat (:\n\n) every 30 seconds via setInterval. This satisfies proxy idle timeouts while adding minimal overhead. The Cache-Control and X-Accel-Buffering headers also disable buffering that could delay event delivery.
How does Pi Web handle events that arrive before the stream is ready?
Events emitted during session initialization are captured in a bufferedEvents array. Once the AgentSessionWrapper resolves and the snapshot publishes, the buffer flushes in order before new events forward. This guarantees the client receives complete, sequentially-ordered event history.
What triggers reconnection in the client-side SSE handler?
The AgentEventConnection class reconnects on any source.onerror event, when readyState remains CONNECTING past the readiness timeout, or when shouldMaintain returns true but the connection closed unexpectedly. Reconnection uses a configurable delay (EVENT_STREAM_RECONNECT_DELAY_MS, default 1000ms).
How does the UI know when to close the SSE connection?
The shouldMaintain predicate in useAgentSession evaluates multiple conditions: hook mount status, current session ID match, agentRunning state, grace period activity, and sidebar running status. When all conditions fail—typically after agent_end or agent_settled plus the idle grace period—the connection closes to free resources.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →