How OpenMAIC Handles Real-Time SSE Delivery: A Three-Pipeline Architecture
OpenMAIC delivers real-time updates to the browser using three independent Server-Sent Events (SSE) pipelines that leverage the native EventSource API, with built-in replay buffers, health monitoring, and periodic reconciliation to ensure reliable live data synchronization.
OpenMAIC (THU-MAIC/OpenMAIC) streams live data from the back-end to the front-end using a low-latency "push-for-latency, pull-for-correctness" model built on Server-Sent Events (SSE). The architecture splits real-time SSE delivery across three distinct pipelines—session chat, course freshness, and owner-session streams—each optimized for specific data consistency requirements and fault tolerance.
Session Chat Stream: Real-Time Agent Communication
The session chat pipeline feeds every chat-related event—including user questions, tool execution updates, and lifecycle signals—to the workbench UI via the useWorkbenchStream hook in lib/workbench/use-workbench-session.ts.
EventSource Lifecycle and Resume Capability
When a session ID is present, the hook creates an EventSource pointing at /api/agent/sessions/:id/events?lastEventId=<resumePoint>. The lastEventId parameter is read from the Redux-style useWorkbenchStore, ensuring the client resumes exactly where it left off after a disconnect:
const source = new EventSource(
`/api/agent/sessions/${encodeURIComponent(sessionId)}/events?lastEventId=${from}`,
);
In lib/workbench/use-workbench-session.ts (lines 53-66), the hook subscribes to all event types enumerated in WORKBENCH_EVENT_TYPES—a combination of lifecycle constants, legacy names, and PI-specific types—by attaching listeners via source.addEventListener(type, onAny).
Replay Buffer and Deterministic Ordering
The server emits a custom caught_up event once the backlog has been transmitted. Until this signal arrives, incoming events accumulate in a backlog buffer. When caught_up fires, the hook compacts the buffer using compactReplayEvents and applies the changes in bulk (lines 119-131), guaranteeing that the UI sees a deterministic, ordered replay rather than partial updates.
Error Handling and State Guards
The hook monitors source.onerror to detect when the underlying connection closes (readyState === EventSource.CLOSED). When closed, it surfaces a hard error via setError('event stream closed'); otherwise, it relies on the native auto-reconnect logic. Stale closure guards using isCurrent() checks compare the current session ID to the store’s state before applying updates (lines 144-152).
Course Freshness Stream: Incremental Canvas Updates
Rather than pulling entire course documents on every change, the canvas uses useStageFreshnessSync in lib/workbench/use-stage-freshness-sync.ts to receive incremental updates via the /api/stages/:id/freshness endpoint.
Selective Synchronization with Manifest Diffing
The stream emits stage_freshness frames whenever the back-end detects a manifest change. Each "pass" fetches the manifest, diffs it against the last rendered version, and narrow-fetches only the changed scenes. The pass aborts if a non-empty batch fetch returns no scenes, preventing the system from silently marking changes as already applied (lines 85-94).
Resilience and Fallback Mechanisms
The sync triggers on multiple conditions: SSE messages, tab visibility changes (visibilitychange → focus), and a fallback timer firing every 30 seconds with jitter. A retry budget (STAGE_SYNC_MAX_RETRIES) with exponential back-off handles transient failures, while the fallback timer guarantees eventual convergence even if the SSE connection drops (lines 158-167).
Owner-Session Stream: Home Rail Synchronization
The OwnerSessionClient class in lib/workbench/owner-session-client.ts manages the owner-level session list, pushing creation events, status updates, and title edits to the "Home" rail via /api/owner/sessions/events.
Closed-Set Event Types for Forward Compatibility
The client defines OWNER_SESSION_EVENT_TYPES as a closed set (session creation, status, deletion, active stage, cancel request, title). Because the native EventSource only delivers events matching registered listeners, any new server-side event name remains invisible until explicitly added to this list, guaranteeing forward compatibility (lines 10-21).
Health Monitoring and Full Reconciliation
Every 5 seconds, the client samples the readyState. If the stream stays in CONNECTING for nine consecutive samples (approximately 45 seconds), it flags the channel as degraded via onStreamHealth. Additionally, every 60-75 seconds (with jitter), the client performs a snapshot fetch (requestFullFetch) to reconcile missed events, ensuring eventual consistency even if the push channel is lost (lines 115-125, 165-174).
Summary
- Three specialized pipelines handle distinct real-time requirements: session chat (
useWorkbenchStream), course freshness (useStageFreshnessSync), and owner-session updates (OwnerSessionClient). - Resume capability via
lastEventIdand replay buffers ensures clients recover state exactly after reconnections, with thecaught_upevent signaling transition from backlog to live mode. - Resilience patterns include native EventSource reconnection, retry budgets with exponential back-off, health monitoring (5-second sampling), and periodic full reconciliation (60-75 second intervals).
- Closed-set event typing in the owner-session client prevents breaking changes when the server introduces new event types.
Frequently Asked Questions
How does OpenMAIC handle SSE reconnections without losing data?
OpenMAIC uses the lastEventId query parameter to resume streams from the last received event, stored in the global workbench store. The server transmits a backlog of missed events upon reconnection, which the client accumulates in a replay buffer until receiving the caught_up signal, ensuring no messages are lost during transient disconnects.
What triggers a course freshness sync in OpenMAIC?
The useStageFreshnessSync hook initiates synchronization on three conditions: receiving a stage_freshness SSE frame from the server, detecting a tab visibility change to focus, or a fallback timer firing every 30 seconds with random jitter. These triggers ensure the canvas converges to the latest state even if the SSE connection stalls.
Why does OpenMAIC use three separate SSE pipelines instead of a single connection?
Each pipeline serves a distinct consistency model and scale requirement. The session chat requires ordered replay buffers for message history, the course freshness needs manifest diffing and selective scene fetching, and the owner-session stream requires health monitoring and periodic reconciliation for the home dashboard. Separating concerns prevents back-pressure in one domain from affecting others.
How does OpenMAIC detect when an SSE connection is permanently degraded?
The OwnerSessionClient samples the EventSource.readyState every 5 seconds. If the state remains CONNECTING for nine consecutive samples (approximately 45 seconds), the client flags the channel as degraded via the onStreamHealth callback. This triggers fallback mechanisms while the native EventSource continues automatic reconnection attempts.
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 →