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:

  1. API route (app/api/agent/[id]/events/route.ts) – opens the SSE endpoint and hands a ReadableStream to the HTTP response
  2. Stream factory (lib/agent-event-stream.ts) – manages heartbeats, session readiness, and event buffering
  3. Client connection manager (lib/agent-event-connection.ts + hooks/useAgentSession.ts) – handles EventSource lifecycle, 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 AgentSessionWrapper already exists via getRpcSession(id), the route reuses it
  • Cold start – otherwise startRpcSession loads the .jsonl file and creates a new wrapper
  • SSE headers – X-Accel-Buffering: no disables nginx buffering; Cache-Control prevents 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 bufferedEvents and 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_error event 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 connected envelope resolves the internal ready promise, signaling that the agent is live
  • Exponential backoff reconnection – on any error, the class schedules retry after reconnectDelayMs
  • Graceful shutdown – when agent_end or agent_settled arrives, 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:

  • shouldMaintain predicate – mirrors the server-side "running-state" poll to keep streams alive only during active sessions
  • handleAgentEvent dispatch – translates raw events (agent_start, message_start, message_update, etc.) into React state updates
  • State reconciliation – background polls every AGENT_STATE_RECONCILE_MS plus 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 headers ReadableStream
  • Stream factory (lib/agent-event-stream.ts) – implements heartbeats, event buffering, and ordered delivery guarantees
  • Client manager (lib/agent-event-connection.ts) – wraps EventSource with 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →