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

> Explore Server-Sent Events SSE architecture in pi-web. Learn how Next.js, ReadableStream, and client-side management power real-time agent streaming with automatic reconnection.

- Repository: [Alex Yang/pi-web](https://github.com/agegr/pi-web)
- Tags: architecture
- Published: 2026-08-13

---

**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.

```typescript
// 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`](https://github.com/agegr/pi-web/blob/main/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.

```typescript
// 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`](https://github.com/agegr/pi-web/blob/main/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.

```typescript
// 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`](https://github.com/agegr/pi-web/blob/main/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.

```typescript
// 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`](https://github.com/agegr/pi-web/blob/main/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.

```typescript
// 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`](https://github.com/agegr/pi-web/blob/main/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`](https://github.com/agegr/pi-web/blob/main/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.