Skip to content

Run Event Stream (IRunEventStream) ​

Feature: 016-run-event-stream · updated for cross-replica streaming

Agentweaver treats a run's event stream as a durable, ordered log first and a live SSE feed second. The current implementation mirrors every RunStreamEntry append into the shared RunEvents table, and the EF/Postgres stream reads from that table by cursor so any replica can stream any run. The pod-local RunStreamStore remains a same-replica compatibility and low-latency path; it is no longer the only place an SSE client can see a run's live history. Source: apps/Agentweaver.Api/Infrastructure/RunStreamStore.cs:98, apps/Agentweaver.Api/Infrastructure/RunStreamStore.cs:115, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:15, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:77, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:416.

For the scaling story, see Distributed execution & scaling. For the event taxonomy, see Events reference.

Architecture — shared store, cursor stream ​

The horizontal-scale invariant is simple: the database log is the source of truth, and the cursor is the replay boundary. EfRunEventStream.AppendAsync writes through before acknowledging (WriteThroughAsync), and SubscribeAsync repeatedly loads rows whose sequence is greater than the caller's last seen cursor, yielding them in sequence order until a terminal event appears. It drains the full replay batch before stopping, so a diagnostic row persisted immediately after a terminal row is still delivered before the SSE subscription closes. Source: apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:63, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:71, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:77, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:84, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:111, apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs:180.

Delivery sequence — durable-first, then cursor polling ​

The EF/Postgres path does not switch to an in-process live channel after replay: each replica continues reading the shared log. SQLite's compatibility implementation registers a local channel before replay and then tails it, dropping already replayed sequences. These are distinct delivery implementations, not a shared cross-replica bus.

What changed from the old in-memory-only stream ​

ConcernCurrent behaviorSource
Event productionRunStreamEntry obtains the authoritative sequence from the durable append before exposing the event to local history/waiters. Local notification is not the cross-replica relay.RunStreamStore.cs:183-212
Cross-replica readsEfRunEventStream.SubscribeAsync polls the shared RunEvents table every 250 ms when no new rows were emitted, so a subscriber on a different replica catches up without sticky sessions.EfRunEventStream.cs:33, EfRunEventStream.cs:77, EfRunEventStream.cs:96, EfRunEventStream.cs:180
Late terminal diagnosticsDiagnostics can remain durable after terminalization. EF's late message-delta suppression is process-local; it is not a database-enforced cross-replica terminal fence. A row outside the drained terminal-containing batch may require a later replay.EfRunEventStream.cs:59-86,144-189
Replay terminal semanticsReplay yields all rows in the loaded batch, then stops if the batch contained a terminal event. coordinator.assembly_failed is terminal; retryable coordinator.assembly_blocked is not, so subscribers can stay attached across recovery.EfRunEventStream.cs:35, EfRunEventStream.cs:84, EfRunEventStream.cs:111, SqliteRunEventStream.cs:34, SqliteRunEventStream.cs:153
Sequence safetyPostgreSQL allocates per-run MAX(Sequence) + 1 under an advisory transaction lock before commit. An explicit duplicate sequence is idempotent only when its contents match; a conflicting payload is rejected.EfRunEventStream.cs:223-284,316-321
Terminal safety netTerminal persistence re-appends the full in-memory history through IRunEventStream; duplicate (RunId, Sequence) rows are skipped, so missed mirrors are reconciled without duplication.RunWorkflowFactory.cs:287, RunWorkflowFactory.cs:296, RunWorkflowFactory.cs:298, RunWorkflowFactory.cs:301
Execution pod badgeAgentHost pod bindings append sandbox.execution_pod.bound and releases append sandbox.execution_pod.unbound to the same shared log, so graph pod badges resolve and clear after refresh and across replicas.RunEventExecutionPodNameStore.cs, KubernetesSandboxExecutor.cs
SSE fallbackIf the current replica has no local stream entry, /api/runs/{id}/stream subscribes to IRunEventStream from the Last-Event-ID cursor and writes those events as SSE frames.RunEndpoints.cs:416, RunEndpoints.cs:423, RunEndpoints.cs:429, RunEndpoints.cs:431, RunEndpoints.cs:443

EfRunEventStreamTests proves the cross-replica behavior by creating two stream instances over the same database: one appends events and the other receives them through SubscribeAsync. The tests also prove that RunStreamStore.RecordNext mirrors into the shared stream. Source: tests/Agentweaver.Tests/EfRunEventStreamTests.cs:27, tests/Agentweaver.Tests/EfRunEventStreamTests.cs:30, tests/Agentweaver.Tests/EfRunEventStreamTests.cs:37, tests/Agentweaver.Tests/EfRunEventStreamTests.cs:42, tests/Agentweaver.Tests/EfRunEventStreamTests.cs:50, tests/Agentweaver.Tests/EfRunEventStreamTests.cs:56, tests/Agentweaver.Tests/EfRunEventStreamTests.cs:63.

SSE wire protocol ​

The /api/runs/{id}/stream endpoint still emits the same Server-Sent Event shape: an id line with the run-event sequence, an event line with the event type, and a JSON data line, followed by a done frame when the server closes the stream. Source: apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:315, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:448, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:468, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:489.

