File size: 10,429 Bytes
f0634fb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
/**
 * acp-server bootstrap β€” wires `@moonshot-ai/agent-core-v2` (the DI Γ— Scope
 * engine) into an ACP (Agent Client Protocol) stdio server.
 *
 * Composition root: `bootstrap()` builds the App `Scope`; a `@moonshot-ai/
 * klient` facade over the in-memory transport is created on top of it, and
 * every ACP method handler drives the engine through that facade. The
 * ACP-backed `IHostFileSystem` (./acp-fs) is imported for its Session-scope
 * registration side effect (see the import below) β€” being registered on the
 * same scope the memory transport dispatches against, it keeps working
 * unchanged.
 */

import { Readable, Writable } from 'node:stream';

import { ndJsonStream, type AgentConnection, type Stream } from '@agentclientprotocol/sdk';
import {
  bootstrap,
  drainLogCloses,
  drainQueryStoreDisposals,
  drainSessionIndexMirror,
  drainSessionMetadataWrites,
  ensureMainAgent,
  getLiveSessionById,
  IAgentLifecycleService,
  IAgentRuntimeBindingService,
  IAppendLogStore,
  IHostEnvironment,
  IHostProcessService,
  ISessionContext,
  ISessionIndexMirror,
  IWorkspaceInstanceManager,
  logSeed,
  resolveConfigPath,
  resolveKimiHome,
  resolveLoggingConfig,
  type Scope,
  type ScopeSeed,
  sessionMediaOriginalsDir,
} from '@moonshot-ai/agent-core-v2';
import type { Klient } from '@moonshot-ai/klient';
import { createKlient } from '@moonshot-ai/klient/memory';

import { acpClientFromContext } from './acp-client';
// Importing the `acp-fs` barrel also registers the ACP-backed Session-scope
// `IHostFileSystem` and the App-scope `IAcpConnection` holder via the barrel's
// module side effects. `IAcpConnection` is used below to bind the ACP client
// connection.
import { IAcpConnection } from './acp-fs';
import { AcpRuntimeProviderFactory } from './acp-terminal';
import { AcpServer, type AcpServerOptions, createAcpAgentApp } from './server';

export interface RunAcpServerOptions extends AcpServerOptions {
  readonly homeDir?: string;
  readonly configPath?: string;
  readonly input?: NodeJS.ReadableStream;
  readonly output?: NodeJS.WritableStream;
  /**
   * Extra App-scope service seeds forwarded to `bootstrap()`. Intended for
   * tests β€” e.g. seeding a scripted `IProtocolAdapterRegistry` to drive a
   * deterministic turn without a real LLM. Seeds shadow any registered
   * binding with the same service identifier.
   */
  readonly extraSeeds?: ScopeSeed;
}

export interface RunningAcpServer {
  readonly core: Scope;
  readonly klient: Klient;
  readonly conn: AgentConnection;
  close(): Promise<void>;
}

/**
 * Redirect `console.*` to stderr. Stdout is the ACP JSON-RPC channel; any stray
 * write from a dependency would corrupt the protocol stream.
 */
function redirectConsoleToStderr(): void {
  const sink = (...args: unknown[]): void => {
    process.stderr.write(`${args.map(String).join(' ')}\n`);
  };
  globalThis.console.log = sink;
  globalThis.console.info = sink;
  globalThis.console.warn = sink;
  globalThis.console.debug = sink;
}

/**
 * Drive an {@link AcpServer} over an arbitrary ACP {@link Stream}.
 *
 * Boots `agent-core-v2`, creates the in-memory `Klient` facade over the app
 * scope, binds the ACP client connection into {@link IAcpConnection} (so the
 * `acp` `IHostFileSystem` can reverse-RPC file IO), and resolves when the
 * connection closes.
 */
