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.
Implements
Section titled “Implements”Constructors
Section titled “Constructors”Constructor
Section titled “Constructor”new PipeStreamSession(opts): PipeStreamSession;Defined in: src/client/pipe.ts:144
Parameters
Section titled “Parameters”| 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 |
Returns
Section titled “Returns”PipeStreamSession
Accessors
Section titled “Accessors”header
Section titled “header”Get Signature
Section titled “Get Signature”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.
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/pipe.ts:173
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/pipe.ts:485
Iterate over producer stream batches (lockstep).
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/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.
Returns
Section titled “Returns”Promise<void>
Implementation of
Section titled “Implementation of”close()
Section titled “close()”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.
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/pipe.ts:363
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/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.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
input |
RawBatch |
Returns
Section titled “Returns”Promise<RawBatch | null>
Implementation of
Section titled “Implementation of”nextWithTokenRaw()
Section titled “nextWithTokenRaw()”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.
Returns
Section titled “Returns”Promise<
| RawBatchWithToken
| null>
Implementation of
Section titled “Implementation of”RawStreamSession.nextWithTokenRaw
tick()
Section titled “tick()”tick(metadata?): Promise<Record<string, any>[]>;Defined in: src/client/pipe.ts:225
Send one producer tick, preserving application message 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/pipe.ts:231
Send one producer tick and return the server’s batch undecoded.
Parameters
Section titled “Parameters”| Parameter | Type |
|---|---|
metadata? |
ReadonlyMap<string, string> |
Returns
Section titled “Returns”Promise<RawBatch | null>
