// Process supervisor manages long-running child and PTY process lifecycles. import crypto from "node:crypto"; import { performance } from "node:perf_hooks"; import { expectDefined } from "@openclaw/normalization-core"; import { resolveTimerTimeoutMs } from "@openclaw/normalization-core/number-coercion"; import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import { sliceUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; import { createDeferredCore } from "../../shared/deferred.js"; import { createLazyRuntimeModule } from "../../shared/lazy-runtime.js"; import { createChildAdapter } from "./adapters/child.js"; import { createPtyAdapter } from "./adapters/pty.js"; import { GRACEFUL_CANCEL_TIMEOUT_MS } from "./cancellation-policy.js"; import type { ManagedRun, ProcessSupervisor, ProcessScopeCleanupPolicy, RunExit, SpawnInput, SpawnProcessAdapter, TerminationReason, } from "./types.js"; type OwnedRun = { runId: string; scopeKey?: string; terminationReason?: TerminationReason; cancel: (reason: TerminationReason) => void; pending?: Promise; waitForExtinction?: () => Promise; cleanupOwners: ScopeCleanupOwner[]; }; type ScopeCleanupOwner = { processTree: ProcessScopeCleanupPolicy; failure?: { error: unknown } }; function requiresProcessTree(scope: ScopeCleanupOwner, external: boolean): boolean { return scope.processTree === "required-all" || (scope.processTree === "owned-only" && !external); } function recordScopeCleanupFailure(owner: OwnedRun, error: unknown): void { for (const cleanupOwner of owner.cleanupOwners) { cleanupOwner.failure ??= { error }; } } type StartingScope = { runs: Set>; replacement?: Promise; }; const DEFAULT_MAX_CAPTURED_OUTPUT_CHARS = 1024 * 1024; const loadSupervisorLogRuntime = createLazyRuntimeModule( () => import("./supervisor-log.runtime.js"), ); function normalizeTimeoutDuration(value?: number): number | undefined { if (typeof value !== "number" || !Number.isFinite(value) || value <= 0) { return undefined; } return Math.max(1, Math.floor(value)); } function clampCapturedOutputChars(value?: number): number { if (typeof value !== "number" || !Number.isFinite(value) || value <= 0) { return DEFAULT_MAX_CAPTURED_OUTPUT_CHARS; } return Math.max(256, Math.floor(value)); } function appendCapturedOutput( current: string, chunk: string, stream: "stdout" | "stderr", maxChars: number, ) { const next = current + chunk; if (next.length <= maxChars) { return next; } const marker = `[openclaw: captured ${stream} truncated to last ${maxChars} chars]\n`; const tailChars = Math.max(0, maxChars - marker.length); return `${marker}${sliceUtf16Safe(next, -tailChars)}`; } function isTimeoutReason(reason: TerminationReason) { return reason === "overall-timeout" || reason === "no-output-timeout"; } function resolveElapsedTimeoutReason(params: { nowMs: number; overallTimeoutDeadlineMs: number | null; noOutputTimeoutDeadlineMs: number | null; }): TerminationReason | null { if ( params.overallTimeoutDeadlineMs !== null && params.nowMs >= params.overallTimeoutDeadlineMs && (params.noOutputTimeoutDeadlineMs === null || params.nowMs < params.noOutputTimeoutDeadlineMs || params.overallTimeoutDeadlineMs <= params.noOutputTimeoutDeadlineMs) ) { return "overall-timeout"; } return params.noOutputTimeoutDeadlineMs !== null && params.nowMs >= params.noOutputTimeoutDeadlineMs ? "no-output-timeout" : null; } export function createProcessSupervisor(): ProcessSupervisor & { shutdown: () => Promise; } { // Retries share a run ID while an older command can still own descendants. // Keep each admission until its own cleanup completes. const ownedRuns = new Set(); const scopeCleanupOwners = new Map>(); const startingScopes = new Map(); let shuttingDown = false; let shutdownPromise: Promise | null = null; let cleanupFailure: { error: unknown } | undefined; const cancel = (runId: string, reason: TerminationReason = "manual-cancel") => { for (const current of ownedRuns) { if (current.runId === runId) { current.cancel(reason); } } }; const cancelActiveScope = (scopeKey: string, reason: TerminationReason) => { for (const current of ownedRuns) { if (current.waitForExtinction && current.scopeKey === scopeKey) { current.cancel(reason); } } }; const cancelScope = (scopeKey: string, reason: TerminationReason = "manual-cancel") => { if (!scopeKey.trim()) { return; } for (const current of ownedRuns) { if (current.scopeKey === scopeKey) { current.cancel(reason); } } }; const waitForRuns = async ( scopeKey: string | null, ignoreStartupFailures = false, ): Promise => { let firstFailure: PromiseRejectedResult | undefined; const observed = new Set(); while (true) { const selected = Array.from(ownedRuns).filter( (current) => !observed.has(current) && (scopeKey === null || current.scopeKey === scopeKey), ); const starts = selected.flatMap((current) => (current.pending ? [current.pending] : [])); const owned = selected.flatMap((current) => { if (!current.waitForExtinction) { return []; } observed.add(current); return [current.waitForExtinction()]; }); if (starts.length === 0 && owned.length === 0) { if (firstFailure) { throw firstFailure.reason; } return; } // Startup can become active while the snapshot settles; recheck admissions // so shutdown cannot outrun an admitted command or retained descendants. const results = await Promise.allSettled([...owned, ...starts]); firstFailure ??= results .slice(0, ignoreStartupFailures ? owned.length : undefined) .find((result): result is PromiseRejectedResult => result.status === "rejected"); } }; const acquireScopeCleanup = ( scopeKey: string, options: { processTree: ProcessScopeCleanupPolicy }, ): (() => Promise) => { const cleanupOwner: ScopeCleanupOwner = { processTree: options.processTree }; const owners = scopeCleanupOwners.get(scopeKey) ?? new Set(); owners.add(cleanupOwner); scopeCleanupOwners.set(scopeKey, owners); let closing: Promise | undefined; return () => (closing ??= (async () => { try { cancelScope(scopeKey); await waitForRuns(scopeKey); } catch (error) { cleanupOwner.failure ??= { error }; } finally { owners.delete(cleanupOwner); if (owners.size === 0) { scopeCleanupOwners.delete(scopeKey); } } if (cleanupOwner.failure) { throw cleanupOwner.failure.error; } })()); }; const startRun = async (input: SpawnInput, owner: OwnedRun): Promise => { const external = input.cleanupOwnership === "external"; const requireProcessTree = owner.cleanupOwners.some((scope) => requiresProcessTree(scope, external), ); // A queued replacement must still own authority before stopping the surviving run. if (!owner.terminationReason) { input.assertCurrent?.(); input.beforeSpawn?.(); // Native PTY has no tree-extinction owner. Reject before spawning so exec's // existing PTY-unavailable fallback can run once under the child anchor. if (input.mode === "pty" && requireProcessTree) { throw new Error("PTY is unavailable when execution requires process-tree cleanup"); } } const { runId, scopeKey } = owner; const startedAtMs = Date.now(); const startingTerminationReason = owner.terminationReason; const settleConstructionResult = ( reason: TerminationReason, cleanup?: Promise, output?: { stdout: string; stderr: string; lastOutputAtMs: number }, ): ManagedRun => { const exit: RunExit = { reason, exitCode: null, exitSignal: null, durationMs: Date.now() - startedAtMs, stdout: output?.stdout ?? "", stderr: output?.stderr ?? "", timedOut: isTimeoutReason(reason), noOutputTimedOut: reason === "no-output-timeout", }; return { runId, startedAtMs, activity: Object.freeze({ resultSettled: true, lastOutputAtMs: output?.lastOutputAtMs ?? startedAtMs, }), wait: async () => exit, ...(cleanup && { waitForExtinction: () => cleanup }), cancel: () => undefined, }; }; if (startingTerminationReason) { // A replacement can be cancelled behind its scope fence. Never launch // its command or terminate the surviving scope after that cancellation. return settleConstructionResult(startingTerminationReason); } // Finish fallible argument preparation before affecting a surviving scope or arming cancellation. if (input.mode !== "anchored-shell" && input.argv.length === 0) { throw new Error("spawn argv cannot be empty"); } const resolvedArgs = input.mode === "child" ? input.resolveArgs?.() : undefined; if (owner.terminationReason) { return settleConstructionResult(owner.terminationReason); } input.assertCurrent?.(); input.beforeSpawn?.(); if (input.replaceExistingScope && scopeKey) { // Scope admission already waited for predecessor startups. Do not // cancel this replacement or later runs reserved behind its fence. cancelActiveScope(scopeKey, "manual-cancel"); } let forcedReason: TerminationReason | null = owner.terminationReason ?? null; let resultSettled = false; let lastOutputAtMs = startedAtMs; let cleanupSettled = false; const captured = { stdout: "", stderr: "" }; // Forced settlement (kill-wait fallback, Windows forced close) resolves the // result while inherited pipes stay open, and callers finalize their own // output state from that terminal result. One fence closes every output path // together: a late chunk reaches no listener, capture buffer, or output clock. let outputDetached = false; const detachOutput = () => { outputDetached = true; }; let forceKillTimer: NodeJS.Timeout | null = null; let cancelRequested = false; const captureOutput = input.captureOutput !== false; const maxCapturedOutputChars = clampCapturedOutputChars(input.maxCapturedOutputChars); const setForcedReason = (reason: TerminationReason) => { if (forcedReason || resultSettled) { return; } forcedReason = reason; }; let cancelAdapter: ((reason: TerminationReason) => void) | null = null; const constructionAbort = new AbortController(); const constructionAbortError = new Error("adapter construction aborted"); const constructionAbortPromise = new Promise((_, reject) => { const rejectConstruction = () => reject(constructionAbortError); if (constructionAbort.signal.aborted) { rejectConstruction(); } else { constructionAbort.signal.addEventListener("abort", rejectConstruction, { once: true }); } }); const requestCancel = (reason: TerminationReason) => { setForcedReason(reason); input.onCancel?.(reason); cancelAdapter?.(reason); // Any cancel must abort construction: the relay may already be spawned // and waiting for ready, and a later deadline must not replace this reason. if (!cancelAdapter) { constructionAbort.abort(); } }; owner.cancel = requestCancel; const createDeadline = (reason: "overall-timeout" | "no-output-timeout", value?: number) => { const durationMs = normalizeTimeoutDuration(value); let deadlineMs: number | null = null; let timer: NodeJS.Timeout | undefined; // Re-arm bounded intervals: a long deadline must not overflow Node's timer cap. const schedule = (remainingMs: number, deadline: number) => { const intervalMs = resolveTimerTimeoutMs(remainingMs, 1); timer = setTimeout(() => { if (resultSettled) { return; } const remaining = Math.min(remainingMs - intervalMs, deadline - performance.now()); if (remaining <= 0) { requestCancel(reason); } else { schedule(remaining, deadline); } }, intervalMs); }; return { get deadlineMs() { return deadlineMs; }, reset: () => { if (!durationMs || resultSettled) { return; } clearTimeout(timer); deadlineMs = performance.now() + durationMs; schedule(durationMs, deadlineMs); }, clear: () => clearTimeout(timer), }; }; const overallDeadline = createDeadline("overall-timeout", input.timeoutMs); const outputDeadline = createDeadline("no-output-timeout", input.noOutputTimeoutMs); const touchOutput = () => { lastOutputAtMs = Date.now(); outputDeadline.reset(); }; try { // Reserve the join before construction: a timeout result does not release // resources acquired later, or hide cleanup when readiness rejects after spawn. const cleanup = createDeferredCore(); owner.waitForExtinction = () => cleanup.promise; void cleanup.promise.catch(() => undefined); let constructionCleanup: Promise | undefined; let ownedAdapter: SpawnProcessAdapter | undefined; const onSpawnCleanup = (promise: Promise) => { constructionCleanup = promise; void promise.catch(() => undefined); }; overallDeadline.reset(); outputDeadline.reset(); const startupPromise = input.mode === "pty" ? createPtyAdapter({ assertCurrent: input.assertCurrent, beforeSpawn: input.beforeSpawn, shell: expectDefined(input.argv[0], "spawn executable"), args: input.argv.slice(1), cwd: input.cwd, env: input.env, abortSignal: constructionAbort.signal, onSpawnCleanup, }).then((adapter) => ({ adapter, ready: Promise.resolve() })) : input.mode === "anchored-shell" ? createChildAdapter({ assertCurrent: input.assertCurrent, beforeSpawn: input.beforeSpawn, anchoredShellCommand: input.command, cwd: input.cwd, env: input.env, abortSignal: constructionAbort.signal, onSpawnCleanup, }) : createChildAdapter({ assertCurrent: input.assertCurrent, beforeSpawn: input.beforeSpawn, ...(requireProcessTree && !external ? { ownProcessTree: true as const } : {}), argv: resolvedArgs ? [...input.argv, ...resolvedArgs] : input.argv, argv0: input.argv0, cwd: input.cwd, env: input.env, exactEnv: input.exactEnv, windowsVerbatimArguments: input.windowsVerbatimArguments, input: input.input, stdinMode: input.stdinMode, secretInput: input.secretInput, abortSignal: constructionAbort.signal, onSpawnCleanup, }); const extinctionPromise = startupPromise .then( async ({ adapter: started, ready }) => { ownedAdapter = started; if (external || !started.waitForExtinction) { for (const scope of owner.cleanupOwners) { if (requiresProcessTree(scope, external)) { scope.failure ??= { error: new Error( "process cleanup cannot confirm owned execution-tree settlement", ), }; } } } if (constructionAbort.signal.aborted) { started.kill("SIGKILL"); // Drain a late adapter's output without reopening the terminal result. void started.wait().catch(() => undefined); } // Child close can precede a descendant's private-input consumption. // Readiness failure is separate from the cleanup owner's outcome. await Promise.allSettled([ready]); await (constructionCleanup ?? started.waitForExtinction?.() ?? started.wait()); }, async () => { await constructionCleanup; }, ) .finally(() => { cleanupSettled = true; if (forceKillTimer) { clearTimeout(forceKillTimer); forceKillTimer = null; } if (resultSettled) { ownedAdapter?.dispose(); } }); void extinctionPromise.then( () => { ownedRuns.delete(owner); cleanup.resolve(); }, (error: unknown) => { recordScopeCleanupFailure(owner, error); cleanupFailure ??= { error }; ownedRuns.delete(owner); cleanup.reject(error); }, ); const settleAbortedConstruction = (reason: TerminationReason) => { resultSettled = true; overallDeadline.clear(); outputDeadline.clear(); detachOutput(); if (cleanupSettled) { ownedAdapter?.dispose(); } return settleConstructionResult(reason, cleanup.promise, { ...captured, lastOutputAtMs }); }; let startup: Awaited; try { startup = await Promise.race([startupPromise, constructionAbortPromise]); } catch (err) { if (err !== constructionAbortError || !forcedReason) { throw err; } return settleAbortedConstruction(forcedReason); } const adapter = startup.adapter; const settleResult = () => { resultSettled = true; overallDeadline.clear(); outputDeadline.clear(); detachOutput(); if (cleanupSettled) { adapter.dispose(); } }; const withOutputFence = (deliver?: (chunk: Chunk) => void, recordsOutput = true) => (chunk: Chunk) => { if (outputDetached) { return; } if (recordsOutput) { touchOutput(); } deliver?.(chunk); }; const rawInput = input.mode === "child" ? input : undefined; // Byte transports can flush decoded text at EOF without fresh activity. // PTYs and Windows Job transports report only text. for (const [stream, subscribe, onText, onRaw] of [ ["stdout", adapter.onStdout, input.onStdout, rawInput?.onStdoutRaw], ["stderr", adapter.onStderr, input.onStderr, rawInput?.onStderrRaw], ] as const) { subscribe( withOutputFence((chunk: string) => { if (captureOutput) { captured[stream] = appendCapturedOutput( captured[stream], chunk, stream, maxCapturedOutputChars, ); } onText?.(chunk); }, !adapter.supportsRawOutput), withOutputFence(onRaw), ); } try { await Promise.race([startup.ready, constructionAbortPromise]); } catch (error) { if (error === constructionAbortError && forcedReason) { return settleAbortedConstruction(forcedReason); } settleResult(); throw error; } cancelAdapter = (reason: TerminationReason) => { if ( cleanupSettled || (cancelRequested && (requireProcessTree || !(resultSettled && forceKillTimer))) ) { return; } cancelRequested = true; if (resultSettled && !requireProcessTree) { if (forceKillTimer) { clearTimeout(forceKillTimer); forceKillTimer = null; } // Root completion closes its terminal record, not ownership of // descendants still retained by the authoritative group or Job. adapter.kill("SIGKILL"); return; } // Windows has no catchable SIGTERM equivalent: the adapter implements it // with asynchronous taskkill, so waiting the cleanup grace only delays an // already-expired deadline before the same forced tree termination. if ( process.platform === "win32" && (reason === "overall-timeout" || reason === "no-output-timeout") ) { adapter.kill("SIGKILL"); return; } adapter.kill("SIGTERM"); forceKillTimer = setTimeout(() => { if (!cleanupSettled) { adapter.kill("SIGKILL"); } }, GRACEFUL_CANCEL_TIMEOUT_MS); forceKillTimer.unref?.(); }; const waitPromise = (async (): Promise => { const result = await adapter.wait(); const deadlineReason = resolveElapsedTimeoutReason({ nowMs: performance.now(), overallTimeoutDeadlineMs: overallDeadline.deadlineMs, noOutputTimeoutDeadlineMs: outputDeadline.deadlineMs, }); const terminalReason = forcedReason ?? deadlineReason; settleResult(); const reason: TerminationReason = terminalReason ?? (result.signal != null ? ("signal" as const) : ("exit" as const)); const exit: RunExit = { reason, exitCode: result.code, exitSignal: result.signal, oomScoreWrapperSelected: adapter.oomScoreWrapperSelected === true, durationMs: Date.now() - startedAtMs, ...captured, timedOut: isTimeoutReason(reason), noOutputTimedOut: terminalReason === "no-output-timeout", }; return exit; })().catch((err: unknown) => { if (!resultSettled) { settleResult(); } throw err; }); const managedRun: ManagedRun = { activity: Object.freeze({ get deadlineAtMs() { return overallDeadline.deadlineMs === null ? undefined : Date.now() + overallDeadline.deadlineMs - performance.now(); }, get resultSettled() { return resultSettled; }, get lastOutputAtMs() { return lastOutputAtMs; }, }), runId, pid: adapter.pid, startedAtMs, stdin: adapter.stdin, wait: async () => await waitPromise, ...(adapter.waitForExtinction && { waitForExtinction: () => cleanup.promise }), cancel: (reason = "manual-cancel") => { requestCancel(reason); }, detachOutput, }; if (forcedReason) { managedRun.cancel(forcedReason); } return managedRun; } catch (err) { resultSettled = true; overallDeadline.clear(); outputDeadline.clear(); detachOutput(); const { warnProcessSupervisorSpawnFailure } = await loadSupervisorLogRuntime(); warnProcessSupervisorSpawnFailure(`spawn failed: runId=${runId} reason=${String(err)}`); throw err; } }; const spawn = (input: SpawnInput): Promise => { if (shuttingDown) { return Promise.reject(new Error("process supervisor is shut down")); } const scopeKey = normalizeOptionalString(input.scopeKey); const runId = normalizeOptionalString(input.runId) ?? crypto.randomUUID(); const owner: OwnedRun = { runId, scopeKey, cancel: (reason) => { owner.terminationReason ??= reason; input.onCancel?.(reason); }, cleanupOwners: scopeKey ? [...(scopeCleanupOwners.get(scopeKey) ?? [])] : [], }; // Reserve cancellation before either adapter startup or a replacement // fence, so stopping a run cannot silently leave a late child alive. ownedRuns.add(owner); const starting = scopeKey ? (startingScopes.get(scopeKey) ?? { runs: new Set>() }) : undefined; if (scopeKey && starting) { startingScopes.set(scopeKey, starting); } // Ordinary runs start together, but replacements fence later arrivals so // delayed cancellation cannot accidentally terminate a newer scoped run. const previous = starting ? input.replaceExistingScope ? Array.from(starting.runs) : starting.replacement ? [starting.replacement] : [] : []; const pending = previous.length > 0 ? Promise.allSettled(previous).then(() => startRun(input, owner)) : startRun(input, owner); owner.pending = pending; starting?.runs.add(pending); if (starting && input.replaceExistingScope) { starting.replacement = pending; } const clearPendingStart = () => { delete owner.pending; if (!owner.waitForExtinction) { ownedRuns.delete(owner); } starting?.runs.delete(pending); if (starting?.replacement === pending) { delete starting.replacement; } if (scopeKey && starting?.runs.size === 0 && startingScopes.get(scopeKey) === starting) { startingScopes.delete(scopeKey); } }; void pending.then(clearPendingStart, clearPendingStart); return pending; }; const shutdown = (): Promise => { // Publish the admission fence before cancellation can invoke owner callbacks. shuttingDown = true; return (shutdownPromise ??= Promise.resolve().then(async () => { while (ownedRuns.size) { for (const owner of ownedRuns) { owner.cancel("manual-cancel"); } // A failed startup owns no live process; only failed owner extinction // must keep the process-wide supervisor fenced for operator recovery. await waitForRuns(null, true); } if (cleanupFailure) { throw cleanupFailure.error; } })); }; return { acquireScopeCleanup, spawn, cancel, cancelScope, shutdown, }; }