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.
Implements
Section titled “Implements”Constructors
Section titled “Constructors”Constructor
Section titled “Constructor”new HttpStreamSession(opts): HttpStreamSession;Defined in: src/client/stream.ts:125
Parameters
Section titled “Parameters”| 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 |
- |
Returns
Section titled “Returns”HttpStreamSession
Accessors
Section titled “Accessors”header
Section titled “header”Get Signature
Section titled “Get Signature”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.
Returns
Section titled “Returns”Record<string, any> | null
The method’s header row (returned once at stream start), or null if the method declares no header.
Implementation of
Section titled “Implementation of”rawHeader
Section titled “rawHeader”Get Signature
Section titled “Get Signature”get rawHeader(): RawBatch | null;Defined in: src/client/stream.ts:215
The stream’s header batch and its custom metadata, undecoded.
Returns
Section titled “Returns”RawBatch | null
The stream’s header batch, or null when the method declares none.
Implementation of
Section titled “Implementation of”Methods
Section titled “Methods”[asyncIterator]()
Section titled “[asyncIterator]()”asyncIterator: AsyncIterableIterator<Record<string, any>[]>;Defined in: src/client/stream.ts:498
Iterate over producer stream batches.
Returns
Section titled “Returns”AsyncIterableIterator<Record<string, any>[]>
Implementation of
Section titled “Implementation of”cancel()
Section titled “cancel()”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.
Returns
Section titled “Returns”Promise<void>
Implementation of
Section titled “Implementation of”close()
Section titled “close()”close(): void;Defined in: src/client/stream.ts:695
No-op: the HTTP transport is stateless, so there is nothing to tear down.
Returns
Section titled “Returns”void
Implementation of
Section titled “Implementation of”exchange()
Section titled “exchange()”exchange(input): Promise<Record<string, any>[]>;Defined in: src/client/stream.ts:288
Send an exchange request and return the data rows.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
input |
ExchangeInput |
Returns
Section titled “Returns”Promise<Record<string, any>[]>
Implementation of
Section titled “Implementation of”exchangeRaw()
Section titled “exchangeRaw()”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.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
input |
RawBatch |
Returns
Section titled “Returns”Promise<RawBatch | null>
Implementation of
Section titled “Implementation of”nextWithToken()
Section titled “nextWithToken()”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.
Returns
Section titled “Returns”Promise<RowsWithToken | null>
nextWithTokenRaw()
Section titled “nextWithTokenRaw()”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.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
metadata? |
ReadonlyMap<string, string> |
Returns
Section titled “Returns”Promise<
| RawBatchWithToken
| null>
Implementation of
Section titled “Implementation of”RawStreamSession.nextWithTokenRaw
seekToToken()
Section titled “seekToToken()”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.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
token |
string |
Returns
Section titled “Returns”void
tick()
Section titled “tick()”tick(metadata?): Promise<Record<string, any>[]>;Defined in: src/client/stream.ts:357
Send one producer continuation tick with application custom metadata.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
metadata? |
ReadonlyMap<string, string> |
Returns
Section titled “Returns”Promise<Record<string, any>[]>
Implementation of
Section titled “Implementation of”tickRaw()
Section titled “tickRaw()”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.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
metadata? |
ReadonlyMap<string, string> |
Returns
Section titled “Returns”Promise<RawBatch | null>
