# How Pi Web Manages SSE Streams for Data Integrity: A Deep Dive into the Event Streaming Architecture

> Discover how Pi Web ensures SSE stream data integrity with its three-layer architecture. Learn about buffering, handshake protocols, and automatic reconnection for reliable event streaming.

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

---

**Pi Web guarantees data integrity in SSE streams through a three-layer architecture that buffers events until snapshots are ready, enforces ordering with handshake protocols, and implements automatic reconnection with server-side reconciliation.**

The Pi Web repository (`agegr/pi-web`) implements a battle-tested Server-Sent Events (SSE) pipeline for streaming real-time agent events to the frontend. This article examines how the codebase prevents lost messages, eliminates duplicates, and maintains strict event ordering—even when connections drop or proxies interfere.

## The Three-Layer SSE Architecture

Pi Web's streaming stack consists of cooperating components across the server, connection manager, and UI:

| Layer | File | Core Responsibility |
|-------|------|---------------------|
| Transport | [`lib/agent-event-stream.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-stream.ts) | Creates `ReadableStream`, manages heartbeats, buffers pre-snapshot events |
| Connection | [`lib/agent-event-connection.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-connection.ts) | Wraps `EventSource`, handles ready-state handshake, retries with backoff |
| Consumer | [`hooks/useAgentSession.ts`](https://github.com/agegr/pi-web/blob/main/hooks/useAgentSession.ts) | Applies events to React state, runs grace periods, reconciles with server |

This separation ensures **no single point of failure** can compromise data integrity.

## Server-Side: Ordered Event Delivery with Snapshot Isolation

### The Buffered Event Queue

The `createAgentEventStream` function in [`lib/agent-event-stream.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-stream.ts) solves a critical race condition: events arriving before the session snapshot is ready.

```ts
// lib/agent-event-stream.ts (lines 73-81)
const bufferedEvents: AgentEventLike[] = [];
let snapshotPublished = false;

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

```

Events received before the snapshot resolve are **captured in `bufferedEvents`** rather than dropped or sent prematurely. Only after the `message_start` snapshot flushes does the server replay the buffer in original order:

```ts
// lib/agent-event-stream.ts (lines 95-100)
for (const event of bufferedEvents) forwardEvent(event, snapshot);
if (snapshot !== undefined && snapshot !== null) {
  encode({ type: "message_start", message: snapshot });
}
snapshotPublished = true;

```

This guarantees clients always receive **snapshot-first ordering**—the foundational invariant for consistent UI state.

### Connection Handshake and Heartbeats

Every SSE stream begins with a protocol handshake:

1. **Immediate heartbeat**: `:\n\n` sent to satisfy proxy keep-alive requirements
2. **Connected event**: `{"type":"connected",...}` signals readiness
3. **Snapshot + replayed buffer**: Complete initial state

```ts
// lib/agent-event-stream.ts (line 22)
const HEARTBEAT_INTERVAL_MS = 30_000;
// ...
heartbeat = setInterval(() => enqueueText(":\n\n"), HEARTBEAT_INTERVAL_MS);

```

The 30-second heartbeat interval prevents intermediate proxies (NGINX, Cloudflare, etc.) from closing idle connections.

## Connection Manager: Resilient Reconnection Without State Loss

`AgentEventConnection` in [`lib/agent-event-connection.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-connection.ts) transforms raw `EventSource` into a **reliability-layer** with typed errors and automatic recovery.

### Readiness Handshake

The connection manager rejects the promise until it receives the server's `connected` event:

```ts
// lib/agent-event-connection.ts (lines 49-55)
source.onmessage = (message) => {
  // ...
  if (event.type === "connected") {
    attempt.ready = true;
    attempt.succeed();
    this.stopRetrying();
  } else if (event.type === "startup_error") {
    // Permanent failure—don't retry
    this.stopRetrying();
  }
  this.options.onEvent(event);
};

```

### Configurable Timeouts and Backoff

```ts
// lib/agent-event-connection.ts (lines 36-41)
export interface AgentEventConnectionOptions {
  readinessTimeoutMs: number,  // default: 60_000 ms
  reconnectDelayMs: number,    // default: 1_000 ms
  // ...
}

```

- **Readiness timeout**: Aborts connection attempts that stall before handshake
- **Reconnect delay**: Prevents thundering herd scenarios during server recovery

On failure, the manager discards the dead connection and schedules retry—preserving the same `sessionId` so the server can resume the session:

```ts
// lib/agent-event-connection.ts (lines 66-73)
private fail(connection: Connection, error: AgentEventConnectionError): void {
  if (this.current !== connection) return; // Ignore stale failures
  this.discard(connection, error);
  if (error.status === "startup_error") this.stopRetrying();
  else this.scheduleRetry(connection.sessionId);
}

```

## UI Layer: Grace Periods and Server Reconciliation

The `useAgentSession` hook implements **client-side data integrity guards** that catch edge cases the transport layer cannot handle alone.

### The Grace Period Pattern

When `agent_end` arrives, the UI does **not** immediately terminate the stream. Instead, it starts a 30-second grace timer:

```ts
// hooks/useAgentSession.ts (lines 998-1004)
const scheduleEventStreamClose = useCallback((sid: string) => {
  cancelEventStreamGrace();
  eventStreamGraceActiveRef.current = true;
  // ...
  eventStreamGraceTimerRef.current = setTimeout(
    () => void checkServerIdle(), 
    EVENT_STREAM_IDLE_GRACE_MS  // 30_000 ms
  );
}, ...);

```

If the server later signals continued activity, the timer cancels. This prevents premature closure during race conditions between agent completion and final event delivery.

### Periodic State Reconciliation

While `agentRunning` is true, the hook polls the canonical server state every 15 seconds:

```ts
// hooks/useAgentSession.ts (lines 973-981)
useEffect(() => {
  if (!agentRunning) return;
  const reconcile = () => {
    const sid = sessionIdRef.current;
    if (sid) void reconcileAgentState(sid);
  };
  const interval = setInterval(reconcile, AGENT_STATE_RECONCILE_MS); // 15_000 ms
  return () => clearInterval(interval);
}, [agentRunning, reconcileAgentState]);

```

If polling reveals the server considers the agent idle while the UI still shows it running, the hook triggers the same completion path used for normal termination—**recovering from any missing SSE messages**.

### Stray Event Protection

Every handler validates against `agentRunningRef` before applying state changes:

```ts
// hooks/useAgentSession.ts (lines 1105-1110)
case "message_start":
case "message_update": {
  if (!agentRunningRef.current) break; // Discard stale events
  // ...apply to state
}

```

This prevents **late-arriving events** from corrupting a session the UI has already settled.

## End-to-End Event Flow

```

Client Request ──► Server creates ReadableStream
                      │
                      ├──► Immediate heartbeat (30s interval)
                      │
                      ├──► Await sessionPromise
                      │      │
                      │      ├──► Buffer incoming events
                      │      └──► On resolve: emit "connected"
                      │             │
                      │             ├──► Flush buffered events (ordered)
                      │             └──► Emit snapshot (message_start)
                      │
Client ◄────────── Stream established, ordered events flowing

     [Connection drops]
        │
        └──► AgentEventConnection detects readyState change
               │
               ├──► Discard connection
               ├──► Schedule retry (1s delay)
               └──► Reconnect with same sessionId
                      │
                      └──► Server resumes from session file

```

## Key Implementation Patterns for SSE Data Integrity

**Buffer-then-replay for ordering**: Events arriving before infrastructure readiness are queued and emitted only after the snapshot

**Handshake protocol for synchronization**: The `connected` event creates a clear boundary between transport establishment and application data

**Heartbeat for proxy compatibility**: Empty SSE comments prevent middlebox timeout without polluting the event stream

**Grace period for completion races**: Delayed closure accommodates out-of-order delivery of terminal events

**Reconciliation for eventual consistency**: Periodic polling provides a fallback when the persistent connection fails silently

## Summary

- **`createAgentEventStream`** ([`lib/agent-event-stream.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-stream.ts)) enforces **snapshot-first ordering** through pre-snapshot event buffering and in-order replay
- **`AgentEventConnection`** ([`lib/agent-event-connection.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-connection.ts)) provides **resilient reconnection** with readiness handshake, typed errors, and backoff
- **`useAgentSession`** ([`hooks/useAgentSession.ts`](https://github.com/agegr/pi-web/blob/main/hooks/useAgentSession.ts)) implements **UI-side integrity checks** including grace periods, reconciliation polling, and stale-event guards
- The 30-second heartbeat, 60-second readiness timeout, and 15-second reconciliation interval are tuned for **production proxy environments**

## Frequently Asked Questions

### How does Pi Web prevent events from arriving out of order?

The `createAgentEventStream` function in [`lib/agent-event-stream.ts`](https://github.com/agegr/pi-web/blob/main/lib/agent-event-stream.ts) maintains a `bufferedEvents` array that captures all events received before the session snapshot resolves. Only after emitting `message_start` does it iterate through the buffer in original insertion order, ensuring clients never receive incremental updates before their base state.

### What happens when an SSE connection drops mid-stream?

`AgentEventConnection` detects the failure through `EventSource` `onerror` handlers and `readyState` monitoring. It discards the dead connection, waits the configured `reconnectDelayMs` (default 1000ms), then opens a new `EventSource` with the identical `sessionId`. The server resumes streaming from its persisted session state, and any missed events are recovered through the reconciliation polling in `useAgentSession`.

### Why does Pi Web use both SSE and periodic HTTP polling?

SSE provides low-latency push delivery for the happy path, while reconciliation polling (every 15 seconds in `useAgentSession`) acts as a **watchdog** for silent failures. If the SSE connection drops without triggering error events—common with certain proxy configurations or background tab throttling—the polling discovers the discrepancy and forces state synchronization with the server.

### What is the purpose of the `agent_end` grace period?

The 30-second `EVENT_STREAM_IDLE_GRACE_MS` timer in [`hooks/useAgentSession.ts`](https://github.com/agegr/pi-web/blob/main/hooks/useAgentSession.ts) handles race conditions where `agent_end` arrives but the server later indicates continued activity. Without this grace period, the UI would immediately close the stream and potentially miss resurrection events. The timer cancels if `checkServerIdle` confirms activity, keeping the connection alive.