Server-Sent Events (SSE) Architecture in pi-web: Real-Time Agent Streaming Explained

The pi-web repository implements Server-Sent Events through a resilient streaming pipeline that combines a Next.js API route managing AgentSession lifecycle, a ReadableStream-based emitter with heartbeat keep-alives, and a client-side connection manager handling automatic reconnection and event filtering.

Pi-web streams real-time updates from Pi agents to the browser using Server-Sent Events (SSE). The architecture separates concerns between session management, wire protocol optimization, and connection resilience to deliver low-latency updates that survive network glitches and browser tab backgrounding. This implementation leverages native browser EventSource APIs while adding a robust handshake and retry layer.

Server-Side SSE Implementation

The server architecture centers on a dynamic API route that either resurrects existing agent sessions or initializes new ones before streaming events through a carefully managed ReadableStream.

API Route and Session Management

The entry point at app/api/agent/[id]/events/route.ts handles HTTP upgrade negotiations and session lifecycle decisions. Rather than creating redundant sessions, it first attempts to reuse an existing AgentSession via getRpcSession(id), falling back to startRpcSession() only when necessary.

// app/api/agent/[id]/events/route.ts
export async function GET(req: Request, { params }: { params: Promise<{ id: string }> }) {
  const { id } = await params;
  if (req.signal.aborted) return new Response(null, { status: 204 });

  // Reuse existing session or start fresh RPC session
  const session = getRpcSession(id);
  const sessionPromise = session?.isAlive()
    ? Promise.resolve(session)
    : startRpcSession(id, await resolveSessionPath(id), undefined).then(r => r.session);

  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",
    },
  });
}

The route returns a streaming Response with critical SSE headers including X-Accel-Buffering: no to prevent proxy buffering and Connection: keep-alive to maintain the TCP socket.

Streaming Helper and Heartbeat Mechanism

The createAgentEventStream function in lib/agent-event-stream.ts constructs a ReadableStream<Uint8Array> that manages the complex dance of event buffering, snapshot delivery, and connection hygiene. It implements a heartbeat mechanism sending :\n\n (SSE comment frames) every few seconds to prevent timeout closures by intermediate proxies.

The stream buffers events that arrive before the initial snapshot is ready, ensuring no data loss during the handshake window. Once the session snapshot publishes, buffered events flush to the client followed by live updates.

// lib/agent-event-stream.ts
const publishSession = async () => {
  const session = await sessionPromise;
  const bufferedEvents: AgentEventLike[] = [];
  let snapshotPublished = false;

  const handleEvent = (event: AgentEventLike) => {
    if (!snapshotPublished) {
      bufferedEvents.push(event);
      return;
    }
    forwardEvent(event, snapshot);
  };

  const stopListening = session.onEvent(handleEvent);
  
  // Send handshake and flush buffer
  const snapshot = session.streamingMessage;
  encode({ type: "connected", sessionId, isStreaming: session.isStreaming });
  for (const ev of bufferedEvents) forwardEvent(ev, snapshot);
  if (snapshot != null) encode({ type: "message_start", message: snapshot });
  snapshotPublished = true;
};

Client-Side Architecture

The browser implementation abstracts raw EventSource complexity behind a connection manager that implements exponential backoff reconnection and strict event filtering.

Connection Manager and Handshake

lib/agent-event-connection.ts exports the AgentEventConnection class, which wraps the native EventSource and enforces a protocol-level handshake. The client waits for a connected event before resolving the connection promise, enabling the UI to distinguish between transport-layer connection and application-layer readiness.

// lib/agent-event-connection.ts
export class AgentEventConnection {
  async ensureConnected(sessionId: string): Promise<void> {
    while (true) {
      let connection = this.current;
      if (!connection || connection.sessionId !== sessionId) {
        connection = this.open(sessionId);
      }
      await connection.attempt.promise;
      if (this.current?.sessionId === sessionId && 
          connection.source.readyState === EVENT_SOURCE_OPEN) return;
      this.discard(connection, new AgentEventConnectionError("closed"));
    }
  }

  private open(sessionId: string): Connection {
    const source = this.options.createSource(sessionId);
    source.onmessage = (msg) => {
      const event: AgentEventLike = JSON.parse(msg.data);
      if (event.type === "connected") {
        attempt.ready = true;
        attempt.succeed();
        this.stopRetrying();
      } else if (event.type === "startup_error") {
        this.fail(connection, new AgentEventConnectionError("startup_error"));
        return;
      }
      this.options.onEvent(event);
    };
    source.onerror = () => this.fail(connection, new AgentEventConnectionError("closed"));
    return { sessionId, source, attempt };
  }
}

