refactor(pi-utils): migrated process spawning from cspawn to spawnGroup/spawnAttached with async resource management
- Migrated process spawning from `cspawn` to `spawnGroup` and `spawnAttached` APIs with explicit resource management using TypeScript's `using` statement. - Refactored `ChildProcess` class to support process group management with new `isProcessGroup` getter and `[Symbol.dispose]()` method for automatic cleanup. - Removed `cspawn` from public API exports in pi-utils, replacing it with `spawnGroup` and `spawnAttached` functions. - Simplified `ChildProcess.kill()` method signature by removing signal parameter and eliminated `killAndWait()` method in favor of async disposable pattern. - Converted `killChild` function to async implementation using Bun's `$` template and sleep-based polling for process termination.
This commit is contained in:
@@ -3,7 +3,7 @@
|
||||
*
|
||||
* Provides unified bash execution for AgentSession.executeBash() and direct calls.
|
||||
*/
|
||||
import { cspawn, Exception, ptree } from "@oh-my-pi/pi-utils";
|
||||
import { Exception, ptree } from "@oh-my-pi/pi-utils";
|
||||
import { OutputSink } from "../session/streaming-output";
|
||||
import { getShellConfig } from "../utils/shell";
|
||||
import { getOrCreateSnapshot, getSnapshotSourceCommand } from "../utils/shell-snapshot";
|
||||
@@ -50,7 +50,7 @@ export async function executeBash(command: string, options?: BashExecutorOptions
|
||||
artifactId: options?.artifactId,
|
||||
});
|
||||
|
||||
const child = cspawn([shell, ...args, finalCommand], {
|
||||
using child = ptree.spawnGroup([shell, ...args, finalCommand], {
|
||||
cwd: options?.cwd,
|
||||
env: finalEnv,
|
||||
signal: options?.signal,
|
||||
|
||||
@@ -35,7 +35,7 @@ export async function execCommand(
|
||||
cwd: string,
|
||||
options?: ExecOptions,
|
||||
): Promise<ExecResult> {
|
||||
const proc = ptree.cspawn([command, ...args], {
|
||||
using proc = ptree.spawnAttached([command, ...args], {
|
||||
cwd,
|
||||
signal: options?.signal,
|
||||
timeout: options?.timeout,
|
||||
|
||||
@@ -112,7 +112,7 @@ export class RpcClient {
|
||||
args.push(...this.options.args);
|
||||
}
|
||||
|
||||
this.process = ptree.cspawn(["bun", cliPath, ...args], {
|
||||
this.process = ptree.spawnAttached(["bun", cliPath, ...args], {
|
||||
cwd: this.options.cwd,
|
||||
env: { ...process.env, ...this.options.env },
|
||||
stdin: "pipe",
|
||||
@@ -154,11 +154,11 @@ export class RpcClient {
|
||||
/**
|
||||
* Stop the RPC agent process.
|
||||
*/
|
||||
async stop(): Promise<void> {
|
||||
stop() {
|
||||
if (!this.process) return;
|
||||
|
||||
this.lineReader?.cancel();
|
||||
await this.process.killAndWait();
|
||||
this.process.kill();
|
||||
|
||||
this.process = null;
|
||||
this.lineReader = null;
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { cspawn, logger, ptree } from "@oh-my-pi/pi-utils";
|
||||
import { logger, ptree } from "@oh-my-pi/pi-utils";
|
||||
import { OutputSink } from "../session/streaming-output";
|
||||
import { buildRemoteCommand, ensureConnection, ensureHostInfo, type SSHConnectionTarget } from "./connection-manager";
|
||||
import { hasSshfs, mountRemote } from "./sshfs-mount";
|
||||
@@ -76,7 +76,7 @@ export async function executeSSH(
|
||||
}
|
||||
}
|
||||
|
||||
const child = cspawn(["ssh", ...(await buildRemoteCommand(host, resolvedCommand))], {
|
||||
using child = ptree.spawnAttached(["ssh", ...(await buildRemoteCommand(host, resolvedCommand))], {
|
||||
signal: options?.signal,
|
||||
timeout: options?.timeout,
|
||||
});
|
||||
|
||||
@@ -111,7 +111,7 @@ async function exec(
|
||||
args: string[],
|
||||
options?: { timeout?: number; input?: string | Buffer },
|
||||
): Promise<{ stdout: string; stderr: string; ok: boolean }> {
|
||||
const proc = ptree.cspawn([cmd, ...args], {
|
||||
using proc = ptree.spawnGroup([cmd, ...args], {
|
||||
stdin: options?.input ? "pipe" : null,
|
||||
timeout: options?.timeout ? options.timeout * 1000 : undefined,
|
||||
});
|
||||
|
||||
@@ -78,7 +78,7 @@ export async function runRg(
|
||||
args: string[],
|
||||
options?: { signal?: AbortSignal; timeoutMs?: number },
|
||||
): Promise<RgResult> {
|
||||
const child = ptree.cspawn([rgPath, ...args], { signal: options?.signal, timeout: options?.timeoutMs });
|
||||
using child = ptree.spawnAttached([rgPath, ...args], { signal: options?.signal, timeout: options?.timeoutMs });
|
||||
const timeoutSeconds = options?.timeoutMs ? Math.max(1, Math.round(options.timeoutMs / 1000)) : undefined;
|
||||
const timeoutMessage = timeoutSeconds ? `rg timed out after ${timeoutSeconds}s` : "rg timed out";
|
||||
|
||||
@@ -370,7 +370,7 @@ export class GrepTool implements AgentTool<typeof grepSchema, GrepToolDetails> {
|
||||
|
||||
args.push("--", normalizedPattern, searchPath);
|
||||
|
||||
const child = ptree.cspawn([rgPath, ...args], { signal });
|
||||
using child = ptree.spawnAttached([rgPath, ...args], { signal });
|
||||
|
||||
let matchCount = 0;
|
||||
let matchLimitReached = false;
|
||||
@@ -583,7 +583,6 @@ export class GrepTool implements AgentTool<typeof grepSchema, GrepToolDetails> {
|
||||
if (maxMatches !== undefined && nextIndex > maxMatches) {
|
||||
matchLimitReached = true;
|
||||
killedDueToLimit = true;
|
||||
child.kill("SIGKILL");
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
@@ -296,7 +296,7 @@ async function convertWithMarkitdown(
|
||||
return { content: "", ok: false, error: "markitdown not found (uv/pip unavailable)" };
|
||||
}
|
||||
|
||||
const child = ptree.cspawn([cmd, filePath], { signal });
|
||||
using child = ptree.spawnGroup([cmd, filePath], { signal });
|
||||
let stdout: string;
|
||||
try {
|
||||
stdout = await child.nothrow().text();
|
||||
|
||||
@@ -49,8 +49,8 @@ export async function convertWithMarkitdown(
|
||||
|
||||
try {
|
||||
await Bun.write(tmpFile, content);
|
||||
const result = await ptree.cspawn([markitdown, tmpFile], { timeout });
|
||||
const [stdout, stderr, exitCode] = await Promise.all([result.stdout.text(), result.stderr.text(), result.exited]);
|
||||
using child = await ptree.spawnGroup([markitdown, tmpFile], { timeout });
|
||||
const [stdout, stderr, exitCode] = await Promise.all([child.stdout.text(), child.stderr.text(), child.exited]);
|
||||
if (exitCode !== 0) {
|
||||
return {
|
||||
content: stdout,
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import * as fs from "node:fs/promises";
|
||||
import * as os from "node:os";
|
||||
import path from "node:path";
|
||||
import { cspawn } from "@oh-my-pi/pi-utils";
|
||||
import { ptree } from "@oh-my-pi/pi-utils";
|
||||
import { nanoid } from "nanoid";
|
||||
import { throwIfAborted } from "../../tools/tool-errors";
|
||||
import { ensureTool } from "../../utils/tools-manager";
|
||||
@@ -16,7 +16,7 @@ async function exec(
|
||||
args: string[],
|
||||
options?: { timeout?: number; input?: string | Buffer; signal?: AbortSignal },
|
||||
): Promise<{ stdout: string; stderr: string; ok: boolean; exitCode: number | null }> {
|
||||
const proc = cspawn([cmd, ...args], {
|
||||
using proc = ptree.spawnGroup([cmd, ...args], {
|
||||
signal: options?.signal,
|
||||
timeout: options?.timeout,
|
||||
stdin: options?.input ? Buffer.from(options.input) : undefined,
|
||||
|
||||
@@ -55,7 +55,7 @@ async function main() {
|
||||
rl.on("line", async line => {
|
||||
if (isWaiting) return;
|
||||
if (line.trim() === "exit") {
|
||||
await client.stop();
|
||||
client.stop();
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
|
||||
@@ -43,7 +43,7 @@ describe.skipIf(!process.env.ANTHROPIC_API_KEY && !process.env.ANTHROPIC_OAUTH_T
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
await client.stop();
|
||||
client.stop();
|
||||
if (sessionDir && fs.existsSync(sessionDir)) {
|
||||
fs.rmSync(sessionDir, { recursive: true });
|
||||
}
|
||||
|
||||
@@ -5,6 +5,6 @@ export * from "./glob";
|
||||
export * as logger from "./logger";
|
||||
export * as postmortem from "./postmortem";
|
||||
export * as ptree from "./ptree";
|
||||
export { AbortError, ChildProcess, cspawn, Exception, NonZeroExitError } from "./ptree";
|
||||
export { AbortError, ChildProcess, Exception, NonZeroExitError } from "./ptree";
|
||||
export * from "./stream";
|
||||
export * from "./temp";
|
||||
|
||||
@@ -6,14 +6,14 @@
|
||||
* - Windows and Unix support for proper tree killing.
|
||||
* - ChildProcess wrapper for capturing output, errors, and kill/detach.
|
||||
*/
|
||||
import { type FileSink, type Spawn, type Subprocess, spawn, spawnSync } from "bun";
|
||||
import { $, type FileSink, type Spawn, type Subprocess, spawn } from "bun";
|
||||
import { postmortem } from ".";
|
||||
|
||||
// Platform detection: process tree kill behavior differs.
|
||||
const isWindows = process.platform === "win32";
|
||||
|
||||
// Set of live children for managed termination/cleanup on shutdown.
|
||||
const managedChildren = new Set<PipedSubprocess>();
|
||||
const managedChildren = new Set<ChildProcess>();
|
||||
|
||||
class AsyncQueue<T> {
|
||||
#items: T[] = [];
|
||||
@@ -73,53 +73,53 @@ function createProcessStream(queue: AsyncQueue<Uint8Array>): ReadableStream<Uint
|
||||
* - Windows: uses taskkill for tree and forceful kill (/T /F)
|
||||
* - Unix: negative PID sends signal to process group (tree kill)
|
||||
*/
|
||||
function killChild(child: PipedSubprocess, signal: NodeJS.Signals = "SIGTERM"): void {
|
||||
async function killChild(child: ChildProcess) {
|
||||
const pid = child.pid;
|
||||
if (!pid) return;
|
||||
if (!pid || child.killed) return;
|
||||
|
||||
try {
|
||||
if (isWindows) {
|
||||
// /T (tree), /F (force): ensure entire tree is killed.
|
||||
spawnSync(["taskkill", ...(signal === "SIGKILL" ? ["/F"] : []), "/T", "/PID", pid.toString()], {
|
||||
stdout: "ignore",
|
||||
stderr: "ignore",
|
||||
timeout: 1000,
|
||||
});
|
||||
} else {
|
||||
// Send signal to process group (negative PID).
|
||||
process.kill(-pid, signal);
|
||||
}
|
||||
const waitForExit = (timeout = 1000) =>
|
||||
Promise.race([Bun.sleep(timeout).then(() => false), child.exited.then(() => true)]);
|
||||
|
||||
// If killed, remove from managed set and clean up.
|
||||
if (child.killed) {
|
||||
managedChildren.delete(child);
|
||||
child.unref();
|
||||
const sendSignal = async (signal?: NodeJS.Signals) => {
|
||||
try {
|
||||
process.kill(pid, signal);
|
||||
} catch {}
|
||||
|
||||
if (child.isProcessGroup) {
|
||||
if (await waitForExit(1000)) return;
|
||||
|
||||
try {
|
||||
if (isWindows) {
|
||||
// /T (tree), /F (force): ensure entire tree is killed.
|
||||
await $`taskkill ${signal === "SIGKILL" ? "/F" : ""} /T /PID ${pid}`.quiet().nothrow();
|
||||
} else {
|
||||
// Send signal to process group (negative PID).
|
||||
process.kill(-pid, signal);
|
||||
}
|
||||
} catch {}
|
||||
}
|
||||
} catch {
|
||||
// Ignore: process may already be dead.
|
||||
}
|
||||
return await waitForExit(1000);
|
||||
};
|
||||
|
||||
if (await sendSignal()) return;
|
||||
await sendSignal("SIGKILL");
|
||||
}
|
||||
|
||||
postmortem.register("managed-children", () => {
|
||||
for (const child of [...managedChildren]) {
|
||||
killChild(child, "SIGKILL");
|
||||
managedChildren.delete(child);
|
||||
}
|
||||
postmortem.register("managed-children", async () => {
|
||||
const children = Array.from(managedChildren);
|
||||
managedChildren.clear();
|
||||
await Promise.all(children.map(killChild));
|
||||
});
|
||||
|
||||
/**
|
||||
* Register a subprocess for managed cleanup.
|
||||
* Will attach to exit Promise so removal happens even if child exits "naturally".
|
||||
*/
|
||||
function registerManaged(child: PipedSubprocess): void {
|
||||
function registerManaged(child: ChildProcess): void {
|
||||
if (child.exitCode !== null) return;
|
||||
if (managedChildren.has(child)) return;
|
||||
child.ref();
|
||||
managedChildren.add(child);
|
||||
|
||||
child.exited.then(() => {
|
||||
child.exited.finally(() => {
|
||||
managedChildren.delete(child);
|
||||
child.unref();
|
||||
});
|
||||
}
|
||||
|
||||
@@ -133,6 +133,7 @@ type PipedSubprocess = Subprocess<"pipe" | "ignore" | null, "pipe", "pipe">;
|
||||
export class ChildProcess {
|
||||
#proc: PipedSubprocess;
|
||||
#detached = false;
|
||||
#group = false;
|
||||
#nothrow = false;
|
||||
#stderrBuffer = "";
|
||||
#stdoutQueue = new AsyncQueue<Uint8Array>();
|
||||
@@ -144,8 +145,9 @@ export class ChildProcess {
|
||||
#exited: Promise<number>;
|
||||
#resolveExited: (ex?: PromiseLike<Exception> | Exception) => void;
|
||||
|
||||
constructor(proc: PipedSubprocess) {
|
||||
registerManaged(proc);
|
||||
constructor(proc: PipedSubprocess, group: boolean) {
|
||||
this.#group = group;
|
||||
registerManaged(this);
|
||||
|
||||
const exitSettled = proc.exited.then(
|
||||
() => {},
|
||||
@@ -241,6 +243,9 @@ export class ChildProcess {
|
||||
this.#proc = proc;
|
||||
}
|
||||
|
||||
get isProcessGroup(): boolean {
|
||||
return this.#group;
|
||||
}
|
||||
get pid(): number | undefined {
|
||||
return this.#proc.pid;
|
||||
}
|
||||
@@ -271,6 +276,9 @@ export class ChildProcess {
|
||||
}
|
||||
return this.#stderrStream;
|
||||
}
|
||||
get proc(): PipedSubprocess {
|
||||
return this.#proc;
|
||||
}
|
||||
|
||||
/**
|
||||
* Peek at the stderr buffer.
|
||||
@@ -286,7 +294,7 @@ export class ChildProcess {
|
||||
detach(): void {
|
||||
if (this.#detached || this.#proc.killed) return;
|
||||
this.#detached = true;
|
||||
if (managedChildren.delete(this.#proc)) {
|
||||
if (managedChildren.delete(this)) {
|
||||
this.#proc.unref();
|
||||
}
|
||||
}
|
||||
@@ -303,25 +311,12 @@ export class ChildProcess {
|
||||
* Kill the process tree.
|
||||
* Optionally set an exit reason (for better error propagation on cancellation).
|
||||
*/
|
||||
kill(signal: NodeJS.Signals = "SIGTERM", reason?: Exception) {
|
||||
kill(reason?: Exception) {
|
||||
if (this.#proc.killed) return;
|
||||
if (reason) {
|
||||
this.#exitReasonPending = reason;
|
||||
}
|
||||
killChild(this.#proc, signal);
|
||||
}
|
||||
|
||||
async killAndWait(): Promise<void> {
|
||||
// Try killing with SIGTERM, then SIGKILL if it doesn't exit within 1 second
|
||||
this.kill("SIGTERM");
|
||||
const exitedOrTimeout = await Promise.race([
|
||||
this.exited.then(() => "exited" as const),
|
||||
Bun.sleep(1000).then(() => "timeout" as const),
|
||||
]);
|
||||
if (exitedOrTimeout === "timeout") {
|
||||
this.kill("SIGKILL");
|
||||
await this.exited.catch(() => {});
|
||||
}
|
||||
killChild(this);
|
||||
}
|
||||
|
||||
// Output utilities (aliases for easy chaining)
|
||||
@@ -356,7 +351,7 @@ export class ChildProcess {
|
||||
attachSignal(signal: AbortSignal): void {
|
||||
const onAbort = () => {
|
||||
const cause = new AbortError(signal.reason, "<cancelled>");
|
||||
this.kill("SIGKILL", cause);
|
||||
this.kill(cause);
|
||||
if (this.#proc.killed) {
|
||||
queueMicrotask(() => {
|
||||
try {
|
||||
@@ -385,7 +380,7 @@ export class ChildProcess {
|
||||
attachTimeout(timeout: number): void {
|
||||
if (timeout <= 0) return;
|
||||
const timeoutId = setTimeout(() => {
|
||||
this.kill("SIGKILL", new TimeoutError(timeout, this.#stderrBuffer));
|
||||
this.kill(new TimeoutError(timeout, this.#stderrBuffer));
|
||||
}, timeout);
|
||||
// Use .finally().catch() to avoid unhandled rejection when #exited rejects
|
||||
this.#exited
|
||||
@@ -394,6 +389,10 @@ export class ChildProcess {
|
||||
})
|
||||
.catch(() => {});
|
||||
}
|
||||
|
||||
[Symbol.dispose](): void {
|
||||
this.kill(new AbortError("process disposed", this.#stderrBuffer));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -468,7 +467,7 @@ type ChildSpawnOptions = Omit<
|
||||
* - Always pipes stdout/stderr, launches in new session/process group (detached).
|
||||
* - Optional AbortSignal integrates with kill-on-abort.
|
||||
*/
|
||||
export function cspawn(cmd: string[], options?: ChildSpawnOptions): ChildProcess {
|
||||
export function spawnGroup(cmd: string[], options?: ChildSpawnOptions): ChildProcess {
|
||||
const { timeout, ...rest } = options ?? {};
|
||||
const child = spawn(cmd, {
|
||||
stdin: "ignore",
|
||||
@@ -478,7 +477,30 @@ export function cspawn(cmd: string[], options?: ChildSpawnOptions): ChildProcess
|
||||
// Windows: new console/pgroup; Unix: setsid for process group.
|
||||
detached: true,
|
||||
});
|
||||
const cproc = new ChildProcess(child);
|
||||
const cproc = new ChildProcess(child, true);
|
||||
if (options?.signal) {
|
||||
cproc.attachSignal(options.signal);
|
||||
}
|
||||
if (timeout && timeout > 0) {
|
||||
cproc.attachTimeout(timeout);
|
||||
}
|
||||
return cproc;
|
||||
}
|
||||
|
||||
/**
|
||||
* Spawn a subprocess as a managed child process.
|
||||
* - Always pipes stdout/stderr, launches in new session/process group (detached).
|
||||
* - Optional AbortSignal integrates with kill-on-abort.
|
||||
*/
|
||||
export function spawnAttached(cmd: string[], options?: ChildSpawnOptions): ChildProcess {
|
||||
const { timeout, ...rest } = options ?? {};
|
||||
const child = spawn(cmd, {
|
||||
stdin: "ignore",
|
||||
...rest,
|
||||
stdout: "pipe",
|
||||
stderr: "pipe",
|
||||
});
|
||||
const cproc = new ChildProcess(child, false);
|
||||
if (options?.signal) {
|
||||
cproc.attachSignal(options.signal);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user