refactor(pi-utils): simplified ptree API and consolidated spawn variants into single function

- Simplified ptree API by removing spawnGroup and spawnAttached variants, consolidating into single spawn function.
- Replaced AsyncQueue class with simpler pushStream utility function for stream management.
- Refactored ChildProcess class to use AbortController instead of callback-based signal handling.
- Renamed captureText method to wait and execText function to exec for clearer API semantics.
- Removed mode parameter from exec function, eliminating distinction between spawn modes.
- Updated all call sites across coding-agent to use simplified ptree API.
This commit is contained in:
can1357
2026-01-28 23:42:59 +01:00
parent e8f63f4502
commit 6260ad6216
11 changed files with 385 additions and 513 deletions
@@ -50,11 +50,12 @@ export async function executeBash(command: string, options?: BashExecutorOptions
artifactId: options?.artifactId,
});
using child = ptree.spawnGroup([shell, ...args, finalCommand], {
using child = ptree.spawn([shell, ...args, finalCommand], {
cwd: options?.cwd,
env: finalEnv,
signal: options?.signal,
timeout: options?.timeout,
detached: true,
});
// Pump streams - errors during abort/timeout are expected
+1 -2
View File
@@ -35,8 +35,7 @@ export async function execCommand(
cwd: string,
options?: ExecOptions,
): Promise<ExecResult> {
const result = await ptree.execText([command, ...args], {
mode: "attached",
const result = await ptree.exec([command, ...args], {
cwd,
signal: options?.signal,
timeout: options?.timeout,
+11 -13
View File
@@ -1,9 +1,9 @@
import { createServer } from "node:net";
import * as path from "node:path";
import { logger } from "@oh-my-pi/pi-utils";
import { $, type Subprocess } from "bun";
import { logger, ptree } from "@oh-my-pi/pi-utils";
import { $ } from "bun";
import { nanoid } from "nanoid";
import { getShellConfig, killProcessTree } from "../utils/shell";
import { getShellConfig } from "../utils/shell";
import { getOrCreateSnapshot } from "../utils/shell-snapshot";
import { time } from "../utils/timings";
import { htmlToBasicMarkdown } from "../web/scrapers/types";
@@ -452,7 +452,7 @@ export function serializeWebSocketMessage(msg: JupyterMessage): ArrayBuffer {
export class PythonKernel {
readonly id: string;
readonly kernelId: string;
readonly gatewayProcess: Subprocess | null;
readonly gatewayProcess: ptree.ChildProcess | null;
readonly gatewayUrl: string;
readonly sessionId: string;
readonly username: string;
@@ -469,7 +469,7 @@ export class PythonKernel {
private constructor(
id: string,
kernelId: string,
gatewayProcess: Subprocess | null,
gatewayProcess: ptree.ChildProcess | null,
gatewayUrl: string,
sessionId: string,
username: string,
@@ -633,14 +633,14 @@ export class PythonKernel {
kernelEnv.PYTHONPATH = pythonPathParts;
}
let gatewayProcess: Subprocess | null = null;
let gatewayProcess: ptree.ChildProcess | null = null;
let gatewayUrl: string | null = null;
let lastError: string | null = null;
for (let attempt = 0; attempt < GATEWAY_STARTUP_ATTEMPTS; attempt += 1) {
const gatewayPort = await allocatePort();
const candidateUrl = `http://127.0.0.1:${gatewayPort}`;
const candidateProcess = Bun.spawn(
const candidateProcess = ptree.spawn(
[
runtime.pythonPath,
"-m",
@@ -653,10 +653,8 @@ export class PythonKernel {
],
{
cwd: options.cwd,
stdin: "ignore",
stdout: "pipe",
stderr: "pipe",
env: kernelEnv,
detached: true,
},
);
@@ -687,7 +685,7 @@ export class PythonKernel {
if (gatewayProcess && gatewayUrl) break;
await killProcessTree(candidateProcess.pid);
candidateProcess.kill();
lastError = exited ? "Kernel gateway process exited during startup" : "Kernel gateway failed to start";
}
@@ -702,7 +700,7 @@ export class PythonKernel {
});
if (!createResponse.ok) {
await killProcessTree(gatewayProcess.pid);
gatewayProcess.kill();
throw new Error(`Failed to create kernel: ${await createResponse.text()}`);
}
@@ -1118,7 +1116,7 @@ export class PythonKernel {
await releaseSharedGateway();
} else if (this.gatewayProcess) {
try {
await killProcessTree(this.gatewayProcess.pid);
this.gatewayProcess.kill();
} catch (err: unknown) {
logger.warn("Failed to terminate gateway process", {
error: err instanceof Error ? err.message : String(err),
@@ -112,7 +112,7 @@ export class RpcClient {
args.push(...this.options.args);
}
this.process = ptree.spawnAttached(["bun", cliPath, ...args], {
this.process = ptree.spawn(["bun", cliPath, ...args], {
cwd: this.options.cwd,
env: { ...process.env, ...this.options.env },
stdin: "pipe",
@@ -76,7 +76,7 @@ export async function executeSSH(
}
}
using child = ptree.spawnAttached(["ssh", ...(await buildRemoteCommand(host, resolvedCommand))], {
using child = ptree.spawn(["ssh", ...(await buildRemoteCommand(host, resolvedCommand))], {
signal: options?.signal,
timeout: options?.timeout,
});
+2 -2
View File
@@ -410,7 +410,7 @@ async function renderHtmlToText(
if (lynx) {
const normalizedPath = tmpFile.replace(/\\/g, "/");
const fileUrl = normalizedPath.startsWith("/") ? `file://${normalizedPath}` : `file:///${normalizedPath}`;
const result = await ptree.execText(["lynx", "-dump", "-nolist", "-width", "120", fileUrl], execOptions);
const result = await ptree.exec(["lynx", "-dump", "-nolist", "-width", "120", fileUrl], execOptions);
if (result.ok) {
return { content: result.stdout, ok: true, method: "lynx" };
}
@@ -419,7 +419,7 @@ async function renderHtmlToText(
// Fall back to html2text (auto-install via uv/pip)
const html2text = await ensureTool("html2text", true);
if (html2text) {
const result = await ptree.execText([html2text, tmpFile], execOptions);
const result = await ptree.exec([html2text, tmpFile], execOptions);
if (result.ok) {
return { content: result.stdout, ok: true, method: "html2text" };
}
+3 -3
View File
@@ -77,7 +77,7 @@ export async function runRg(
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";
const result = await ptree.execText([rgPath, ...args], {
const result = await ptree.exec([rgPath, ...args], {
signal: options?.signal,
timeout: options?.timeoutMs,
allowNonZero: true,
@@ -149,7 +149,7 @@ export class GrepTool implements AgentTool<typeof grepSchema, GrepToolDetails> {
// Run ripgrep against the null device with the pattern - this validates regex syntax
// without searching any files
const nullDevice = process.platform === "win32" ? "NUL" : "/dev/null";
const result = await ptree.execText([rgPath, "--no-config", "--quiet", "--", pattern, nullDevice], {
const result = await ptree.exec([rgPath, "--no-config", "--quiet", "--", pattern, nullDevice], {
allowNonZero: true,
allowAbort: true,
});
@@ -314,7 +314,7 @@ export class GrepTool implements AgentTool<typeof grepSchema, GrepToolDetails> {
args.push("--", normalizedPattern, searchPath);
using child = ptree.spawnAttached([rgPath, ...args], { signal });
using child = ptree.spawn([rgPath, ...args], { signal });
let matchCount = 0;
let matchLimitReached = false;
+2 -2
View File
@@ -296,12 +296,12 @@ async function convertWithMarkitdown(
return { content: "", ok: false, error: "markitdown not found (uv/pip unavailable)" };
}
const result = await ptree.execText([cmd, filePath], {
mode: "group",
const result = await ptree.exec([cmd, filePath], {
signal,
allowNonZero: true,
allowAbort: true,
stderr: "buffer",
detached: true,
});
if (result.exitError?.aborted) {
@@ -49,11 +49,11 @@ export async function convertWithMarkitdown(
try {
await Bun.write(tmpFile, content);
const result = await ptree.execText([markitdown, tmpFile], {
mode: "group",
const result = await ptree.exec([markitdown, tmpFile], {
timeout,
allowNonZero: true,
stderr: "full",
detached: true,
});
if (!result.ok) {
return {
@@ -146,7 +146,7 @@ export const handleYouTube: SpecialHandler = async (
// Fetch video metadata
throwIfAborted(signal);
const metaResult = await ptree.execText(
const metaResult = await ptree.exec(
[ytdlp, "--dump-json", "--no-warnings", "--no-playlist", "--skip-download", videoUrl],
execOptions,
);
@@ -191,7 +191,7 @@ export const handleYouTube: SpecialHandler = async (
// First, list available subtitles
throwIfAborted(signal);
const listResult = await ptree.execText(
const listResult = await ptree.exec(
[ytdlp, "--list-subs", "--no-warnings", "--no-playlist", "--skip-download", videoUrl],
execOptions,
);
@@ -208,7 +208,7 @@ export const handleYouTube: SpecialHandler = async (
// Try manual subtitles first (English preferred)
if (hasManualSubs) {
throwIfAborted(signal);
const subResult = await ptree.execText(
const subResult = await ptree.exec(
[
ytdlp,
"--write-sub",
@@ -243,7 +243,7 @@ export const handleYouTube: SpecialHandler = async (
// Fall back to auto-generated captions
if (!transcript && hasAutoSubs) {
throwIfAborted(signal);
const autoResult = await ptree.execText(
const autoResult = await ptree.exec(
[
ytdlp,
"--write-auto-sub",
+356 -482
View File
@@ -1,126 +1,131 @@
/**
* Process tree management utilities for Bun subprocesses.
*
* Provides:
* - Managed tracking of child subprocesses for cleanup on exit/signals.
* - Windows and Unix support for proper tree killing.
* - ChildProcess wrapper for capturing output, errors, and kill/detach.
* Exposes the same public interface as the original implementation, but with
* much less code:
* - Track managed child processes for cleanup on shutdown (postmortem).
* - Drain stdout/stderr to avoid subprocess pipe deadlocks.
* - Cross-platform tree kill for process groups (Windows taskkill, Unix -pid).
* - Convenience helpers: captureText / execText, AbortSignal, timeouts.
*/
import { $, type FileSink, type Spawn, type Subprocess, spawn } from "bun";
import { $, type FileSink, type Spawn, type Subprocess } 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<ChildProcess>();
class AsyncQueue<T> {
#items: T[] = [];
#resolvers: Array<(result: IteratorResult<T>) => void> = [];
#closed = false;
/** A Bun subprocess with stdout/stderr always piped (stdin may vary). */
type PipedSubprocess = Subprocess<"pipe" | "ignore" | null, "pipe", "pipe">;
push(item: T): void {
if (this.#closed) return;
const resolver = this.#resolvers.shift();
if (resolver) {
resolver({ value: item, done: false });
return;
}
this.#items.push(item);
}
/** Minimal push-based ReadableStream that buffers unboundedly (like the old queue). */
function pushStream<T>() {
let controller!: ReadableStreamDefaultController<T>;
let closed = false;
close(options?: { discard?: boolean }): void {
if (this.#closed) {
if (options?.discard) {
this.#items = [];
}
return;
}
this.#closed = true;
if (options?.discard) {
this.#items = [];
}
while (this.#resolvers.length > 0) {
const resolver = this.#resolvers.shift();
if (resolver) {
resolver({ value: undefined, done: true });
}
}
}
async next(): Promise<IteratorResult<T>> {
if (this.#items.length > 0) {
return { value: this.#items.shift() as T, done: false };
}
if (this.#closed) {
return { value: undefined, done: true };
}
return await new Promise<IteratorResult<T>>(resolve => {
this.#resolvers.push(resolve);
});
}
}
function createProcessStream(queue: AsyncQueue<Uint8Array>, onCancel?: () => void): ReadableStream<Uint8Array> {
const stream = new ReadableStream<Uint8Array>({
pull: async controller => {
const result = await queue.next();
if (result.done) {
controller.close();
return;
}
controller.enqueue(result.value);
const stream = new ReadableStream<T>({
start(c) {
controller = c;
},
cancel: () => {
onCancel?.();
queue.close({ discard: true });
cancel() {
closed = true; // consumer no longer cares; keep draining but drop
},
});
return stream;
return {
stream,
push(value: T) {
if (closed) return;
try {
controller.enqueue(value);
} catch {
closed = true;
}
},
close() {
if (closed) return;
closed = true;
try {
controller.close();
} catch {}
},
};
}
const DONE = { done: true, value: undefined } as const;
function abortRead(signal: AbortSignal) {
if (signal.aborted) return Promise.resolve(DONE);
const { promise, resolve } = Promise.withResolvers<typeof DONE>();
signal.addEventListener("abort", () => resolve(DONE), { once: true });
return promise;
}
/** Drain a ReadableStream into a pushStream, optionally tapping each chunk. */
async function pump(
src: ReadableStream<Uint8Array>,
dst: ReturnType<typeof pushStream<Uint8Array>>,
opts?: { signal?: AbortSignal; onChunk?: (chunk: Uint8Array) => void; onFinally?: () => void },
) {
const reader = src.getReader();
const stop = opts?.signal ? abortRead(opts.signal) : null;
try {
while (true) {
const r = stop ? await Promise.race([reader.read(), stop]) : await reader.read();
if (r.done) break;
if (!r.value) continue;
opts?.onChunk?.(r.value);
dst.push(r.value);
}
} catch {
// ignore; this module is "best effort" for streaming/cleanup
} finally {
try {
await reader.cancel();
} catch {}
try {
reader.releaseLock();
} catch {}
dst.close();
opts?.onFinally?.();
}
}
/**
* Kill a child process and its descendents.
* - Windows: uses taskkill for tree and forceful kill (/T /F)
* - Unix: negative PID sends signal to process group (tree kill)
* - Windows: taskkill /T, add /F on SIGKILL
* - Unix: negative PID signals the process group
*/
async function killChild(child: ChildProcess) {
const pid = child.pid;
if (!pid || child.killed) return;
const waitForExit = (timeout = 1000) =>
Promise.race([
Bun.sleep(timeout).then(() => false),
child.proc.exited.then(
() => true,
() => true,
),
]);
const exited = child.proc.exited.then(
() => true,
() => true,
);
const waitForExit = (timeout = 1000) => Promise.race([Bun.sleep(timeout).then(() => false), exited]);
const sendSignal = async (signal?: NodeJS.Signals) => {
// Give it a moment to exit gracefully first.
try {
child.proc.kill();
} catch {}
if (await waitForExit(1000)) return true;
if (child.isProcessGroup) {
try {
child.proc.kill(signal);
if (isWindows) {
await $`taskkill /F /T /PID ${pid}`.quiet().nothrow();
} else {
process.kill(-pid);
}
} catch {}
}
try {
child.proc.kill("SIGKILL");
} 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 {}
}
return await waitForExit(1000);
};
if (await sendSignal()) return;
await sendSignal("SIGKILL");
return await waitForExit(1000);
}
postmortem.register("managed-children", async () => {
@@ -129,27 +134,19 @@ postmortem.register("managed-children", async () => {
await Promise.all(children.map(killChild));
});
// A Bun subprocess with stdin=Writable/ignore, stdout/stderr=pipe (for tracking/cleanup).
type PipedSubprocess = Subprocess<"pipe" | "ignore" | null, "pipe", "pipe">;
type StreamReadResult = { done: boolean; value: Uint8Array | undefined };
/**
* Options for capturing process output as text.
* Options for waiting for process exit and capturing output.
*/
export interface CaptureTextOptions {
/** Allow non-zero exit codes without throwing. */
export interface WaitOptions {
allowNonZero?: boolean;
/** Allow abort/timeout without throwing. */
allowAbort?: boolean;
/** Select stderr source: full stream or bounded buffer. */
stderr?: "full" | "buffer";
}
/**
* Result from captureText/execText.
* Result from wait and captureText.
*/
export interface CaptureTextResult {
export interface ExecResult {
stdout: string;
stderr: string;
exitCode: number | null;
@@ -157,323 +154,6 @@ export interface CaptureTextResult {
exitError?: Exception;
}
/**
* ChildProcess wraps a managed subprocess, capturing output, errors, and providing
* cross-platform kill/detach logic plus AbortSignal integration.
*/
export class ChildProcess {
#nothrow = false;
#stderrBuffer = "";
#stdoutQueue = new AsyncQueue<Uint8Array>();
#stderrQueue = new AsyncQueue<Uint8Array>();
#stderrDone!: Promise<void>;
#requestStreamStop: () => void;
#stdoutActive = true;
#stdoutStream?: ReadableStream<Uint8Array>;
#stderrStream?: ReadableStream<Uint8Array>;
#exitReason?: Exception;
#exitReasonPending?: Exception;
#exited: Promise<number>;
constructor(
public readonly proc: PipedSubprocess,
public readonly isProcessGroup: boolean,
) {
const { promise: stopStreaming, resolve: resolveStopStreaming } = Promise.withResolvers<StreamReadResult>();
this.#requestStreamStop = () => void resolveStopStreaming({ done: true, value: undefined });
const { promise: stderrDone, resolve: resolveStderrDone } = Promise.withResolvers<void>();
this.#stderrDone = stderrDone;
// Capture stdout while active. Buffering starts enabled and is disabled when the
// stream is cancelled. The underlying process stdout is always drained to prevent
// the process from blocking on a full pipe buffer.
void (async () => {
const reader = proc.stdout.getReader();
try {
while (this.#stdoutActive) {
const result = await Promise.race([reader.read(), stopStreaming]);
if (result.done) break;
if (!result.value) continue;
this.#stdoutQueue.push(result.value);
}
} catch {
// ignore
} finally {
try {
await reader.cancel();
} catch {}
try {
reader.releaseLock();
} catch {}
this.#stdoutQueue.close();
}
})().catch(() => {
this.#stdoutQueue.close();
});
// Capture stderr at all times, with a capped buffer for errors.
const decoder = new TextDecoder();
void (async () => {
const reader = proc.stderr.getReader();
try {
while (true) {
const result = await Promise.race([reader.read(), stopStreaming]);
if (result.done) break;
if (!result.value) continue;
this.#stderrQueue.push(result.value);
this.#stderrBuffer += decoder.decode(result.value, { stream: true });
if (this.#stderrBuffer.length > NonZeroExitError.MAX_TRACE) {
this.#stderrBuffer = this.#stderrBuffer.slice(-NonZeroExitError.MAX_TRACE);
}
}
} catch {
// ignore
} finally {
this.#stderrBuffer += decoder.decode();
if (this.#stderrBuffer.length > NonZeroExitError.MAX_TRACE) {
this.#stderrBuffer = this.#stderrBuffer.slice(-NonZeroExitError.MAX_TRACE);
}
try {
await reader.cancel();
} catch {}
try {
reader.releaseLock();
} catch {}
this.#stderrQueue.close();
resolveStderrDone();
}
})().catch(() => {
this.#stderrQueue.close();
resolveStderrDone();
});
const { promise, resolve, reject } = Promise.withResolvers<number>();
this.#exited = promise;
// On exit, resolve with a ChildError if nonzero code.
if (this.proc.exitCode === null) {
managedChildren.add(this);
}
proc.exited
.catch(() => null)
.then(async exitCode => {
// If we have an exit reason pending (e.g., kill() was called), use it immediately.
if (this.#exitReasonPending) {
this.#exitReason = this.#exitReasonPending;
reject(this.#exitReasonPending);
return;
}
// If successful, resolve as 0.
if (exitCode === 0) {
resolve(0);
return;
}
// Wait for stderr capture to complete before creating error with stderr content.
await this.#stderrDone;
let ex: Exception;
if (exitCode !== null) {
this.#exitReason = new NonZeroExitError(exitCode, this.#stderrBuffer);
resolve(exitCode);
return;
} else if (this.proc.killed) {
ex = new AbortError(new Error("process killed"), this.#stderrBuffer);
} else {
ex = new NonZeroExitError(-1, this.#stderrBuffer);
}
this.#exitReason = ex;
reject(ex);
})
.finally(() => {
managedChildren.delete(this);
});
}
get pid(): number | undefined {
return this.proc.pid;
}
get exited(): Promise<number> {
return this.#exited;
}
get exitedCleanly(): Promise<number> {
if (this.#nothrow) return this.exited;
return this.exited.then(code => {
if (code !== 0) {
throw new NonZeroExitError(code, this.#stderrBuffer);
}
return code;
});
}
get exitCode(): number | null {
return this.proc.exitCode;
}
get exitReason(): Exception | undefined {
return this.#exitReason;
}
get killed(): boolean {
return this.proc.killed;
}
get stdin(): FileSink | undefined {
return this.proc.stdin;
}
get stdout(): ReadableStream<Uint8Array> {
if (!this.#stdoutStream) {
this.#stdoutStream = createProcessStream(this.#stdoutQueue, () => {
this.#stdoutActive = false;
});
}
return this.#stdoutStream;
}
get stderr(): ReadableStream<Uint8Array> {
if (!this.#stderrStream) {
// stderr cancellation doesn't affect the internal buffer used for error context
this.#stderrStream = createProcessStream(this.#stderrQueue, () => {});
}
return this.#stderrStream;
}
/**
* Peek at the stderr buffer.
* @returns The stderr buffer.
*/
peekStderr(): string {
return this.#stderrBuffer;
}
/**
* Prevents thrown ChildError on nonzero exit code, for optional error handling.
*/
nothrow(): this {
this.#nothrow = true;
return this;
}
/**
* Kill the process tree.
* Optionally set an exit reason (for better error propagation on cancellation).
*/
kill(reason?: Exception) {
if (reason && !this.#exitReasonPending) {
this.#exitReasonPending = reason;
}
this.#requestStreamStop();
if (this.proc.killed) return;
killChild(this);
}
// Output utilities (aliases for easy chaining)
async text(): Promise<string> {
return (await this.blob()).text();
}
async json(): Promise<unknown> {
return (await this.blob()).json();
}
async arrayBuffer(): Promise<ArrayBuffer> {
return (await this.blob()).arrayBuffer();
}
async bytes() {
return (await this.blob()).bytes();
}
async blob() {
const { promise, resolve, reject } = Promise.withResolvers<Blob>();
const blob = this.stdout.blob();
if (!this.#nothrow) {
this.exitedCleanly.catch(reject);
}
blob.then(resolve, reject);
return promise;
}
/**
* Capture stdout/stderr as text with optional exit handling.
*/
async captureText(options?: CaptureTextOptions): Promise<CaptureTextResult> {
const stderrMode = options?.stderr ?? "buffer";
const stdoutPromise = this.stdout.text();
const stderrPromise =
stderrMode === "full"
? this.stderr.text()
: (async () => {
await Promise.allSettled([stdoutPromise, this.exited, this.#stderrDone]);
return this.peekStderr();
})();
const [stdout, stderr] = await Promise.all([stdoutPromise, stderrPromise]);
let exitError: Exception | undefined;
try {
await this.exited;
} catch (err) {
if (err instanceof Exception) {
exitError = err;
} else {
throw err;
}
}
const exitCode = this.exitCode ?? (exitError && !exitError.aborted ? exitError.exitCode : null);
const ok = exitCode === 0;
if (exitError) {
const allowAbort = options?.allowAbort ?? false;
const allowNonZero = options?.allowNonZero ?? false;
if ((exitError.aborted && !allowAbort) || (!exitError.aborted && !allowNonZero)) {
throw exitError;
}
}
return { stdout, stderr, exitCode, ok, exitError };
}
/**
* Attach an AbortSignal to this process. Will kill tree with SIGKILL if aborted.
*/
attachSignal(signal: AbortSignal): void {
const onAbort = () => {
const cause = new AbortError(signal.reason, "<cancelled>");
this.kill(cause);
};
if (signal.aborted) {
return void onAbort();
}
signal.addEventListener("abort", onAbort, { once: true });
// Use .finally().catch() to avoid unhandled rejection when #exited rejects
this.#exited
.catch(() => {})
.finally(() => {
signal.removeEventListener("abort", onAbort);
});
}
/**
* Attach a timeout to this process. Will kill the process with SIGKILL if the timeout is reached.
*/
attachTimeout(timeout: number): void {
if (timeout <= 0) return;
if (this.proc.killed) return;
void (async () => {
const result = await Promise.race([
Bun.sleep(timeout).then(() => true),
this.proc.exited.then(
() => false,
() => false,
),
]);
if (result) {
this.kill(new TimeoutError(timeout, this.#stderrBuffer));
}
});
}
[Symbol.dispose](): void {
this.kill(new AbortError("process disposed", this.#stderrBuffer));
}
}
/**
* Base for all exceptions representing child process nonzero exit, killed, or cancellation.
*/
@@ -531,81 +211,275 @@ export class TimeoutError extends AbortError {
}
}
/**
* ChildProcess wraps a managed subprocess, capturing stderr tail, providing
* cross-platform kill/detach logic plus AbortSignal integration.
*/
export class ChildProcess {
#nothrow = false;
#stderrBuffer = "";
#exitReason?: Exception;
#exitReasonPending?: Exception;
#stop = new AbortController();
#stdoutOut = pushStream<Uint8Array>();
#stderrOut = pushStream<Uint8Array>();
#stderrDone: Promise<void>;
#exited: Promise<number>;
constructor(
public readonly proc: PipedSubprocess,
public readonly isProcessGroup: boolean,
) {
const { promise: stderrDone, resolve: resolveStderrDone } = Promise.withResolvers<void>();
this.#stderrDone = stderrDone;
// Drain stdout always -> expose our buffered stream to the user.
void pump(proc.stdout, this.#stdoutOut, { signal: this.#stop.signal }).catch(() => this.#stdoutOut.close());
// Drain stderr always -> expose stream + keep a bounded tail buffer.
const decoder = new TextDecoder();
const trim = () => {
if (this.#stderrBuffer.length > NonZeroExitError.MAX_TRACE) {
this.#stderrBuffer = this.#stderrBuffer.slice(-NonZeroExitError.MAX_TRACE);
}
};
void pump(proc.stderr, this.#stderrOut, {
signal: this.#stop.signal,
onChunk: chunk => {
this.#stderrBuffer += decoder.decode(chunk, { stream: true });
trim();
},
onFinally: () => {
this.#stderrBuffer += decoder.decode();
trim();
resolveStderrDone();
},
}).catch(() => {
try {
this.#stderrBuffer += decoder.decode();
trim();
} catch {}
this.#stderrOut.close();
resolveStderrDone();
});
const { promise, resolve, reject } = Promise.withResolvers<number>();
this.#exited = promise;
if (this.proc.exitCode === null) managedChildren.add(this);
// Normalize Bun's exited promise into our "exitReason / exitedCleanly" model.
proc.exited
.catch(() => null)
.then(async exitCode => {
if (this.#exitReasonPending) {
this.#exitReason = this.#exitReasonPending;
reject(this.#exitReasonPending);
return;
}
if (exitCode === 0) {
resolve(0);
return;
}
await this.#stderrDone;
if (exitCode !== null) {
this.#exitReason = new NonZeroExitError(exitCode, this.#stderrBuffer);
resolve(exitCode);
return;
}
const ex = this.proc.killed
? new AbortError(new Error("process killed"), this.#stderrBuffer)
: new NonZeroExitError(-1, this.#stderrBuffer);
this.#exitReason = ex;
reject(ex);
})
.finally(() => {
managedChildren.delete(this);
});
}
get pid(): number | undefined {
return this.proc.pid;
}
get exited(): Promise<number> {
return this.#exited;
}
get exitedCleanly(): Promise<number> {
if (this.#nothrow) return this.exited;
return this.exited.then(code => {
if (code !== 0) throw new NonZeroExitError(code, this.#stderrBuffer);
return code;
});
}
get exitCode(): number | null {
return this.proc.exitCode;
}
get exitReason(): Exception | undefined {
return this.#exitReason;
}
get killed(): boolean {
return this.proc.killed;
}
get stdin(): FileSink | undefined {
return this.proc.stdin;
}
get stdout(): ReadableStream<Uint8Array> {
return this.#stdoutOut.stream;
}
get stderr(): ReadableStream<Uint8Array> {
return this.#stderrOut.stream;
}
peekStderr(): string {
return this.#stderrBuffer;
}
nothrow(): this {
this.#nothrow = true;
return this;
}
kill(reason?: Exception) {
if (reason && !this.#exitReasonPending) this.#exitReasonPending = reason;
this.#stop.abort();
if (this.proc.killed) return;
void killChild(this);
}
// Output helpers
async blob(): Promise<Blob> {
const blobPromise = new Response(this.stdout).blob();
if (this.#nothrow) return await blobPromise;
const [blob] = await Promise.all([blobPromise, this.exitedCleanly]);
return blob;
}
async text(): Promise<string> {
return (await this.blob()).text();
}
async json(): Promise<unknown> {
return await new Response(await this.blob()).json();
}
async arrayBuffer(): Promise<ArrayBuffer> {
return (await this.blob()).arrayBuffer();
}
async bytes(): Promise<Uint8Array> {
return new Uint8Array(await this.arrayBuffer());
}
async wait(options?: WaitOptions): Promise<ExecResult> {
const { allowNonZero = false, allowAbort = false, stderr: stderrMode = "buffer" } = options ?? {};
const stdoutPromise = new Response(this.stdout).text();
const stderrPromise =
stderrMode === "full"
? new Response(this.stderr).text()
: (async () => {
await Promise.allSettled([stdoutPromise, this.exited, this.#stderrDone]);
return this.peekStderr();
})();
const [stdout, stderr] = await Promise.all([stdoutPromise, stderrPromise]);
let exitError: Exception | undefined;
try {
await this.exited;
} catch (err) {
if (err instanceof Exception) exitError = err;
else throw err;
}
const exitCode = this.exitCode ?? (exitError && !exitError.aborted ? exitError.exitCode : null);
const ok = exitCode === 0;
if (exitError) {
if ((exitError.aborted && !allowAbort) || (!exitError.aborted && !allowNonZero)) {
throw exitError;
}
}
return { stdout, stderr, exitCode, ok, exitError };
}
attachSignal(signal: AbortSignal): void {
const onAbort = () => this.kill(new AbortError(signal.reason, "<cancelled>"));
if (signal.aborted) return void onAbort();
signal.addEventListener("abort", onAbort, { once: true });
this.#exited
.catch(() => {})
.finally(() => {
signal.removeEventListener("abort", onAbort);
});
}
attachTimeout(timeout: number): void {
if (timeout <= 0 || this.proc.killed) return;
void (async () => {
const timedOut = await Promise.race([
Bun.sleep(timeout).then(() => true),
this.proc.exited.then(
() => false,
() => false,
),
]);
if (timedOut) this.kill(new TimeoutError(timeout, this.#stderrBuffer));
})();
}
[Symbol.dispose](): void {
this.kill(new AbortError("process disposed", this.#stderrBuffer));
}
}
/**
* Options for cspawn (child spawn). Always pipes stdout/stderr, allows signal.
*/
type ChildSpawnOptions = Omit<
Spawn.SpawnOptions<"pipe" | "ignore" | Buffer | null, "pipe", "pipe">,
Spawn.SpawnOptions<"pipe" | "ignore" | Buffer | Uint8Array | null, "pipe", "pipe">,
"stdout" | "stderr"
> & {
signal?: AbortSignal;
};
> & { signal?: AbortSignal; detached?: boolean };
function spawnManaged(
cmd: string[],
options: ChildSpawnOptions | undefined,
config: { detached: boolean; processGroup: boolean },
): ChildProcess {
const { timeout, ...rest } = options ?? {};
const child = spawn(cmd, {
/**
* Spawn a child process.
* @param cmd - The command to spawn.
* @param options - The options for the spawn.
* @returns A ChildProcess instance.
*/
export function spawn(cmd: string[], options?: ChildSpawnOptions): ChildProcess {
const { detached = false, timeout, signal, ...rest } = options ?? {};
const child = Bun.spawn(cmd, {
stdin: "ignore",
...rest,
stdout: "pipe",
stderr: "pipe",
...(config.detached ? { detached: true } : {}),
detached,
...rest,
});
const cproc = new ChildProcess(child, config.processGroup);
if (options?.signal) {
cproc.attachSignal(options.signal);
}
if (timeout && timeout > 0) {
cproc.attachTimeout(timeout);
}
const cproc = new ChildProcess(child, detached);
if (signal) cproc.attachSignal(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 spawnGroup(cmd: string[], options?: ChildSpawnOptions): ChildProcess {
return spawnManaged(cmd, options, { detached: true, processGroup: true });
}
/**
* Spawn a subprocess as a managed child process.
* - Always pipes stdout/stderr, inherits the current session (not detached).
* - Optional AbortSignal integrates with kill-on-abort.
*/
export function spawnAttached(cmd: string[], options?: ChildSpawnOptions): ChildProcess {
return spawnManaged(cmd, options, { detached: false, processGroup: false });
}
/**
* Options for execText.
*/
export interface ExecTextOptions extends Omit<ChildSpawnOptions, "stdin">, CaptureTextOptions {
/** Spawn mode (process group or attached). */
mode?: "group" | "attached";
/** Input to write to stdin (Buffer or UTF-8 string). */
export interface ExecOptions extends Omit<ChildSpawnOptions, "stdin">, WaitOptions {
input?: string | Buffer | Uint8Array;
}
function toStdinBuffer(input: string | Buffer | Uint8Array): Buffer {
if (typeof input === "string") {
return Buffer.from(input);
}
return Buffer.isBuffer(input) ? input : Buffer.from(input);
}
/**
* Spawn a process and capture stdout/stderr as text.
*/
export async function execText(cmd: string[], options?: ExecTextOptions): Promise<CaptureTextResult> {
const { mode = "attached", input, stderr, allowAbort, allowNonZero, ...spawnOptions } = options ?? {};
const stdin = input === undefined ? undefined : toStdinBuffer(input);
export async function exec(cmd: string[], options?: ExecOptions): Promise<ExecResult> {
const { input, stderr, allowAbort, allowNonZero, ...spawnOptions } = options ?? {};
const stdin = typeof input === "string" ? Buffer.from(input) : input;
const resolvedOptions: ChildSpawnOptions = stdin === undefined ? { ...spawnOptions } : { ...spawnOptions, stdin };
using child = mode === "group" ? spawnGroup(cmd, resolvedOptions) : spawnAttached(cmd, resolvedOptions);
return await child.captureText({ stderr, allowAbort, allowNonZero });
using child = spawn(cmd, resolvedOptions);
return await child.wait({ stderr, allowAbort, allowNonZero });
}