crouter
SDK

Streaming

Follow assistant text, tool activity, reports, and outcomes as an agent runs.

Watch a run as it works: assistant text as it is produced, tool calls as they start and finish, reports as they are pushed, and the settled outcome.

Forwarding a run to a browser

The app server keeps the runtime token and connection; the browser connects only to the app. The app server calls forwardRunEvents(client, runId, {after: cursorFromRequest(request), signal: request.signal}) and returns the Fetch Response. After a host cuts the browser connection, EventSource sends its last SSE id as Last-Event-ID; cursorFromRequest prefers it over the URL's after query parameter.

Every forwarded run event is an unnamed SSE record with its sequence_number as id and the complete event object as data. Receive all event types with one handler:

const source = new EventSource(`/api/runs/${runId}/events?after=${cursor}`);
source.onmessage = ({data}) => {
  const event = JSON.parse(data);
  if (event.type === 'stream.refused') { showRefusal(event.error); source.close(); return; }
  if (event.type === 'stream.reconnecting') { showConnectionFailure(event.error); return; }
  if (event.type === 'stream.reconnected') { clearConnectionFailure(); return; }
  showRunEvent(event);
};

The three stream.* records are forwarded-stream notices, not daemon run events; they have no sequence_number or SSE id. A stream.refused carries the full error envelope, including stream_gap's details.latest_sequence and details.earliest_sequence. The helper ends that response after a refusal. The SDK does not recover a gap by skipping missing events.

When the page reloads mid-reply or gets stream_gap, the app reads runs.trace(runId) through its server and shows each node's in_progress_message ({thinking, text} or null) as that node's open message. The trace's latest_sequence is read consistently with its conversations and open messages. Resume runs.events(runId, {after: trace.latest_sequence}) (or forward with that cursor); do not use a separately read runs.get cursor. Continue appending new deltas to the open message, but replace its thinking or text with the complete content repeated by node.thinking.done or node.output_text.done. This shows the message from its beginning without duplicating content. runs.get and every runs.list entry expose latest_sequence and last_activity_at; a cursor equal to latest_sequence is not a gap even when older events have expired.

Following the reply to one message

A chat app sends a message and shows the answer to it. runs.reply(runId, sent, {signal?}) takes what runs.message returned and follows only the turn that delivers that message:

const sent = await client.runs.message(runId, 'What did the report say?');
for await (const event of client.runs.reply(runId, sent)) {
  if (event.type === 'text.delta') append(event.delta);
  if (event.type === 'turn.completed' && event.error) showError(event.error); // {code: 'content_refused' | 'model_error', message}
}
// or collect it: text as it streams, then the whole text and how the turn ended
const {text, ended, error, outcome} = await client.runs.reply(runId, sent).collect({onDelta: (delta) => append(delta)});
// ended: 'completed' | 'error' (error set) | 'waiting_on_user' | 'settled' (outcome set)

It reads runs.events(runId, {after: sent.sequence_number}), skips child nodes' events, starts at the root's node.turn.started whose input.message_ids holds sent.message_id, and yields:

EventFieldsFrom
text.deltadeltanode.output_text.delta
message.donetext (the message's whole text, possibly empty), stop_reasonnode.output_text.done, once per assistant message
turn.completedrun_status, stop_reason, error?, message_idsnode.turn.completed; always last
run.settledoutcomenode.settled, when the run settles before the turn completes; always last

Every event also carries its sequence_number. The reply ends after turn.completed or run.settled and closes its event stream. A resident run never settles when it goes idle, so a plain runs.events stream stays open, but a reply does not. .text() consumes the reply and resolves to the text of its message.done events joined with a blank line, so a turn that speaks, calls a tool and speaks again returns both texts. .collect({onDelta?}) consumes the reply too: it passes each streamed piece to onDelta, starting a later assistant message's first piece with a blank line, and resolves to RunReplyResult {text, ended, error?, outcome?}, where text is the trimmed messages joined with a blank line. ended is completed (text may be empty when the model wrote nothing), error (the turn's last message failed; error says why), waiting_on_user (the turn ended on a question to the person), or settled (the run settled first; outcome says how). .text({onDelta?}) is collect's text. A reply can be consumed once.

Several messages delivered in one turn, including a start_turn: false message that joined it, all resolve to that turn. To have a reply to the run's first message, pass it as runs.start({prompt, message}) and follow runs.reply(run.run_id, run.message!): a message stored with the run always joins the first turn. A refused or failed model call ends the turn with an empty message.done (stop_reason: 'error') and turn.completed.error; the run stays usable, and runs.get shows the same error as last_turn_error until a later turn completes without one. A provider failure the runtime retries on its own (connection, rate limit, overloaded) does not end the reply: the recovery attempt continues the same turn, so the reply yields its answer and one turn.completed without an error, and only a turn whose retries run out ends with model_error. An interrupt, runs.cancel, or a closed node ends the turn with stop_reason: 'aborted' and no error. An expired cursor throws APIError with code stream_gap; aborting signal rejects with APIUserAbortError.

Typed run events

RunEvent, exported by the SDK, is the union of every run event, discriminated on type. Its types are exactly the daemon's run event vocabulary, declared once. RunEventOf<'node.turn.completed'> picks one event, and RunStream.on(type, listener) passes the listener that event. An event type a newer daemon adds and this SDK does not declare is still delivered as-is.

The SDK surface

const stream = client.nodes.stream({ prompt: 'Fix the failing test', cwd }, { headers: { 'x-trace': 'run-9' } });

stream.on('node.output_text.delta', (event) => process.stdout.write(event.delta));

for await (const event of stream) {
  // the same events, typed by `type`
}

const outcome = await stream.finalOutcome(); // resolves on node.settled
stream.abort(); // ends the HTTP response; the node keeps running
MemberBehaviour
client.nodes.stream(params, options?)Creates the node, then opens its event stream. Returns synchronously; stream.node is a Promise<NodeDetailDTO> for the created node. The request options apply to creation and to opening the stream, except timeout, because a stream has no wall-clock timeout.
client.nodes.events(id, options?)Validates and retrieves an existing node, then streams it. options is NodeEventsOptions: after plus request headers, signal, and maxRetries; it deliberately has no timeout.
.on(type, listener) / .off(type, listener)Typed event emitter; each returns the NodeStream.
[Symbol.asyncIterator]()Yields the same NodeEventDTO union.
.finalOutcome()Returns Promise<NodeOutcomeDTO>, resolving on node.settled and rejecting on a terminal error event or a connection failure.
.abort()Aborts the underlying fetch. An in-flight iterator and finalOutcome() reject with APIUserAbortError. Passing your own signal in the options aborts the stream the same way; NodeStream exposes no signal property.

The method is called events, not subscribe: a subscription is already a push-delivery edge between two nodes on the canvas, and one word must not mean two things.

A stream is an observer, never a lifecycle hold. Disconnecting drops your subscriber and leaves the node running. Opening a stream never revives a dormant node — call client.nodes.revive(id) first if that is what you want.

Activity helper

followActivity(stream, { describe? }) turns streamed tool-call events into immutable ActivityStep[] snapshots for a plain activity feed. describeToolDefault(tool, summary) names the common tools; pass describe to use labels for your application or return null to hide a tool call. Tool summaries describe the argument shape, never argument values or raw arguments.

import { followActivity } from '@crouter/sdk';

for await (const steps of followActivity(stream)) {
  render(steps);
}

A step starts as running; a node.tool_call.completed event changes it to done when status is ok or failed when status is error.

Events

Every event except error carries node_id and sequence_number in addition to the fields listed. error carries only its error object.

Eventdata
node.output_text.delta{ node_id, sequence_number, delta }
node.output_text.done{ node_id, sequence_number, text }
node.tool_call.started{ node_id, sequence_number, tool_call_id, tool, summary }
node.tool_call.completed{ node_id, sequence_number, tool_call_id, tool, status: 'ok' | 'error', summary }
node.turn.started{ node_id, sequence_number }
node.turn.completed{ node_id, sequence_number }
node.report.pushed{ node_id, sequence_number, report }
node.status.changed{ node_id, sequence_number, status }, where status can also be stream-only 'dormant'
node.settled{ node_id, sequence_number, outcome } — terminal; the daemon ends the response after it
error{ error: { code: 'stream_gap' | 'stream_dropped' | 'stream_error', message, details?: { earliest_sequence?: number } } } — stream_gap is recoverable; the other codes terminate the response