Reconnects use the browser's Last-Event-ID header as the fromSequence cursor. The endpoint parses that header on the replay path and on the local path, so refreshes resume from the last emitted sequence instead of starting over. Source: apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:423, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:424, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:429, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:452, apps/Agentweaver.Api/Endpoints/RunEndpoints.cs:453.

Read-through coordinator gates ​

Browser refreshes around coordinator gates no longer surface transient 404 or 409 responses just because the spec or plan row is being created. GET /api/runs/{id}/outcome-spec waits up to three seconds for the persisted OutcomeSpec, and GET /api/runs/{coordinatorRunId}/work-plan waits up to five seconds for the persisted work plan before returning not found. Confirming an already-confirmed outcome spec also attempts to return the confirmed spec instead of a conflict. Source: apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:53, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:88, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:90, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:167, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:571, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:574, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:580, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:591, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:594, apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs:600.

Source ​

ConcernFile
Local stream entry and mirror to shared streamapps/Agentweaver.Api/Infrastructure/RunStreamStore.cs
EF/Postgres shared event stream and cursor pollingapps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs
SSE endpoint, Last-Event-ID, replay fallbackapps/Agentweaver.Api/Endpoints/RunEndpoints.cs
Coordinator outcome/work-plan read-through waitsapps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs
Terminal backfill and recording writer integrationapps/Agentweaver.Api/Runs/RunWorkflowFactory.cs
Cross-instance teststests/Agentweaver.Tests/EfRunEventStreamTests.cs

See also ​

Diagram details and constraints
ElementContract
titlePostgres is the event relay
subtitleAny API replica can serve a cursor over durable RunEvents—no sticky session required.
group-title0Write path · replica A
group-title1Read path · replica B
Run producerRun producer
Run producerAppend a structured event
Run producerrunId + type + payload
EF event streamEF event stream
EF event streamSerialize writes per run
EF event streampg_advisory_xact_lock
RunEventsRunEvents
RunEventsShared PostgreSQL table
RunEvents(RunId, Sequence)
Web / MCP watcherWeb / MCP watcher
Web / MCP watcherConsume ordered events
Web / MCP watcherlast delivered cursor
SSE endpointSSE endpoint
SSE endpointEmit id + event + data
SSE endpointordered response frames
EF subscriberEF subscriber
EF subscriberRead Sequence > cursor
EF subscriberidle poll: 250 ms
e1append
e2commit
e3ordered batch
e4yield
e5SSE frames
assurance-titlePOSTGRES LANE ONLY
assurance-line1SQLite register-channel / replay / tail is a separate implementation—not this architecture.
assurance-line2Late-delta suppression is process-local; do not read it as a database-wide terminal fence.
Run producerInput
Run producerRunStreamEntry
Run producerIdentity
Run producerrunId + event type
Run producerBody
Run producerStructured payload
Run producerAck
Run producerAfter durable commit
EF event streamLock
EF event streamPer-run advisory lock
EF event streamNext
EF event streamMAX(Sequence) + 1
EF event streamWrite
EF event streamSave transaction
EF event streamCommit
EF event streamBefore acknowledgement
RunEventsTable
RunEventsKey
RunEventsRunId + Sequence
RunEventsOrder
RunEventsAscending sequence
RunEventsReuse
RunEventsSame type / payload
Web / MCP watcherClient
Web / MCP watcherWeb or MCP
Web / MCP watcherResume
Web / MCP watcherLast delivered cursor
Web / MCP watcherReplica
Web / MCP watcherNo sticky requirement
Web / MCP watcherHistory
Web / MCP watcherDurable ordered events
SSE endpointFrame
SSE endpointid + event + data
SSE endpointCursor
SSE endpointLast-Event-ID
SSE endpointDelivery
SSE endpointYield ordered events
SSE endpointClose
SSE endpointAfter batch is drained
EF subscriberQuery
EF subscriberSequence > cursor
EF subscriberIdle
EF subscriberPoll after 250 ms
EF subscriberState
EF subscriberShared durable table
EF subscriberBlocked
EF subscriberRetryable: keep open
producerCoordinator or run execution; Acknowledgement follows commit
appendAllocate MAX(Sequence) + 1; Save and commit transaction
storeCross-replica ordered history; Explicit duplicates must match payload
clientReconnect from the cursor; No local channel dependency
sseCursor advances after delivery; Drain batch before terminal close
readerQuery the shared durable table; Retryable assembly_blocked stays open
notesPOSTGRES LANE ONLY; SQLite register-channel / replay / tail is a separate implementation—not this architecture.; Late-delta suppression is process-local; do not read it as a database-wide terminal fence.
groupsWrite path · replica A; Read path · replica B
Diagram details and constraints
ElementContract
notesLOOP · repeat durable reads; idle wait = 250 ms; Drain the whole batch before terminal close. Retryable assembly_blocked is not terminal.; Explicit-sequence reuse is idempotent only for matching type/payload. SQLite live channels are a separate lane.