export async function runAcpServerWithStream(
  stream: Stream,
  opts: RunAcpServerOptions = {},
): Promise<RunningAcpServer> {
  const homeDir = resolveKimiHome(opts.homeDir);
  const configPath = resolveConfigPath({ homeDir, configPath: opts.configPath });
  // `ILogOptions` (logSeed) is required by the Session-scoped log writer; any
  // session creation would otherwise fail to instantiate the Session scope.
  const logging = resolveLoggingConfig({ homeDir, env: process.env });
  // `bootstrap()` seeds `IFileSystemStorageService` with a `FileStorageService`
  // rooted at `homeDir`, so session metadata, wire records, blobs, and the
  // session index all persist to disk. `clientIdentity` is required by the
  // engine: reuse the advertised ACP `agentInfo` (the embedding CLI's
  // name/version) with the CLI platform β€” the literal matches
  // `KIMI_CODE_PLATFORM` from `@moonshot-ai/kimi-code-oauth`, which this
  // package does not depend on.
  const { app: core } = bootstrap(
    {
      homeDir,
      configPath,
      clientIdentity: {
        productName: opts.agentInfo?.name ?? 'kimi-code-acp',
        version: opts.agentInfo?.version ?? '0.0.0',
        platform: 'kimi_code_cli',
      },
    },
    [...logSeed(logging), ...(opts.extraSeeds ?? [])],
  );

  // The klient dispatches against the same app scope β€” calls and events stay
  // in-process but observe wire-shaped (JSON-cloned) data. The klient does
  // NOT own the scope: lifecycle stays with this composition root.
  const klient = createKlient({ scope: core });
  const acpConnection = core.accessor.get(IAcpConnection);

  // Route every inbound ACP method to the `AcpServer`. The app must be
  // connected before the outbound client surface (`conn.client`) β€” and thus
  // the server β€” exists, so handlers dereference `server` lazily. No inbound
  // message can be dispatched before this synchronous block yields (the
  // connection's reader only runs on later microtasks), so `server` is always
  // assigned by the time a handler fires.
  let server: AcpServer;
  const app = createAcpAgentApp(() => server);
  const conn = app.connect(stream);
  const client = acpClientFromContext(conn.client);
  // Bind the process-wide ACP client connection before any session performs
  // file IO. The `acp` `IHostFileSystem` reads it lazily via
  // `IAcpConnection.get()`.
  acpConnection.bind(client);
  const workspaceManager = core.accessor.get(IWorkspaceInstanceManager);
  const acpRuntimeProvider = new AcpRuntimeProviderFactory(acpConnection, core.accessor.get(IHostEnvironment), core.accessor.get(IHostProcessService));
  const acpProviderRegistration = await workspaceManager.addProvider(acpRuntimeProvider);
  const sessionWorkspaces = new Map<string, string>();
  server = new AcpServer(client, klient, acpConnection, {
    agentInfo: opts.agentInfo,
    disableAuth: opts.disableAuth,
    terminalAuthEnv: opts.terminalAuthEnv,
    terminalAuthLegacyCommand: opts.terminalAuthLegacyCommand,
    slashCommands: opts.slashCommands,
    bindSessionRuntime: async (sessionId) => {
      const handle = getLiveSessionById(core.accessor, sessionId);
      if (handle === undefined) throw new Error(`session ${sessionId} is not live`);
      const context = handle.accessor.get(ISessionContext);
      const runtimeId = acpRuntimeProvider.bindSession(context.workspaceId, sessionId, context.cwd);
      sessionWorkspaces.set(sessionId, context.workspaceId);
      const agentContext = await ensureMainAgent(handle, { runtimeId });
      handle.accessor
        .get(IAgentLifecycleService)
        .handleOf(agentContext.agentId)!
        .accessor.get(IAgentRuntimeBindingService)
        .switch(runtimeId);
    },
    unbindSessionRuntime: async (sessionId) => {
      const workspaceId = sessionWorkspaces.get(sessionId);
      if (workspaceId === undefined) return;
      sessionWorkspaces.delete(sessionId);
      await acpRuntimeProvider.unbindSession(workspaceId, sessionId);
    },
    // Prompt-image compression persists originals into the session's own
    // media-originals dir (same resolution as kap-server's prompt route):
    // live session scope β†’ `ISessionContext.sessionDir`. A session that is
    // not live in this process yields undefined β†’ temp-dir fallback.
    resolveOriginalsDir: (sessionId) => {
      const handle = getLiveSessionById(core.accessor, sessionId);
      return handle === undefined
        ? undefined
        : sessionMediaOriginalsDir(handle.accessor.get(ISessionContext).sessionDir);
    },
  });

  let closePromise: Promise<void> | undefined;
  const close = async (): Promise<void> => {
    if (closePromise !== undefined) return closePromise;
    closePromise = (async () => {
      // Detach the klient's event subscriptions first so disposal below cannot
      // deliver into a torn-down scope.
      await klient.close();
      // Flush the append-log write-behind before disposing, so a clean shutdown
      // never races a pending drain against teardown (and doesn't drop the last
      // persisted ops). Best-effort: a flush failure must not block disposal.
      const appendLogStore = core.accessor.get(IAppendLogStore);
      try {
        await appendLogStore.flush();
      } catch {
        // ignore β€” disposal proceeds regardless
      }
      // Same shutdown order as kap-server: settle queued session-metadata
      // writes, then drain the session-index mirror while the query store is
      // still open, so a queued summary lands in the read model.
      await drainSessionMetadataWrites();
      await core.accessor.get(ISessionIndexMirror).drain();
      await acpProviderRegistration.dispose();
      core.dispose();
      // `core.dispose()` runs the mirror's and the query store's synchronous
      // `dispose()`, whose drains/closes are asynchronous β€” await them so an
      // embedding host that removes homeDir right after close() never races
      // an in-flight shard close (ENOTEMPTY on teardown). The same window
      // exists for the append-log retirement flushes released by disposal.
      await appendLogStore.drainRetirements();
      await drainSessionIndexMirror();
      await drainQueryStoreDisposals();
      await drainSessionMetadataWrites();
      await drainLogCloses();
    })();
    return closePromise;
  };

  void conn.closed.then(() => {
    void close();
  });

  return { core, klient, conn, close };
}

/**
 * Drive an {@link AcpServer} over Node stdio (or the supplied streams).
 *
 * The ACP SDK speaks Web `ReadableStream` / `WritableStream`, so Node stdio is
 * bridged through `Readable.toWeb` / `Writable.toWeb`.
 */
export async function runAcpServer(opts: RunAcpServerOptions = {}): Promise<void> {
  redirectConsoleToStderr();
  const input = (opts.input ?? process.stdin) as Readable;
  const output = (opts.output ?? process.stdout) as Writable;
  const stream = ndJsonStream(Writable.toWeb(output), Readable.toWeb(input));
  const server = await runAcpServerWithStream(stream, opts);
  await server.conn.closed;
  await server.close();
}