How the OpenMAIC Agent Runtime Handles Event Replay for Session History
The OpenMAIC agent runtime replays session history by reading immutable events from PostgreSQL, compacting them to remove no-ops, and folding them into session state using a monotonic seq ordering key.
Every interaction in an OpenMAIC session—user messages, tool calls, thinking events, assistant replies—is durably logged as an immutable row. When a client reconnects or a session resumes, the runtime reconstructs the current view by replaying these events through a deterministic fold operation. This article explains the complete replay pipeline, from storage to state reconstruction.
Event Storage and Ordering Guarantees
OpenMAIC persists session events to a PostgreSQL table managed by PgAgentSessionStore. Each row carries a monotonically increasing seq field that serves as the sole ordering key for replay, as defined in packages/@openmaic/storage/src/agent-session/types.ts.
This design guarantees:
- Deterministic ordering – Events are always processed in the sequence they were written
- Idempotent folding – Re-applying an already-processed event is a no-op because
seqensures unique positioning - Gap detection – Clients can resume from any point using a cursor value
The Replay Pipeline: Fetch, Compact, Fold
The runtime reconstructs session state through three distinct phases.
1. Fetching Events from Durable Storage
The PgAgentSessionStore.readEventsAfterForReplay(sessionId, cursor) method pulls all rows with seq greater than the provided cursor. This interface, located in the storage package, enables both cold replays (full session reconstruction) and mid-gap replays (resuming from a disconnection point).
2. Compacting the Event Stream
Before folding, events pass through compactReplayEvents to remove no-op duplicates. For example, a re-sent thinking_end event that carries no new information is eliminated while preserving the original order of meaningful events.
3. Folding into Session State
The foldEvents function, implemented in packages/@openmaic/dsl/src/runtime.ts, walks the compacted stream and applies each record to an in-memory session state object. This state includes:
- Chat transcript accumulation
- Generating order tracking
- Page and scene context
The fold operation is pure and deterministic—given the same event sequence, it always produces identical state.
Replay Consistency Across Connection Scenarios
The system handles two primary reconnection patterns, validated in tests/workbench/session-fold.test.ts.
Cold Replay
After a full disconnect, the client reconstructs state from scratch:
// Rebuild from cursor 0 (complete history)
const { events } = await store.readEventsAfterForReplay(sessionId, 0);
const replayable = compactReplayEvents(events);
const sessionState = foldEvents(undefined, replayable);
As verified in tests/workbench/session-fold.test.ts lines 227-238, this produces a state matching the live session exactly.
Mid-Gap Replay
When reconnecting with a Last-Event-ID header, the runtime resumes from the known cursor:
// Resume from last known sequence number
const { events, scanned } = await store.readEventsAfterForReplay(
sessionId,
lastSeq // cursor from client
);
The overlapping prefix is folded once; duplicates are ignored. Test lines 166-165 validate that the cursor stays synchronized without state corruption.
Event Notification and Real-Time Delivery
Every durable append triggers a PostgreSQL NOTIFY within the same transaction that writes the row. In lib/server/agent-runtime/store.ts lines 71-80, the notifyDurableAgentEvent call wakes the per-session Server-Sent-Events (SSE) tail.
The notification flow:
- Transaction commits event to
agent-sessiontable NOTIFYpayload emitted with session ID and newseqevent-notify-bus.tsroutes to active SSE connections- Clients trigger replay for any missed events
If notifications are missed due to network partitions, SSE clients fall back to periodic polling. Correctness is never compromised because the replay mechanism is source-of-truth based.
Key Source Files and Responsibilities
| Component | File Path | Role in Replay |
|---|---|---|
| Persistent event store | lib/server/agent-runtime/store.ts |
Transaction wrapper, readEventsAfterForReplay |
| Notification dispatch | lib/server/agent-runtime/event-notify-bus.ts |
PostgreSQL LISTEN/ NOTIFY bridge to SSE |
| Resume orchestration | lib/server/agent-runtime/resume.ts |
Entry point when sessions (re)attach |
| Fold/compact logic | packages/@openmaic/dsl/src/runtime.ts |
foldEvents, compactReplayEvents implementations |
| Type definitions | packages/@openmaic/storage/src/agent-session/types.ts |
Event shapes, seq ordering contract |
| Determinism tests | tests/workbench/session-fold.test.ts |
Validation of cold and mid-gap replay |
Complete Replay Example
The following pattern, derived from the runtime's implementation in resume.ts, demonstrates full session reconstruction:
import { getAgentSessionStore } from '@/lib/server/agent-runtime/store';
import { foldEvents, compactReplayEvents } from '@openmaic/dsl';
interface SessionState {
messages: Array<{ role: string; content: string }>;
generating: boolean;
seq: number;
}
async function rebuildSession(
sessionId: string,
cursor: number
): Promise<SessionState> {
// 1. Initialize storage layer
const store = await getAgentSessionStore();
// 2. Fetch events after known cursor
const { events, scanned } = await store.readEventsAfterForReplay(
sessionId,
cursor
);
// 3. Remove redundant events while preserving order
const replayable = compactReplayEvents(events);
// 4. Fold into fresh state (undefined = initial state)
const sessionState = foldEvents(undefined, replayable);
// 5. Return reconstructed state with new cursor
return {
...sessionState,
seq: scanned // Latest sequence for next resume
};
}
In production, resume.ts automates this workflow—detecting whether to perform cold replay or resume from a partial state based on client-provided cursors.
Summary
- Immutable durable log: All events stored in PostgreSQL with monotonic
seqordering - Three-phase replay: Fetch via
readEventsAfterForReplay, compact viacompactReplayEvents, fold viafoldEvents - Idempotent guarantee:
seqkey ensures safe reprocessing of overlapping events - Dual consistency modes: Cold replay from cursor 0, mid-gap replay from
Last-Event-ID - Notification-driven real-time: PostgreSQL
NOTIFYtriggers SSE updates with polling fallback
Frequently Asked Questions
How does OpenMAIC handle duplicate events during replay?
The seq field provides unique positioning for every event. When foldEvents processes a stream, events with seq values already incorporated into state are ignored. This makes the fold idempotent—replaying the same prefix multiple times produces identical results.
What happens if the PostgreSQL notification is lost?
Correctness is maintained through polling fallback. The SSE client periodically queries readEventsAfterForReplay with its last known cursor, discovering any events missed due to notification failures. Notifications optimize latency; the durable log guarantees consistency.
Can session replay work across server restarts?
Yes. All event state resides in PostgreSQL, not server memory. A new server process can reconstruct any session's state by reading the persistent log and applying the standard fold operation, as implemented in resume.ts.
Where is the event ordering contract defined?
The seq field and event type definitions are located in packages/@openmaic/storage/src/agent-session/types.ts. The runtime in packages/@openmaic/dsl/src/runtime.ts implements the fold logic that consumes this contract.
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 →