# How SSE Event Streaming Works for Real-Time Agent Communication in Pi Web

> Discover how Pi Web's SSE event streaming pipeline enables real-time agent communication. Learn about its API route, stream factory, and client connection manager.

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

---

**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`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-stream.ts)) – manages heartbeats, session readiness, and event buffering
3. **Client connection manager** ([`lib/agent-event-connection.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-connection.ts) + [`hooks/useAgentSession.ts`](https://github.com/agegr/pi-web/blob/main/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:

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

```ts
// 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`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-connection.ts) wraps the native `EventSource` with robust retry logic:

```ts
// 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`](https://github.com/agegr/pi-web/blob/main/useAgentSession.ts) hook instantiates a single `AgentEventConnection` and orchestrates its lifecycle with UI state:

```ts
// 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:

```tsx
// 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`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-stream.ts)) – implements heartbeats, event buffering, and ordered delivery guarantees
- **Client manager** ([`lib/agent-event-connection.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-connection.ts)) – wraps `EventSource` with handshake, retry, and graceful shutdown logic
- **React integration** ([`hooks/useAgentSession.ts`](https://github.com/agegr/pi-web/blob/main/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.