/** * Wire protocol — envelope shapes and the transport-agnostic {@link HostEndpoint}/{@link GuestEndpoint} * correlation engines built on top of {@link PortLike}. * * @remarks * Generalizes the pattern hand-rolled per-battery in * `docs/.vitepress/theme/components/agent/litert_lm_worker_proxy.ts` + * `litert_lm_worker.ts` (id-correlated one-shot calls, a persistent stream-sink map fed by unsolicited * events, string-only error crossing) into a reusable, method/stream/event-name-generic core. Both * endpoint classes are PURE over {@link PortLike} — no `Worker`/`postMessage`/`process.send` reference * anywhere in this module; the browser and Node transports supply concrete `PortLike` adapters, while * this file is exercised only against linked in-memory fake ports (see the unit specs). * * `HostEndpoint` queues every outbound call/stream-start made before the guest's `ready` envelope * arrives, then flushes the queue in order once it does. `GuestEndpoint` requires no such queueing (it * only ever reacts to inbound envelopes). */ import type { PortLike } from "./types"; /** * A single argument/result value as it crosses the wire — the codec's (`codec.ts`) output shape. `enc: * 'raw'` ships the value (mostly) untouched; `enc: 'nhtio'` ships an `@nhtio/encoder`-encoded string * (or a BYO-codec-encoded string). `transfer` is a pass-through marker the browser transport unwraps into a `postMessage` transfer list; * Node transports ignore it. */ export type WireValue = { enc: 'raw'; v: unknown; transfer?: unknown[]; } | { enc: 'nhtio'; v: string; }; /** * An error as it crosses the wire. `message`/`name`/`stack` are ALWAYS present baseline string fields * (never omitted, regardless of encoder availability) so error-classification-by-message-signature * always works even when the encoder is unavailable or fails to decode. `nhtio` carries the * `@nhtio/encoder`-encoded original `Error` instance when BOTH sides advertised `encoderAvailable` on * `ready` — the receiving side decodes it for full fidelity (custom error subclasses, extra * properties) and falls back silently to the baseline fields on any decode failure. */ export interface WireError { /** Baseline error message — always populated, even when `nhtio` is absent or fails to decode. */ message: string; /** Baseline error name (e.g. `'TypeError'`) — always populated. */ name: string; /** Baseline stack trace, when the original error had one — always forwarded as-is (never re-derived). */ stack?: string; /** `@nhtio/encoder`-encoded original `Error`, when both sides have the encoder. */ nhtio?: string; } /** Host → guest envelopes. */ export type HostToGuestEnvelope = { t: 'call'; id: string; method: string; args: WireValue[]; } | { t: 'hostresult'; id: string; ok: true; value: WireValue | string; } | { t: 'hostresult'; id: string; ok: false; error?: WireError; value?: string; } | { t: 'abort'; id: string; } | { t: 'stream:start'; id: string; stream: string; args: WireValue[]; } | { t: 'stream:cancel'; id: string; reason?: WireValue; } | { t: 'shutdown'; }; /** Guest → host envelopes. */ export type GuestToHostEnvelope = { t: 'ready'; encoderAvailable: boolean; } | { t: 'hostcall'; id: string; method: string; args: WireValue[]; } | { t: 'result'; id: string; ok: true; value: WireValue; } | { t: 'result'; id: string; ok: false; error: WireError; } | { t: 'stream:delta'; id: string; delta: WireValue; } | { t: 'stream:end'; id: string; } | { t: 'stream:error'; id: string; error: WireError; } | { t: 'event'; channel: string; payload: WireValue; }; /** Either direction's envelope — used by generic wire-tracing hooks. */ export type WireEnvelope = HostToGuestEnvelope | GuestToHostEnvelope; /** Monotonic id generator shared by both endpoints (module-scoped counter — fine across many * instances in one realm since ids are only ever compared within a single connection). */ export declare const nextCorrelationId: () => string; interface StreamSink { push: (delta: WireValue) => void; end: () => void; error: (err: WireError) => void; } /** Hooks {@link HostEndpoint} invokes on protocol-level events; `host.ts` wires these to the * observability layer + guest-event fan-out. All optional. */ export interface HostcallQuotas { /** Per-request deadline in milliseconds. */ hostcallTimeoutMs: number; /** Maximum accepted requests for one evaluation. */ maxHostcallsPerEvaluation: number; /** Maximum concurrently running requests. */ maxConcurrentHostcalls: number; } /** Host-side capability registry. The handler receives decoded wire arguments. */ export type HostcallHandler = (args: WireValue[], signal: AbortSignal) => WireValue | Promise; /** UTF-8 producer-side measurement used by both RPC realms. */ export declare const measureHostcallBytes: (value: unknown) => number; /** * Callback surface for observing host-endpoint lifecycle and guest-originated events. * * @remarks Hooks are notifications only; dispatch and correlation remain owned by the endpoint. */ export interface HostEndpointHooks { /** The guest's `ready` envelope arrived. */ onReady?: (info: { encoderAvailable: boolean; }) => void; /** An `event` envelope arrived for `channel`. */ onEvent?: (channel: string, payload: WireValue) => void; /** A guest-to-host capability request arrived. It is deliberately independent of `call`. */ onHostcall?: (id: string, method: string, args: WireValue[]) => void; /** Any envelope was sent (`dir: 'out'`) or received (`dir: 'in'`) — for wire tracing. */ onEnvelope?: (dir: 'out' | 'in', envelope: WireEnvelope) => void; } /** * Host-side correlation engine over a {@link PortLike}. Queues calls/stream-starts made before `ready` * and flushes them in order once it arrives; tracks in-flight calls (one-shot, resolved/rejected by a * `result` envelope) and open streams (persistent, fed by `stream:delta`/`stream:end`/`stream:error` * until closed). `terminate()` rejects every in-flight call and errors every open stream with a * caller-supplied reason (the message text `host.ts` uses is `E_ISOLATED_TERMINATED`'s message). */ export declare class HostEndpoint { #private; constructor(port: PortLike, hooks?: HostEndpointHooks, hostcalls?: { handlers?: ReadonlyMap; quotas?: HostcallQuotas; maxHostcallBytes?: number; }); /** Whether the guest has signaled `ready` yet. */ get isReady(): boolean; /** Number of calls currently awaiting a `result` envelope. Used by `host.ts` to report an accurate * `inFlight` count on a crash before `terminate()` clears the pending map. */ get pendingCallCount(): number; /** Number of streams currently open (started, not yet ended/errored). Used by `host.ts` alongside * {@link pendingCallCount} to report an accurate `inFlight` count on a crash. */ get openStreamCount(): number; /** * Issue a request/response call. Resolves with the guest's returned {@link WireValue}, rejects with a * reconstructed `Error` (see {@link wireErrorToError}) on failure or on `terminate()`. */ call(method: string, args: WireValue[]): { id: string; promise: Promise; }; /** Post a guest capability result. Unknown/late ids are harmlessly ignored by the guest. */ hostresult(id: string, result: { ok: true; value: WireValue | string; } | { ok: false; error?: WireError; value?: string; }): void; /** Send an `abort` envelope for an in-flight call's id. Does not itself reject the call — the guest * is expected to respond with a `result` (ok:false) once it observes the abort. */ abort(id: string): void; /** * Start a fire-and-forward stream. Returns the correlation id immediately (before the guest * necessarily even exists, if not yet `ready`) and a `sink` the caller wires to a `ReadableStream` * controller. */ startStream(stream: string, args: WireValue[], sink: StreamSink): string; /** Send a `stream:cancel` envelope and stop tracking the stream locally. */ cancelStream(id: string, reason?: WireValue): void; /** Send a `shutdown` envelope (graceful-exit request; does not itself tear down the port). */ shutdown(): void; /** * Reject every in-flight call and error every open stream with `reason`, clear all queued-but-unsent * envelopes, and unsubscribe from the port. Idempotent. */ terminate(reason: string): void; } /** Reconstruct an `Error` from a {@link WireError} baseline (name/message/stack only — the `nhtio`-rich * path is decoded separately by the caller when an encoder is available; see `host.ts`). */ export declare const wireErrorToError: (wireError: WireError) => Error; /** Hooks {@link GuestEndpoint} invokes for the guest server (`serve.ts`) to react to. */ export interface GuestEndpointHooks { /** A `call` envelope arrived — resolve/reject `settle` with the method's outcome. */ onCall?: (id: string, method: string, args: WireValue[], signal: AbortSignal) => void; /** A host capability result arrived. */ onHostResult?: (id: string, result: { ok: true; value: WireValue | string; } | { ok: false; error?: WireError; value?: string; }) => void; /** A `stream:start` envelope arrived — the handler pushes deltas via the returned sink. */ onStreamStart?: (id: string, stream: string, args: WireValue[], signal: AbortSignal) => void; /** A `stream:cancel` envelope arrived for an open stream id. */ onStreamCancel?: (id: string, reason?: WireValue) => void; /** A `shutdown` envelope arrived. */ onShutdown?: () => void; /** Any envelope was sent (`dir: 'out'`) or received (`dir: 'in'`) — for wire tracing. */ onEnvelope?: (dir: 'out' | 'in', envelope: WireEnvelope) => void; } /** * Guest-side correlation engine over a {@link PortLike}. Owns per-call `AbortController`s (aborted on * an inbound `abort`/`stream:cancel` envelope) and exposes `settleCall`/`pushDelta`/`endStream`/ * `errorStream` for `serve.ts` to report outcomes back across the wire. */ export declare class GuestEndpoint { #private; constructor(port: PortLike, hooks?: GuestEndpointHooks); /** Issue a guest-to-host capability request using the separate hostcall id space. */ hostcall(method: string, args: WireValue[], maxBytes?: number): { id: string; promise: Promise; }; /** Announce readiness. Must be sent exactly once, before any `result`/`stream:*`/`event` envelope. */ ready(encoderAvailable: boolean): void; /** Report a successful call outcome and release the call's abort controller. */ settleOk(id: string, value: WireValue): void; /** Report a failed call outcome and release the call's abort controller. */ settleError(id: string, error: WireError): void; /** Push a stream delta. */ pushDelta(id: string, delta: WireValue): void; /** Signal a stream's clean end and release its abort controller. */ endStream(id: string): void; /** Signal a stream's terminal error and release its abort controller. */ errorStream(id: string, error: WireError): void; /** Reject all guest capability requests when this endpoint is stopped. */ terminate(reason?: string): void; /** Emit an unsolicited event on `channel`. */ emit(channel: string, payload: WireValue): void; } export {};