From f96acad43829e6e7461418fa738265cdd8398396 Mon Sep 17 00:00:00 2001 From: pku-xht Date: Tue, 4 Aug 2026 19:55:39 +0800 Subject: [PATCH] fix(subagent): preserve Codex fatal and grace semantics --- packages/subagent/subagent-codex/src/run.ts | 50 ++++++++++++++- packages/subagent/subagent-codex/src/wire.ts | 26 ++++---- .../tests/subagent-codex.spec.ts | 48 +++++++++++++- .../subprocess/subprocess-local/src/spawn.ts | 63 +++++++++++++++++-- .../subprocess-local/tests/spawn.spec.ts | 43 ++++++++++++- 5 files changed, 210 insertions(+), 20 deletions(-) diff --git a/packages/subagent/subagent-codex/src/run.ts b/packages/subagent/subagent-codex/src/run.ts index 21f22d8a1b..811f7c8f98 100644 --- a/packages/subagent/subagent-codex/src/run.ts +++ b/packages/subagent/subagent-codex/src/run.ts @@ -24,6 +24,47 @@ import { CodexAppServerWire } from './wire.ts' /** Default POSIX grace between subprocess termination tiers. */ export const DEFAULT_DISPOSE_GRACE_MS = 3_000 +/** Largest delay Node schedules without collapsing it to one millisecond. */ +const MAX_TIMER_DELAY_MS = 2_147_483_647n + +/** + * Bound final exit observation at twice a positive finite grace without + * narrowing the public config to Node's single-timer integer range. + */ +function doubledGraceWindow(graceMs: number): { + readonly signal: AbortSignal + readonly cancel: () => void +} { + const whole = Math.floor(graceMs) + let remaining = BigInt(whole) * 2n + + BigInt(Math.ceil((graceMs - whole) * 2)) + const controller = new AbortController() + let timer: ReturnType | undefined + const arm = (): void => { + const chunk = remaining > MAX_TIMER_DELAY_MS + ? MAX_TIMER_DELAY_MS + : remaining + remaining -= chunk + timer = setTimeout(() => { + timer = undefined + if (remaining === 0n) { + controller.abort() + } else { + arm() + } + }, Number(chunk)) + } + arm() + return { + signal: controller.signal, + cancel: () => { + if (timer === undefined) return + clearTimeout(timer) + timer = undefined + }, + } +} + /** Fully resolved inputs for one Codex app-server run. */ export interface CodexRunSpec { /** Parent Session workspace, also supplied to `thread/start`. */ @@ -88,8 +129,13 @@ export async function disposeCodexChild( // A concurrently closed stdin does not change tree ownership below. } child.terminate() - if (!(await child.waitForExit(AbortSignal.timeout(graceMs * 2)))) { - throw new Error('subagent-codex: app-server process tree did not exit within its dispose window') + const exitWindow = doubledGraceWindow(graceMs) + try { + if (!(await child.waitForExit(exitWindow.signal))) { + throw new Error('subagent-codex: app-server process tree did not exit within its dispose window') + } + } finally { + exitWindow.cancel() } await child.done } diff --git a/packages/subagent/subagent-codex/src/wire.ts b/packages/subagent/subagent-codex/src/wire.ts index 4e920113e3..304c5eadb4 100644 --- a/packages/subagent/subagent-codex/src/wire.ts +++ b/packages/subagent/subagent-codex/src/wire.ts @@ -17,12 +17,17 @@ type JsonObject = Record interface Deferred { readonly promise: Promise readonly resolve: (value: T) => void + readonly reject: (reason?: unknown) => void } function deferred(): Deferred { let resolve!: (value: T) => void - const promise = new Promise((settle) => { resolve = settle }) - return { promise, resolve } + let reject!: (reason?: unknown) => void + const promise = new Promise((settle, fail) => { + resolve = settle + reject = fail + }) + return { promise, resolve, reject } } function object(value: unknown, label: string): JsonObject { @@ -93,7 +98,7 @@ async function raceAbort(pending: Promise, signal: AbortSignal): Promise() + private readonly fatal = deferred() private threadId: string | undefined private turnId: string | undefined private pendingTurnId: string | undefined @@ -111,6 +116,10 @@ export class CodexAppServerWire { output: Writable, ) { this.transport = new JsonRpcLineTransport(input, output) + // Fatal protocol state can arrive after the current guarded operation has + // already settled. Keep the shared rejection observed without inserting + // another promise-adoption hop into active races. + void this.fatal.promise.catch(() => {}) this.transport.onRequest((method, params) => this.handleServerRequest(method, params)) this.transport.onNotification((method, params) => { try { @@ -157,9 +166,8 @@ export class CodexAppServerWire { * Create the run's private ephemeral thread and retain its identity. * @param cwd - parent Session workspace. * @param signal - unpublished-start cancellation. - * @returns the app-server thread id. */ - async startThread(cwd: string, signal: AbortSignal): Promise { + async startThread(cwd: string, signal: AbortSignal): Promise { const response = object(await this.guarded(this.transport.request('thread/start', { cwd, ephemeral: true, @@ -170,7 +178,6 @@ export class CodexAppServerWire { throw new Error('subagent-codex: app-server did not create an ephemeral thread') } this.threadId = id - return id } /** @@ -249,15 +256,12 @@ export class CodexAppServerWire { } private async guarded(pending: Promise, signal: AbortSignal): Promise { - const withFatal = Promise.race([ - pending, - this.fatal.promise.then((error): Promise => Promise.reject(error)), - ]) + const withFatal = Promise.race([pending, this.fatal.promise]) return raceAbort(withFatal, signal) } private fail(error: Error): void { - this.fatal.resolve(error) + this.fatal.reject(error) } private readonly onInputError = (error: Error): void => { diff --git a/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts b/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts index 6cc4461e06..18e28cc7cc 100644 --- a/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts +++ b/packages/subagent/subagent-codex/tests/subagent-codex.spec.ts @@ -210,7 +210,7 @@ async function initializeWire(): Promise<{ const starting = wire.startThread(process.cwd(), new AbortController().signal) const threadStart = await child.peer.nextMethod('thread/start') child.peer.respond(threadStart, { thread: { id: 'thread-1', ephemeral: true } }) - await expect(starting).resolves.toBe('thread-1') + await starting return { child, wire } } @@ -516,6 +516,21 @@ describe('CodexAppServerWire', () => { } }) + it('keeps an earlier fatal frame authoritative over later completion in the same chunk', async () => { + const { child, wire } = await initializeWire() + const result = wire.runTurn(['task'], new AbortController().signal, () => false) + const turnStart = await child.peer.nextMethod('turn/start') + child.peer.respond(turnStart, { turn: { id: 'turn-1' } }) + await nextTask() + child.peer.send( + agentMessage('invalid', 'future_phase'), + agentMessage('late answer', 'final_answer'), + turnCompleted('completed'), + ) + await expect(result).rejects.toThrow('unknown agent message phase') + wire.close() + }) + it('gives local cancellation precedence over a remote completed turn', async () => { const { child, wire } = await initializeWire() let cancelled = false @@ -1038,6 +1053,37 @@ describe('disposeCodexChild', () => { expect(child.waitForExit).toHaveBeenCalledTimes(1) }) + it('accepts fractional and larger-than-Node grace windows', async () => { + for (const graceMs of [0.25, Number.MAX_VALUE]) { + const child = fakeChild() + const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!) + await expect(disposeCodexChild(wire, child.handle, graceMs)) + .resolves.toBeUndefined() + const signal = vi.mocked(child.waitForExit).mock.calls[0]?.[0] + expect(signal?.aborted).toBe(false) + } + }) + + it('chains a doubled grace window beyond one Node timer segment', async () => { + vi.useFakeTimers() + try { + const child = fakeChild({ exitOnTerminate: false }) + const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!) + const disposal = disposeCodexChild( + wire, + child.handle, + 1_073_741_823.75, + ) + const rejected = expect(disposal) + .rejects.toThrow('did not exit within its dispose window') + await vi.advanceTimersByTimeAsync(2_147_483_647) + await vi.advanceTimersByTimeAsync(1) + await rejected + } finally { + vi.useRealTimers() + } + }) + it('contains a concurrently closed stdin error', async () => { const child = fakeChild() const wire = new CodexAppServerWire(child.handle.stdout!, child.handle.stdin!) diff --git a/packages/subprocess/subprocess-local/src/spawn.ts b/packages/subprocess/subprocess-local/src/spawn.ts index 90d460c2c5..d3cbb0cf55 100644 --- a/packages/subprocess/subprocess-local/src/spawn.ts +++ b/packages/subprocess/subprocess-local/src/spawn.ts @@ -55,6 +55,47 @@ function sleepTick(): Promise { return sleepMs(15) } +/** Largest delay Node schedules without collapsing it to one millisecond. */ +const MAX_TIMER_DELAY_MS = 2_147_483_647n + +/** + * Schedule a positive finite millisecond delay across as many Node-safe timer + * segments as necessary. Fractional milliseconds round up so a grace never + * expires earlier than configured. + * @param delayMs - positive finite delay in milliseconds. + * @param callback - work to run after the complete delay. + * @returns a handle that cancels the active segment and all future segments. + */ +export function scheduleFiniteTimeout( + delayMs: number, + callback: () => void, +): { cancel(): void } { + let remaining = BigInt(Math.ceil(delayMs)) + let timer: ReturnType | undefined + const arm = (): void => { + const chunk = remaining > MAX_TIMER_DELAY_MS + ? MAX_TIMER_DELAY_MS + : remaining + remaining -= chunk + timer = setTimeout(() => { + timer = undefined + if (remaining === 0n) { + callback() + } else { + arm() + } + }, Number(chunk)) + } + arm() + return { + cancel(): void { + if (timer === undefined) return + clearTimeout(timer) + timer = undefined + }, + } +} + let spillCounter = 0 let defaultSpillDir: string | undefined @@ -341,7 +382,7 @@ export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInter const stdoutCollector = collectStream(outMode, child.stdout, 'stdout') const stderrCollector = collectStream(errMode, child.stderr, 'stderr') - let graceTimer: NodeJS.Timeout | undefined + let graceTimer: ReturnType | undefined let settled = false // Failed spawns use pid -1 so signalling remains a no-op. @@ -377,6 +418,8 @@ export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInter // child and must stay signalable, while a fully-dead tree (possible pid // reuse) must not be re-signalled by a later tier. const kill = (sig: NodeJS.Signals): void => { + /* v8 ignore next -- the exit monitor cancels the ordinary dead-tree timer; + this remains the timer/death race guard and cannot be staged deterministically. */ if (!treeAlive()) return signalTree(platform, pid, sig, child, taskkill) } @@ -390,7 +433,15 @@ export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInter // kill() re-probes tree liveness before force-killing. It stays ref'd: // the pending SIGKILL is a commitment, and a parent exiting before it // fires would orphan a trapped survivor. Self-bounds at graceMs. - graceTimer = setTimeout(() => { kill('SIGKILL') }, spec.graceMs) + const timer = scheduleFiniteTimeout(spec.graceMs, () => { kill('SIGKILL') }) + graceTimer = timer + // A very large configured grace must not pin the parent after TERM already + // removed the whole tree. Keep the escalation armed only while its target + // remains alive; direct-child settlement alone is not sufficient. + void waitForExit().then(() => { + timer.cancel() + graceTimer = undefined + }) } // The caller owns timeout classification; this layer only reacts to abort. @@ -405,7 +456,7 @@ export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInter } const done = new Promise((resolve, reject) => { - let pipeDrainTimer: NodeJS.Timeout | undefined + let pipeDrainTimer: ReturnType | undefined const settle = (exitCode: number | null, signal: NodeJS.Signals | null): void => { if (settled) return settled = true @@ -428,13 +479,15 @@ export function spawnSubprocess(spec: SubprocessSpawnSpec, internals: SpawnInter // A surviving descendant that inherited a pipe must not hold the // outcome open indefinitely: after exit, the same bounded grace that // governs kills also bounds the close wait. - pipeDrainTimer = setTimeout(() => { settle(exitCode, signal) }, spec.graceMs) + pipeDrainTimer = scheduleFiniteTimeout(spec.graceMs, () => { + settle(exitCode, signal) + }) }) child.on('close', settle) function cleanup(): void { // graceTimer deliberately NOT cleared: the SIGKILL escalation must be // able to reach tree survivors after the direct child settles. - if (pipeDrainTimer !== undefined) clearTimeout(pipeDrainTimer) + pipeDrainTimer?.cancel() spec.signal?.removeEventListener('abort', onAbort) } }) diff --git a/packages/subprocess/subprocess-local/tests/spawn.spec.ts b/packages/subprocess/subprocess-local/tests/spawn.spec.ts index 491756f01f..87c81116ff 100644 --- a/packages/subprocess/subprocess-local/tests/spawn.spec.ts +++ b/packages/subprocess/subprocess-local/tests/spawn.spec.ts @@ -2,7 +2,13 @@ import { mkdtempSync, readFileSync, statSync, unlinkSync } from 'node:fs' import { tmpdir } from 'node:os' import { dirname, join } from 'node:path' import { describe, expect, it, vi } from 'vitest' -import { killGroup, OutputCollector, spawnSubprocess, taskkillProcessTree } from '../src/spawn.ts' +import { + killGroup, + OutputCollector, + scheduleFiniteTimeout, + spawnSubprocess, + taskkillProcessTree, +} from '../src/spawn.ts' import type { SubprocessHandle, SubprocessOutputReader } from '@deepseek-ai/dsh-subprocess' const { failNextClose, failNextUnlink } = vi.hoisted(() => ({ @@ -101,6 +107,30 @@ async function waitForPidFile(path: string, timeoutMs = 5_000): Promise throw new Error(`pid file ${path} was not written after ${timeoutMs}ms`) } +describe('scheduleFiniteTimeout', () => { + it('rounds fractions up, chains Node-safe segments, and cancels idempotently', async () => { + vi.useFakeTimers() + try { + const fired = vi.fn() + const chained = scheduleFiniteTimeout(2_147_483_647.25, fired) + await vi.advanceTimersByTimeAsync(2_147_483_647) + expect(fired).not.toHaveBeenCalled() + await vi.advanceTimersByTimeAsync(1) + expect(fired).toHaveBeenCalledOnce() + chained.cancel() + + const cancelled = vi.fn() + const timer = scheduleFiniteTimeout(0.25, cancelled) + timer.cancel() + timer.cancel() + await vi.advanceTimersByTimeAsync(1) + expect(cancelled).not.toHaveBeenCalled() + } finally { + vi.useRealTimers() + } + }) +}) + describe('spawnSubprocess', () => { it('captures stdout on success', async () => { const result = await finish(spawnSubprocess(spec('echo hello'))) @@ -164,6 +194,17 @@ describe('spawnSubprocess', () => { expect(result.signal).toBe('SIGKILL') }) + it('cancels a larger-than-Node escalation timer once SIGTERM removes the tree', async () => { + const running = spawnSubprocess(spec('echo ready; sleep 60', { + graceMs: Number.MAX_VALUE, + })) + await waitForStdout(running, 'ready\n') + running.terminate() + const result = await running.done + expect(result.signal).toBe('SIGTERM') + await expect(running.waitForExit()).resolves.toBe(true) + }) + it('terminates the whole process group (grandchildren die too)', async () => { // The subshell writes the sleep's pid then waits on it; terminating the // group must take the sleep down with bash.