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
| Concern | Current behavior | Source |
|---|---|---|
| Event production | RunStreamEntry 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 reads | EfRunEventStream.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 diagnostics | Diagnostics 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 semantics | Replay 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 safety | PostgreSQL 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 net | Terminal 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 badge | AgentHost 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 fallback | If 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
| Concern | File |
|---|---|
| Local stream entry and mirror to shared stream | apps/Agentweaver.Api/Infrastructure/RunStreamStore.cs |
| EF/Postgres shared event stream and cursor polling | apps/Agentweaver.Api/Infrastructure/EfRunEventStream.cs |
SSE endpoint, Last-Event-ID, replay fallback | apps/Agentweaver.Api/Endpoints/RunEndpoints.cs |
| Coordinator outcome/work-plan read-through waits | apps/Agentweaver.Api/Endpoints/CoordinatorEndpoints.cs |
| Terminal backfill and recording writer integration | apps/Agentweaver.Api/Runs/RunWorkflowFactory.cs |
| Cross-instance tests | tests/Agentweaver.Tests/EfRunEventStreamTests.cs |
See also
- Distributed execution & scaling — why the shared event store is required for multi-replica deployments.
- Events & observability — event taxonomy and observability model.
- Token usage monitoring — one UI surface that consumes the same live stream and usage projections.
Diagram details and constraints
| Element | Contract |
|---|---|
| title | Postgres is the event relay |
| subtitle | Any API replica can serve a cursor over durable RunEvents—no sticky session required. |
| group-title0 | Write path · replica A |
| group-title1 | Read path · replica B |
| Run producer | Run producer |
| Run producer | Append a structured event |
| Run producer | runId + type + payload |
| EF event stream | EF event stream |
| EF event stream | Serialize writes per run |
| EF event stream | pg_advisory_xact_lock |
| RunEvents | RunEvents |
| RunEvents | Shared PostgreSQL table |
| RunEvents | (RunId, Sequence) |
| Web / MCP watcher | Web / MCP watcher |
| Web / MCP watcher | Consume ordered events |
| Web / MCP watcher | last delivered cursor |
| SSE endpoint | SSE endpoint |
| SSE endpoint | Emit id + event + data |
| SSE endpoint | ordered response frames |
| EF subscriber | EF subscriber |
| EF subscriber | Read Sequence > cursor |
| EF subscriber | idle poll: 250 ms |
| e1 | append |
| e2 | commit |
| e3 | ordered batch |
| e4 | yield |
| e5 | SSE frames |
| assurance-title | POSTGRES LANE ONLY |
| assurance-line1 | SQLite register-channel / replay / tail is a separate implementation—not this architecture. |
| assurance-line2 | Late-delta suppression is process-local; do not read it as a database-wide terminal fence. |
| Run producer | Input |
| Run producer | RunStreamEntry |
| Run producer | Identity |
| Run producer | runId + event type |
| Run producer | Body |
| Run producer | Structured payload |
| Run producer | Ack |
| Run producer | After durable commit |
| EF event stream | Lock |
| EF event stream | Per-run advisory lock |
| EF event stream | Next |
| EF event stream | MAX(Sequence) + 1 |
| EF event stream | Write |
| EF event stream | Save transaction |
| EF event stream | Commit |
| EF event stream | Before acknowledgement |
| RunEvents | Table |
| RunEvents | Key |
| RunEvents | RunId + Sequence |
| RunEvents | Order |
| RunEvents | Ascending sequence |
| RunEvents | Reuse |
| RunEvents | Same type / payload |
| Web / MCP watcher | Client |
| Web / MCP watcher | Web or MCP |
| Web / MCP watcher | Resume |
| Web / MCP watcher | Last delivered cursor |
| Web / MCP watcher | Replica |
| Web / MCP watcher | No sticky requirement |
| Web / MCP watcher | History |
| Web / MCP watcher | Durable ordered events |
| SSE endpoint | Frame |
| SSE endpoint | id + event + data |
| SSE endpoint | Cursor |
| SSE endpoint | Last-Event-ID |
| SSE endpoint | Delivery |
| SSE endpoint | Yield ordered events |
| SSE endpoint | Close |
| SSE endpoint | After batch is drained |
| EF subscriber | Query |
| EF subscriber | Sequence > cursor |
| EF subscriber | Idle |
| EF subscriber | Poll after 250 ms |
| EF subscriber | State |
| EF subscriber | Shared durable table |
| EF subscriber | Blocked |
| EF subscriber | Retryable: keep open |
| producer | Coordinator or run execution; Acknowledgement follows commit |
| append | Allocate MAX(Sequence) + 1; Save and commit transaction |
| store | Cross-replica ordered history; Explicit duplicates must match payload |
| client | Reconnect from the cursor; No local channel dependency |
| sse | Cursor advances after delivery; Drain batch before terminal close |
| reader | Query the shared durable table; Retryable assembly_blocked stays open |
| notes | POSTGRES 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. |
| groups | Write path · replica A; Read path · replica B |
Diagram details and constraints
| Element | Contract |
|---|---|
| notes | LOOP · 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. |
