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

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 Creates ReadableStream, manages heartbeats, buffers pre-snapshot events
Connection lib/agent-event-connection.ts Wraps EventSource, handles ready-state handshake, retries with backoff
Consumer 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 solves a critical race condition: events arriving before the session snapshot is ready.

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

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

// 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

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

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

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

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

// 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) enforces snapshot-first ordering through pre-snapshot event buffering and in-order replay
  • AgentEventConnection (lib/agent-event-connection.ts) provides resilient reconnection with readiness handshake, typed errors, and backoff
  • useAgentSession (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 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 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.

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 →