diff --git a/packages/coding-agent/src/config/settings-manager.ts b/packages/coding-agent/src/config/settings-manager.ts index 7c44bc784..2eea4ee42 100644 --- a/packages/coding-agent/src/config/settings-manager.ts +++ b/packages/coding-agent/src/config/settings-manager.ts @@ -1,6 +1,6 @@ import * as fs from "node:fs/promises"; import * as path from "node:path"; -import { isEnoent, logger } from "@oh-my-pi/pi-utils"; +import { isEnoent, logger, procmgr } from "@oh-my-pi/pi-utils"; import { YAML } from "bun"; import { type Settings as SettingsItem, settingsCapability } from "../capability/settings"; import { getAgentDbPath, getAgentDir } from "../config"; @@ -461,6 +461,8 @@ export class SettingsManager { private settings!: Settings; private persist: boolean; + static #lastInstance: SettingsManager | null = null; + /** * Private constructor - use static factory methods instead. * @param storage - SQLite storage instance for auth/cache, or null for in-memory mode @@ -477,6 +479,7 @@ export class SettingsManager { initialSettings: Settings, persist: boolean, projectSettings: Settings, + private agentDir: string | null, ) { this.storage = storage; this.configPath = configPath; @@ -517,6 +520,9 @@ export class SettingsManager { * @returns Configured SettingsManager with merged global and user settings */ static async create(cwd: string = process.cwd(), agentDir: string = getAgentDir()): Promise { + cwd = path.normalize(cwd); + agentDir = path.normalize(agentDir); + const configPath = path.join(agentDir, "config.yml"); const storage = await AgentStorage.open(getAgentDbPath(agentDir)); @@ -541,7 +547,9 @@ export class SettingsManager { // Load project settings before construction (constructor is sync) const projectSettings = await SettingsManager.loadProjectSettingsStatic(cwd); - return new SettingsManager(storage, configPath, cwd, globalSettings, true, projectSettings); + const instance = new SettingsManager(storage, configPath, cwd, globalSettings, true, projectSettings, agentDir); + SettingsManager.#lastInstance = instance; + return instance; } /** @@ -550,7 +558,7 @@ export class SettingsManager { * @returns SettingsManager that won't persist changes to disk */ static inMemory(settings: Partial = {}): SettingsManager { - return new SettingsManager(null, null, null, settings, false, {}); + return new SettingsManager(null, null, null, settings, false, {}, null); } /** @@ -1784,4 +1792,54 @@ export class SettingsManager { await this.save(); } } + + _compareUniqueCtorKeys(cwd: string, agentDir: string): boolean { + if (this.cwd !== cwd) { + cwd = path.normalize(cwd); + if (this.cwd !== cwd) { + return false; + } + } + if (this.agentDir !== agentDir) { + agentDir = path.normalize(agentDir); + if (this.agentDir !== agentDir) { + return false; + } + } + return true; + } + + /** + * Acquire the last created SettingsManager instance. + * If no instance exists, create a new one. + * @returns The SettingsManager instance + */ + static acquire( + cwd: string = process.cwd(), + agentDir: string = getAgentDir(), + ): SettingsManager | Promise { + const prev = SettingsManager.#lastInstance; + if (prev?._compareUniqueCtorKeys(cwd, agentDir)) { + return prev; + } + return SettingsManager.create(cwd, agentDir); + } + + /** + * Gets the shell configuration + * @returns The shell configuration + */ + async getShellConfig() { + const shell = this.getShellPath(); + return procmgr.getShellConfig(shell); + } + + /** + * Gets the shell configuration from the last created SettingsManager instance. + * @returns The shell configuration + */ + static async getGlobalShellConfig() { + const settings = await SettingsManager.acquire(); + return settings.getShellConfig(); + } } diff --git a/packages/coding-agent/src/exec/bash-executor.ts b/packages/coding-agent/src/exec/bash-executor.ts index dc79b3aa3..9dcfee108 100644 --- a/packages/coding-agent/src/exec/bash-executor.ts +++ b/packages/coding-agent/src/exec/bash-executor.ts @@ -4,8 +4,8 @@ * Provides unified bash execution for AgentSession.executeBash() and direct calls. */ import { Exception, ptree } from "@oh-my-pi/pi-utils"; +import { SettingsManager } from "../config/settings-manager"; import { OutputSink } from "../session/streaming-output"; -import { getShellConfig } from "../utils/shell"; import { getOrCreateSnapshot, getSnapshotSourceCommand } from "../utils/shell-snapshot"; export interface BashExecutorOptions { @@ -33,7 +33,7 @@ export interface BashResult { } export async function executeBash(command: string, options?: BashExecutorOptions): Promise { - const { shell, args, env, prefix } = await getShellConfig(); + const { shell, args, env, prefix } = await SettingsManager.getGlobalShellConfig(); // Merge additional env vars if provided const finalEnv = options?.env ? { ...env, ...options.env } : env; diff --git a/packages/coding-agent/src/extensibility/plugins/installer.ts b/packages/coding-agent/src/extensibility/plugins/installer.ts index 12ea10fc4..0cfb29fb9 100644 --- a/packages/coding-agent/src/extensibility/plugins/installer.ts +++ b/packages/coding-agent/src/extensibility/plugins/installer.ts @@ -45,7 +45,7 @@ export async function installPlugin(packageName: string): Promise { await ensurePluginsDir(); - const proc = Bun.spawn(["npm", "uninstall", name], { + const proc = Bun.spawn(["bun", "uninstall", name], { cwd: PLUGINS_DIR, stdin: "ignore", stdout: "pipe", diff --git a/packages/coding-agent/src/extensibility/plugins/manager.ts b/packages/coding-agent/src/extensibility/plugins/manager.ts index 2e10c4517..21b3e12a9 100644 --- a/packages/coding-agent/src/extensibility/plugins/manager.ts +++ b/packages/coding-agent/src/extensibility/plugins/manager.ts @@ -154,7 +154,7 @@ export class PluginManager { } // Run npm install - const proc = Bun.spawn(["npm", "install", spec.packageName], { + const proc = Bun.spawn(["bun", "install", spec.packageName], { cwd: getPluginsDir(), stdin: "ignore", stdout: "pipe", @@ -234,7 +234,7 @@ export class PluginManager { validatePackageName(name); await this.ensurePackageJson(); - const proc = Bun.spawn(["npm", "uninstall", name], { + const proc = Bun.spawn(["bun", "uninstall", name], { cwd: getPluginsDir(), stdin: "ignore", stdout: "pipe", @@ -616,7 +616,7 @@ export class PluginManager { private async fixMissingPlugin(): Promise { try { - const proc = Bun.spawn(["npm", "install"], { + const proc = Bun.spawn(["bun", "install"], { cwd: getPluginsDir(), stdin: "ignore", stdout: "pipe", diff --git a/packages/coding-agent/src/index.ts b/packages/coding-agent/src/index.ts index f49d9c0e0..00f897ed9 100644 --- a/packages/coding-agent/src/index.ts +++ b/packages/coding-agent/src/index.ts @@ -276,4 +276,3 @@ export { truncateTail, type WriteToolDetails, } from "./tools"; -export { getShellConfig } from "./utils/shell"; diff --git a/packages/coding-agent/src/ipy/gateway-coordinator.ts b/packages/coding-agent/src/ipy/gateway-coordinator.ts index 21cf837a2..12dfee17e 100644 --- a/packages/coding-agent/src/ipy/gateway-coordinator.ts +++ b/packages/coding-agent/src/ipy/gateway-coordinator.ts @@ -1,10 +1,10 @@ import * as fs from "node:fs"; import { createServer } from "node:net"; import * as path from "node:path"; -import { isEnoent, logger } from "@oh-my-pi/pi-utils"; +import { isEnoent, logger, procmgr } from "@oh-my-pi/pi-utils"; import type { Subprocess } from "bun"; import { getAgentDir } from "../config"; -import { getShellConfig, killProcessTree } from "../utils/shell"; +import { SettingsManager } from "../config/settings-manager"; import { getOrCreateSnapshot } from "../utils/shell-snapshot"; import { time } from "../utils/timings"; @@ -304,7 +304,7 @@ async function withGatewayLock(handler: () => Promise): Promise { const lockPid = lockInfo?.pid; const lockAgeMs = lockInfo?.startedAt ? Date.now() - lockInfo.startedAt : Date.now() - lockStat.mtimeMs; const staleByTime = lockAgeMs > GATEWAY_LOCK_STALE_MS; - const staleByPid = lockPid !== undefined && !isPidRunning(lockPid); + const staleByPid = lockPid !== undefined && !procmgr.isPidRunning(lockPid); const staleByMissingPid = lockPid === undefined && staleByTime; if (staleByPid || staleByMissingPid) { await fs.promises.unlink(lockPath); @@ -365,15 +365,6 @@ async function clearGatewayInfo(): Promise { } } -function isPidRunning(pid: number): boolean { - try { - process.kill(pid, 0); - return true; - } catch { - return false; - } -} - async function isGatewayHealthy(url: string): Promise { try { const response = await fetch(`${url}/api/kernelspecs`, { @@ -386,14 +377,14 @@ async function isGatewayHealthy(url: string): Promise { } async function isGatewayAlive(info: GatewayInfo): Promise { - if (!isPidRunning(info.pid)) return false; + if (!procmgr.isPidRunning(info.pid)) return false; return await isGatewayHealthy(info.url); } async function startGatewayProcess( cwd: string, ): Promise<{ url: string; pid: number; pythonPath: string; venvPath: string | null }> { - const { shell, env } = await getShellConfig(); + const { shell, env } = await SettingsManager.getGlobalShellConfig(); const filteredEnv = filterEnv(env); const runtime = await resolvePythonRuntime(cwd, filteredEnv); const snapshotPath = await getOrCreateSnapshot(shell, env).catch((err: unknown) => { @@ -428,17 +419,16 @@ async function startGatewayProcess( stdin: "ignore", stdout: "pipe", stderr: "pipe", + detached: true, env: kernelEnv, }, ); let exited = false; gatewayProcess.exited + .catch(() => {}) .then(() => { exited = true; - }) - .catch(() => { - exited = true; }); const startTime = Date.now(); @@ -459,13 +449,13 @@ async function startGatewayProcess( await Bun.sleep(100); } - await killProcessTree(gatewayProcess.pid); + await procmgr.terminate({ target: gatewayProcess, group: true }); throw new Error("Gateway startup timeout"); } async function killGateway(pid: number, context: string): Promise { try { - await killProcessTree(pid); + await procmgr.terminate({ target: pid, group: true }); } catch (err) { logger.warn("Failed to kill shared gateway process", { error: err instanceof Error ? err.message : String(err), @@ -495,7 +485,7 @@ export async function acquireSharedGateway(cwd: string): Promise { venvPath: null, }; } - const active = isPidRunning(info.pid); + const active = procmgr.isPidRunning(info.pid); return { active, url: info.url, @@ -573,7 +563,7 @@ export async function shutdownSharedGateway(): Promise { await withGatewayLock(async () => { const info = await readGatewayInfo(); if (!info) return; - if (isPidRunning(info.pid)) { + if (procmgr.isPidRunning(info.pid)) { await killGateway(info.pid, "shutdown"); } await clearGatewayInfo(); diff --git a/packages/coding-agent/src/ipy/kernel.ts b/packages/coding-agent/src/ipy/kernel.ts index 90adb7591..5bcb5f6ee 100644 --- a/packages/coding-agent/src/ipy/kernel.ts +++ b/packages/coding-agent/src/ipy/kernel.ts @@ -3,7 +3,7 @@ import * as path from "node:path"; import { logger, ptree } from "@oh-my-pi/pi-utils"; import { $ } from "bun"; import { nanoid } from "nanoid"; -import { getShellConfig } from "../utils/shell"; +import { SettingsManager } from "../config/settings-manager"; import { getOrCreateSnapshot } from "../utils/shell-snapshot"; import { time } from "../utils/timings"; import { htmlToBasicMarkdown } from "../web/scrapers/types"; @@ -279,7 +279,7 @@ export async function checkPythonKernelAvailability(cwd: string): Promise { - const { shell, env } = await getShellConfig(); + const { shell, env } = await SettingsManager.getGlobalShellConfig(); const filteredEnv = filterEnv(env); const runtime = await resolvePythonRuntime(options.cwd, filteredEnv); const snapshotPath = await getOrCreateSnapshot(shell, env).catch((err: unknown) => { diff --git a/packages/coding-agent/src/lsp/client.ts b/packages/coding-agent/src/lsp/client.ts index e36352010..62d3ea849 100644 --- a/packages/coding-agent/src/lsp/client.ts +++ b/packages/coding-agent/src/lsp/client.ts @@ -1,4 +1,4 @@ -import { isEnoent, logger } from "@oh-my-pi/pi-utils"; +import { isEnoent, logger, ptree } from "@oh-my-pi/pi-utils"; import { ToolAbortError, throwIfAborted } from "../tools/tool-errors"; import { applyWorkspaceEdit } from "./edits"; import { getLspmuxCommand, isLspmuxSupported } from "./lspmux"; @@ -206,7 +206,7 @@ function concatBuffers(a: Uint8Array, b: Uint8Array): Uint8Array { } async function writeMessage( - sink: import("bun").FileSink, + sink: Bun.FileSink, message: LspJsonRpcRequest | LspJsonRpcNotification | LspJsonRpcResponse, ): Promise { const content = JSON.stringify(message); @@ -230,7 +230,7 @@ async function startMessageReader(client: LspClient): Promise { if (client.isReading) return; client.isReading = true; - const reader = (client.process.stdout as ReadableStream).getReader(); + const reader = (client.proc.stdout as ReadableStream).getReader(); try { while (true) { @@ -364,7 +364,7 @@ async function sendResponse( }; try { - await writeMessage(client.process.stdin as import("bun").FileSink, response); + await writeMessage(client.proc.stdin, response); } catch (err) { logger.error("LSP failed to respond.", { method, error: String(err) }); } @@ -409,18 +409,17 @@ export async function getOrCreateClient(config: ServerConfig, cwd: string, initT ? await getLspmuxCommand(baseCommand, baseArgs) : { command: baseCommand, args: baseArgs }; - const proc = Bun.spawn([command, ...args], { + const proc = ptree.spawn([command, ...args], { cwd, + detached: true, stdin: "pipe", - stdout: "pipe", - stderr: "pipe", env: env ? { ...process.env, ...env } : undefined, }); const client: LspClient = { name: key, cwd, - process: proc, + proc, config, requestId: 0, diagnostics: new Map(), @@ -686,7 +685,7 @@ export function shutdownClient(key: string): void { sendRequest(client, "shutdown", null).catch(() => {}); // Kill process - client.process.kill(); + client.proc.kill(); clients.delete(key); } @@ -773,7 +772,7 @@ export async function sendRequest( }); // Write request - writeMessage(client.process.stdin as import("bun").FileSink, request).catch(err => { + writeMessage(client.proc.stdin, request).catch(err => { if (timeout) clearTimeout(timeout); client.pendingRequests.delete(id); cleanup(); @@ -793,26 +792,33 @@ export async function sendNotification(client: LspClient, method: string, params }; client.lastActivity = Date.now(); - await writeMessage(client.process.stdin as import("bun").FileSink, notification); + await writeMessage(client.proc.stdin, notification); } /** * Shutdown all LSP clients. */ export function shutdownAll(): void { - for (const client of Array.from(clients.values())) { - // Reject all pending requests - for (const pending of Array.from(client.pendingRequests.values())) { - pending.reject(new Error("LSP client shutdown")); - } - client.pendingRequests.clear(); - - // Send shutdown request (best effort, don't wait) - sendRequest(client, "shutdown", null).catch(() => {}); - - client.process.kill(); - } + const clientsToShutdown = Array.from(clients.values()); clients.clear(); + + const err = new Error("LSP client shutdown"); + for (const client of clientsToShutdown) { + /// Reject all pending requests + const reqs = Array.from(client.pendingRequests.values()); + client.pendingRequests.clear(); + for (const pending of reqs) { + pending.reject(err); + } + + void (async () => { + // Send shutdown request (best effort, don't wait) + const timeout = Bun.sleep(5_000); + const result = sendRequest(client, "shutdown", null).catch(() => {}); + await Promise.race([result, timeout]); + client.proc.kill(); + })().catch(() => {}); + } } /** Status of an LSP server */ diff --git a/packages/coding-agent/src/lsp/types.ts b/packages/coding-agent/src/lsp/types.ts index 23514db4c..cd8237f25 100644 --- a/packages/coding-agent/src/lsp/types.ts +++ b/packages/coding-agent/src/lsp/types.ts @@ -1,6 +1,6 @@ import { StringEnum } from "@oh-my-pi/pi-ai"; +import type { ptree } from "@oh-my-pi/pi-utils"; import { type Static, Type } from "@sinclair/typebox"; -import type { Subprocess } from "bun"; // ============================================================================= // Tool Schema @@ -400,7 +400,7 @@ export interface LspClient { name: string; cwd: string; config: ServerConfig; - process: Subprocess; + proc: ptree.ChildProcess<"pipe">; requestId: number; diagnostics: Map; diagnosticsVersion: number; diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index 00e189a93..eca9a1621 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-client.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-client.ts @@ -524,7 +524,7 @@ export class RpcClient { // Write to stdin after registering the handler const stdin = this.process!.stdin as import("bun").FileSink; - stdin.write(new TextEncoder().encode(`${JSON.stringify(fullCommand)}\n`)); + stdin.write(`${JSON.stringify(fullCommand)}\n`); // flush() returns number | Promise - handle both cases const flushResult = stdin.flush(); if (flushResult instanceof Promise) { diff --git a/packages/coding-agent/test/tools.test.ts b/packages/coding-agent/test/tools.test.ts index 043f935c5..2bb49db8e 100644 --- a/packages/coding-agent/test/tools.test.ts +++ b/packages/coding-agent/test/tools.test.ts @@ -2,6 +2,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; +import { SettingsManager } from "@oh-my-pi/pi-coding-agent/config/settings-manager"; import { EditTool } from "@oh-my-pi/pi-coding-agent/patch"; import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools"; import { BashTool } from "@oh-my-pi/pi-coding-agent/tools/bash"; @@ -11,7 +12,6 @@ import { LsTool } from "@oh-my-pi/pi-coding-agent/tools/ls"; import { wrapToolWithMetaNotice } from "@oh-my-pi/pi-coding-agent/tools/output-meta"; import { ReadTool } from "@oh-my-pi/pi-coding-agent/tools/read"; import { WriteTool } from "@oh-my-pi/pi-coding-agent/tools/write"; -import * as shellModule from "@oh-my-pi/pi-coding-agent/utils/shell"; import { nanoid } from "nanoid"; // Helper to extract text from content blocks @@ -425,7 +425,7 @@ function b() { }); it("should handle process spawn errors", async () => { - const getShellConfigSpy = vi.spyOn(shellModule, "getShellConfig").mockResolvedValueOnce({ + const getShellConfigSpy = vi.spyOn(SettingsManager, "getGlobalShellConfig").mockResolvedValueOnce({ shell: "/nonexistent-shell-path-xyz123", args: ["-c"], env: {}, diff --git a/packages/utils/src/index.ts b/packages/utils/src/index.ts index 89b1713ed..39bd3e1f5 100644 --- a/packages/utils/src/index.ts +++ b/packages/utils/src/index.ts @@ -4,6 +4,7 @@ export * from "./fs-error"; export * from "./glob"; export * as logger from "./logger"; export * as postmortem from "./postmortem"; +export * as procmgr from "./procmgr"; export * as ptree from "./ptree"; export { AbortError, ChildProcess, Exception, NonZeroExitError } from "./ptree"; export * from "./stream"; diff --git a/packages/coding-agent/src/utils/shell.ts b/packages/utils/src/procmgr.ts similarity index 58% rename from packages/coding-agent/src/utils/shell.ts rename to packages/utils/src/procmgr.ts index cff262bf2..f959f85ab 100644 --- a/packages/coding-agent/src/utils/shell.ts +++ b/packages/utils/src/procmgr.ts @@ -1,6 +1,6 @@ import * as fs from "node:fs"; -import { $ } from "bun"; -import { SettingsManager } from "../config/settings-manager"; +import * as timers from "node:timers"; +import type { Subprocess } from "bun"; export interface ShellConfig { shell: string; @@ -11,6 +11,9 @@ export interface ShellConfig { let cachedShellConfig: ShellConfig | null = null; +const IS_WINDOWS = process.platform === "win32"; +const TERM_SIGNAL = IS_WINDOWS ? undefined : "SIGTERM"; + /** * Check if a shell binary is executable. */ @@ -87,14 +90,11 @@ function buildConfig(shell: string): ShellConfig { * 3. On Unix: $SHELL if bash/zsh, then fallback paths * 4. Fallback: sh */ -export async function getShellConfig(): Promise { +export async function getShellConfig(customShellPath?: string): Promise { if (cachedShellConfig) { return cachedShellConfig; } - const settings = await SettingsManager.create(); - const customShellPath = settings.getShellPath(); - // 1. Check user-specified shell path if (customShellPath) { if (await Bun.file(customShellPath).exists()) { @@ -176,127 +176,142 @@ export async function getShellConfig(): Promise { return cachedShellConfig; } -let pgrepAvailable: string | null | undefined; - /** - * Check if pgrep is available on this system (cached). + * Options for terminating a process and all its descendants. */ -function hasPgrep(): string | null { - if (pgrepAvailable === undefined) { - try { - pgrepAvailable = Bun.which("pgrep") ?? null; - } catch { - pgrepAvailable = null; - } - } - return pgrepAvailable; +export interface TerminateOptions { + /** The process to terminate */ + target: Subprocess | number; + /** Whether to terminate the process group (Windows only) */ + group?: boolean; + /** Timeout in milliseconds */ + timeout?: number; + /** Abort signal */ + signal?: AbortSignal; } /** - * Get direct children of a PID using pgrep. + * Check if a process is running. */ -async function getChildrenViaPgrep(pid: number): Promise { - const result = await $`pgrep -P ${pid}`.quiet().nothrow(); - if (result.exitCode !== 0) return []; - const output = result.stdout.toString().trim(); - if (!output) return []; - - const children: number[] = []; - for (const line of output.split("\n")) { - const childPid = parseInt(line, 10); - if (!Number.isNaN(childPid)) children.push(childPid); - } - return children; -} - -/** - * Get direct children of a PID using /proc (Linux only). - */ -async function getChildrenViaProc(pid: number): Promise { +export function isPidRunning(pid: number | Subprocess): boolean { try { - const script = `for p in /proc/[0-9]*/stat; do cat "$p" 2>/dev/null; done | awk -v ppid=${pid} '$4 == ppid { print $1 }'`; - const result = await $`sh -c ${script}`.quiet().nothrow(); - if (result.exitCode !== 0) return []; - const output = result.stdout.toString().trim(); - if (!output) return []; - - const children: number[] = []; - for (const line of output.split("\n")) { - const childPid = parseInt(line, 10); - if (!Number.isNaN(childPid)) children.push(childPid); + if (typeof pid === "number") { + process.kill(pid, 0); + } else { + if (pid.killed) return false; + if (pid.exitCode !== null) return false; } - return children; - } catch { - return []; - } -} - -/** - * Collect all descendant PIDs breadth-first. - * Returns deepest descendants first (reverse BFS order) for proper kill ordering. - */ -async function getDescendantPids(pid: number): Promise { - const getChildren = hasPgrep() ? getChildrenViaPgrep : getChildrenViaProc; - const descendants: number[] = []; - const queue = [pid]; - - while (queue.length > 0) { - const current = queue.shift()!; - const children = await getChildren(current); - for (const child of children) { - descendants.push(child); - queue.push(child); - } - } - - // Reverse so deepest children are killed first - return descendants.reverse(); -} - -function tryKill(pid: number, signal: NodeJS.Signals): boolean { - try { - process.kill(pid, signal); return true; } catch { return false; } } -/** - * Kill a process and all its descendants. - * @param gracePeriodMs - Time to wait after SIGTERM before SIGKILL (0 = immediate SIGKILL) - */ -export async function killProcessTree(pid: number, gracePeriodMs = 0): Promise { - if (process.platform === "win32") { - await $`taskkill /F /T /PID ${pid}`.quiet().nothrow(); - return; - } - - const signal = gracePeriodMs > 0 ? "SIGTERM" : "SIGKILL"; - - // Fast path: process group kill (works if pid is group leader) - try { - process.kill(-pid, signal); - if (gracePeriodMs > 0) { - await Bun.sleep(gracePeriodMs); - try { - process.kill(-pid, "SIGKILL"); - } catch { - // Already dead - } - } - return; - } catch { - // Not a process group leader, fall through - } - - // Collect descendants BEFORE killing to minimize race window - const allPids = [...(await getDescendantPids(pid)), pid]; - - if (gracePeriodMs > 0) { - for (const p of allPids) tryKill(p, "SIGTERM"); - await Bun.sleep(gracePeriodMs); - } - - for (const p of allPids) tryKill(p, "SIGKILL"); +function joinSignals(...sigs: (AbortSignal | null | undefined)[]): AbortSignal | undefined { + const nn = sigs.filter(Boolean) as AbortSignal[]; + if (nn.length === 0) return undefined; + if (nn.length === 1) return nn[0]; + return AbortSignal.any(nn); +} + +export function onProcessExit(proc: Subprocess | number, abortSignal?: AbortSignal): Promise { + if (typeof proc !== "number") { + return proc.exited.then( + () => true, + () => true, + ); + } + + if (!isPidRunning(proc)) { + return Promise.resolve(true); + } + + const { promise, resolve, reject } = Promise.withResolvers(); + const localAbortController = new AbortController(); + + const timer = timers.promises.setInterval(300, null, { + signal: joinSignals(abortSignal, localAbortController.signal), + }); + void (async () => { + try { + for await (const _ of timer) { + if (!isPidRunning(proc)) { + resolve(true); + break; + } + } + } catch (error) { + return reject(error); + } finally { + localAbortController.abort(); + } + resolve(false); + })(); + + return promise; +} + +/** + * Terminate a process and all its descendants. + */ +export async function terminate(options: TerminateOptions): Promise { + const { target, group = false, timeout = 5000, signal } = options; + + const abortController = new AbortController(); + try { + const abortSignal = joinSignals(signal, abortController.signal); + + // Determine PID + let pid: number | undefined; + const exitPromise = onProcessExit(target, abortSignal); + if (typeof target === "number") { + pid = target; + } else { + pid = target.pid; + if (target.killed) return true; + } + + // Give it a moment to exit gracefully first. + try { + if (typeof target === "number") { + process.kill(target, TERM_SIGNAL); + } else { + target.kill(TERM_SIGNAL); + } + + if (exitPromise) { + const exited = await Promise.race([Bun.sleep(1000).then(() => false), exitPromise]); + if (exited) return true; + } + } catch {} + + if (group) { + try { + if (IS_WINDOWS) { + const taskkill = Bun.spawn({ + cmd: ["taskkill", "/F", "/T", "/PID", pid.toString()], + stdin: "ignore", + stdout: "ignore", + stderr: "ignore", + timeout: 5000, + }); + void taskkill.exited.catch(() => {}); + taskkill.unref(); + } else { + process.kill(-pid, "SIGKILL"); + } + } catch {} + } + try { + if (typeof target === "number") { + process.kill(target, "SIGKILL"); + } else { + target.kill("SIGKILL"); + } + } catch {} + + return await Promise.race([Bun.sleep(timeout).then(() => false), exitPromise]); + } finally { + abortController.abort(); + } } diff --git a/packages/utils/src/ptree.ts b/packages/utils/src/ptree.ts index 7ae7294cb..72400c5b3 100644 --- a/packages/utils/src/ptree.ts +++ b/packages/utils/src/ptree.ts @@ -8,14 +8,15 @@ * - Cross-platform tree kill for process groups (Windows taskkill, Unix -pid). * - Convenience helpers: captureText / execText, AbortSignal, timeouts. */ -import { $, type FileSink, type Spawn, type Subprocess } from "bun"; -import { postmortem } from "."; -const isWindows = process.platform === "win32"; +import type { Spawn, Subprocess } from "bun"; +import { postmortem } from "."; +import { terminate } from "./procmgr"; + const managedChildren = new Set(); /** A Bun subprocess with stdout/stderr always piped (stdin may vary). */ -type PipedSubprocess = Subprocess<"pipe" | "ignore" | null, "pipe", "pipe">; +type PipedSubprocess = Subprocess; /** Minimal push-based ReadableStream that buffers unboundedly (like the old queue). */ function pushStream() { @@ -97,35 +98,7 @@ async function pump( * - Unix: negative PID signals the process group */ async function killChild(child: ChildProcess) { - const pid = child.pid; - if (!pid || child.killed) return; - - const exited = child.proc.exited.then( - () => true, - () => true, - ); - const waitForExit = (timeout = 1000) => Promise.race([Bun.sleep(timeout).then(() => false), exited]); - - // Give it a moment to exit gracefully first. - try { - child.proc.kill(); - } catch {} - if (await waitForExit(1000)) return true; - - if (child.isProcessGroup) { - try { - if (isWindows) { - await $`taskkill /F /T /PID ${pid}`.quiet().nothrow(); - } else { - process.kill(-pid); - } - } catch {} - } - try { - child.proc.kill("SIGKILL"); - } catch {} - - return await waitForExit(1000); + await terminate({ target: child.proc, group: child.isProcessGroup }); } postmortem.register("managed-children", async () => { @@ -211,11 +184,13 @@ export class TimeoutError extends AbortError { } } +type InMask = "pipe" | "ignore" | Buffer | Uint8Array | null; + /** * ChildProcess wraps a managed subprocess, capturing stderr tail, providing * cross-platform kill/detach logic plus AbortSignal integration. */ -export class ChildProcess { +export class ChildProcess { #nothrow = false; #stderrBuffer = ""; @@ -231,7 +206,7 @@ export class ChildProcess { #exited: Promise; constructor( - public readonly proc: PipedSubprocess, + public readonly proc: PipedSubprocess, public readonly isProcessGroup: boolean, ) { const { promise: stderrDone, resolve: resolveStderrDone } = Promise.withResolvers(); @@ -329,7 +304,7 @@ export class ChildProcess { get killed(): boolean { return this.proc.killed; } - get stdin(): FileSink | undefined { + get stdin(): Bun.SpawnOptions.WritableToIO { return this.proc.stdin; } get stdout(): ReadableStream { @@ -443,8 +418,8 @@ export class ChildProcess { /** * Options for cspawn (child spawn). Always pipes stdout/stderr, allows signal. */ -type ChildSpawnOptions = Omit< - Spawn.SpawnOptions<"pipe" | "ignore" | Buffer | Uint8Array | null, "pipe", "pipe">, +type ChildSpawnOptions = Omit< + Spawn.SpawnOptions, "stdout" | "stderr" > & { signal?: AbortSignal; detached?: boolean }; @@ -454,7 +429,7 @@ type ChildSpawnOptions = Omit< * @param options - The options for the spawn. * @returns A ChildProcess instance. */ -export function spawn(cmd: string[], options?: ChildSpawnOptions): ChildProcess { +export function spawn(cmd: string[], options?: ChildSpawnOptions): ChildProcess { const { detached = false, timeout, signal, ...rest } = options ?? {}; const child = Bun.spawn(cmd, { stdin: "ignore",