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