Skip to content

Event streams

Streaming methods return an EventStream: an async iterator over Server-Sent Events. It reconnects after a dropped connection and resumes from the last event it saw (Last-Event-ID), and it stops when you break or abort.

Exported from @morevoice/sdk (source: packages/sdk-ts/src/lib/streams.ts).

interface

transcript.partial / transcript.final payload (WS05 W3-16). Partials are superseded by the final.

interface TranscriptObject {
call_id?: string;
speaker: "customer" | "assistant" | "agent";
text: string;
start_ms?: number;
end_ms?: number;
language?: string;
[field: string]: unknown;
}
interface

call.state_changed payload.

interface CallStateObject {
call_id?: string;
status: string;
[field: string]: unknown;
}
interface

call.tool_called payload.

interface ToolCalledObject {
call_id?: string;
tool: string;
arguments?: Record<string, unknown>;
result?: unknown;
[field: string]: unknown;
}
interface

call.ended payload: the call object.

interface CallObject {
id: string;
object?: "call";
status?: string;
end_reason?: string | null;
duration_ms?: number;
[field: string]: unknown;
}
type

One member per type, so Extract<…, { type: K }> and switch (ev.type) narrow per type.

type TranscriptStreamEvent = EventEnvelope<"transcript.partial", TranscriptObject> | EventEnvelope<"transcript.final", TranscriptObject>;
type

What GET /v1/calls/{id}/events carries (WS05 W3-15). W2-03 may replace these with generated types.

type CallStreamEvent =
| TranscriptStreamEvent
| EventEnvelope<"call.state_changed", CallStateObject>
| EventEnvelope<"call.tool_called", ToolCalledObject>
| EventEnvelope<"call.ended", CallObject>;
type
type CallStreamEventType = CallStreamEvent["type"];
interface

What the stream methods need from the client: where, how to authenticate, which fetch.

interface StreamClient {
/** API origin, e.g. "https://api.morevoice.ai" (paths start with `pathPrefix`). */
baseUrl: string;
/** Prefix of every stream path (default "/v1"); "" when `baseUrl` already ends with the version, as the client's does. [WS06-W2-01] */
pathPrefix?: string;
fetch?: FetchLike;
/** Auth and version headers, or a function returning them (called on every reconnect). */
headers?: HeadersProvider;
/** Default for every stream (5). */
maxRetries?: number;
}
interface
interface CallStreamOptions<K extends string> extends CommonStreamOptions {
/** Only these event types (sent to the server, enforced client-side too). */
types?: readonly K[];
/** End the iteration after `call.ended` (default true). */
untilEnded?: boolean;
}
interface
interface EventsStreamOptions extends CommonStreamOptions {
/** Event types to receive; `call.*` matches a group. Default: all. */
types?: readonly string[];
}
interface
interface Streams {
calls: {
stream<K extends CallStreamEventType = CallStreamEventType>(callId: string, options?: CallStreamOptions<K>): EventStream<Extract<CallStreamEvent, { type: K }>>;
transcriptStream(callId: string, options?: CommonStreamOptions & { untilEnded?: boolean }): EventStream<TranscriptStreamEvent>;
};
events: {
stream<T extends EventEnvelope = EventEnvelope>(options?: EventsStreamOptions): EventStream<T>;
};
}
function

Build the stream methods for a client.

function createStreams(client: StreamClient): Streams;
interface

The parts of the generated @hey-api/client-fetch client the transport uses (structural, so no import of src/api).

interface StreamClientLike {
buildUrl(options: { url: string; path?: Record<string, unknown>; query?: object }): string;
getConfig(): { fetch?: FetchLike; headers?: unknown; auth?: unknown };
}
interface

One SSE operation as the generated facade describes it (W1-03 StreamRequest).

interface StreamRequestLike {
operationId: string;
/** Path template, e.g. `/v1/calls/{id}/events`. */
url: string;
path?: Record<string, string>;
query?: object;
options?: { headers?: Record<string, string>; signal?: AbortSignal; lastEventId?: string };
}
function

The facade’s stream transport (FacadeRuntime.stream, W1-03): opens the operation’s SSE endpoint through streamEvents with the client’s base URL, fetch, headers and auth (re-read on every reconnect), the caller’s signal and extra headers, and a client-side types filter. Call-scoped streams (/calls/{…}/…) end on call.ended.

new MoreVoice(...) → bindFacade({ client, paginate: autoPaginate, stream: createStreamTransport() }) // W2-01
function createStreamTransport(defaults: { maxRetries?: number } = {}): (client: StreamClientLike, request: StreamRequestLike) => EventStream<unknown>;

