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>;
}