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:
- Immediate heartbeat:
:\n\nsent to satisfy proxy keep-alive requirements - Connected event:
{"type":"connected",...}signals readiness - 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 replayAgentEventConnection(lib/agent-event-connection.ts) provides resilient reconnection with readiness handshake, typed errors, and backoffuseAgentSession(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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →