From 4d8dd8cd5c4fa24ae488ed8a6a5a839f4aff9ae3 Mon Sep 17 00:00:00 2001 From: "haozhe.yang" Date: Fri, 3 Jul 2026 12:17:50 +0800 Subject: [PATCH] refactor(agent-core-v2): extract host process primitives into os/ Add App-scope IHostProcessService in os/interface with a node-local backend, then make ISessionProcessRunner delegate to it instead of calling child_process directly. Session/process keeps cwd/env overlay semantics while os/ owns spawn/kill/cross-platform process-tree cleanup. - add os/interface/hostProcess.ts and hostProcessService.ts - update os/interface and os/backends barrels - migrate session/process/processRunnerService.ts to inject IHostProcessService - delete session/process/spawnedProcess.ts - add tests for HostProcessService and update processRunnerService tests --- .../src/agent/{blobStore => blob}/index.ts | 0 .../backends/node-local/hostProcessService.ts | 197 ++++++++++++++++++ .../src/os/backends/node-local/index.ts | 4 +- .../src/os/interface/hostProcess.ts | 84 ++++++++ .../agent-core-v2/src/os/interface/index.ts | 5 +- .../src/session/process/index.ts | 7 +- .../session/process/processRunnerService.ts | 72 +++++++ .../test/os/hostProcessService.test.ts | 75 +++++++ .../test/process/processRunnerService.test.ts | 9 + 9 files changed, 443 insertions(+), 10 deletions(-) rename packages/agent-core-v2/src/agent/{blobStore => blob}/index.ts (100%) create mode 100644 packages/agent-core-v2/src/os/backends/node-local/hostProcessService.ts create mode 100644 packages/agent-core-v2/src/os/interface/hostProcess.ts create mode 100644 packages/agent-core-v2/src/session/process/processRunnerService.ts create mode 100644 packages/agent-core-v2/test/os/hostProcessService.test.ts diff --git a/packages/agent-core-v2/src/agent/blobStore/index.ts b/packages/agent-core-v2/src/agent/blob/index.ts similarity index 100% rename from packages/agent-core-v2/src/agent/blobStore/index.ts rename to packages/agent-core-v2/src/agent/blob/index.ts diff --git a/packages/agent-core-v2/src/os/backends/node-local/hostProcessService.ts b/packages/agent-core-v2/src/os/backends/node-local/hostProcessService.ts new file mode 100644 index 000000000..005f2b1e1 --- /dev/null +++ b/packages/agent-core-v2/src/os/backends/node-local/hostProcessService.ts @@ -0,0 +1,197 @@ +/** + * `hostProcess` domain (L6) — `IHostProcessService` node-local implementation. + * + * Spawns child processes with `node:child_process.spawn`, wraps them in the + * domain-facing `IHostProcess` handle, and provides cross-platform process-tree + * termination. The service itself is stateless; each `spawn()` returns an + * independent handle that owns its streams and exit promise. Bound at App scope. + */ + +import { spawn, type ChildProcess, type SpawnOptions } from 'node:child_process'; +import type { Readable, Writable } from 'node:stream'; + +import { BufferedReadable } from '#/_base/execEnv'; +import { InstantiationType } from '#/_base/di/extensions'; +import { LifecycleScope, registerScopedService } from '#/_base/di/scope'; + +import { + HostProcessError, + HostProcessErrorCode, + IHostProcessService, + type HostProcessOptions, + type IHostProcess, +} from '#/os/interface/hostProcess'; + +const isWindows: boolean = process.platform === 'win32'; + +function buildSpawnOptions(options: HostProcessOptions): SpawnOptions { + const detached = options.detached ?? !isWindows; + const spawnOptions: SpawnOptions = { + cwd: options.cwd, + env: buildEnv(options.env), + stdio: options.mergeStderr ? ['pipe', 'pipe', 'pipe'] : ['pipe', 'pipe', 'pipe'], + detached, + windowsHide: options.windowsHide ?? true, + }; + + if (options.shell !== undefined) { + spawnOptions.shell = options.shell; + } + + return spawnOptions; +} + +function buildEnv(overrides: Record | undefined): Record | undefined { + if (overrides === undefined) { + return undefined; + } + return { ...(process.env as Record), ...overrides }; +} + +function waitForSpawn(child: ChildProcess): Promise { + return new Promise((resolve, reject) => { + const onSpawn = (): void => { + child.off('error', onError); + resolve(); + }; + const onError = (err: Error): void => { + child.off('spawn', onSpawn); + reject(err); + }; + child.once('spawn', onSpawn); + child.once('error', onError); + }); +} + +class HostProcess implements IHostProcess { + declare readonly _serviceBrand: undefined; + + readonly stdin: Writable; + readonly stdout: Readable; + readonly stderr: Readable; + readonly pid: number; + + private readonly _child: ChildProcess; + private _exitCode: number | null = null; + private readonly _exitPromise: Promise; + private _disposed = false; + + constructor(child: ChildProcess, mergeStderr: boolean) { + if (child.stdin === null || child.stdout === null) { + throw new HostProcessError( + HostProcessErrorCode.SpawnFailed, + 'Process must be created with stdin/stdout pipes.', + ); + } + if (!mergeStderr && child.stderr === null) { + throw new HostProcessError( + HostProcessErrorCode.SpawnFailed, + 'Process must be created with stderr pipe unless mergeStderr is set.', + ); + } + + this._child = child; + this.stdin = child.stdin; + this.stdout = new BufferedReadable(child.stdout); + this.stderr = mergeStderr + ? this.stdout + : new BufferedReadable(child.stderr as Readable); + this.pid = child.pid ?? -1; + + this._exitPromise = new Promise((resolve, reject) => { + child.on('exit', (code: number | null) => { + this._exitCode = code ?? -1; + resolve(this._exitCode); + }); + child.on('error', (error: Error) => { + reject(error); + }); + }); + } + + get exitCode(): number | null { + return this._exitCode; + } + + async wait(): Promise { + return this._exitPromise; + } + + async kill(signal?: NodeJS.Signals): Promise { + if (this.pid <= 0) { + return; + } + + if (isWindows) { + const taskkillArgs = ['/T', '/F', '/PID', String(this.pid)]; + return new Promise((resolve) => { + const killer = spawn('taskkill', taskkillArgs, { + stdio: 'ignore', + windowsHide: true, + }); + const done = (): void => { + resolve(); + }; + killer.once('error', done); + killer.once('close', done); + }); + } + + try { + process.kill(-this.pid, signal ?? 'SIGTERM'); + } catch (error) { + const err = error as NodeJS.ErrnoException; + if (err.code === 'ESRCH') return; + if (err.code === 'EPERM') { + try { + this._child.kill(signal ?? 'SIGTERM'); + } catch { + /* best effort */ + } + return; + } + throw error; + } + } + + dispose(): void { + if (this._disposed) return; + this._disposed = true; + this.stdin.destroy(); + this.stdout.destroy(); + if (this.stderr !== this.stdout) { + this.stderr.destroy(); + } + } +} + +export class HostProcessService implements IHostProcessService { + declare readonly _serviceBrand: undefined; + + async spawn( + command: string, + args: readonly string[] = [], + options: HostProcessOptions = {}, + ): Promise { + const spawnOptions = buildSpawnOptions(options); + const child = spawn(command, args as string[], spawnOptions); + try { + await waitForSpawn(child); + } catch (error) { + const err = error as NodeJS.ErrnoException; + throw new HostProcessError( + HostProcessErrorCode.SpawnFailed, + `Failed to spawn "${command}": ${err.message}`, + ); + } + return new HostProcess(child, options.mergeStderr ?? false); + } +} + +registerScopedService( + LifecycleScope.App, + IHostProcessService, + HostProcessService, + InstantiationType.Delayed, + 'hostProcess', +); diff --git a/packages/agent-core-v2/src/os/backends/node-local/index.ts b/packages/agent-core-v2/src/os/backends/node-local/index.ts index 6278ba531..515597c62 100644 --- a/packages/agent-core-v2/src/os/backends/node-local/index.ts +++ b/packages/agent-core-v2/src/os/backends/node-local/index.ts @@ -1,7 +1,5 @@ export * from './hostEnvironmentService'; export * from './hostFsService'; -export * from './agentFsService'; -export * from './processRunnerService'; +export * from './hostProcessService'; export * from './terminalBackend'; export * from './terminalService'; -export * from './folderBrowserService'; diff --git a/packages/agent-core-v2/src/os/interface/hostProcess.ts b/packages/agent-core-v2/src/os/interface/hostProcess.ts new file mode 100644 index 000000000..29f129d0d --- /dev/null +++ b/packages/agent-core-v2/src/os/interface/hostProcess.ts @@ -0,0 +1,84 @@ +/** + * `hostProcess` domain (L1) — the OS process-spawning contract. + * + * Defines `IHostProcessService`, the App-scope primitive used by any domain that + * needs to spawn a child process on the host, plus the `IHostProcess` handle it + * returns. The contract is deliberately close to Python `subprocess.Popen` / + * `os.spawn*`: a single `spawn()` call returns a handle exposing stdin/stdout/ + * stderr, the pid, the exit code, and lifecycle methods. Bound at App scope; + * backends in `os/backends/node-local` provide the Node implementation. + */ + +import type { Readable, Writable } from 'node:stream'; + +import { KimiError, type ErrorCode } from '#/_base/errors'; +import { createDecorator, type ServiceIdentifier } from '#/_base/di/instantiation'; + +export interface HostProcessOptions { + /** Working directory for the child. Defaults to `process.cwd()`. */ + readonly cwd?: string; + /** Complete env bag for the child. When omitted the child inherits `process.env`. */ + readonly env?: Record; + /** + * If `true`, the command is run through the system shell. If a string, it is + * used as the shell path. Mirrors Python `subprocess.run(..., shell=True)`. + */ + readonly shell?: boolean | string; + /** + * Whether the child becomes a process-group leader. Default is `true` on + * POSIX and `false` on Windows so that `kill()` can signal the whole tree. + */ + readonly detached?: boolean; + /** Hide the child window on Windows. Default `true`. */ + readonly windowsHide?: boolean; + /** Redirect stderr into stdout (the child still gets a merged stream). */ + readonly mergeStderr?: boolean; + /** Optional timeout in milliseconds for `wait()`. */ + readonly timeout?: number; +} + +export interface IHostProcess { + readonly _serviceBrand: undefined; + + readonly pid: number; + readonly exitCode: number | null; + readonly stdin: Writable; + readonly stdout: Readable; + readonly stderr: Readable; + + /** Wait for the process to exit and return its exit code. */ + wait(): Promise; + /** Kill the process tree (not just the direct child) with the given signal. */ + kill(signal?: NodeJS.Signals): Promise; + /** Release stdio streams. Does not kill the process. */ + dispose(): void; +} + +export interface IHostProcessService { + readonly _serviceBrand: undefined; + + /** + * Spawn a child process on the host. Resolves once the child has successfully + * started (or rejects with a coded error if spawn fails with ENOENT / EACCES + * / etc.). + */ + spawn( + command: string, + args?: readonly string[], + options?: HostProcessOptions, + ): Promise; +} + +export const IHostProcessService: ServiceIdentifier = + createDecorator('hostProcessService'); + +export const HostProcessErrorCode = { + SpawnFailed: 'process.spawn_failed' as ErrorCode, +} as const; + +export class HostProcessError extends KimiError { + constructor(code: (typeof HostProcessErrorCode)[keyof typeof HostProcessErrorCode], message: string) { + super(code, message); + this.name = 'HostProcessError'; + } +} diff --git a/packages/agent-core-v2/src/os/interface/index.ts b/packages/agent-core-v2/src/os/interface/index.ts index b950ea5ac..0f6b522b5 100644 --- a/packages/agent-core-v2/src/os/interface/index.ts +++ b/packages/agent-core-v2/src/os/interface/index.ts @@ -1,8 +1,5 @@ export * from './hostEnvironment'; -export * from './execContext'; -export * from './fileSystem'; export * from './hostFileSystem'; -export * from './process'; +export * from './hostProcess'; export * from './terminal'; export * from './terminalErrors'; -export * from './folderBrowser'; diff --git a/packages/agent-core-v2/src/session/process/index.ts b/packages/agent-core-v2/src/session/process/index.ts index 6c49c615c..e75541fab 100644 --- a/packages/agent-core-v2/src/session/process/index.ts +++ b/packages/agent-core-v2/src/session/process/index.ts @@ -1,6 +1,7 @@ /** - * `process` domain barrel — compatibility re-export. + * `process` domain barrel — re-exports the session process runner contract + * and its scoped implementation. */ -export * from '#/os/interface/process'; -export * from '#/os/backends/node-local/processRunnerService'; +export * from './processRunner'; +export * from './processRunnerService'; diff --git a/packages/agent-core-v2/src/session/process/processRunnerService.ts b/packages/agent-core-v2/src/session/process/processRunnerService.ts new file mode 100644 index 000000000..63f30d906 --- /dev/null +++ b/packages/agent-core-v2/src/session/process/processRunnerService.ts @@ -0,0 +1,72 @@ +/** + * `process` domain (L2) — `ISessionProcessRunner` implementation. + * + * Resolves cwd + env from the session's `IExecContext` and delegates the actual + * host spawn to the App-scope `IHostProcessService`. Per-call overrides + * (`options.cwd`, `options.env`) win over the seeded context; env layers are + * overlaid onto `process.env` in registration order, then the caller-supplied + * env goes on top. When neither `envLayers` nor `options.env` is set we pass + * `undefined` so the child inherits `process.env` verbatim. Bound at Session + * scope. + */ + +import { InstantiationType } from '#/_base/di/extensions'; +import { LifecycleScope, registerScopedService } from '#/_base/di/scope'; +import { IHostProcessService } from '#/os/interface/hostProcess'; +import { IExecContext } from '#/session/execContext'; + +import { type IProcess, ISessionProcessRunner, type ProcessExecOptions } from './processRunner'; + +export class SessionProcessRunner implements ISessionProcessRunner { + declare readonly _serviceBrand: undefined; + + constructor( + @IExecContext private readonly ctx: IExecContext, + @IHostProcessService private readonly hostProcess: IHostProcessService, + ) {} + + async exec(args: readonly string[], options?: ProcessExecOptions): Promise { + const command = args[0]; + if (command === undefined) { + throw new Error( + 'SessionProcessRunner.exec(): at least one argument (the command to run) is required.', + ); + } + const restArgs = args.slice(1); + + const cwd = options?.cwd ?? this.ctx.cwd; + const env = this._buildExecEnv(options?.env); + + return this.hostProcess.spawn(command, restArgs, { cwd, env }); + } + + private _buildExecEnv( + invocationEnv: Record | undefined, + ): Record | undefined { + // No overrides at all — inherit process.env verbatim by passing `undefined` + // to the host process service. Mirrors the pre-refactor behaviour when + // neither the session context nor the caller wanted to touch the child's + // environment. + if (this.ctx.envLayers.length === 0 && invocationEnv === undefined) { + return undefined; + } + const merged: Record = { + ...(process.env as Record), + }; + for (const layer of this.ctx.envLayers) { + Object.assign(merged, layer); + } + if (invocationEnv !== undefined) { + Object.assign(merged, invocationEnv); + } + return merged; + } +} + +registerScopedService( + LifecycleScope.Session, + ISessionProcessRunner, + SessionProcessRunner, + InstantiationType.Delayed, + 'process', +); diff --git a/packages/agent-core-v2/test/os/hostProcessService.test.ts b/packages/agent-core-v2/test/os/hostProcessService.test.ts new file mode 100644 index 000000000..c906ad898 --- /dev/null +++ b/packages/agent-core-v2/test/os/hostProcessService.test.ts @@ -0,0 +1,75 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; + +import { Readable } from 'node:stream'; + +import { DisposableStore } from '#/_base/di/lifecycle'; +import { createServices, type TestInstantiationService } from '#/_base/di/test'; +import { + HostProcessError, + HostProcessErrorCode, + IHostProcessService, +} from '#/os/interface/hostProcess'; +import { HostProcessService } from '#/os/backends/node-local/hostProcessService'; + +async function collect(stream: Readable): Promise { + const chunks: Buffer[] = []; + for await (const chunk of stream) { + chunks.push(chunk as Buffer); + } + return Buffer.concat(chunks).toString('utf8'); +} + +describe('HostProcessService', () => { + let disposables: DisposableStore; + let ix: TestInstantiationService; + + beforeEach(() => { + disposables = new DisposableStore(); + ix = createServices(disposables, { + additionalServices: (reg) => { + reg.define(IHostProcessService, HostProcessService); + }, + }); + }); + + afterEach(() => { + disposables.dispose(); + }); + + it('spawns a process and captures stdout + exit code', async () => { + const svc = ix.get(IHostProcessService); + const proc = await svc.spawn('node', ['-e', 'process.stdout.write("ok")']); + const out = await collect(proc.stdout); + expect(out).toBe('ok'); + expect(await proc.wait()).toBe(0); + expect(proc.exitCode).toBe(0); + }); + + it('passes env overrides to the child', async () => { + const svc = ix.get(IHostProcessService); + const proc = await svc.spawn('node', ['-e', 'process.stdout.write(process.env.FOO ?? "")'], { + env: { FOO: 'bar' }, + }); + const out = await collect(proc.stdout); + expect(out).toBe('bar'); + expect(await proc.wait()).toBe(0); + }); + + it('throws a coded error when the command does not exist', async () => { + const svc = ix.get(IHostProcessService); + await expect(svc.spawn('definitely-not-a-real-command-42')).rejects.toSatisfy((err: unknown) => { + expect(err).toBeInstanceOf(HostProcessError); + expect((err as HostProcessError).code).toBe(HostProcessErrorCode.SpawnFailed); + return true; + }); + }); + + it('terminates a running process with kill()', async () => { + const svc = ix.get(IHostProcessService); + const proc = await svc.spawn('node', ['-e', 'setTimeout(() => {}, 30000)']); + expect(proc.pid).toBeGreaterThan(0); + await proc.kill('SIGTERM'); + const code = await proc.wait(); + expect(code).not.toBe(0); + }); +}); diff --git a/packages/agent-core-v2/test/process/processRunnerService.test.ts b/packages/agent-core-v2/test/process/processRunnerService.test.ts index 0bfa770df..59aec90b1 100644 --- a/packages/agent-core-v2/test/process/processRunnerService.test.ts +++ b/packages/agent-core-v2/test/process/processRunnerService.test.ts @@ -12,6 +12,8 @@ import { registerScopedService, } from '#/_base/di/scope'; import { createScopedTestHost, stubPair } from '#/_base/di/test'; +import { IHostProcessService } from '#/os/interface/hostProcess'; +import { HostProcessService } from '#/os/backends/node-local/hostProcessService'; import { createExecContext, IExecContext } from '#/session/execContext'; import { ISessionProcessRunner, SessionProcessRunner } from '#/session/process'; @@ -28,6 +30,13 @@ describe('SessionProcessRunner (backed by IExecContext)', () => { beforeEach(async () => { _clearScopedRegistryForTests(); + registerScopedService( + LifecycleScope.App, + IHostProcessService, + HostProcessService, + InstantiationType.Delayed, + 'hostProcess', + ); registerScopedService( LifecycleScope.Session, ISessionProcessRunner,