File size: 2,867 Bytes
4e23b01 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 | /**
* Transport SPI — the single abstraction every klient transport implements.
*
* A `KlientChannel` carries service calls and event subscriptions for one
* scope triple. The facade above it never knows which transport is underneath
* (http, ipc, or in-memory); transports never know which facade method
* triggered a frame. `ScopeRef` already carries session/agent coordinates so
* future session/agent facades plug in without changing this interface.
*/
export interface IDisposable {
dispose(): void;
}
/** Optional per-call knobs a transport may honor. */
export interface CallOptions {
/**
* Per-call deadline (ms). A transport with a default call timeout (ipc)
* takes it as an override — long-poll calls pass a deadline covering the
* engine-side wait; transports without a timeout (memory) ignore it.
*/
readonly timeoutMs?: number;
}
/** Scope coordinates of a call/subscription. Empty object = core (app) scope. */
export interface ScopeRef {
readonly workspaceId?: string;
readonly sessionId?: string;
readonly agentId?: string;
}
/**
* Where an event subscription reads from:
* - `stream` — a scope's named event stream, mirroring kap-server's WS
* `eventMap`: core `events` (the global `IEventService` bus), session
* `interactions` / `interactions:resolved`, agent `events` (the per-agent
* `IEventBus`). The scope coordinates disambiguate which scope's stream.
* - `emitter` — one service's `onDid*` `Event<T>` property, addressed by the
* service's wire name and the property name (e.g. `onDidChangeModels`).
*/
export type EventSourceRef =
| { readonly kind: 'stream'; readonly name: string }
| { readonly kind: 'emitter'; readonly service: string; readonly event: string };
export interface KlientChannel {
/** Invoke `service.method(...args)` in the given scope; resolves with the raw wire result. */
call(
scope: ScopeRef,
service: string,
method: string,
args: unknown[],
options?: CallOptions,
): Promise<unknown>;
/**
* Invoke `service.method(...args)` in the given scope and return a streaming
* result. The callee must return an `AsyncIterable`; each yielded chunk is
* surfaced as-is (after the transport's serialization round-trip).
*/
stream(scope: ScopeRef, service: string, method: string, args: unknown[]): AsyncIterable<unknown>;
/**
* Subscribe to an event source; `handler` receives raw wire payloads.
* `onError` reports asynchronous subscription failures (bad source, dropped
* remote subscription) — synchronous validation may also throw.
*/
listen(
scope: ScopeRef,
source: EventSourceRef,
handler: (data: unknown) => void,
onError?: (error: Error) => void,
): IDisposable;
/** Tear the transport down (sockets, lazy bridges). Rejects in-flight calls. */
close(): Promise<void>;
}
|