Exported from @morevoice/sdk (source: packages/sdk-ts/src/lib/sse.ts).

interface

One dispatched SSE event.

interface SseMessage {
/** The `event:` field; "message" when the block had none. */
event: string;
/** The `data:` lines joined with "\n". */
data: string;
/** The `id:` field of this block, when it had one. */
id: string | undefined;
/** The last event ID in force at dispatch (an `id:` persists until the next one). What a reconnect sends. */
lastEventId: string;
}
class

Incremental parser for the text/event-stream format (HTML Living Standard §9.2.6): UTF-8 text in, events out. Handles CRLF / LF / CR line ends (also split across chunks), a leading BOM, comments (: heartbeats), multi-line data, id (ignored when it contains NUL) and retry (digits only). An unterminated last block is discarded.

class SseParser {
/** The last event ID buffer (persists across events, as the spec requires). */
lastEventId: string;
/** The latest valid `retry:` value in milliseconds, if the server sent one. */
retry: number | undefined;
constructor(lastEventId = "");
/** Feed decoded text; returns the events completed by it. */
feed(chunk: string): SseMessage[];
/** End of stream: any incomplete block is discarded (spec). Resets the line state for reuse. */
end(): void;
}
function

Parse a complete text/event-stream document (tests, recorded fixtures).

function parseSse(text: string): SseMessage[];
type
type SseErrorCode = "http_error" | "network_error" | "invalid_content_type" | "invalid_json";
class
class SseError extends Error {
readonly name = "SseError";
readonly code: SseErrorCode;
/** HTTP status, for `http_error`. */
readonly status: number | undefined;
/** The start of the error response body, for `http_error`. */
readonly body: string | undefined;
constructor(code: SseErrorCode, message: string, extra: { status?: number; body?: string; cause?: unknown } = {});
}
type

Request headers, or a function returning them (called on every (re)connect, e.g. to refresh a token).

type HeadersProvider = Readonly<Record<string, string>> | (() => Readonly<Record<string, string>> | Promise<Readonly<Record<string, string>>>);
interface
interface StreamEventsOptions<T> {
headers?: HeadersProvider;
/** Abort the stream: the iteration rejects with the signal's reason. (`stream.close()` / `break` end it quietly.) */
signal?: AbortSignal;
/** Resume after this event ID (sent as `Last-Event-ID` on the first request too). */
lastEventId?: string;
/** Consecutive failed (re)connection attempts before giving up (default 5). A connection that delivered data resets the count. */
maxRetries?: number;
/** Delay before reconnect attempt `attempt` (1-based); `retryMs` is the server's `retry:` hint or the default. */
backoff?: (attempt: number, retryMs: number) => number;
/** Reconnection delay before the server sends `retry:` (default 1000 ms). */
retryMs?: number;
/** Stop after an event whose type is in `endEvents` (default false). */
untilEnded?: boolean;
/** The types `untilEnded` stops on (default ["call.ended"]). */
endEvents?: readonly string[];
/** Only yield these event types; `call.*` matches a prefix. Also sent to the server by the stream helpers. */
types?: readonly string[];
/** Turn an SSE message into a value; return undefined to skip it. Default: JSON.parse(data), skipping `ping` / `heartbeat` events. */
parse?: (message: SseMessage) => T | undefined;
/** How many recent event IDs to remember for deduplication (default 1000). */
dedupeWindow?: number;
}
constant
const DEFAULT_SSE_MAX_RETRIES = 5;
function

True when type matches one of patterns (exact, or prefix.*).

function matchesEventType(type: string | undefined, patterns: readonly string[]): boolean;
class

A live event stream: iterate it once with for await. break (or close()) aborts the underlying request. lastEventId is the resume point, e.g. to persist and pass back as lastEventId later.

class EventStream<T> implements AsyncIterable<T> {
constructor(start: (stream: EventStream<T>, signal: AbortSignal) => AsyncGenerator<T>, lastEventId?: string);
/** The ID to resume after (the last event ID the server set), or undefined. */
get lastEventId(): string | undefined;
/** @internal */
_setLastEventId(id: string): void;
/** True once `close()` was called. */
get closed(): boolean;
/** Stop the stream and abort the request; a running `for await` ends without an error. */
close(): void;
[Symbol.asyncIterator](): AsyncIterator<T>;
}
function

Open an SSE stream as a typed async iterable. fetchFn defaults to the global fetch. Reconnects with Last-Event-ID after drops; see the file header for the delivery rules.

function streamEvents<T = unknown>(fetchFn: FetchLike | undefined, url: string | URL, options: StreamEventsOptions<T> = {}): EventStream<T>;