summary on a tool-call event states whether arguments are null, an array and its item count, an object and its field count, or a primitive type. It never includes argument values or tool output.

Deliberately not carried: thinking deltas, the model's tool-call construction deltas, raw tool output, and the system prompt. Those belong to the owner's viewer, not to an application's run stream.

The route

GET /v1/nodes/{id}/events        Accept: text/event-stream
    ?after=<sequence_number>     resume from a cursor (optional)

Server-sent events: one record per event as event: <type>, data: <JSON>, blank line, plus : keepalive comment lines every 15 s. Authentication is the daemon's existing rule — filesystem permission on the unix socket, bearer token over TCP.

What you get depending on the node's state

Node state when you callWhat the stream does
Already settledWrites retained events when available, ending in node.settled; otherwise writes node.settled with the outcome and ends. Nothing is revived.
RunningStreams live, seeded as described below.
Dormant and not settledWrites node.status.changed { status: 'dormant' } and holds the response open with keepalives. It starts streaming if and when the daemon brings a broker up for that node.

A node can settle between your create and your events call, so both branches are ordinary. Either way the terminal event is node.settled, read from the same outcome row nodes.outcome reads — a client that joins after settlement and a client that was live at settlement see the same outcome.

Sequence, resume, and failure

sequence_number is monotonic per node while that node's in-memory stream hub exists. The daemon keeps one in-memory ring per streamed node holding the last 512 events or 256 KiB, whichever binds first.

SituationBehaviour
You join a run already in progressFirst replays already-pushed reports oldest first, then seeds from the broker's snapshot: one node.output_text.done per completed assistant message, then one node.output_text.delta carrying the accumulated partial, then live. No content is lost — only delta granularity.
A second subscriber joinsSeeded from the ring, so both subscribers see the same sequence numbers.
after=N within the retained rangeReplays events after N.
after=N outside the retained range, including a future cursorAn error event with code stream_gap carries the earliest available sequence, then the stream continues from that point. It never silently skips.
One event exceeds the 256 KiB replay budgetLive subscribers receive it, but the daemon clears the resume ring. A cursor before that event gets stream_gap; a cursor on it resumes at the next retained event.
Your HTTP reader is too slowThe daemon ends the response with an error event whose code is stream_dropped. One slow reader cannot stall the engine.
The node's broker is replaced or its observer socket closesThe daemon reconnects and emits node.status.changed. Sequence continues while the stream hub remains in memory.
The daemon restartsThe ring is gone. A resuming client gets stream_gap with earliest_sequence: 1. There is no durable event log.
You disconnectYour subscriber is dropped. The node keeps running.

Resuming

import { APIError } from '@crouter/sdk';

let cursor: number | undefined;

for (;;) {
  const stream = client.nodes.events(id, { after: cursor });
  try {
    for await (const event of stream) {
      if (event.type === 'error') {
        if (event.error.code === 'stream_gap') {
          console.warn('missed events before', event.error.details?.earliest_sequence);
          continue; // the stream continues from the floor
        }
        throw new APIError(0, event.error.code, event.error.message, event.error.details);
      }
      cursor = event.sequence_number;
      handle(event);
    }
    return; // ended on node.settled
  } catch (error) {
    if (error instanceof APIError && error.code === 'stream_dropped') continue;
    throw error;
  }
}

stream_gap is recoverable and the stream continues after it. stream_dropped is terminal for that HTTP response — reconnect with the cursor you last saw.

Identifier validation

nodes.events(id, options?) validates id before it requests the node or opens SSE. Because it returns a NodeStream synchronously, an invalid node identifier rejects stream.node, stream.finalOutcome(), and an iteration with TypeError; no request is sent. See Errors.

Event shapes depend on the pinned pi version

The delta and tool-call shapes come from the installed @earendil-works/pi-agent-core and pi-ai packages, not from crouter. The daemon's translator is the one place that depends on them. A pi version bump must be checked against it.

On this page