Skip to content

HttpStreamSession

Defined in: src/client/stream.ts:98

StreamSession implementation for the HTTP transport. Stream state is carried statelessly across requests via an HMAC state token: each HttpStreamSession.exchange or producer-continuation POST sends the current token and receives the next one in the response metadata.

new HttpStreamSession(opts): HttpStreamSession;

Defined in: src/client/stream.ts:125

Parameter Type Description
opts { acceptedMaxResponseBytes?: number; authorization?: string; baseUrl: string; callStateToken?: string | null; compressFn?: CompressFn; compressionLevel?: number; decompressFn?: DecompressFn; externalConfig?: ExternalLocationConfig; finished: boolean; header: Record<string, any> | null; inputSchema?: Schema<any>; method: string; onLog?: (msg) => void; outputSchema: Schema; pendingBatches: RecordBatch<any>[]; postFn?: PostFn; prefix: string; rawHeader?: RawBatch | null; stateToken: string | null; } -
opts.acceptedMaxResponseBytes? number -
opts.authorization? string -
opts.baseUrl string -
opts.callStateToken? string | null -
opts.compressFn? CompressFn -
opts.compressionLevel? number -
opts.decompressFn? DecompressFn -
opts.externalConfig? ExternalLocationConfig -
opts.finished boolean -
opts.header Record<string, any> | null -
opts.inputSchema? Schema<any> -
opts.method string -
opts.onLog? (msg) => void -
opts.outputSchema Schema -
opts.pendingBatches RecordBatch<any>[] -
opts.postFn? PostFn -
opts.prefix string The already-namespaced {prefix}/{protocol} an /exchange path hangs off. Folded by the caller, which learns the protocol from the server’s description.
opts.rawHeader? RawBatch | null -
opts.stateToken string | null -

HttpStreamSession

get header(): Record<string, any> | null;

Defined in: src/client/stream.ts:210

The stream’s one-time header row, or null if the method declares no header.

Record<string, any> | null

The method’s header row (returned once at stream start), or null if the method declares no header.

StreamSession.header


get rawHeader(): RawBatch | null;

Defined in: src/client/stream.ts:215

The stream’s header batch and its custom metadata, undecoded.

RawBatch | null

The stream’s header batch, or null when the method declares none.

RawStreamSession.rawHeader

asyncIterator: AsyncIterableIterator<Record<string, any>[]>;

Defined in: src/client/stream.ts:498

Iterate over producer stream batches.

AsyncIterableIterator<Record<string, any>[]>

StreamSession.[asyncIterator]


cancel(): Promise<void>;

Defined in: src/client/stream.ts:405

Ask the server to discard this stream’s state and stop producing.

Sends POST {prefix}/{method}/exchange carrying vgi_rpc.cancel alongside the current tokens, so the server runs the state’s cancel hook and releases it. Idempotent and best-effort: a transport failure here is swallowed, because the session is finished either way. Mirrors Python’s HttpStreamSession.cancel.

Promise<void>

RawStreamSession.cancel


close(): void;

Defined in: src/client/stream.ts:695

No-op: the HTTP transport is stateless, so there is nothing to tear down.

void

RawStreamSession.close


exchange(input): Promise<Record<string, any>[]>;

Defined in: src/client/stream.ts:288

Send an exchange request and return the data rows.

Parameter Type
input ExchangeInput

Promise<Record<string, any>[]>

StreamSession.exchange


exchangeRaw(input): Promise<RawBatch | null>;

Defined in: src/client/stream.ts:387

Send one encoded batch with its custom metadata and read the reply.

The declared-batch branch of HttpStreamSession.exchange without the row decoding: the caller’s schema, buffers and metadata cross verbatim (plus this turn’s stream tokens), and the reply comes back as the server encoded it.

Parameter Type
input RawBatch

Promise<RawBatch | null>

RawStreamSession.exchangeRaw


nextWithToken(): Promise<RowsWithToken | null>;

Defined in: src/client/stream.ts:575

Read one producer batch and surface the worker’s continuation token.

Reads exactly one data batch and returns it paired with the resume token that continues the stream AFTER that batch — the worker’s own serialized producer state. A fresh session positioned at that token (see HttpStreamSession.seekToToken / the client’s resumeStream) resumes on any node, which is the basis for stateless, load-balanced relays that must not pin a scan to one process.

Returns null at end-of-stream. Throws a ProtocolError if a peer violates the lock-step contract by returning more than one data batch in a turn.

Drives the same wire protocol as async iteration but yields one { rows, token } per call instead of auto-following the token. Do not interleave with iteration/exchange on the same session.

Mirrors Python’s HttpStreamSession.next_with_token.

Promise<RowsWithToken | null>


nextWithTokenRaw(metadata?): Promise<
| RawBatchWithToken
| null>;

Defined in: src/client/stream.ts:586

Read one producer batch, undecoded, together with its resume token.

The batch-level twin of HttpStreamSession.nextWithToken; see there for the resume-token contract.

Parameter Type
metadata? ReadonlyMap<string, string>

Promise< | RawBatchWithToken | null>

RawStreamSession.nextWithTokenRaw


seekToToken(token): void;

Defined in: src/client/stream.ts:661

Reposition a freshly-initialised session to resume from token.

Discards any init-preloaded batches and points the session at the given resume token (as returned by HttpStreamSession.nextWithToken), so the next nextWithToken() continues from exactly there. Used to resume a scan on a new process/node — which is why the call token travels inside the blob too: that node may never have seen this stream’s /init. Mirrors Python’s seek_to_token.

Parameter Type
token string

void


tick(metadata?): Promise<Record<string, any>[]>;

Defined in: src/client/stream.ts:357

Send one producer continuation tick with application custom metadata.

Parameter Type
metadata? ReadonlyMap<string, string>

Promise<Record<string, any>[]>

StreamSession.tick


tickRaw(metadata?): Promise<RawBatch | null>;

Defined in: src/client/stream.ts:374

Send one producer tick and return the batch undecoded.

A tick is one batch forward, whatever it took to get there: /init may already have buffered the first one, in which case this consumes that rather than issuing a continuation for a batch the client is holding. Refusing outright — which this did — made tick() unusable on the HTTP transport, because the first tick of every producer is the buffered one. Only explicit metadata is still refused while a batch is buffered, and for a reason that survives: that metadata belongs on a request this turn does not make.

Parameter Type
metadata? ReadonlyMap<string, string>

Promise<RawBatch | null>

RawStreamSession.tickRaw