Skip to content

Streaming in JavaScript

View as Markdown

K-Agent streams with standard Server-Sent Events, but you read them with fetch rather than EventSource: EventSource can’t send an Authorization header or a POST body. The code below works in browsers and in Node 18+, has no dependencies, and is tested against streams that split lines, UTF-8 characters and events across network chunks.

/** Parses a text/event-stream response into { id, event, data } objects. */
export async function* readEvents(response) {
const reader = response.body.pipeThrough(new TextDecoderStream()).getReader();
let buffer = '';
let ev = { id: undefined, event: 'message', data: [] };
for (;;) {
const { value, done } = await reader.read();
if (done) return;
buffer += value;
const lines = buffer.split('\n');
buffer = lines.pop(); // keep the unfinished last line for the next chunk
for (const raw of lines) {
const line = raw.endsWith('\r') ? raw.slice(0, -1) : raw;
if (line === '') {
// A blank line ends an event.
if (ev.data.length) yield { id: ev.id, event: ev.event, data: JSON.parse(ev.data.join('\n')) };
ev = { id: undefined, event: 'message', data: [] };
} else if (!line.startsWith(':')) {
// Lines starting with ":" are heartbeats; everything else is "field: value".
const i = line.indexOf(':');
const field = i === -1 ? line : line.slice(0, i);
const value = i === -1 ? '' : line.slice(i + 1).replace(/^ /, '');
if (field === 'id') ev.id = value;
else if (field === 'event') ev.event = value;
else if (field === 'data') ev.data.push(value);
}
}
}
}

The rules from Runs and streaming turn into a few lines:

  • keep the text per message_id — one run can produce more than one assistant message;
  • append message.delta text as a preview;
  • on message.completed, replace the preview with the authoritative text — or drop it when discarded: true;
  • remember the last event id, so you can resume;
  • stop at the terminal run.* event.
const BASE = 'https://api.k-agent.kerneltics.com/v1';
const TERMINAL = new Set(['run.completed', 'run.failed', 'run.cancelled', 'run.superseded', 'run.requires_action']);
const textOf = (message) =>
message.content.filter((part) => part.type === 'text').map((part) => part.text).join('');
/**
* Sends a message and streams the reply. onText receives the reply text so far
* (all assistant messages of the run). Resumes with Last-Event-ID if the connection drops.
* Returns the terminal event: run.completed, run.failed, run.requires_action, …
*/
export async function streamReply({ session, input, token, onText, signal }) {
const headers = { Authorization: `Bearer ${token}`, Accept: 'text/event-stream' };
const texts = new Map(); // message_id → text
let runId;
let lastId;
async function consume(response) {
if (!response.ok) {
const { error } = await response.json();
throw Object.assign(new Error(error.message), { code: error.code });
}
for await (const { id, event, data } of readEvents(response)) {
if (id) lastId = id;
if (event === 'run.created') runId = data.run.id;
else if (event === 'message.delta') {
texts.set(data.message_id, (texts.get(data.message_id) ?? '') + data.delta);
} else if (event === 'message.completed') {
if (data.discarded) texts.delete(data.message_id ?? data.message?.id);
else texts.set(data.message.id, textOf(data.message)); // authoritative text
} else if (event === 'error') {
throw Object.assign(new Error(data.code), { code: data.code });
} else if (TERMINAL.has(event)) {
onText([...texts.values()].join('\n\n'));
return { event, data };
} else continue; // other events: ignore what you don't use
onText([...texts.values()].join('\n\n'));
}
return undefined; // the connection ended early
}
let response = await fetch(`${BASE}/sessions/${encodeURIComponent(session)}/messages`, {
method: 'POST',
headers: { ...headers, 'Content-Type': 'application/json' },
body: JSON.stringify({ input, stream: true, client_message_id: crypto.randomUUID() }),
signal,
});
for (let attempt = 0; ; attempt++) {
try {
const end = await consume(response);
if (end) return end;
} catch (err) {
if (err.code || signal?.aborted) throw err; // an API error, or you aborted
// otherwise a network error: resume below
}
if (!runId || attempt === 3) throw new Error('The stream was interrupted.');
await new Promise((resolve) => setTimeout(resolve, 500 * 2 ** attempt));
response = await fetch(`${BASE}/runs/${runId}/events`, {
headers: lastId ? { ...headers, 'Last-Event-ID': lastId } : headers,
signal,
});
}
}

Use it:

const output = document.querySelector('#reply');
const end = await streamReply({
session: 'order-8812',
input: 'Can I change the delivery address?',
token, // a client token in the browser; a secret key only on your server
onText: (text) => {
output.textContent = text; // plain text: never innerHTML for model output
},
});
if (end.event === 'run.completed' && end.data.run.outcome === 'handed_off') {
showBanner('A member of our team will reply here.');
}

Secret keys never belong in a browser, and the API sends no CORS headers for them. In front-end code, use a client token minted by your server for the signed-in user (POST /v1/client_tokens). Client tokens may call the message and event routes of their own sessions from the browser, including resuming with Last-Event-ID. See End users and identity.

When a client token expires during a stream, the stream sends event: error with {"code": "token_expired"} and closes. streamReply throws an error with that code: fetch a fresh token from your server and resume with GET /v1/runs/{run}/events and the last id.

A run stream ends with its run. To show everything in a conversation — the agent’s replies, your team’s replies during a handoff, mode changes — keep a session stream open:

async function followSession({ session, token, after, onEvent, signal }) {
const url = new URL(`${BASE}/sessions/${encodeURIComponent(session)}/events`);
if (after) url.searchParams.set('after', after); // replay what you missed
const response = await fetch(url, {
headers: { Authorization: `Bearer ${token}`, Accept: 'text/event-stream' },
signal,
});
if (!response.ok) throw new Error((await response.json()).error.code);
for await (const ev of readEvents(response)) {
onEvent(ev); // message.created, message.delta, message.completed, session.updated, …
}
}
  • A session stream never closes on its own; abort it with an AbortController when the chat closes.
  • Without after it starts live; ?after=0 replays everything still retained (7 days). Store the last id and pass it back after a reload.
  • When a stream is open, send messages with "background": true and render only from the session stream — that way no message appears twice. This is how the K-Agent widget works.
  • session.updated with mode: "human" is your cue to show “you’re talking to a person now”.
  • Heartbeats (: ping) arrive every 15 seconds; the parser skips them. If you see nothing for much longer, reconnect.
  • Stopping the stream doesn’t stop the run. Closing the connection only stops your reading; the run finishes and its reply is stored. To stop the run itself, call POST /v1/runs/{run}/cancel.
  • Errors before streaming come back as normal JSON with an HTTP status (409 session_busy, 401 invalid_api_key…), which consume throws with the error code.
  • New events appear over time. Ignore types you don’t handle, as above.