Streaming API
Cortex uses Server-Sent Events (SSE) to stream long-running execution state to clients in real time. The bulk of the API is a normal request/response HTTP surface; streaming is reserved for surfaces where progressive delivery is the point — agent runs, message streaming, tool-call status, and run lifecycle events.
When to use streaming
Reach for the streaming surface when you need to render output as it is produced rather than waiting for a final result. There are two streaming endpoints:
GET /v1/threads/:threadId/runs/:runId/events
GET /v1/threads/:threadId/runs/:runId/turn/events forwards the worker's full envelope stream verbatim — every channel (lifecycle, acts, updates, tokens, interrupts, artifacts, errors). /turn is a smaller, purpose-built stream for rendering a single "what is this run doing right now" view: the API folds the same underlying envelopes server-side into one small versioned document and streams that instead. See Turn document stream below. Use /events when you need the raw envelopes (token deltas, artifact announcements, interrupt prompts and options); use /turn when you just need an activity list and a status to render. Both are read-only — actions (sending messages, answering an interrupt, cancelling a run) go through the normal request/response endpoints.
Calling /events opens an SSE stream. The first frame is always a snapshot event carrying the run's full durable state (the run, its act DAG, any in-flight partial text, artifacts, and the final message if one exists); live events follow until the run reaches a terminal state. For everything else (creating threads, posting messages, reading run state once it has settled), use the normal request/response endpoints — the same snapshot is available as GET /v1/threads/:threadId/runs/:runId.
Durable delivery contract
Live SSE events are latency optimizations, not the source of truth. PostgreSQL owns the transcript, run status, act DAG, artifacts, and final output message. A run is delivered only when its snapshot meets all three conditions:
statusiscompleted,failed, orcancelled.outputMessageIdis non-null.outputMessageis present and readable.
Do not treat a terminal SSE frame by itself as delivery. Fetch the run snapshot after the frame and keep reconnecting or polling until the durable predicate above is true. This also covers the small interval in which a worker publishes a terminal hint while its database transaction is still committing.
Execution path
- The client posts a user message. In one database transaction the API assigns its thread sequence, creates the running run, and inserts a pending
run.startrow inrun_outbox. - An API or worker outbox relay claims due rows with
FOR UPDATE SKIP LOCKED, publishes the command to a Redis Stream, and marks the row sent. Failed publishes stay pending with capped exponential backoff; they are never discarded for exceeding an attempt count. - A worker consumer receives
run.start, materializes and claims the initial act, then acknowledges the Redis command. Commands that fail before durable ownership remain in the consumer group's pending list forXAUTOCLAIMrecovery. - A single-step run writes its assistant message, output link, terminal status, and next queued turn in one transaction. A decomposed run persists each act and dependency, dispatches ready acts through the same transactional outbox, and advances the DAG with row locks and compare-and-set claims.
- Before multi-act synthesis starts, the finalizer reserves a durable deterministic reply and links it through
outputMessageIdwhile the run remains running. Synthesis updates that message in place. If the process dies, stale-run recovery completes the reserved reply instead of leaving a terminal run without output. - Redis Pub/Sub carries tokens, activity, act, artifact, and lifecycle hints to the SSE connection. On every connection the API subscribes first and then reads PostgreSQL for the snapshot, so events published during snapshot construction are forwarded after it.
- The client reconciles the durable output, then reloads all paginated transcript pages.
GET /v1/threads/:idalso returnsactiveRunId, allowing a new tab or device to discover and resume an in-flight run without browser-local state.
Recovery paths
| Failure | Recovery |
|---|---|
| Redis publish fails | The durable outbox row retries with backoff until sent. |
| Worker dies before claiming work | The Redis pending entry is reclaimed; no explicit acknowledgement was sent. |
| Worker dies after claiming an act | A stale heartbeat reclaims the act until its bounded act-attempt limit is exhausted. |
act.ready delivery is lost | The stale ready-act sweep re-enqueues it; the claim compare-and-set makes duplicates harmless. |
| Finalizer dies before or during synthesis | The all-terminal run sweep starts finalization or completes its reserved fallback. |
| SSE disconnects or loses a terminal hint | The client rebuilds from a new snapshot and polls until the durable delivery predicate is true. |
Run events
Events are grouped by channel: lifecycle events mark run state, acts events carry the run's step timeline, updates events carry visible tool activity and narration, tokens events stream assistant text, interrupts events pause the run on a human question, and artifacts events announce generated files.
| Event | Family | Purpose |
|---|---|---|
snapshot | — | The catch-up frame sent once when the stream opens: full durable run state. Rebuild your view from it, then apply live events. |
lifecycle.started | Lifecycle | The worker has started executing the run. |
acts.updated | Acts | The run's act list after a transition (an act was created, claimed, completed, failed, retried, paused, or cancelled). The payload carries every act's id, key, status, description, dependsOn, and error — replace your act state wholesale. |
updates.activity | Updates | A visible activity started, completed, failed (with a short detail reason), or was cancelled. The envelope's actId links it to its act. |
updates.narration | Updates | Model narration emitted before a round of tool calls. |
interrupts.clarification / interrupts.approval | Interrupts | The run paused on a human question. The payload carries the prompt and options; answer via the resume/approval endpoints. It also carries actKey, the act's key, which matches acts.updated's keys and the turn document's activity entry ids for correlation. |
tokens.delta | Tokens | Append text to the active assistant draft. |
tokens.completed | Tokens | The assistant text stream for the run is complete. |
artifacts.created | Artifacts | A document or artifact produced by the run is now retrievable. |
errors.raised | Errors | A run-level error report accompanying a failure. |
lifecycle.completed / lifecycle.failed / lifecycle.cancelled | Lifecycle | Terminal run state. The server closes the stream after one of these arrives. |
Turn document stream
GET /v1/threads/:threadId/runs/:runId/turn streams a small versioned turn document instead of raw envelopes: a status, an ordered activity list, and the streamed reply text, enough to render a run's live progress surface and its chat bubble without reassembling either from the /events channels yourself. The API folds the run's worker envelopes into this document server-side, so the wire content is different, but auth/visibility, keepalives, the server-side stream deadline, and the terminal-lifecycle close all match /events exactly.
The first frame is always turn.snapshot: the full document, seeded from the run's durable state. Every frame after that is turn.patch: a small versioned diff. Not every worker envelope produces a patch — only the channels the fold owns (acts.updated, lifecycle.*, updates.activity, updates.narration, errors.raised, interrupts.clarification / interrupts.approval, tokens.delta, artifacts.created) advance the document. Envelopes outside that set (deck-authoring events and anything unrecognized) produce no turn.patch — use /events if you need those.
A worker envelope that folds to no patch still writes a turn.live frame: empty data, no document change. It signals only that the run is alive — an event arrived and changed nothing in the document — so a client should treat it the same as a turn.snapshot/turn.patch for liveness purposes (resetting a stall timer, for example) without otherwise acting on it. The 15-second SSE keepalive comment is a transport-level no-op and is not a turn.live frame: only a real worker envelope produces one, so a stalled worker behind a healthy connection still shows as quiet.
event: turn.snapshot
data: {"v":0,"runId":"9b2f…","threadId":"7ac1…","status":"running","activity":[{"id":"initial","label":"Working on your request","state":"active","detail":null}],"message":"","interrupt":null,"artifacts":[]}
event: turn.patch
data: {"v":1,"ops":[{"p":"/activity/0/label","o":"set","v":"Designing slide 1"},{"p":"/activity/0/state","o":"set","v":"active"}]}
event: turn.patch
data: {"v":2,"ops":[{"p":"/message","o":"append","v":"Here's "}]}
event: turn.patch
data: {"v":3,"ops":[{"p":"/message","o":"append","v":"your deck."}]}
event: turn.patch
data: {"v":4,"ops":[{"p":"/status","o":"set","v":"completed"}]}
The turn document shape:
v— the document's version. Starts at0on the snapshot; each patch carries the version it produces.runId/threadId— the run and thread this stream is for.status— one ofrunning,awaiting,completed,failed,cancelled.activity— an ordered list of{id, label, state, detail}entries.labelis always authored, human copy — never a raw key, tool name, or machinery string.stateis one ofpending,active,done,failed,cancelled.detailcarries verbatim machinery text (an exit code, a traceback) when present, kept separate fromlabelso a client can render it collapsed.message— the run's chat-bubble reply text, built up as it streams. The snapshot seeds it from durable state: the persisted output message's text once the run has settled, or the in-flight partial for a still-running reply, so a reconnect resumes with the complete text so far rather than an empty bubble. Live, eachtokens.deltaenvelope folds to a/messageappend op. Oncestatusreaches a terminal value the reply is frozen — any stragglingtokens.deltafolds to no patch.interrupt—null, or the run's current paused question:{type, actId, actKey, prompt, options, multiSelect}. Set by aninterrupts.clarification/interrupts.approvalenvelope, cleared once the run resumes (no act left awaiting) or reaches a terminal status.artifacts— an ordered list of the run's produced files, each projected to{id, filename, title, mimeType, sizeBytes, status, createdByTool}. The snapshot seeds it from durable state; eachartifacts.createdenvelope appends one, idempotently — re-delivery of an id already present is a no-op.
A patch carries only the ops needed to bring the document forward:
{"p":"/status","o":"set","v":<status>}— replace the document's status.{"p":"/activity","o":"append","v":<entry>}— append a new activity entry.{"p":"/activity/{i}/{field}","o":"set","v":<value>}— set one field (label,state, ordetail) on the activity entry at indexi.{"p":"/message","o":"append","v":<text>}— appendtextto the document'smessage.{"p":"/interrupt","o":"set","v":<object|null>}— replace the document'sinterrupt, or clear it withnull.{"p":"/artifacts","o":"append","v":<object>}— append one artifact to the document'sartifacts.
Versioning and reconnect-on-gap. A patch only applies when its v is exactly one more than the document's current v. Any other version — including a patch arriving before any snapshot — is a gap: stop applying patches and reconnect (re-open the same URL). There is no partial or op-level recovery and no Last-Event-ID support, matching /events. The new connection's turn.snapshot frame is always a full, self-consistent document built from durable state, so a reconnect fully re-bases the client regardless of what was missed.
SSE framing
Cortex follows the standard SSE wire format. Every event is a block of field: value lines terminated by a blank line. All frames — the snapshot included — use camelCase field names.
event: tokens.delta
data: {"channel":"tokens","type":"tokens.delta","actId":"9b2f…","payload":{"text":"Hello, "},"ts":1751879000.12}
event: acts.updated
data: {"channel":"acts","type":"acts.updated","actId":null,"payload":{"acts":[{"id":"9b2f…","key":"initial","status":"running","description":"Answer the question","dependsOn":[],"error":null}]},"ts":1751879001.02}
event— the event type. Always one of the names in the table above.data— a single JSON object. Live events share the envelope shape{channel, type, actId, payload, ts}; the snapshot frame is the run snapshot object itself. Cortex never splits a JSON payload across multipledata:lines.- Cortex emits SSE comment lines (lines starting with
:) as keepalives every ~15 seconds. Standards-compliant SSE parsers ignore these automatically. - The server bounds every stream's lifetime. When the cap is reached it writes a
: stream-timeoutcomment and closes; reconnect if the run is still going.
Reconnection
Long-running streams may drop for reasons unrelated to your client — network blips, load balancer rolls, idle timeouts on intermediate proxies. Clients are expected to reconnect.
Live events are ephemeral: there is no durable replay log and no Last-Event-ID support. Reconnecting simply re-opens the same URL — the new stream's snapshot frame carries everything durable (act statuses, partial text, artifacts, the final message), so rebuild your view from it and continue with the live events that follow. If the run finished while you were away, the snapshot reflects the terminal state and the stream closes.
Ordering guarantees
- The
snapshotframe always precedes live events on a stream, and any event published while the snapshot was being built is delivered after it — nothing is dropped in between. - Events for a single act arrive in order. Acts on a decomposed run execute in parallel, so events from different acts may interleave — use
actIdto attribute them. - A run may emit many
tokens.deltaevents followed by onetokens.completedevent for assistant text. - Exactly one terminal lifecycle event (
lifecycle.completed,lifecycle.failed, orlifecycle.cancelled) terminates the stream. After it arrives, the connection closes cleanly.
Errors on a stream
Authentication and authorization errors are returned as a normal HTTP response with a 4xx status before any event fires. Errors that occur mid-run are surfaced as an errors.raised event followed by lifecycle.failed. The connection then closes.