The ensureConnected method implements an infinite retry loop with configurable delays (reconnectDelayMs), automatically re-establishing the SSE pipe if the server closes the connection or the network drops.

Event Filtering and Wire Protocol

Before events reach React state, lib/agent-event-wire.ts filters raw Pi SDK events through toClientAgentEvent(). This function strips internal events listed in OMITTED_EVENT_TYPES and normalizes message_update events to remove the partial field from delta updates, reducing payload size and simplifying UI logic.

// lib/agent-event-wire.ts
export function toClientAgentEvent(event: AgentEventLike): AgentEventLike | ClientMessageUpdateEvent | null {
  if (OMITTED_EVENT_TYPES.has(event.type)) return null;

  if (event.type === "message_update") {
    const { assistantMessageEvent } = event;
    if (assistantMessageEvent && !("partial" in assistantMessageEvent)) {
      return { type: "message_update", assistantMessageEvent };
    }
    // Strip partial field for delta updates
    const { partial, ...delta } = assistantMessageEvent as any;
    return { type: "message_update", assistantMessageEvent: delta };
  }

  if (event.type === "agent_end") return { type: "agent_end" };
  return event;
}

React Hook Integration

The useAgentSession hook in hooks/useAgentSession.ts instantiates a persistent AgentEventConnection using a useRef pattern to survive React re-renders. It configures the connection with a shouldMaintain predicate that evaluates session state, mounting status, and grace periods to determine whether to keep the SSE alive or allow it to close.

// hooks/useAgentSession.ts
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)
  ),
  readinessTimeoutMs: EVENT_STREAM_READY_TIMEOUT_MS,
  reconnectDelayMs: EVENT_STREAM_RECONNECT_DELAY_MS,
});

End-to-End Event Flow

Understanding the complete data path clarifies how pi-web maintains synchronization between the agent process and browser UI:

  1. Browser instantiates EventSource pointing to /api/agent/:id/events
  2. API route locates or starts an AgentSession, passing the promise to createAgentEventStream
  3. ReadableStream sends the connected handshake event, flushes buffered history, then streams live message_update events
  4. Client connection resolves its handshake promise upon receiving connected, then dispatches filtered events via onEvent
  5. React hook updates local state, triggering UI re-renders with new agent messages
  6. Automatic recovery activates if the connection drops, with AgentEventConnection reinitializing the EventSource after a backoff delay

Summary

  • Session Reuse: The SSE endpoint at app/api/agent/[id]/events/route.ts checks getRpcSession() to avoid spinning duplicate agent processes when refreshing the browser.
  • Reliable Streaming: createAgentEventStream implements heartbeat comments (:\n\n) and event buffering to prevent proxy timeouts and data loss during initialization.
  • Handshake Protocol: AgentEventConnection waits for an explicit connected event before resolving, distinguishing transport connectivity from application readiness.
  • Wire Optimization: toClientAgentEvent in lib/agent-event-wire.ts filters SDK internals and strips partial fields to minimize bandwidth.
  • Resilient Clients: Exponential backoff reconnection and the shouldMaintain lifecycle predicate ensure connections survive network interruptions without zombie processes.

Frequently Asked Questions

How does pi-web handle SSE reconnection when the network drops?

The AgentEventConnection class implements an infinite retry loop in ensureConnected() that detects EventSource errors via the onerror handler. Upon disconnection, it discards the dead connection, waits for a configurable delay (reconnectDelayMs), and instantiates a new EventSource. The client buffers no events locally during disconnection; instead, it relies on the server-side session state and periodic reconciliation to recover missed updates.

What is the purpose of the connected event in pi-web's SSE implementation?

The connected event serves as an application-layer handshake. While the browser's EventSource may report OPEN state immediately, pi-web's server only emits connected after the AgentSession is fully initialized and ready to stream. This prevents race conditions where the client might process message_update events before receiving the initial message_start snapshot, ensuring UI consistency.

How does pi-web filter events before sending them to the browser?

Raw events from the Pi SDK pass through toClientAgentEvent() in lib/agent-event-wire.ts, which blocks events listed in OMITTED_EVENT_TYPES and transforms message_update payloads. Specifically, it removes the partial field from delta updates to reduce JSON size and filters out internal system events that lack UI relevance. This occurs server-side, minimizing bandwidth and client processing overhead.

What keeps the SSE connection alive during idle periods?

The server sends periodic heartbeat frames consisting of SSE comment lines (:\n\n) every few seconds via the createAgentEventStream heartbeat interval. These comments keep the TCP connection warm and prevent intermediate proxies or load balancers from closing idle sockets. Additionally, the response headers include X-Accel-Buffering: no to disable Nginx buffering, ensuring frames flush immediately to the wire rather than accumulating in server buffers.

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 →