Skip to content

PipeStreamSession

Defined in: src/client/pipe.ts:122

StreamSession implementation for the pipe/subprocess transport. Drives lockstep streaming over a single bidirectional pipe: each PipeStreamSession.exchange or iteration step writes one input batch and reads one output batch. Holds the connection’s single-threaded busy lock until closed.

new PipeStreamSession(opts): PipeStreamSession;

Defined in: src/client/pipe.ts:144

Parameter Type
opts { externalConfig?: ExternalLocationConfig; header: Record<string, any> | null; onLog?: (msg) => void; outputSchema: Schema; rawHeader?: RawBatch | null; reader: IpcStreamReader | (() => Promise<IpcStreamReader>); releaseBusy: () => void; setDrainPromise: (p) => void; writeFn: WriteFn; }
opts.externalConfig? ExternalLocationConfig
opts.header Record<string, any> | null
opts.onLog? (msg) => void
opts.outputSchema Schema
opts.rawHeader? RawBatch | null
opts.reader IpcStreamReader | (() => Promise<IpcStreamReader>)
opts.releaseBusy () => void
opts.setDrainPromise (p) => void
opts.writeFn WriteFn

PipeStreamSession

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

Defined in: src/client/pipe.ts:168

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/pipe.ts:173

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/pipe.ts:485

Iterate over producer stream batches (lockstep).

AsyncIterableIterator<Record<string, any>[]>

StreamSession.[asyncIterator]


cancel(): Promise<void>;

Defined in: src/client/pipe.ts:326

Signal the server to stop processing and discard the stream’s state.

Writes a zero-row batch carrying vgi_rpc.cancel, closes the input stream, and drains whatever the server still had queued. Idempotent and best-effort: a transport that has already failed is not worth a second failure during teardown. Mirrors Python’s StreamSession.cancel.

Promise<void>

RawStreamSession.cancel


close(): void;

Defined in: src/client/pipe.ts:525

End the stream: close the input side (or send an empty stream if nothing was sent yet) and drain the server’s remaining output in the background, releasing the connection’s busy lock once the drain completes.

void

RawStreamSession.close


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

Defined in: src/client/pipe.ts:363

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/pipe.ts:290

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

The declared-batch branch of PipeStreamSession.exchange without the row decoding: input schema and buffers cross verbatim, and the server’s answer comes back as it was encoded.

Parameter Type
input RawBatch

Promise<RawBatch | null>

RawStreamSession.exchangeRaw


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

Defined in: src/client/pipe.ts:240

A byte-stream transport carries no resumable stream state, so the token is always null. Declared so one caller can drive either transport.

Promise< | RawBatchWithToken | null>

RawStreamSession.nextWithTokenRaw


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

Defined in: src/client/pipe.ts:225

Send one producer tick, preserving application message metadata.

Parameter Type
metadata? ReadonlyMap<string, string>

Promise<Record<string, any>[]>

StreamSession.tick


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

Defined in: src/client/pipe.ts:231

Send one producer tick and return the server’s batch undecoded.

Parameter Type
metadata? ReadonlyMap<string, string>

Promise<RawBatch | null>

RawStreamSession.tickRaw