fix(subagent): preserve Codex fatal and grace semantics

This commit is contained in:
pku-xht
2026-08-04 19:55:39 +08:00
parent c09b20f96b
commit f96acad438
5 changed files with 210 additions and 20 deletions
+48 -2
View File
@@ -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<typeof setTimeout> | 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
}
+15 -11
View File
@@ -17,12 +17,17 @@ type JsonObject = Record<string, unknown>
interface Deferred<T> {
readonly promise: Promise<T>
readonly resolve: (value: T) => void
readonly reject: (reason?: unknown) => void
}
function deferred<T>(): Deferred<T> {
let resolve!: (value: T) => void
const promise = new Promise<T>((settle) => { resolve = settle })
return { promise, resolve }
let reject!: (reason?: unknown) => void
const promise = new Promise<T>((settle, fail) => {
resolve = settle
reject = fail
})
return { promise, resolve, reject }
}
function object(value: unknown, label: string): JsonObject {
@@ -93,7 +98,7 @@ async function raceAbort<T>(pending: Promise<T>, signal: AbortSignal): Promise<T
*/
export class CodexAppServerWire {
private readonly transport: JsonRpcLineTransport
private readonly fatal = deferred<Error>()
private readonly fatal = deferred<never>()
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<string> {
async startThread(cwd: string, signal: AbortSignal): Promise<void> {
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<T>(pending: Promise<T>, signal: AbortSignal): Promise<T> {
const withFatal = Promise.race([
pending,
this.fatal.promise.then((error): Promise<never> => 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 => {