mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 08:02:21 +00:00
Share the managed process lifecycle for structured providers (#25204)
* fix(native-chat): a child's root exit is reported even during its close, and bookkeeping after it never reads as unproven - Both connections report the root process's exit once, with `expected` set when a close had begun. A close that came back unproven and whose root exits later is finished by the adapter, and its end reaches the host like any other. - A Claude close whose resume-point write fails after the exit was proven, and a Codex close whose terminal row is refused, now end the session and report the failure, instead of keeping a dead child indexed as if its exit were unproven. - A Codex close whose forced tree kill can't prove the descendants gone but saw the root exit reports the descendants and counts the root exit. - Every child exit with an identity, expected or not, is forwarded to the host. * fix(native-chat): the exit ends the child's record; an unfinished stop is the child's own close, which everyone joins - The host keeps no stored "stop still owed" record any more. A stop begins the child's close (`child.close`), which lives on the child and ends with it. A second Stop, the idle reaper, quit, a send and an option/answer/goal/rewind all join that close instead of retrying a separate obligation. - A caller waits on the close only as long as the step deadline; the close itself is never abandoned. A proof that lands after every caller stopped waiting reaches the host as the adapter's report of that exit, which ends the record through the same handler. - Once the exit is proven, draining, settling, the lease release and the adapter's acknowledgement are each attempted and reported on failure; none keeps the child on record. A start, and the handle's close, write a release that failed from this host's proof of that exit, so a failed write never refuses a send. - A start that meets a close still unverifiable is refused with `previousExitUnverifiable`, so the queued message is rejected with a send-again reason; nothing is held and nothing starts beside the old process. - The idle sweep goes back to idle reaping only. - Removes #24333's retry entry points, the wait row and its hold rule, the ask/failure cursors on the stored record, and the stop's own wake. Tests replace the #24333 unproven-stop test: a send joining an unproven close and an in-flight one, a late proof past the caller's bound, a root exiting after its close gave up, a proven exit whose resume-point write and lease release both failed, an unverifiable close rejecting the send and refusing an option change, a surviving descendant, quit and the idle reaper; and Codex's unverifiable, late-exit and joined-close cases. * fix(native-chat): a message refused because the old process's exit is unverifiable says so, and to send again The start failure for a refusal with reason `previousExitUnverifiable` reads "Orca couldn't confirm Claude's previous process ended. Send your message to try again." instead of "Claude couldn't restart." The status-row kind and the refusal reason stay in the shared lists for rows and hosts that still carry them; the catalogs keep one sentence for both. * fix(native-chat): a close's verdict is the root's exit alone, and what follows it is logged - A Claude close resolves as soon as the root's exit is proven: the session ends and its `ended` report goes out then. Saving the resume point runs afterwards and a failure is logged, so a slow or hung write never reads as an unproven exit or keeps a dead child on record. - A root that exits after its close came back unproven finishes that close through the same path as any close, so the session's child work is published as ended (background tasks and subagents no longer stay shown running for a dead agent), and a failure there is logged. - Codex logs a refused final row, and reports a root exit whose forced tree kill could not prove the rest of the tree gone the way Claude does, so the host logs it and blocks nothing. - Both adapters take the host's logger for this bookkeeping. * fix(native-chat): one handler ends every child's exit, and a join waits on the adapter's own close - One exit handler (`structured-agent-session-child-exit`) ends a child's record for an exit expected or not. `expected` only changes what the chat is told: the stop's cause, its end at the stop's ask, the settlement id, and no crash outcome row. The lease release keeps the exit's evidence; the handoff guard, lifecycle barrier, sink release and adapter acknowledgement apply to both. A Claude journal-sink failure ends in the same step as its stop, as Orca's own fault. - Joining a close is asking the adapter, whose close is memoized while it runs and bounded by its own kill escalation; the host keeps no attempt of its own and no 10 s caller bound. An ask after a close came back unproven runs the stop again. - A close's end is stamped where its stop was asked for (a repeated ask moves it), so the closed chat and failed start checks order a message accepted meanwhile after it. - A start refused because the old exit is unverifiable rejects what was queued in the same step. - The end of a close the host asked for no longer waits on the cross-session recovery chain. - The kill no longer waits for the stop event's write; the journal writes rows in order. * fix(native-chat): an exit's lease release lands whatever the length of its reason A crash's reason can carry kilobytes of the provider's stderr, and a lease whose death detail is over 512 characters fails the store's own check. The exit handler cut it, but the release a start or the chat handle's close re-derives did not, so after a crash whose own release failed every message was refused as not resumable until restart. The record's builder now cuts the detail to the record's bound, so no writer can hand it one too long. * fix(claude): a proven close waits at most 2 s for the output it already wrote Once the root's exit is proven, the close still waited for the SDK's output reader to end. Something outside the process tree that holds the output open would keep that close, and every send, Stop and quit joining it, waiting with no bound. The wait is now bounded; past it the close resolves as proven and the open output is logged. * fix(codex): an exit reported inside Orca's close keeps the reason Orca closed it for The connection reports the app-server's exit inside the close that ends it, so that report ended every Codex close and replaced the close's own reason (for example, a provider frame that could not be recorded) with the connection's stderr text in the ended record and the lease's exit evidence. The session now records Orca's close with its reason, and the exit it ends keeps that reason. The test connection reports its exit inside close the way the real one does. * fix(native-chat): quit stops delivery before it drains exit recovery Every exit now wakes delivery, and teardown drained exit recovery before it stopped delivery, so an exit settled in that window could start a fresh agent that teardown then killed. Teardown stops delivery first; queued messages wait for the next launch. * docs(native-chat): the unverifiable-exit refusal no longer names a caller's wait The caller's bounded wait was removed; the comment describes the close as it is now. * fix(native-chat): a stop whose kill did not take is logged, and the next ask kills again When a close's kill leaves the agent's root running, the host now logs it. Tests pin what a later ask does: each connection runs its whole stop again (Codex sends SIGKILL a second time), refuses input meanwhile, and proves the exit once the kill takes. * fix(native-chat): a start refused over the old process says Orca couldn't stop it The host reaches an unverifiable verdict only after its own kill left the agent's root running, on the machine that runs the agent, so the sentence now says that: "Orca couldn't stop {agent}'s previous process." The refusal reason, failure kind and wire shapes are unchanged. The host test also checks the failed kill is logged. * docs(native-chat): an unverifiable close verdict is a root that survived the kill The host's close runs where the agent runs, so lost contact never yields this verdict; the comment no longer says it does. * fix(native-chat): a kill that did not take is reported once, by whoever met it The log added at the close fired beside a Stop's own failure report for the same event. A stop still reports it through its failure; a send or option change refused over it now logs it at the refusal, the only place it is otherwise invisible. * test(native-chat): a second Stop joins a close the first could not prove and retries its kill * fix(native-chat): say a start refused beside an unstopped process plainly The rejection now reads "Couldn't stop {{agent}} from before. Send your message again to try once more." This kind has its own send-again step; every other failure keeps "Send your message to try again." * Move provider process supervision and stream reading out of Codex * Preserve teardown behavior with checked mock types after move * Apply provider launch environment and caller teardown labels * Share managed provider launch, exit observation, and close * Preserve synchronous stdin closure before the grace wait * Keep the stacked Claude adapter within the module size limit * Handle nullable provider stdin and processless fixture identity * Preserve close-time cleanup diagnostics after a late provider exit * Separate provider root and descendant exit observations * Give provider child env one owner and gate Codex contract on the shared reader resolveProviderChildEnv is now the only place that overlays and strips a provider's environment; the spawn spec and the request-scoped Codex session both call it. supervisedPosixLaunch only accepts a launch without env fields, so an override can no longer be silently ignored there. Edits to the shared stream reader or the env rule now run the real-binary Codex contract job. * Give every provider one close result, stderr tail and root-only default - A close reports the root and the descendants with the one verdict vocabulary (live / unverifiable / exited); `tree` is null when the close made no descendant observation, and the fallback teardown's outcome is written into it instead of a side flag a caller could miss. - The managed process drains stderr and keeps the 8 KiB tail, so no provider can forget to drain the pipe. - Root-only providers take the default close policy and completion rule; only Claude overrides them. The one already-exited guard lives in the managed close and is recorded when no earlier close ran. - A failed spawn is never read as an observed root exit. * Report only observed descendant exits from the fallback teardown The fallback teardown answered "accepted", and the close turned that into `tree: 'exited'`, though a Windows tree kill's outcome is never read and an unreadable process table observes nothing. It now returns what it observed: `exited` only when the captured descendants were verified gone (or the kernel reported the group empty), `live` / `unverifiable` as verification found them, and null when it signalled without observing. Codex's process-tree diagnostic fires in exactly the cases it did before. * Give two test fixtures' casts a SAFETY rationale for the changed-code gate * Claim no observation from an empty process group ESRCH from the dedicated-group signal says only that the group is empty: a descendant that left it, or a root that never led it, may still be running. It now reports no observation instead of `exited`. The unreadable-table comment says why that case stays "no observation" for root-only providers. * fix(native-chat): a person's close joining a failed one still binds the turn its child end cuts On main every close of the chat writes its own Stop and settle. Here a later close joins the first and writes no row, and the first's settle closed when its kill failed, so a turn that opened in between and was cut by the next close read as failed. A person's close joining a person's close whose Stop opened a settle now reopens that settle until its attempt is done. Tests: a close whose kill failed still closes its settle; a turn opened between a failed close and the next reads as the person's cancellation (each fails without its half of the fix). * test(native-chat): build the Stop-opened-turn test's identity with the opaque handle The test (#25056) landed before the opaque provider handle (#24991), so main still built the old {kind, threadId} handle. * test(ratchet): require src/main/provider-process now that it has landed * test(native-chat): keep main's opaque-handle import in the Stop-opened-turn test Main's #25706 made the same fix as this branch at a different line; the merge kept both imports.
This commit is contained in:
@@ -11,6 +11,7 @@ import {
|
||||
proveClaudeChildExit,
|
||||
type ClaudeChildTreeReaper
|
||||
} from './claude-agent-sdk-exit-proof'
|
||||
import { managedChild } from './claude-child-exit-proof-fixture'
|
||||
import { GRACEFUL_EXIT_MS } from './claude-child-exit-proof-ladder'
|
||||
|
||||
// The descendant models an MCP server: it either cooperates or, when it traps
|
||||
@@ -99,11 +100,13 @@ async function proveExitWithRetries(
|
||||
}
|
||||
|
||||
function spawnScript(script: string): ReturnType<typeof spawnProcess> {
|
||||
return spawnProcess({
|
||||
const child = spawnProcess({
|
||||
program: process.execPath,
|
||||
args: ['-e', script],
|
||||
stdio: ['pipe', 'pipe', 'pipe']
|
||||
})
|
||||
managedChild(child)
|
||||
return child
|
||||
}
|
||||
|
||||
function firstStdoutLine(child: ReturnType<typeof spawnProcess>): Promise<string> {
|
||||
@@ -112,28 +115,19 @@ function firstStdoutLine(child: ReturnType<typeof spawnProcess>): Promise<string
|
||||
})
|
||||
}
|
||||
|
||||
function observeExit(child: EventEmitter): { exitPromise: Promise<void>; exited: () => boolean } {
|
||||
let exited = false
|
||||
const exitPromise = new Promise<void>((resolve) => {
|
||||
child.once('exit', () => {
|
||||
exited = true
|
||||
resolve()
|
||||
})
|
||||
})
|
||||
return { exitPromise, exited: () => exited }
|
||||
}
|
||||
|
||||
/** `null` models a spawn that failed before a pid existed. */
|
||||
function mockChild(
|
||||
pid: number | null = 424242
|
||||
): EventEmitter &
|
||||
Pick<SpawnedProcess, 'pid' | 'kill' | 'stdin'> & { kill: ReturnType<typeof vi.fn> } {
|
||||
const child = new EventEmitter()
|
||||
return Object.assign(child, {
|
||||
Pick<SpawnedProcess, 'pid' | 'kill' | 'stdin' | 'stderr'> & { kill: ReturnType<typeof vi.fn> } {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: pid ?? undefined,
|
||||
stdin: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
kill: vi.fn(() => true)
|
||||
}) as never
|
||||
})
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The proof reads only pid, events, kill, stdin and stderr from this fixture.
|
||||
return child as never
|
||||
}
|
||||
|
||||
/** A tree whose verdict is scripted per reap, recording when it was armed. */
|
||||
@@ -200,7 +194,7 @@ describe('claude child exit proof', () => {
|
||||
await ageDescendantPastTheCaptureSecond()
|
||||
|
||||
try {
|
||||
const proven = await proveExitWithRetries({ child, ...observeExit(child) })
|
||||
const proven = await proveExitWithRetries({ managed: managedChild(child) })
|
||||
// Evaluated AT the boundary, not by polling until a deferred sweep timer
|
||||
// wins: true releases the lease, so a descendant still running here is
|
||||
// exactly the orphan the proof exists to prevent. False would be the
|
||||
@@ -237,7 +231,7 @@ describe('claude child exit proof', () => {
|
||||
await ageDescendantPastTheCaptureSecond()
|
||||
|
||||
try {
|
||||
const proven = await proveExitWithRetries({ child, ...observeExit(child) })
|
||||
const proven = await proveExitWithRetries({ managed: managedChild(child) })
|
||||
expect({ proven, descendant: descendantState(descendantPid) }).toEqual({
|
||||
proven: true,
|
||||
descendant: 'exited'
|
||||
@@ -263,7 +257,7 @@ describe('claude child exit proof', () => {
|
||||
descendantPid: number
|
||||
}
|
||||
try {
|
||||
const proven = await proveExitWithRetries({ child, ...observeExit(child) })
|
||||
const proven = await proveExitWithRetries({ managed: managedChild(child) })
|
||||
expect({ proven, descendant: descendantState(descendantPid) }).toEqual({
|
||||
proven: true,
|
||||
descendant: 'exited'
|
||||
@@ -282,26 +276,26 @@ describe('claude child exit proof', () => {
|
||||
it('arms the snapshot before stdin closes and verifies it after a clean exit', async () => {
|
||||
const child = spawnScript(COOPERATIVE_CHILD)
|
||||
expect(await firstStdoutLine(child)).toBe('ready')
|
||||
const exit = observeExit(child)
|
||||
const managed = managedChild(child)
|
||||
const tree = mockTree(['exited'])
|
||||
let exitedWhenArmed: boolean | null = null
|
||||
tree.capture.mockImplementation(async () => {
|
||||
exitedWhenArmed = exit.exited()
|
||||
exitedWhenArmed = managed.rootVerdict === 'exited'
|
||||
})
|
||||
|
||||
await expect(proveClaudeChildExit({ child, ...exit, tree })).resolves.toBe(true)
|
||||
await expect(proveClaudeChildExit({ managed, tree })).resolves.toBe(true)
|
||||
// The snapshot is the only proof that survives the root: taken while it lived,
|
||||
// verified once it left. A reap before the exit would have been the forced ladder.
|
||||
expect(exitedWhenArmed).toBe(false)
|
||||
expect(tree.reap).toHaveBeenCalledTimes(1)
|
||||
expect(exit.exited()).toBe(true)
|
||||
expect(managed.rootVerdict).toBe('exited')
|
||||
}, 20_000)
|
||||
|
||||
it('proves a clean close of a childless root with one snapshot and no signal', async () => {
|
||||
const child = spawnScript(COOPERATIVE_CHILD)
|
||||
expect(await firstStdoutLine(child)).toBe('ready')
|
||||
|
||||
await expect(proveClaudeChildExit({ child, ...observeExit(child) })).resolves.toBe(true)
|
||||
await expect(proveClaudeChildExit({ managed: managedChild(child) })).resolves.toBe(true)
|
||||
}, 20_000)
|
||||
|
||||
it('reports an unprovable exit as false rather than assuming the child died', async () => {
|
||||
@@ -309,9 +303,7 @@ describe('claude child exit proof', () => {
|
||||
try {
|
||||
const tree = mockTree(['exited'])
|
||||
const proof = proveClaudeChildExit({
|
||||
child: mockChild(),
|
||||
exitPromise: new Promise<void>(() => {}),
|
||||
exited: () => false,
|
||||
managed: managedChild(mockChild()),
|
||||
tree
|
||||
})
|
||||
await new Promise((resolve) => setImmediate(resolve))
|
||||
@@ -330,18 +322,18 @@ describe('claude child exit proof', () => {
|
||||
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
|
||||
try {
|
||||
const child = mockChild()
|
||||
const exit = observeExit(child)
|
||||
const managed = managedChild(child)
|
||||
const tree = mockTree(['live'])
|
||||
tree.reap.mockImplementation(async () => {
|
||||
child.emit('exit', null, 'SIGKILL')
|
||||
return 'live'
|
||||
})
|
||||
const proof = proveClaudeChildExit({ child, ...exit, tree })
|
||||
const proof = proveClaudeChildExit({ managed, tree })
|
||||
await new Promise((resolve) => setImmediate(resolve))
|
||||
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS)
|
||||
|
||||
await expect(proof).resolves.toBe(false)
|
||||
expect(exit.exited()).toBe(true)
|
||||
expect(managed.rootVerdict).toBe('exited')
|
||||
// One verification per attempt: the retried close re-verifies, this one does not.
|
||||
expect(tree.reap).toHaveBeenCalledTimes(1)
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
@@ -352,16 +344,18 @@ describe('claude child exit proof', () => {
|
||||
|
||||
it('re-verifies an unproven tree on a retried close instead of trusting the dead root', async () => {
|
||||
const child = mockChild()
|
||||
const managed = managedChild(child)
|
||||
child.emit('exit', 0, null)
|
||||
const tree = mockTree(['exited'])
|
||||
|
||||
await expect(
|
||||
proveClaudeChildExit({ child, exitPromise: Promise.resolve(), exited: () => true, tree })
|
||||
).resolves.toBe(true)
|
||||
await expect(proveClaudeChildExit({ managed, tree })).resolves.toBe(true)
|
||||
expect(tree.reap).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('stays unproven for a root that left before any snapshot could be armed', async () => {
|
||||
const child = mockChild()
|
||||
const managed = managedChild(child)
|
||||
child.emit('exit', 0, null)
|
||||
const captureDescendants = vi.fn(async () => snapshotOf(4243))
|
||||
const terminateDescendants = vi.fn()
|
||||
const tree = createClaudeChildTreeReaper(child, {
|
||||
@@ -371,9 +365,7 @@ describe('claude child exit proof', () => {
|
||||
terminateDescendants
|
||||
})
|
||||
|
||||
await expect(
|
||||
proveClaudeChildExit({ child, exitPromise: Promise.resolve(), exited: () => true, tree })
|
||||
).resolves.toBe(false)
|
||||
await expect(proveClaudeChildExit({ managed, tree })).resolves.toBe(false)
|
||||
// A dead root's descendants have reparented: walking its pid now could only
|
||||
// sweep a stranger, so no walk is attempted and nothing is proven.
|
||||
expect(captureDescendants).not.toHaveBeenCalled()
|
||||
|
||||
@@ -361,6 +361,8 @@ export function createClaudeChildTreeReaper(
|
||||
*/
|
||||
export function proveClaudeChildExit(input: ClaudeChildExitProofInput): Promise<boolean> {
|
||||
return proveClaudeChildExitWithReaper(input, () =>
|
||||
createClaudeChildTreeReaper(input.child, { exited: input.exited })
|
||||
createClaudeChildTreeReaper(input.managed.child, {
|
||||
exited: () => input.managed.rootVerdict === 'exited'
|
||||
})
|
||||
)
|
||||
}
|
||||
|
||||
@@ -121,32 +121,21 @@ describe('claude agent SDK process spawn', () => {
|
||||
spawn.spawn(sdkOptions())
|
||||
expect(spawn.supervised).toBe(specSupervised)
|
||||
|
||||
let exited = false
|
||||
let settle = (): void => {}
|
||||
const exitPromise = new Promise<void>((resolve) => {
|
||||
settle = resolve
|
||||
})
|
||||
// Claude leaves shortly after stdin ends, so the ladder never needs its forced rung.
|
||||
const managed = spawn.managed
|
||||
if (!managed) {
|
||||
throw new Error('Claude spawner did not retain its managed child')
|
||||
}
|
||||
// Claude leaves shortly after stdin ends, before the forced stop.
|
||||
process.child.stdin.on('finish', () =>
|
||||
setTimeout(() => {
|
||||
exited = true
|
||||
settle()
|
||||
}, 10)
|
||||
setTimeout(() => process.child.emit('exit', 0, null), 10)
|
||||
)
|
||||
const tree = {
|
||||
capture: vi.fn(async () => {}),
|
||||
reap: vi.fn(async () => 'exited' as const),
|
||||
treeVerdict: 'exited' as const
|
||||
}
|
||||
await proveClaudeChildExitWithReaper(
|
||||
{
|
||||
child: process.child,
|
||||
exitPromise,
|
||||
exited: () => exited,
|
||||
tree,
|
||||
supervised: spawn.supervised
|
||||
},
|
||||
() => tree
|
||||
await expect(proveClaudeChildExitWithReaper({ managed, tree }, () => tree)).resolves.toBe(
|
||||
true
|
||||
)
|
||||
// SIGTERM to an unsupervised Claude on Windows is TerminateProcess; a skipped one leaves it running.
|
||||
if (specSupervised) {
|
||||
@@ -176,8 +165,8 @@ describe('claude agent SDK process spawn', () => {
|
||||
process.child.stderr.write('claude: not signed in')
|
||||
await new Promise((resolve) => setImmediate(resolve))
|
||||
|
||||
expect(spawn.stderrTail).toMatch(/claude: not signed in$/)
|
||||
expect(spawn.stderrTail.length).toBe(8192)
|
||||
expect(spawn.managed?.stderrTail()).toMatch(/claude: not signed in$/)
|
||||
expect(spawn.managed?.stderrTail().length).toBe(8192)
|
||||
})
|
||||
|
||||
it('hands a Windows .cmd shim to Orca\u2019s argument encoder', () => {
|
||||
|
||||
@@ -1,17 +1,20 @@
|
||||
import type { SpawnOptions as ClaudeAgentSdkSpawnOptions } from '@anthropic-ai/claude-agent-sdk'
|
||||
import { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { createProviderSpawnSpec } from '../provider-process/provider-process-supervisor'
|
||||
import {
|
||||
spawnManagedProviderProcess,
|
||||
type ManagedProviderProcess
|
||||
} from '../provider-process/managed-provider-process'
|
||||
import { claudeChildClosePolicy } from './claude-child-exit-proof-ladder'
|
||||
|
||||
/** Derived rather than imported: only src/shared/child-process may name node:child_process. */
|
||||
type ClaudeCodeChild = ReturnType<typeof spawnProcess>
|
||||
|
||||
const STDERR_TAIL_MAX_BYTES = 8192
|
||||
|
||||
export type ClaudeCodeProcessSpawn = {
|
||||
/** Pass as the SDK's `spawnClaudeCodeProcess`; the SDK never learns the pid because it never owns it. */
|
||||
spawn: (options: ClaudeAgentSdkSpawnOptions) => ClaudeCodeChild
|
||||
/** The retained child, so Orca keeps its own tree-kill and exit-proof ladder. Null until the SDK spawns. */
|
||||
readonly child: ClaudeCodeChild | null
|
||||
readonly managed: ManagedProviderProcess | null
|
||||
/**
|
||||
* Ownership proof: the durable lease adjudicates on this pid plus start time plus the spawn
|
||||
* token. On POSIX it is the provider supervisor's, which outlives Claude by construction.
|
||||
@@ -19,7 +22,6 @@ export type ClaudeCodeProcessSpawn = {
|
||||
readonly pid: number | undefined
|
||||
/** The spawn spec's verdict, so the close ladder never re-decides it. False until the SDK spawns. */
|
||||
readonly supervised: boolean
|
||||
readonly stderrTail: string
|
||||
}
|
||||
|
||||
function definedEnv(env: Record<string, string | undefined>): Record<string, string> {
|
||||
@@ -46,50 +48,39 @@ export function createClaudeCodeProcessSpawn(
|
||||
spawnImpl: typeof spawnProcess = spawnProcess,
|
||||
platform: NodeJS.Platform = process.platform
|
||||
): ClaudeCodeProcessSpawn {
|
||||
let child: ClaudeCodeChild | null = null
|
||||
let stderrTail = ''
|
||||
let supervised = false
|
||||
let managed: ManagedProviderProcess | null = null
|
||||
return {
|
||||
spawn: (options) => {
|
||||
const spec = createProviderSpawnSpec(
|
||||
// SDK abort would kill the child outside Orca's ladder, losing observed exit proof.
|
||||
managed = spawnManagedProviderProcess(
|
||||
{
|
||||
command: options.command,
|
||||
args: [...options.args],
|
||||
...(options.cwd === undefined ? {} : { cwd: options.cwd })
|
||||
},
|
||||
definedEnv(options.env),
|
||||
platform
|
||||
{
|
||||
spawnImpl,
|
||||
platform,
|
||||
inheritedEnv: definedEnv(options.env),
|
||||
site: 'claude-stream-json-teardown',
|
||||
policy: claudeChildClosePolicy,
|
||||
acceptClose: (result) => result.root === 'exited' && result.tree === 'exited'
|
||||
}
|
||||
)
|
||||
// Why `options.signal` is dropped: it would let the SDK kill the child outside
|
||||
// Orca's ladder, and close() may never report an exit it did not observe.
|
||||
const spawned = spawnImpl({
|
||||
program: spec.program,
|
||||
args: spec.args,
|
||||
cwd: spec.cwd,
|
||||
env: spec.env,
|
||||
detached: spec.detached,
|
||||
stdio: ['pipe', 'pipe', 'pipe']
|
||||
})
|
||||
child = spawned
|
||||
supervised = spec.supervised
|
||||
// The SDK drains stderr only for its own local spawn, so a custom spawner must:
|
||||
// otherwise the child blocks on a full pipe and exit errors lose their tail.
|
||||
spawned.stderr.setEncoding('utf8').on('data', (chunk: string) => {
|
||||
stderrTail = (stderrTail + chunk).slice(-STDERR_TAIL_MAX_BYTES)
|
||||
})
|
||||
return spawned
|
||||
// The SDK drains stderr only for its own local spawn; the managed process drains it here.
|
||||
return managed.child
|
||||
},
|
||||
get child() {
|
||||
return child
|
||||
return managed?.child ?? null
|
||||
},
|
||||
get managed() {
|
||||
return managed
|
||||
},
|
||||
get pid() {
|
||||
return child?.pid
|
||||
return managed?.child.pid
|
||||
},
|
||||
get supervised() {
|
||||
return supervised
|
||||
},
|
||||
get stderrTail() {
|
||||
return stderrTail
|
||||
return managed?.supervised ?? false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
import type { EventEmitter } from 'node:events'
|
||||
import type { spawnProcess, SpawnedProcess } from '../../shared/child-process/run-process'
|
||||
import {
|
||||
spawnManagedProviderProcess,
|
||||
type ManagedProviderProcess
|
||||
} from '../provider-process/managed-provider-process'
|
||||
import { claudeChildClosePolicy } from './claude-child-exit-proof-ladder'
|
||||
|
||||
const managedChildren = new WeakMap<object, ManagedProviderProcess>()
|
||||
|
||||
export function managedChild(
|
||||
child: Pick<SpawnedProcess, 'pid' | 'kill' | 'stdin' | 'stderr'> & EventEmitter
|
||||
): ManagedProviderProcess {
|
||||
const existing = managedChildren.get(child)
|
||||
if (existing) {
|
||||
return existing
|
||||
}
|
||||
const managed = spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The managed process reads only the fixture's owned pid, events, kill, stdin and stderr.
|
||||
spawnImpl: () => child as ReturnType<typeof spawnProcess>,
|
||||
platform: 'win32',
|
||||
site: 'claude-proof-fixture',
|
||||
policy: claudeChildClosePolicy,
|
||||
acceptClose: (result) => result.root === 'exited' && result.tree === 'exited'
|
||||
}
|
||||
)
|
||||
managedChildren.set(child, managed)
|
||||
return managed
|
||||
}
|
||||
@@ -1,6 +1,10 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from '../provider-process/provider-process-supervisor'
|
||||
import type { ClaudeChildTreeReaper } from './claude-agent-sdk-exit-proof'
|
||||
import { createClaudeCodeProcessSpawn } from './claude-agent-sdk-process-spawn'
|
||||
import { proveClaudeChildExitWithReaper } from './claude-child-exit-proof-ladder'
|
||||
|
||||
function fakeTree(): ClaudeChildTreeReaper & { reap: ReturnType<typeof vi.fn> } {
|
||||
@@ -12,61 +16,65 @@ function fakeTree(): ClaudeChildTreeReaper & { reap: ReturnType<typeof vi.fn> }
|
||||
}
|
||||
}
|
||||
|
||||
/** A root that leaves only once a SIGTERM has had `stopMs` to act, the way a supervisor does. */
|
||||
function rootStoppedBySigterm(stopMs: number) {
|
||||
let exited = false
|
||||
let settle = (): void => {}
|
||||
const exitPromise = new Promise<void>((resolve) => {
|
||||
settle = resolve
|
||||
function rootStoppedBySigterm(stopMs: number, platform: NodeJS.Platform) {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: 4321,
|
||||
stdin: new PassThrough(),
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
kill: vi.fn((signal?: NodeJS.Signals | number) => {
|
||||
if (signal === 'SIGTERM') {
|
||||
setTimeout(() => child.emit('exit', 0, 'SIGTERM'), stopMs)
|
||||
}
|
||||
return true
|
||||
})
|
||||
})
|
||||
const kill = vi.fn((signal?: NodeJS.Signals | number) => {
|
||||
if (signal === 'SIGTERM') {
|
||||
setTimeout(() => {
|
||||
exited = true
|
||||
settle()
|
||||
}, stopMs)
|
||||
}
|
||||
return true
|
||||
const spawner = createClaudeCodeProcessSpawn(() => {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The fixture supplies every event, stream and process field used by the spawner and close.
|
||||
return child as unknown as ReturnType<typeof spawnProcess>
|
||||
}, platform)
|
||||
spawner.spawn({
|
||||
command: 'fixture-provider',
|
||||
args: [],
|
||||
env: {},
|
||||
signal: new AbortController().signal
|
||||
})
|
||||
const stdin = { end: vi.fn() }
|
||||
return {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The ladder reads only pid, kill and stdin.end from its child.
|
||||
child: { pid: 4321, kill, stdin } as unknown as Parameters<
|
||||
typeof proveClaudeChildExitWithReaper
|
||||
>[0]['child'],
|
||||
kill,
|
||||
stdin,
|
||||
exitPromise,
|
||||
exited: () => exited
|
||||
const managed = spawner.managed
|
||||
if (!managed) {
|
||||
throw new Error('Fixture did not retain its managed child')
|
||||
}
|
||||
return { child, managed }
|
||||
}
|
||||
|
||||
afterEach(() => vi.useRealTimers())
|
||||
|
||||
describe('Claude child exit proof ladder', () => {
|
||||
it('stops a supervised child with SIGTERM and waits out the supervisor stop before forcing', async () => {
|
||||
// Slower than the unsupervised 1.5 s grace, still inside the supervisor's own bound.
|
||||
const root = rootStoppedBySigterm(PROVIDER_SUPERVISOR_MAX_STOP_MS - 500)
|
||||
vi.useFakeTimers()
|
||||
const root = rootStoppedBySigterm(PROVIDER_SUPERVISOR_MAX_STOP_MS - 500, 'darwin')
|
||||
const tree = fakeTree()
|
||||
|
||||
await expect(
|
||||
proveClaudeChildExitWithReaper({ ...root, supervised: true, tree }, () => tree)
|
||||
).resolves.toBe(true)
|
||||
|
||||
expect(root.stdin.end).toHaveBeenCalled()
|
||||
expect(root.kill).toHaveBeenCalledWith('SIGTERM')
|
||||
// Forcing here would SIGKILL the supervisor mid-stop and orphan Claude in its own group.
|
||||
expect(root.kill).not.toHaveBeenCalledWith('SIGKILL')
|
||||
expect(root.managed.rootVerdict).toBe('live')
|
||||
const proof = proveClaudeChildExitWithReaper({ managed: root.managed, tree }, () => tree)
|
||||
await vi.advanceTimersByTimeAsync(PROVIDER_SUPERVISOR_MAX_STOP_MS)
|
||||
await expect(proof).resolves.toBe(true)
|
||||
expect(root.child.stdin.writableEnded).toBe(true)
|
||||
expect(root.child.kill).toHaveBeenCalledExactlyOnceWith('SIGTERM')
|
||||
expect(root.managed.lastCloseResult).toEqual({ root: 'exited', tree: 'exited' })
|
||||
expect(tree.reap).not.toHaveBeenCalled()
|
||||
}, 10_000)
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
|
||||
it('never signals an unsupervised child for the graceful stop', async () => {
|
||||
const root = rootStoppedBySigterm(0)
|
||||
vi.useFakeTimers()
|
||||
const root = rootStoppedBySigterm(0, 'win32')
|
||||
const tree = fakeTree()
|
||||
|
||||
await proveClaudeChildExitWithReaper({ ...root, tree }, () => tree)
|
||||
|
||||
// On Windows a direct SIGTERM is TerminateProcess: stdin end stays the only graceful rung.
|
||||
expect(root.stdin.end).toHaveBeenCalled()
|
||||
expect(root.kill).not.toHaveBeenCalledWith('SIGTERM')
|
||||
expect(tree.reap).toHaveBeenCalled()
|
||||
}, 10_000)
|
||||
const proof = proveClaudeChildExitWithReaper({ managed: root.managed, tree }, () => tree)
|
||||
await vi.advanceTimersByTimeAsync(2_500)
|
||||
await expect(proof).resolves.toBe(false)
|
||||
expect(root.child.stdin.writableEnded).toBe(true)
|
||||
expect(root.child.kill).not.toHaveBeenCalledWith('SIGTERM')
|
||||
expect(tree.reap).toHaveBeenCalledOnce()
|
||||
expect(root.managed.lastCloseResult).toEqual({ root: 'live', tree: 'exited' })
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { SpawnedProcess } from '../../shared/child-process/run-process'
|
||||
import { waitForProcessExitUntil } from '../provider-process/provider-process-exit-deadline'
|
||||
import type { ManagedProviderProcess } from '../provider-process/managed-provider-process'
|
||||
import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from '../provider-process/provider-process-supervisor'
|
||||
import type { ProviderProcessClosePolicy } from '../provider-process/provider-process-close'
|
||||
import type { ClaudeChildTreeReaper } from './claude-agent-sdk-exit-proof'
|
||||
|
||||
export const GRACEFUL_EXIT_MS = 1_500
|
||||
@@ -8,47 +8,23 @@ export const GRACEFUL_EXIT_MS = 1_500
|
||||
export const SUPERVISED_GRACEFUL_EXIT_MS = PROVIDER_SUPERVISOR_MAX_STOP_MS + 500
|
||||
const FORCED_EXIT_MS = 1_000
|
||||
|
||||
export function claudeChildClosePolicy(supervised: boolean): ProviderProcessClosePolicy {
|
||||
return {
|
||||
gracefulExitMs: supervised ? SUPERVISED_GRACEFUL_EXIT_MS : GRACEFUL_EXIT_MS,
|
||||
forcedExitMs: FORCED_EXIT_MS,
|
||||
signalSupervisorOnClose: true
|
||||
}
|
||||
}
|
||||
|
||||
export type ClaudeChildExitProofInput = {
|
||||
child: Pick<SpawnedProcess, 'pid' | 'kill' | 'stdin'>
|
||||
exitPromise: Promise<void>
|
||||
exited: () => boolean
|
||||
managed: ManagedProviderProcess
|
||||
tree?: ClaudeChildTreeReaper
|
||||
/** The child is the POSIX provider supervisor: SIGTERM stops Claude, which reaps its tools. */
|
||||
supervised?: boolean
|
||||
}
|
||||
|
||||
export async function proveClaudeChildExitWithReaper(
|
||||
input: ClaudeChildExitProofInput,
|
||||
createTree: () => ClaudeChildTreeReaper
|
||||
): Promise<boolean> {
|
||||
const tree = input.tree ?? createTree()
|
||||
// Arm before the stop: only a live root can identify its descendants.
|
||||
await tree.capture()
|
||||
try {
|
||||
input.child.stdin?.end()
|
||||
} catch {
|
||||
// The reap below still owns the process.
|
||||
}
|
||||
// Stdin end alone lets Claude finish its turn, tools and edits included; a close is a stop.
|
||||
// Windows has no supervisor, and a direct SIGTERM there is TerminateProcess.
|
||||
if (input.supervised && !input.exited()) {
|
||||
input.child.kill('SIGTERM')
|
||||
}
|
||||
let reaped = false
|
||||
if (!input.exited()) {
|
||||
await waitForProcessExitUntil(
|
||||
input.exitPromise,
|
||||
input.supervised ? SUPERVISED_GRACEFUL_EXIT_MS : GRACEFUL_EXIT_MS
|
||||
)
|
||||
if (!input.exited()) {
|
||||
reaped = true
|
||||
await tree.refresh?.()
|
||||
await tree.reap()
|
||||
await waitForProcessExitUntil(input.exitPromise, FORCED_EXIT_MS)
|
||||
}
|
||||
}
|
||||
if (!reaped && input.exited() && tree.treeVerdict !== 'exited') {
|
||||
await tree.reap()
|
||||
}
|
||||
return input.exited() && tree.treeVerdict === 'exited'
|
||||
const result = await input.managed.close(input.tree ?? createTree())
|
||||
return result.root === 'exited' && result.tree === 'exited'
|
||||
}
|
||||
|
||||
@@ -706,7 +706,7 @@ describe('Claude stream-json connection', () => {
|
||||
)
|
||||
|
||||
await until(
|
||||
() => (connection.exitVerdict.root === 'processless' ? connection.exitVerdict : null),
|
||||
() => (connection.exitVerdict.processless === true ? connection.exitVerdict : null),
|
||||
'the processless spawn settlement'
|
||||
)
|
||||
expect(connection.pid).toBeUndefined()
|
||||
@@ -717,7 +717,7 @@ describe('Claude stream-json connection', () => {
|
||||
true
|
||||
])
|
||||
await expect(connection.close()).resolves.toBe(true)
|
||||
expect(connection.exitVerdict).toEqual({ root: 'processless', tree: 'exited' })
|
||||
expect(connection.exitVerdict).toEqual({ root: 'exited', tree: 'exited', processless: true })
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
@@ -83,8 +83,9 @@ export type ClaudeStreamJsonConnectionHandlers = {
|
||||
* is never collapsed into either neighbour.
|
||||
*/
|
||||
export type ClaudeChildExitVerdict = {
|
||||
root: 'exited' | 'live' | 'processless'
|
||||
root: DescendantTreeVerdict
|
||||
tree: DescendantTreeVerdict
|
||||
processless?: boolean
|
||||
}
|
||||
|
||||
export type ClaudeStreamJsonConnection = ClaudeControlSurface & {
|
||||
@@ -148,8 +149,9 @@ export async function openClaudeStreamJsonConnection(
|
||||
...(handlers.onUserDialog ? { onUserDialog: handlers.onUserDialog } : {})
|
||||
}
|
||||
})
|
||||
const child = spawner.child
|
||||
if (!child) {
|
||||
const managed = spawner.managed
|
||||
const child = managed?.child
|
||||
if (!child || !managed) {
|
||||
throw new Error('the claude agent SDK returned without spawning a child')
|
||||
}
|
||||
// This child owns the account's credentials for as long as it runs, exactly as a
|
||||
@@ -158,11 +160,8 @@ export async function openClaudeStreamJsonConnection(
|
||||
// path exists.
|
||||
const authGateKey = randomUUID()
|
||||
const releaseAuthGate = (): void => markClaudeStructuredChildExited(authGateKey)
|
||||
let exited = false
|
||||
let exitStatus: ExitStatus | null = null
|
||||
let closing = false
|
||||
let processless = false
|
||||
let prePidSpawnError = false
|
||||
let terminalError: Error | null = null
|
||||
let faultReported = false
|
||||
let exitReported = false
|
||||
@@ -170,7 +169,7 @@ export async function openClaudeStreamJsonConnection(
|
||||
let readingBarrier: Promise<void> | null = null
|
||||
let releaseReadingBarrier: (() => void) | null = null
|
||||
const pauseReading = (): void => {
|
||||
if (closing || exited || terminalError || readingBarrier) {
|
||||
if (closing || managed.rootVerdict === 'exited' || terminalError || readingBarrier) {
|
||||
return
|
||||
}
|
||||
readingBarrier = new Promise<void>((resolve) => {
|
||||
@@ -185,7 +184,7 @@ export async function openClaudeStreamJsonConnection(
|
||||
}
|
||||
const waitUntilReadable = (): Promise<void> => readingBarrier ?? Promise.resolve()
|
||||
// One reaper per child: every close attempt and error-path reap shares its proof.
|
||||
const rootSettled = (): boolean => exited || processless
|
||||
const rootSettled = (): boolean => managed.rootVerdict === 'exited'
|
||||
const tree = createClaudeChildTreeReaper(child, { exited: rootSettled })
|
||||
|
||||
// Arm lazily on actual child output instead of issuing a process-table scan for
|
||||
@@ -203,34 +202,19 @@ export async function openClaudeStreamJsonConnection(
|
||||
// The SDK may synchronously spawn the CLI and consume an early stderr chunk
|
||||
// before this connection can attach its listener; the bounded tail preserves
|
||||
// that observation for the same lazy arm.
|
||||
if (spawner.stderrTail.length > 0) {
|
||||
if (managed.stderrTail().length > 0) {
|
||||
armTreeOnOutput()
|
||||
}
|
||||
|
||||
let settleExit = (): void => {}
|
||||
const exitPromise = new Promise<void>((resolve) => {
|
||||
settleExit = resolve
|
||||
})
|
||||
const markExited = (): void => {
|
||||
exited = true
|
||||
releaseAuthGate()
|
||||
settleExit()
|
||||
}
|
||||
child.on('exit', (code, signal) => {
|
||||
exitStatus = { code, signal }
|
||||
markExited()
|
||||
handleUnexpectedEnd()
|
||||
})
|
||||
|
||||
const handleUnexpectedEnd = (cause?: Error): void => {
|
||||
resumeReading()
|
||||
terminalError ??= exitError(spawner.stderrTail, exitStatus, cause)
|
||||
terminalError ??= exitError(managed.stderrTail(), exitStatus, cause)
|
||||
inbox.fail(terminalError)
|
||||
if (!closing && !faultReported) {
|
||||
faultReported = true
|
||||
handlers.onFault?.(terminalError)
|
||||
}
|
||||
if (exited && !exitReported) {
|
||||
if (managed.rootExitObserved && !exitReported) {
|
||||
exitReported = true
|
||||
handlers.onExit?.(terminalError, { expected: closing })
|
||||
}
|
||||
@@ -262,29 +246,26 @@ export async function openClaudeStreamJsonConnection(
|
||||
} catch (error: unknown) {
|
||||
// The SDK ends its generator in error when the child dies or the transport
|
||||
// fails; a transport failure with a live child still has to reap the tree.
|
||||
if (!closing && !exited) {
|
||||
if (!closing && managed.rootVerdict !== 'exited') {
|
||||
void tree.reap()
|
||||
}
|
||||
handleUnexpectedEnd(error instanceof Error ? error : new Error(String(error)))
|
||||
}
|
||||
})()
|
||||
|
||||
managed.onExit((exit) => {
|
||||
exitStatus = exit
|
||||
releaseAuthGate()
|
||||
handleUnexpectedEnd()
|
||||
})
|
||||
child.on('error', (error) => {
|
||||
if (spawner.pid === undefined) {
|
||||
prePidSpawnError = true
|
||||
}
|
||||
if (!closing && !exited) {
|
||||
if (!closing && managed.rootVerdict !== 'exited') {
|
||||
void tree.reap()
|
||||
}
|
||||
handleUnexpectedEnd(error)
|
||||
})
|
||||
child.on('close', () => {
|
||||
// Covers the spawn-failure path too, where no 'exit' ever arrives.
|
||||
releaseAuthGate()
|
||||
if (prePidSpawnError && spawner.pid === undefined) {
|
||||
processless = true
|
||||
settleExit()
|
||||
}
|
||||
handleUnexpectedEnd()
|
||||
})
|
||||
child.stdin.on('error', (error) => {
|
||||
@@ -302,7 +283,13 @@ export async function openClaudeStreamJsonConnection(
|
||||
markClaudeStructuredChildSpawned(authGateKey)
|
||||
|
||||
const send: ClaudeStreamJsonConnection['send'] = (message, beforeDispatch) => {
|
||||
if (closing || exited || terminalError || child.stdin.destroyed || !child.stdin.writable) {
|
||||
if (
|
||||
closing ||
|
||||
managed.rootVerdict === 'exited' ||
|
||||
terminalError ||
|
||||
child.stdin.destroyed ||
|
||||
!child.stdin.writable
|
||||
) {
|
||||
return Promise.reject(
|
||||
claudeUnwrittenUserMessageError(
|
||||
terminalError ?? new Error('claude stream-json connection is closed')
|
||||
@@ -322,15 +309,12 @@ export async function openClaudeStreamJsonConnection(
|
||||
await (tree.refresh?.() ?? tree.capture())
|
||||
inbox.end()
|
||||
const proven = await proveClaudeChildExit({
|
||||
child,
|
||||
exitPromise,
|
||||
exited: rootSettled,
|
||||
tree,
|
||||
supervised: spawner.supervised
|
||||
managed
|
||||
})
|
||||
inbox.fail(new Error('claude stream-json connection closed'))
|
||||
if (!proven) {
|
||||
if (exited && tree.treeVerdict === 'live') {
|
||||
if (managed.lastCloseResult?.root === 'exited' && managed.lastCloseResult.tree === 'live') {
|
||||
console.warn('[claude-stream-json] root exited but a descendant survived the close:', {
|
||||
pid: spawner.pid
|
||||
})
|
||||
@@ -359,12 +343,13 @@ export async function openClaudeStreamJsonConnection(
|
||||
return spawner.pid
|
||||
},
|
||||
get closed() {
|
||||
return closing || exited || terminalError !== null
|
||||
return closing || managed.rootVerdict === 'exited' || terminalError !== null
|
||||
},
|
||||
get exitVerdict() {
|
||||
return {
|
||||
root: processless ? 'processless' : exited ? 'exited' : 'live',
|
||||
tree: tree.treeVerdict
|
||||
root: managed.rootVerdict,
|
||||
tree: tree.treeVerdict,
|
||||
...(managed.processless ? { processless: true } : {})
|
||||
} as const
|
||||
},
|
||||
pauseReading,
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
import { claudePromptCardWritten } from './claude-child-work-evidence'
|
||||
import type {
|
||||
ClaudeSession,
|
||||
ClaudeStructuredSessionAdapterDeps,
|
||||
ClaudeStructuredSessionEvent
|
||||
} from './claude-structured-session-state'
|
||||
|
||||
type ClaudeEventDelivery = {
|
||||
session: ClaudeSession | null
|
||||
event: ClaudeStructuredSessionEvent
|
||||
deps: Pick<ClaudeStructuredSessionAdapterDeps, 'onEvent'>
|
||||
publishChildWork: (
|
||||
sessionId: string,
|
||||
session?: ClaudeSession | null,
|
||||
message?: Record<string, unknown> | null
|
||||
) => void
|
||||
}
|
||||
|
||||
export function emitClaudeStructuredSessionEvent({
|
||||
session,
|
||||
event,
|
||||
deps,
|
||||
publishChildWork
|
||||
}: ClaudeEventDelivery): void {
|
||||
// Host evidence follows the journal; the tracker's roster remains available to contract tests.
|
||||
if (event.type === 'ended') {
|
||||
session?.childWork.clear()
|
||||
session?.backgroundTasks.clear()
|
||||
} else if (event.type === 'message') {
|
||||
session?.childWork.observe(event.message)
|
||||
session?.backgroundTasks.observe(event.message, event.startsTurn === true)
|
||||
} else if (event.type === 'prompt-cancelled') {
|
||||
// Free the child before its withdrawn card is written closed.
|
||||
publishChildWork(event.sessionId, session)
|
||||
}
|
||||
if (event.type === 'message' && session?.commands.observe(event.message)) {
|
||||
session.events?.publish()
|
||||
}
|
||||
session?.translator?.handle(event)
|
||||
deps.onEvent?.(event)
|
||||
publishChildWork(event.sessionId, session, event.type === 'message' ? event.message : null)
|
||||
// A prompt blocks its child only after the journal has written its card.
|
||||
void claudePromptCardWritten(session, event)?.then(() => publishChildWork(event.sessionId))
|
||||
}
|
||||
@@ -29,7 +29,7 @@ describe('Claude structured processless acquisition', () => {
|
||||
const connection: ClaudeStreamJsonConnection = {
|
||||
pid: undefined,
|
||||
closed: true,
|
||||
exitVerdict: { root: 'processless', tree: 'exited' },
|
||||
exitVerdict: { root: 'exited', tree: 'exited', processless: true },
|
||||
initializationResult: async () => {
|
||||
throw fault
|
||||
},
|
||||
|
||||
@@ -14,7 +14,7 @@ import { supportsClaudeStructuredLocation } from './claude-structured-location-s
|
||||
import { setClaudeStructuredSessionOption } from './claude-structured-options'
|
||||
import { readClaudeStructuredSessionOptions } from './claude-structured-session-options'
|
||||
import {
|
||||
claudeStartupFailureFact,
|
||||
awaitClaudeSessionStarted,
|
||||
claudeStartupSettledWithin
|
||||
} from './claude-structured-session-startup-state'
|
||||
import { CLAUDE_DEFAULT_REQUEST_TIMEOUT_MS } from './claude-agent-sdk-control-requests'
|
||||
@@ -40,7 +40,8 @@ import {
|
||||
} from './claude-structured-session-exit-lifecycle'
|
||||
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
|
||||
import { resolveClaudeProviderHistoryWindow } from './claude-structured-history-window'
|
||||
import { claudePromptCardWritten, drainClaudeChildWork } from './claude-child-work-evidence'
|
||||
import { drainClaudeChildWork } from './claude-child-work-evidence'
|
||||
import { emitClaudeStructuredSessionEvent } from './claude-structured-event-delivery'
|
||||
import {
|
||||
answerClaudeStructuredPrompt,
|
||||
cancelClaudeStructuredTurn,
|
||||
@@ -130,14 +131,8 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
|
||||
/** Resolves once a published session's startup has landed, faulted, or been ended by a close;
|
||||
* with the reason when it did not land. */
|
||||
awaitStarted = async (sessionId: string): Promise<void | SubmissionRejectionFact> => {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (!session) {
|
||||
return
|
||||
}
|
||||
await session.startup.settled
|
||||
return claudeStartupFailureFact(session) ?? undefined
|
||||
}
|
||||
awaitStarted = (sessionId: string): Promise<void | SubmissionRejectionFact> =>
|
||||
awaitClaudeSessionStarted(this.sessions.get(sessionId))
|
||||
|
||||
/** Restart reconciliation reads the transcript a resume replays; these maps track liveness. */
|
||||
providerHistoryWindow: NonNullable<StructuredAgentSessionAdapter['providerHistoryWindow']> = (
|
||||
@@ -151,27 +146,12 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
})
|
||||
|
||||
private emit(session: ClaudeSession | null, event: ClaudeStructuredSessionEvent): void {
|
||||
// The host's child records, fed by the decoder's evidence drained below, are what every surface
|
||||
// and every Stop reads; the tracker's roster is kept only for tests that compare the two.
|
||||
if (event.type === 'ended') {
|
||||
session?.childWork.clear()
|
||||
session?.backgroundTasks.clear()
|
||||
} else if (event.type === 'message') {
|
||||
session?.childWork.observe(event.message)
|
||||
session?.backgroundTasks.observe(event.message, event.startsTurn === true)
|
||||
} else if (event.type === 'prompt-cancelled') {
|
||||
// A withdrawn request frees its child before its card closes: the journal may take that
|
||||
// write, and publish it, as it is submitted.
|
||||
this.publishChildWork(event.sessionId, session)
|
||||
}
|
||||
if (event.type === 'message' && session?.commands.observe(event.message)) {
|
||||
session.events?.publish()
|
||||
}
|
||||
session?.translator?.handle(event)
|
||||
this.deps.onEvent?.(event)
|
||||
this.publishChildWork(event.sessionId, session, event.type === 'message' ? event.message : null)
|
||||
// A subagent's card holds it waiting only once its row is written: its wait goes out after.
|
||||
void claudePromptCardWritten(session, event)?.then(() => this.publishChildWork(event.sessionId))
|
||||
emitClaudeStructuredSessionEvent({
|
||||
session,
|
||||
event,
|
||||
deps: this.deps,
|
||||
publishChildWork: (id, child, message) => this.publishChildWork(id, child, message)
|
||||
})
|
||||
}
|
||||
|
||||
/** After the journal handled the frame, which republished the parent's own row: the host never
|
||||
|
||||
@@ -13,10 +13,21 @@ import {
|
||||
import type { AgentChildWorkEvidence } from '../../shared/agent-status-child-work-evidence'
|
||||
import { AgentSessionAcquisitionRootExitObservedError } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
|
||||
import { ClaudePromptRegistry } from './claude-structured-prompt-replies'
|
||||
import { closeClaudeSession } from './claude-structured-session-close'
|
||||
import { claudeRootExitObserved, closeClaudeSession } from './claude-structured-session-close'
|
||||
import { ClaudeAcquisitionRegistry } from './claude-structured-session-state'
|
||||
|
||||
describe('Claude published session close lifecycle', () => {
|
||||
it('never reads a failed spawn as an observed root exit, whatever checks it first', async () => {
|
||||
const claude = fakeClaude()
|
||||
const adapter = adapterFor(claude)
|
||||
await adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
|
||||
const connection = claude.connections[0]!
|
||||
connection.exitVerdict = { root: 'exited', tree: 'exited', processless: true }
|
||||
expect(claudeRootExitObserved(connection)).toBe(false)
|
||||
connection.exitVerdict = { root: 'exited', tree: 'unverifiable' }
|
||||
expect(claudeRootExitObserved(connection)).toBe(true)
|
||||
})
|
||||
|
||||
it('reports a proven root exit when published-session close cannot prove descendants', async () => {
|
||||
const claude = fakeClaude()
|
||||
const adapter = adapterFor(claude)
|
||||
|
||||
@@ -22,11 +22,12 @@ import { settleClaudeTurnEndWaiters } from './claude-request-end-wait'
|
||||
import type { StructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
|
||||
/** The root's own exit was seen first-hand. The lease follows the root, so a descendant
|
||||
* left unverified or seen alive does not hold it. */
|
||||
* left unverified or seen alive does not hold it. A failed spawn had no process to exit. */
|
||||
export function claudeRootExitObserved(
|
||||
connection: ClaudeStreamJsonConnection | null | undefined
|
||||
): boolean {
|
||||
return connection?.exitVerdict.root === 'exited'
|
||||
const verdict = connection?.exitVerdict
|
||||
return verdict?.root === 'exited' && verdict.processless !== true
|
||||
}
|
||||
|
||||
export function claudeAcquisitionCleanupError(
|
||||
@@ -34,7 +35,7 @@ export function claudeAcquisitionCleanupError(
|
||||
cause: unknown
|
||||
): Error {
|
||||
const verdict = connection?.exitVerdict
|
||||
if (verdict?.root === 'processless') {
|
||||
if (verdict?.processless === true) {
|
||||
return new AgentSessionPreSpawnError(cause)
|
||||
}
|
||||
return claudeRootExitObserved(connection)
|
||||
@@ -57,7 +58,7 @@ export async function resolveClaudeAcquisitionError(input: {
|
||||
prompt.settle(null)
|
||||
}
|
||||
const closed = (await input.attempt.connection?.close()) ?? true
|
||||
if (input.attempt.connection?.exitVerdict.root === 'processless') {
|
||||
if (input.attempt.connection?.exitVerdict.processless === true) {
|
||||
acquisitionError = new AgentSessionPreSpawnError(input.error)
|
||||
} else if (!closed) {
|
||||
acquisitionError = claudeAcquisitionCleanupError(input.attempt.connection, input.error)
|
||||
|
||||
@@ -31,6 +31,16 @@ export function claudeStartupFailureFact(session: ClaudeSession): SubmissionReje
|
||||
: null
|
||||
}
|
||||
|
||||
export async function awaitClaudeSessionStarted(
|
||||
session: ClaudeSession | undefined
|
||||
): Promise<void | SubmissionRejectionFact> {
|
||||
if (!session) {
|
||||
return
|
||||
}
|
||||
await session.startup.settled
|
||||
return claudeStartupFailureFact(session) ?? undefined
|
||||
}
|
||||
|
||||
/** Resolves when startup lands or `timeoutMs` passes; a stuck start then refuses the write as before. */
|
||||
export function claudeStartupSettledWithin(
|
||||
session: ClaudeSession | undefined,
|
||||
|
||||
@@ -140,22 +140,16 @@ async function spawnClaude(env: Record<string, string> = {}) {
|
||||
const options = sdkOptions(env)
|
||||
const child = spawner.spawn(options)
|
||||
recordedPids.add(child.pid!)
|
||||
let exited = false
|
||||
const managed = spawner.managed
|
||||
if (!managed) {
|
||||
throw new Error('Claude spawner did not retain its managed child')
|
||||
}
|
||||
const exit = new Promise<{ code: number | null; signal: NodeJS.Signals | null }>((resolve) =>
|
||||
child.once('exit', (code, signal) => {
|
||||
exited = true
|
||||
resolve({ code, signal })
|
||||
})
|
||||
managed.onExit(resolve)
|
||||
)
|
||||
const pids = await readPids(child, ['claude', 'tool', 'daemon'])
|
||||
const close = (): Promise<boolean> =>
|
||||
proveClaudeChildExit({
|
||||
child,
|
||||
exitPromise: exit.then(() => undefined),
|
||||
exited: () => exited,
|
||||
supervised: spawner.supervised
|
||||
})
|
||||
return { child, exit, pids, close, marker: String(options.env.ORCA_TEST_SIGTERM_MARKER) }
|
||||
const close = (): Promise<boolean> => proveClaudeChildExit({ managed })
|
||||
return { child, managed, exit, pids, close, marker: String(options.env.ORCA_TEST_SIGTERM_MARKER) }
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
@@ -172,31 +166,37 @@ afterEach(() => {
|
||||
|
||||
describe.runIf(process.platform !== 'win32')('Claude under the POSIX provider supervisor', () => {
|
||||
it('stops a mid-turn Claude on close instead of letting stdin end finish its turn', async () => {
|
||||
const { child, exit, pids, close, marker } = await spawnClaude()
|
||||
const { child, managed, exit, pids, close, marker } = await spawnClaude()
|
||||
expect(child.pid).not.toBe(pids.claude)
|
||||
|
||||
const startedAt = Date.now()
|
||||
await expect(close()).resolves.toBe(true)
|
||||
expect(managed.lastCloseResult).toEqual({ root: 'exited', tree: 'exited' })
|
||||
await expect(close()).resolves.toBe(true)
|
||||
|
||||
// Claude's own SIGTERM reap ran at once, not after the supervisor's stdin-end grace.
|
||||
expect(Date.now() - startedAt).toBeLessThan(PROVIDER_STDIN_END_GRACE_MS)
|
||||
expect(existsSync(marker)).toBe(true)
|
||||
await expect(exit).resolves.toEqual({ code: null, signal: 'SIGTERM' })
|
||||
await expect(exit).resolves.toEqual({ code: null, signal: 'SIGTERM', processless: false })
|
||||
expect(alive(pids.claude)).toBe(false)
|
||||
expect(alive(pids.tool)).toBe(false)
|
||||
})
|
||||
|
||||
it('lets the supervisor escalate a Claude that ignores SIGTERM, and exits only after it', async () => {
|
||||
const { exit, pids, close } = await spawnClaude({ ORCA_TEST_CLAUDE_IGNORES_SIGTERM: '1' })
|
||||
const { managed, exit, pids, close } = await spawnClaude({
|
||||
ORCA_TEST_CLAUDE_IGNORES_SIGTERM: '1'
|
||||
})
|
||||
|
||||
const startedAt = Date.now()
|
||||
await expect(close()).resolves.toBe(true)
|
||||
expect(managed.lastCloseResult).toEqual({ root: 'exited', tree: 'exited' })
|
||||
await expect(close()).resolves.toBe(true)
|
||||
|
||||
const elapsed = Date.now() - startedAt
|
||||
expect(elapsed).toBeGreaterThanOrEqual(PROVIDER_SIGTERM_GRACE_MS)
|
||||
expect(elapsed).toBeLessThan(PROVIDER_SUPERVISOR_MAX_STOP_MS + 1_000)
|
||||
// The supervisor's own SIGTERM stop finished the job; nothing forced the supervisor itself.
|
||||
await expect(exit).resolves.toEqual({ code: null, signal: 'SIGTERM' })
|
||||
await expect(exit).resolves.toEqual({ code: null, signal: 'SIGTERM', processless: false })
|
||||
expect(alive(pids.claude)).toBe(false)
|
||||
})
|
||||
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import type { ProviderProcessTeardownVerdict } from '../provider-process/provider-process-teardown'
|
||||
import { openCodexAppServerConnection } from './codex-app-server-connection'
|
||||
|
||||
const teardown = vi.hoisted(() => {
|
||||
const state: { verdict: ProviderProcessTeardownVerdict; rootExits: () => void } = {
|
||||
verdict: null,
|
||||
rootExits: () => {}
|
||||
}
|
||||
return state
|
||||
})
|
||||
vi.mock('../provider-process/provider-process-teardown', () => ({
|
||||
terminateProviderProcessTree: vi.fn(async () => {
|
||||
teardown.rootExits()
|
||||
return teardown.verdict
|
||||
})
|
||||
}))
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
function stubChild() {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: 9_999_999,
|
||||
stdin: new PassThrough(),
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
kill: vi.fn(() => true)
|
||||
})
|
||||
child.stdin.once('data', () => {
|
||||
child.stdout.write(`${JSON.stringify({ id: 1, result: {} })}\n`)
|
||||
})
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The connection reads only events, pid, streams and kill from this stub.
|
||||
const spawnImpl = (() => child) as unknown as typeof spawnProcess
|
||||
return { child, spawnImpl }
|
||||
}
|
||||
|
||||
describe('Codex process-tree diagnostic after a forced close', () => {
|
||||
// The flag keeps its meaning: a forced teardown that did not prove the descendants gone.
|
||||
it.each([
|
||||
['unverifiable', true],
|
||||
['live', true],
|
||||
['exited', false],
|
||||
[null, false]
|
||||
] as const)('teardown observing %s reads unproven=%s', async (verdict, unproven) => {
|
||||
vi.useFakeTimers()
|
||||
teardown.verdict = verdict
|
||||
const { child, spawnImpl } = stubChild()
|
||||
teardown.rootExits = () => child.emit('exit', null, 'SIGKILL')
|
||||
const connection = await openCodexAppServerConnection(
|
||||
{ command: 'codex', args: ['app-server'] },
|
||||
{},
|
||||
spawnImpl
|
||||
)
|
||||
const closing = connection.close()
|
||||
await vi.advanceTimersByTimeAsync(10_000)
|
||||
await expect(closing).resolves.toBe(true)
|
||||
expect(connection.processTreeUnproven).toBe(unproven)
|
||||
// A repeat close answers from the memo and keeps the diagnostic.
|
||||
await expect(connection.close()).resolves.toBe(true)
|
||||
expect(connection.processTreeUnproven).toBe(unproven)
|
||||
})
|
||||
})
|
||||
@@ -431,6 +431,7 @@ describe('openCodexAppServerConnection', () => {
|
||||
|
||||
expect(error.name).toBe('CodexAppServerHandshakeExitUnprovenError')
|
||||
expect(error.connection).toBeDefined()
|
||||
child.emit('exit', 1, null)
|
||||
child.emit('close', 1, null)
|
||||
await expect(error.connection?.close()).resolves.toBe(true)
|
||||
})
|
||||
|
||||
@@ -1,15 +1,9 @@
|
||||
import { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { RetryableProcessExitProof } from '../../shared/child-process/retryable-process-exit-proof'
|
||||
import { spawnManagedProviderProcess } from '../provider-process/managed-provider-process'
|
||||
import type { ProviderProcessLaunch } from '../provider-process/provider-process-launch'
|
||||
import {
|
||||
createProviderSpawnSpec,
|
||||
PROVIDER_SUPERVISOR_MAX_STOP_MS
|
||||
} from '../provider-process/provider-process-supervisor'
|
||||
import { buildCodexAppServerExitError } from './codex-app-server-exit-error'
|
||||
import { initializeCodexAppServerConnection } from './codex-app-server-handshake'
|
||||
import { CodexAppServerHandshakeExitUnprovenError } from './codex-app-server-handshake-exit-proof'
|
||||
import { terminateProviderProcessTree } from '../provider-process/provider-process-teardown'
|
||||
import { waitForProcessExitUntil } from '../provider-process/provider-process-exit-deadline'
|
||||
import {
|
||||
CodexAppServerTimeoutError,
|
||||
CodexAppServerUnsupportedError
|
||||
@@ -31,6 +25,7 @@ export {
|
||||
isCodexAppServerRequestError
|
||||
} from './codex-app-server-request-error'
|
||||
export { CodexAppServerFrameSizeError } from './codex-app-server-frame-size-error'
|
||||
export { ROOT_ONLY_GRACEFUL_EXIT_MS as GRACEFUL_EXIT_MS } from '../provider-process/provider-process-close'
|
||||
|
||||
// Structured chat needs a persistent bidirectional child and per-request deadlines;
|
||||
// the request-scoped app-server runner cannot carry approvals or streamed turns.
|
||||
@@ -38,9 +33,6 @@ export { CodexAppServerFrameSizeError } from './codex-app-server-frame-size-erro
|
||||
export type CodexAppServerLaunch = ProviderProcessLaunch
|
||||
|
||||
const DEFAULT_REQUEST_TIMEOUT_MS = 30_000
|
||||
export const GRACEFUL_EXIT_MS = 1_500
|
||||
const FORCED_EXIT_MS = 1_000
|
||||
const STDERR_TAIL_MAX_BYTES = 8192
|
||||
|
||||
/**
|
||||
* Spawns `codex app-server`, completes the initialize handshake, and returns a
|
||||
@@ -52,46 +44,22 @@ export async function openCodexAppServerConnection(
|
||||
handlers: CodexAppServerConnectionHandlers = {},
|
||||
spawnImpl: typeof spawnProcess = spawnProcess
|
||||
): Promise<CodexAppServerConnection> {
|
||||
const spawnSpec = createProviderSpawnSpec(launch, process.env, process.platform)
|
||||
const child = spawnImpl(spawnSpec)
|
||||
const managed = spawnManagedProviderProcess(launch, {
|
||||
spawnImpl,
|
||||
site: 'codex-app-server-teardown'
|
||||
})
|
||||
const { child, terminateTree: terminateProcessTree } = managed
|
||||
|
||||
function terminateProcessTree(): Promise<boolean> {
|
||||
// The supervisor and provider own separate POSIX groups so the supervisor can prove the
|
||||
// provider group empty before relaying its exit. Forced wrapper teardown uses descendant proof.
|
||||
return terminateProviderProcessTree(child, { site: 'codex-app-server-teardown' })
|
||||
}
|
||||
|
||||
let stderrTail = ''
|
||||
let nextRequestId = 1
|
||||
let exited = false
|
||||
let exitObserved = false
|
||||
let closing = false
|
||||
let exitReported = false
|
||||
let processTreeUnproven = false
|
||||
const exitProof = new RetryableProcessExitProof()
|
||||
/** First terminal cause, or null while the transport is still usable. Set once:
|
||||
* a child that dies reaches us through several listeners, and the specific
|
||||
* first cause is the one worth reporting. */
|
||||
let terminalError: Error | null = null
|
||||
|
||||
let resolveExit = (): void => undefined
|
||||
const exitPromise = new Promise<void>((resolve) => {
|
||||
resolveExit = resolve
|
||||
})
|
||||
|
||||
function observeExit(): void {
|
||||
exited = true
|
||||
exitObserved = true
|
||||
resolveExit()
|
||||
}
|
||||
|
||||
child.on('exit', () => {
|
||||
observeExit()
|
||||
handleUnexpectedEnd()
|
||||
})
|
||||
|
||||
function buildExitError(cause?: Error): Error {
|
||||
return buildCodexAppServerExitError(stderrTail, cause)
|
||||
return buildCodexAppServerExitError(managed.stderrTail(), cause)
|
||||
}
|
||||
|
||||
const dispatcher = createCodexAppServerRecordDispatcher({
|
||||
@@ -115,7 +83,7 @@ export async function openCodexAppServerConnection(
|
||||
// Transport/protocol failures make the connection unusable immediately so
|
||||
// callers do not hang, but recovery must not treat that as a child exit
|
||||
// until the execution host has observed `exit`/`close`.
|
||||
if (exitObserved && !exitReported) {
|
||||
if (managed.rootVerdict === 'exited' && !exitReported) {
|
||||
exitReported = true
|
||||
handlers.onExit?.(terminalError, { expected: closing })
|
||||
}
|
||||
@@ -124,13 +92,7 @@ export async function openCodexAppServerConnection(
|
||||
child.on('error', (error) => {
|
||||
handleUnexpectedEnd(error)
|
||||
})
|
||||
child.on('close', () => {
|
||||
observeExit()
|
||||
handleUnexpectedEnd()
|
||||
})
|
||||
child.stderr.setEncoding('utf8').on('data', (chunk: string) => {
|
||||
stderrTail = (stderrTail + chunk).slice(-STDERR_TAIL_MAX_BYTES)
|
||||
})
|
||||
managed.onExit(() => handleUnexpectedEnd())
|
||||
child.stdin.on('error', (error) => {
|
||||
// A broken pipe is terminal, not one failed write: every later request can
|
||||
// only error or time out, so the session must learn its lease is worthless
|
||||
@@ -172,7 +134,7 @@ export async function openCodexAppServerConnection(
|
||||
}
|
||||
|
||||
function notify(method: string, params?: Record<string, unknown>): void {
|
||||
if (exited || terminalError) {
|
||||
if (managed.rootVerdict === 'exited' || terminalError) {
|
||||
return
|
||||
}
|
||||
try {
|
||||
@@ -193,7 +155,7 @@ export async function openCodexAppServerConnection(
|
||||
if (terminalError) {
|
||||
return Promise.reject(terminalError)
|
||||
}
|
||||
if (exited) {
|
||||
if (managed.rootVerdict === 'exited') {
|
||||
return Promise.reject(buildExitError())
|
||||
}
|
||||
const id = nextRequestId++
|
||||
@@ -217,7 +179,12 @@ export async function openCodexAppServerConnection(
|
||||
}
|
||||
|
||||
function writeResponse(payload: Record<string, unknown>): void {
|
||||
if (exited || terminalError || child.stdin.destroyed || !child.stdin.writable) {
|
||||
if (
|
||||
managed.rootVerdict === 'exited' ||
|
||||
terminalError ||
|
||||
child.stdin.destroyed ||
|
||||
!child.stdin.writable
|
||||
) {
|
||||
return
|
||||
}
|
||||
try {
|
||||
@@ -228,32 +195,10 @@ export async function openCodexAppServerConnection(
|
||||
}
|
||||
|
||||
function close(): Promise<boolean> {
|
||||
if (exitObserved) {
|
||||
return Promise.resolve(true)
|
||||
}
|
||||
closing = true
|
||||
return exitProof.run(async () => {
|
||||
try {
|
||||
child.stdin.end()
|
||||
} catch {
|
||||
// Already destroyed; the reap below still runs.
|
||||
}
|
||||
if (!exited) {
|
||||
// The POSIX supervisor stops its own provider group; forcing it any sooner can orphan it.
|
||||
await waitForProcessExitUntil(
|
||||
exitPromise,
|
||||
process.platform === 'win32' ? GRACEFUL_EXIT_MS : PROVIDER_SUPERVISOR_MAX_STOP_MS
|
||||
)
|
||||
if (!exited) {
|
||||
const treeExited = await terminateProcessTree()
|
||||
await waitForProcessExitUntil(exitPromise, FORCED_EXIT_MS)
|
||||
// The lease follows the root, which is gone: a child left behind is reported by the
|
||||
// owner, and blocks nothing.
|
||||
processTreeUnproven = !treeExited && exitObserved
|
||||
}
|
||||
}
|
||||
closing ||= managed.rootVerdict !== 'exited'
|
||||
return managed.close().then((result) => {
|
||||
dispatcher.failPending(new Error('codex app-server connection closed'))
|
||||
return exitObserved
|
||||
return result.root === 'exited'
|
||||
})
|
||||
}
|
||||
|
||||
@@ -262,10 +207,13 @@ export async function openCodexAppServerConnection(
|
||||
return child.pid
|
||||
},
|
||||
get closed() {
|
||||
return closing || exited || terminalError !== null
|
||||
return closing || managed.rootVerdict === 'exited' || terminalError !== null
|
||||
},
|
||||
get processTreeUnproven() {
|
||||
return processTreeUnproven
|
||||
const tree = managed.lastCloseResult?.tree
|
||||
return (
|
||||
managed.lastCloseResult?.root === 'exited' && (tree === 'unverifiable' || tree === 'live')
|
||||
)
|
||||
},
|
||||
request,
|
||||
notify,
|
||||
|
||||
@@ -10,6 +10,7 @@ const processWork = vi.hoisted(() => {
|
||||
capturedAtMs: 0
|
||||
})),
|
||||
terminateDescendantSnapshotAndWait: vi.fn(never),
|
||||
terminateDescendantSnapshotWithVerdict: vi.fn(never),
|
||||
queryWindowsProcessDescendants: vi.fn(never),
|
||||
terminateWindowsProcessTree: vi.fn(never)
|
||||
}
|
||||
@@ -20,7 +21,8 @@ vi.mock('../pty-descendant-termination', async (importOriginal) => ({
|
||||
}))
|
||||
vi.mock('../pty-descendant-exit-verification', async (importOriginal) => ({
|
||||
...(await importOriginal<Record<string, unknown>>()),
|
||||
terminateDescendantSnapshotAndWait: processWork.terminateDescendantSnapshotAndWait
|
||||
terminateDescendantSnapshotAndWait: processWork.terminateDescendantSnapshotAndWait,
|
||||
terminateDescendantSnapshotWithVerdict: processWork.terminateDescendantSnapshotWithVerdict
|
||||
}))
|
||||
vi.mock('../providers/windows-foreground-process-rows', async (importOriginal) => ({
|
||||
...(await importOriginal<Record<string, unknown>>()),
|
||||
|
||||
@@ -0,0 +1,76 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import type { DescendantSnapshot } from '../pty-descendant-termination'
|
||||
import { spawnManagedProviderProcess } from './managed-provider-process'
|
||||
|
||||
// The real fallback teardown; only the primitives that touch the OS are faked.
|
||||
const os = vi.hoisted(() => ({
|
||||
taskkill: vi.fn(async () => {}),
|
||||
capture: vi.fn(async (): Promise<DescendantSnapshot | null> => null),
|
||||
verifySnapshot: vi.fn(async () => 'exited' as const)
|
||||
}))
|
||||
vi.mock('../windows-process-tree-kill', () => ({ terminateWindowsProcessTree: os.taskkill }))
|
||||
vi.mock('../pty-descendant-termination', () => ({ captureDescendantSnapshot: os.capture }))
|
||||
vi.mock('../pty-descendant-exit-verification', () => ({
|
||||
terminateDescendantSnapshotWithVerdict: os.verifySnapshot
|
||||
}))
|
||||
vi.mock('../crash-reporting/self-initiated-tree-kill-log', () => ({
|
||||
recordSelfInitiatedTreeKill: vi.fn()
|
||||
}))
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
vi.clearAllMocks()
|
||||
})
|
||||
|
||||
function rootOnly(platform: NodeJS.Platform) {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: 4242,
|
||||
stdin: new PassThrough(),
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
// Only the root dies to SIGKILL; nothing in these paths examines a descendant.
|
||||
kill: vi.fn((signal?: NodeJS.Signals | number) => {
|
||||
if (signal === 'SIGKILL') {
|
||||
queueMicrotask(() => child.emit('exit', null, 'SIGKILL'))
|
||||
}
|
||||
return true
|
||||
})
|
||||
})
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The managed lifecycle reads only events, pid, streams and kill from this fixture.
|
||||
const spawnImpl = (() => child) as unknown as typeof spawnProcess
|
||||
return spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{ spawnImpl, platform, site: 'fixture-provider-teardown' }
|
||||
)
|
||||
}
|
||||
|
||||
describe('fallback teardown never claims descendants it did not observe', () => {
|
||||
it('reports no observation after a Windows tree kill, whose outcome is unreadable', async () => {
|
||||
vi.useFakeTimers()
|
||||
const closing = rootOnly('win32').close()
|
||||
await vi.advanceTimersByTimeAsync(1_500)
|
||||
await expect(closing).resolves.toEqual({ root: 'exited', tree: null })
|
||||
expect(os.taskkill).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('reports no observation when the POSIX process table cannot be read', async () => {
|
||||
vi.useFakeTimers()
|
||||
const closing = rootOnly('darwin').close()
|
||||
await vi.advanceTimersByTimeAsync(5_500)
|
||||
await expect(closing).resolves.toEqual({ root: 'exited', tree: null })
|
||||
expect(os.capture).toHaveBeenCalledOnce()
|
||||
expect(os.verifySnapshot).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('reports exited only when the captured descendants were verified gone', async () => {
|
||||
vi.useFakeTimers()
|
||||
os.capture.mockResolvedValueOnce({ rootPgid: 1, descendants: [], capturedAtMs: 1 })
|
||||
const closing = rootOnly('darwin').close()
|
||||
await vi.advanceTimersByTimeAsync(5_500)
|
||||
await expect(closing).resolves.toEqual({ root: 'exited', tree: 'exited' })
|
||||
expect(os.verifySnapshot).toHaveBeenCalledOnce()
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,73 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { spawnManagedProviderProcess } from './managed-provider-process'
|
||||
import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from './provider-process-supervisor'
|
||||
|
||||
function fixture() {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: 9_999_999,
|
||||
stdin: new PassThrough(),
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
kill: vi.fn(() => true)
|
||||
})
|
||||
const spawn = vi.fn<typeof spawnProcess>(() => {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The fixture supplies all process fields the managed lifecycle consumes.
|
||||
return child as unknown as ReturnType<typeof spawnProcess>
|
||||
})
|
||||
return { child, spawn }
|
||||
}
|
||||
|
||||
describe('managed provider supervisor grace', () => {
|
||||
it.each([false, true])(
|
||||
'rejects a short grace before spawning (signal on close: %s)',
|
||||
(signalSupervisorOnClose) => {
|
||||
const { spawn } = fixture()
|
||||
expect(() =>
|
||||
spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{
|
||||
spawnImpl: spawn,
|
||||
platform: 'linux',
|
||||
site: 'fixture',
|
||||
policy: () => ({
|
||||
gracefulExitMs: PROVIDER_SUPERVISOR_MAX_STOP_MS - 1,
|
||||
forcedExitMs: 50,
|
||||
signalSupervisorOnClose
|
||||
}),
|
||||
acceptClose: (result) => result.root === 'exited'
|
||||
}
|
||||
)
|
||||
).toThrow(
|
||||
`Supervised provider graceful exit must wait at least ${PROVIDER_SUPERVISOR_MAX_STOP_MS} ms`
|
||||
)
|
||||
expect(spawn).not.toHaveBeenCalled()
|
||||
}
|
||||
)
|
||||
|
||||
it.each(['linux', 'win32'] as const)(
|
||||
'accepts the floor on POSIX and the shorter Windows grace (%s)',
|
||||
(platform) => {
|
||||
const { child, spawn } = fixture()
|
||||
const managed = spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{
|
||||
spawnImpl: spawn,
|
||||
platform,
|
||||
site: 'fixture',
|
||||
policy: (supervised) => ({
|
||||
gracefulExitMs: supervised ? PROVIDER_SUPERVISOR_MAX_STOP_MS : 100,
|
||||
forcedExitMs: 50
|
||||
}),
|
||||
acceptClose: (result) => result.root === 'exited'
|
||||
}
|
||||
)
|
||||
expect(spawn).toHaveBeenCalledOnce()
|
||||
expect(managed.rootVerdict).toBe('live')
|
||||
child.emit('exit', 0, null)
|
||||
expect(managed.rootVerdict).toBe('exited')
|
||||
}
|
||||
)
|
||||
})
|
||||
@@ -0,0 +1,123 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { spawnManagedProviderProcess } from './managed-provider-process'
|
||||
import { ROOT_ONLY_GRACEFUL_EXIT_MS } from './provider-process-close'
|
||||
import type { ProviderProcessTeardownVerdict } from './provider-process-teardown'
|
||||
|
||||
const teardown = vi.hoisted(() => ({
|
||||
terminate: vi.fn(async (): Promise<ProviderProcessTeardownVerdict> => 'exited')
|
||||
}))
|
||||
vi.mock('./provider-process-teardown', () => ({
|
||||
terminateProviderProcessTree: teardown.terminate
|
||||
}))
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
vi.clearAllMocks()
|
||||
})
|
||||
|
||||
function fakeChild(pid: number | null = 9_999_999) {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: pid ?? undefined,
|
||||
stdin: new PassThrough(),
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
kill: vi.fn(() => true)
|
||||
})
|
||||
const spawn = vi.fn<typeof spawnProcess>(() => {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The managed lifecycle reads only events, pid, streams and kill from this fixture.
|
||||
return child as unknown as ReturnType<typeof spawnProcess>
|
||||
})
|
||||
return { child, spawn }
|
||||
}
|
||||
|
||||
/** A provider with no reaper of its own takes every close default. */
|
||||
function rootOnly(fixture: ReturnType<typeof fakeChild>) {
|
||||
return spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{ spawnImpl: fixture.spawn, platform: 'win32', site: 'fixture-provider-teardown' }
|
||||
)
|
||||
}
|
||||
|
||||
describe('root-only managed provider close', () => {
|
||||
it('reports no descendant observation when the root leaves on stdin end', async () => {
|
||||
const fixture = fakeChild()
|
||||
fixture.child.stdin.once('finish', () => fixture.child.emit('exit', 0, null))
|
||||
const managed = rootOnly(fixture)
|
||||
await expect(managed.close()).resolves.toEqual({ root: 'exited', tree: null })
|
||||
expect(teardown.terminate).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it.each(['exited', 'live', 'unverifiable', null] as const)(
|
||||
'writes what the fallback teardown observed (%s) as the tree verdict',
|
||||
async (tree) => {
|
||||
vi.useFakeTimers()
|
||||
teardown.terminate.mockImplementationOnce(async () => {
|
||||
fixture.child.emit('exit', null, 'SIGKILL')
|
||||
return tree
|
||||
})
|
||||
const fixture = fakeChild()
|
||||
const managed = rootOnly(fixture)
|
||||
const closing = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(ROOT_ONLY_GRACEFUL_EXIT_MS)
|
||||
await expect(closing).resolves.toEqual({ root: 'exited', tree })
|
||||
// The root is gone, so the close is done: a repeat answers from the memo, not a second teardown.
|
||||
await expect(managed.close()).resolves.toEqual({ root: 'exited', tree })
|
||||
expect(teardown.terminate).toHaveBeenCalledOnce()
|
||||
expect(managed.lastCloseResult).toEqual({ root: 'exited', tree })
|
||||
}
|
||||
)
|
||||
|
||||
it('waits the default root-only grace before forcing', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = rootOnly(fixture)
|
||||
void managed.close()
|
||||
await vi.advanceTimersByTimeAsync(ROOT_ONLY_GRACEFUL_EXIT_MS - 1)
|
||||
expect(teardown.terminate).not.toHaveBeenCalled()
|
||||
await vi.advanceTimersByTimeAsync(1)
|
||||
expect(teardown.terminate).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('records the already-exited answer when no close ran, without touching the child', async () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = rootOnly(fixture)
|
||||
fixture.child.emit('exit', 0, null)
|
||||
await expect(managed.close()).resolves.toEqual({ root: 'exited', tree: null })
|
||||
expect(managed.lastCloseResult).toEqual({ root: 'exited', tree: null })
|
||||
expect(fixture.child.stdin.writableEnded).toBe(false)
|
||||
expect(fixture.child.kill).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('drains stderr into a bounded tail', async () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = rootOnly(fixture)
|
||||
fixture.child.stderr.write('x'.repeat(9000))
|
||||
fixture.child.stderr.write('provider: not signed in')
|
||||
await new Promise((resolve) => setImmediate(resolve))
|
||||
expect(managed.stderrTail()).toMatch(/provider: not signed in$/)
|
||||
expect(managed.stderrTail()).toHaveLength(8192)
|
||||
})
|
||||
})
|
||||
|
||||
describe('root exit observation', () => {
|
||||
it('never reads a failed spawn as an observed root exit', () => {
|
||||
const fixture = fakeChild(null)
|
||||
const managed = rootOnly(fixture)
|
||||
fixture.child.emit('error', Object.assign(new Error('spawn ENOENT'), { code: 'ENOENT' }))
|
||||
fixture.child.emit('close', -2, null)
|
||||
expect(managed.rootVerdict).toBe('exited')
|
||||
expect(managed.processless).toBe(true)
|
||||
expect(managed.rootExitObserved).toBe(false)
|
||||
})
|
||||
|
||||
it('reads a real process exit as observed', () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = rootOnly(fixture)
|
||||
expect(managed.rootExitObserved).toBe(false)
|
||||
fixture.child.emit('exit', 1, null)
|
||||
expect(managed.rootExitObserved).toBe(true)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,290 @@
|
||||
import { EventEmitter } from 'node:events'
|
||||
import { PassThrough } from 'node:stream'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from './provider-process-supervisor'
|
||||
import { spawnManagedProviderProcess } from './managed-provider-process'
|
||||
import type { DescendantTreeVerdict } from '../pty-descendant-exit-verification'
|
||||
import type { ProviderProcessTree } from './provider-process-close'
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
capture: vi.fn(async () => null),
|
||||
windowsTree: vi.fn(async () => {})
|
||||
}))
|
||||
vi.mock('../pty-descendant-termination', () => ({ captureDescendantSnapshot: mocks.capture }))
|
||||
vi.mock('../windows-process-tree-kill', () => ({ terminateWindowsProcessTree: mocks.windowsTree }))
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
vi.clearAllMocks()
|
||||
})
|
||||
|
||||
function fakeChild(pid: number | null = 9_999_999) {
|
||||
const child = Object.assign(new EventEmitter(), {
|
||||
pid: pid ?? undefined,
|
||||
stdin: new PassThrough(),
|
||||
stdout: new PassThrough(),
|
||||
stderr: new PassThrough(),
|
||||
kill: vi.fn(() => true)
|
||||
})
|
||||
const spawn = vi.fn<typeof spawnProcess>(() => {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: The managed lifecycle reads only events, pid, streams and kill from this fixture.
|
||||
return child as unknown as ReturnType<typeof spawnProcess>
|
||||
})
|
||||
return { child, spawn }
|
||||
}
|
||||
|
||||
function launch(fixture: ReturnType<typeof fakeChild>, platform: NodeJS.Platform = 'win32') {
|
||||
return spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: ['serve'], cwd: '/workspace' },
|
||||
{
|
||||
spawnImpl: fixture.spawn,
|
||||
platform,
|
||||
site: 'fixture-provider-teardown',
|
||||
acceptClose: (result) => result.root === 'exited',
|
||||
policy: () => ({ gracefulExitMs: 100, forcedExitMs: 50 })
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
function fakeTree(initial: DescendantTreeVerdict = 'unverifiable') {
|
||||
let verdict = initial
|
||||
const tree: ProviderProcessTree = {
|
||||
capture: vi.fn(async () => {}),
|
||||
refresh: vi.fn(async () => {}),
|
||||
reap: vi.fn(async () => verdict),
|
||||
get treeVerdict() {
|
||||
return verdict
|
||||
}
|
||||
}
|
||||
return {
|
||||
tree,
|
||||
setVerdict: (next: DescendantTreeVerdict) => {
|
||||
verdict = next
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
describe('managed provider process', () => {
|
||||
it('applies launch environment and uses the supervisor on its execution platform', () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = spawnManagedProviderProcess(
|
||||
{
|
||||
command: 'fixture-provider',
|
||||
args: [],
|
||||
env: { AGENT_HOME: '/pinned' },
|
||||
envToDelete: ['SECRET']
|
||||
},
|
||||
{
|
||||
spawnImpl: fixture.spawn,
|
||||
platform: 'darwin',
|
||||
inheritedEnv: { SECRET: 'inherited', PATH: '/bin' },
|
||||
site: 'fixture',
|
||||
acceptClose: (result) => result.root === 'exited' && result.tree === 'exited',
|
||||
policy: () => ({ gracefulExitMs: PROVIDER_SUPERVISOR_MAX_STOP_MS, forcedExitMs: 50 })
|
||||
}
|
||||
)
|
||||
expect(managed.supervised).toBe(true)
|
||||
expect(fixture.spawn).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
detached: true,
|
||||
env: expect.objectContaining({ AGENT_HOME: '/pinned', PATH: '/bin' })
|
||||
})
|
||||
)
|
||||
expect(fixture.spawn.mock.calls[0][0].env).not.toHaveProperty('SECRET')
|
||||
expect(managed.rootVerdict).toBe('live')
|
||||
})
|
||||
|
||||
it('observes exit once across exit and close, including a late subscriber', async () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture)
|
||||
const onExit = vi.fn()
|
||||
managed.onExit(onExit)
|
||||
fixture.child.emit('exit', 7, 'SIGTERM')
|
||||
fixture.child.emit('close', 7, 'SIGTERM')
|
||||
fixture.child.emit('exit', 8, 'SIGKILL')
|
||||
await managed.exitPromise
|
||||
expect(onExit).toHaveBeenCalledExactlyOnceWith({
|
||||
code: 7,
|
||||
signal: 'SIGTERM',
|
||||
processless: false
|
||||
})
|
||||
const late = vi.fn()
|
||||
managed.onExit(late)
|
||||
expect(late).toHaveBeenCalledExactlyOnceWith({ code: 7, signal: 'SIGTERM', processless: false })
|
||||
await expect(managed.close()).resolves.toMatchObject({ root: 'exited' })
|
||||
expect(fixture.child.kill).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('proves stdin-close exit without forcing and clears the grace timer', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture)
|
||||
fixture.child.stdin.once('finish', () => fixture.child.emit('exit', 0, null))
|
||||
const closing = managed.close()
|
||||
expect(fixture.child.stdin.writableEnded).toBe(true)
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
await expect(closing).resolves.toMatchObject({ root: 'exited' })
|
||||
expect(fixture.child.kill).not.toHaveBeenCalled()
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
|
||||
it('joins a close, keeps root live after the kill deadline, and retries on the same child', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture)
|
||||
const first = managed.close()
|
||||
expect(managed.close()).toBe(first)
|
||||
await vi.advanceTimersByTimeAsync(150)
|
||||
await expect(first).resolves.toMatchObject({ root: 'live' })
|
||||
expect(fixture.child.kill).toHaveBeenCalledWith('SIGKILL')
|
||||
expect(managed.rootVerdict).toBe('live')
|
||||
fixture.child.kill.mockImplementation(() => {
|
||||
fixture.child.emit('exit', null, 'SIGKILL')
|
||||
return true
|
||||
})
|
||||
const second = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(150)
|
||||
await expect(second).resolves.toMatchObject({ root: 'exited' })
|
||||
expect(fixture.spawn).toHaveBeenCalledOnce()
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
|
||||
it('still reports a root exit once after an unconfirmed close', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture)
|
||||
const report = vi.fn()
|
||||
managed.onExit(report)
|
||||
const close = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(150)
|
||||
await expect(close).resolves.toMatchObject({ root: 'live' })
|
||||
fixture.child.emit('exit', 0, null)
|
||||
fixture.child.emit('close', 0, null)
|
||||
expect(report).toHaveBeenCalledOnce()
|
||||
await expect(managed.close()).resolves.toMatchObject({ root: 'exited' })
|
||||
})
|
||||
|
||||
it('requires error then close for a processless spawn to settle', async () => {
|
||||
const fixture = fakeChild(null)
|
||||
const managed = launch(fixture)
|
||||
expect(managed.rootVerdict).toBe('unverifiable')
|
||||
fixture.child.emit('close', -2, null)
|
||||
expect(managed.rootVerdict).toBe('unverifiable')
|
||||
fixture.child.emit('error', new Error('ENOENT'))
|
||||
expect(managed.rootVerdict).toBe('unverifiable')
|
||||
fixture.child.emit('close', null, null)
|
||||
expect(managed.processless).toBe(true)
|
||||
expect(managed.rootVerdict).toBe('exited')
|
||||
await expect(managed.close()).resolves.toMatchObject({ root: 'exited' })
|
||||
})
|
||||
|
||||
it('does not turn a transport close with an existing pid into Claude exit proof', () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture)
|
||||
fixture.child.emit('error', new Error('EPIPE'))
|
||||
fixture.child.emit('close', 0, null)
|
||||
expect(managed.rootVerdict).toBe('live')
|
||||
expect(managed.processless).toBe(false)
|
||||
})
|
||||
|
||||
it('uses Windows tree teardown and still requires observed root exit', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture, 'win32')
|
||||
expect(managed.supervised).toBe(false)
|
||||
const first = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(150)
|
||||
await expect(first).resolves.toMatchObject({ root: 'live' })
|
||||
expect(mocks.windowsTree).toHaveBeenCalledWith(fixture.child.pid, {
|
||||
site: 'fixture-provider-teardown'
|
||||
})
|
||||
expect(fixture.child.kill).not.toHaveBeenCalledWith('SIGTERM')
|
||||
fixture.child.emit('exit', null, 'SIGKILL')
|
||||
await expect(managed.close()).resolves.toMatchObject({ root: 'exited' })
|
||||
})
|
||||
|
||||
it('preserves a caller requiring tree proof, including live and unverifiable descendants', async () => {
|
||||
const fixture = fakeChild()
|
||||
const managed = spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{
|
||||
spawnImpl: fixture.spawn,
|
||||
platform: 'darwin',
|
||||
site: 'fixture',
|
||||
acceptClose: (result) => result.root === 'exited' && result.tree === 'exited',
|
||||
policy: () => ({ gracefulExitMs: PROVIDER_SUPERVISOR_MAX_STOP_MS, forcedExitMs: 50 })
|
||||
}
|
||||
)
|
||||
const proof = fakeTree()
|
||||
fixture.child.emit('exit', 0, null)
|
||||
await expect(managed.close(proof.tree)).resolves.toEqual({
|
||||
root: 'exited',
|
||||
tree: 'unverifiable'
|
||||
})
|
||||
proof.setVerdict('live')
|
||||
await expect(managed.close(proof.tree)).resolves.toEqual({ root: 'exited', tree: 'live' })
|
||||
proof.setVerdict('exited')
|
||||
await expect(managed.close(proof.tree)).resolves.toEqual({ root: 'exited', tree: 'exited' })
|
||||
await expect(managed.close(proof.tree)).resolves.toEqual({ root: 'exited', tree: 'exited' })
|
||||
expect(proof.tree.capture).toHaveBeenCalledTimes(3)
|
||||
})
|
||||
|
||||
it('reports failed tree cleanup separately after observed root exit', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture, 'win32')
|
||||
mocks.windowsTree.mockImplementationOnce(async () => {
|
||||
fixture.child.emit('exit', null, 'SIGKILL')
|
||||
throw new Error('tree teardown unavailable')
|
||||
})
|
||||
const close = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(150)
|
||||
await expect(close).resolves.toMatchObject({ root: 'exited' })
|
||||
expect(managed.lastCloseResult).toEqual({ root: 'exited', tree: 'unverifiable' })
|
||||
})
|
||||
|
||||
it('keeps the cleanup diagnostic tied to the close that finished before a late root exit', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = launch(fixture, 'win32')
|
||||
mocks.windowsTree.mockRejectedValueOnce(new Error('tree teardown unavailable'))
|
||||
const close = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(150)
|
||||
await expect(close).resolves.toMatchObject({ root: 'live' })
|
||||
fixture.child.emit('exit', null, 'SIGKILL')
|
||||
await expect(managed.close()).resolves.toMatchObject({ root: 'exited' })
|
||||
expect(managed.lastCloseResult).toEqual({ root: 'live', tree: 'unverifiable' })
|
||||
})
|
||||
|
||||
it('lets a supervised caller signal immediately and waits its configured grace before forcing', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fixture = fakeChild()
|
||||
const managed = spawnManagedProviderProcess(
|
||||
{ command: 'fixture-provider', args: [] },
|
||||
{
|
||||
spawnImpl: fixture.spawn,
|
||||
platform: 'darwin',
|
||||
site: 'fixture',
|
||||
acceptClose: (result) => result.root === 'exited' && result.tree === 'exited',
|
||||
policy: () => ({
|
||||
gracefulExitMs: PROVIDER_SUPERVISOR_MAX_STOP_MS + 500,
|
||||
forcedExitMs: 50,
|
||||
signalSupervisorOnClose: true
|
||||
})
|
||||
}
|
||||
)
|
||||
fixture.child.kill.mockImplementation(() => {
|
||||
setTimeout(() => fixture.child.emit('exit', 0, 'SIGTERM'), 550)
|
||||
return true
|
||||
})
|
||||
const close = managed.close()
|
||||
await vi.advanceTimersByTimeAsync(549)
|
||||
expect(fixture.child.kill).toHaveBeenCalledExactlyOnceWith('SIGTERM')
|
||||
expect(managed.rootVerdict).toBe('live')
|
||||
await vi.advanceTimersByTimeAsync(1)
|
||||
await expect(close).resolves.toMatchObject({ root: 'exited' })
|
||||
expect(fixture.child.kill).not.toHaveBeenCalledWith('SIGKILL')
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,168 @@
|
||||
import { spawnProcess } from '../../shared/child-process/run-process'
|
||||
import { RetryableProcessExitProof } from '../../shared/child-process/retryable-process-exit-proof'
|
||||
import type { ProviderProcessLaunch } from './provider-process-launch'
|
||||
import {
|
||||
PROVIDER_SUPERVISOR_MAX_STOP_MS,
|
||||
createProviderSpawnSpec
|
||||
} from './provider-process-supervisor'
|
||||
import {
|
||||
terminateProviderProcessTree,
|
||||
type ProviderProcessTeardownVerdict
|
||||
} from './provider-process-teardown'
|
||||
import type { DescendantTreeVerdict } from '../pty-descendant-exit-verification'
|
||||
import {
|
||||
acceptProviderRootExit,
|
||||
closeProviderProcess,
|
||||
rootOnlyProviderClosePolicy,
|
||||
type ProviderProcessClosePolicy,
|
||||
type ProviderProcessCloseResult,
|
||||
type ProviderProcessTree
|
||||
} from './provider-process-close'
|
||||
|
||||
const STDERR_TAIL_MAX_CHARS = 8192
|
||||
|
||||
export type ProviderProcessExit = {
|
||||
code: number | null
|
||||
signal: NodeJS.Signals | null
|
||||
processless: boolean
|
||||
}
|
||||
|
||||
type ManagedProviderProcessOptions = {
|
||||
site: string
|
||||
/** Defaults to the root-only policy; only a provider with its own reaper overrides it. */
|
||||
policy?: (supervised: boolean) => ProviderProcessClosePolicy
|
||||
spawnImpl?: typeof spawnProcess
|
||||
platform?: NodeJS.Platform
|
||||
inheritedEnv?: NodeJS.ProcessEnv
|
||||
/** Defaults to "the root is gone". */
|
||||
acceptClose?: (result: ProviderProcessCloseResult) => boolean
|
||||
}
|
||||
|
||||
export type ManagedProviderProcess = {
|
||||
child: ReturnType<typeof spawnProcess>
|
||||
supervised: boolean
|
||||
/** The spawn failed before a process existed: absence is proven, but no exit was observed. */
|
||||
readonly processless: boolean
|
||||
/** `exited` covers a processless child too; use `rootExitObserved` for "a real process exited". */
|
||||
readonly rootVerdict: DescendantTreeVerdict
|
||||
/** A process that existed was seen to exit; never true for a processless child. */
|
||||
readonly rootExitObserved: boolean
|
||||
/** The last close that ran the ladder; the already-exited answer only when none did. */
|
||||
readonly lastCloseResult: ProviderProcessCloseResult | null
|
||||
readonly exitPromise: Promise<void>
|
||||
/** The last 8 KiB of stderr, which the managed process drains so the child never blocks on it. */
|
||||
stderrTail(): string
|
||||
onExit(listener: (exit: ProviderProcessExit) => void): void
|
||||
terminateTree(): Promise<ProviderProcessTeardownVerdict>
|
||||
close(tree?: ProviderProcessTree): Promise<ProviderProcessCloseResult>
|
||||
}
|
||||
|
||||
/** One child owns its exit observation and every retry of an unconfirmed close. */
|
||||
export function spawnManagedProviderProcess(
|
||||
launch: ProviderProcessLaunch,
|
||||
options: ManagedProviderProcessOptions
|
||||
): ManagedProviderProcess {
|
||||
const platform = options.platform ?? process.platform
|
||||
const spec = createProviderSpawnSpec(launch, options.inheritedEnv ?? process.env, platform)
|
||||
const policy = (options.policy ?? rootOnlyProviderClosePolicy)(spec.supervised)
|
||||
if (spec.supervised && !(policy.gracefulExitMs >= PROVIDER_SUPERVISOR_MAX_STOP_MS)) {
|
||||
throw new RangeError(
|
||||
`Supervised provider graceful exit must wait at least ${PROVIDER_SUPERVISOR_MAX_STOP_MS} ms; received ${policy.gracefulExitMs} ms`
|
||||
)
|
||||
}
|
||||
const child = (options.spawnImpl ?? spawnProcess)({
|
||||
program: spec.program,
|
||||
args: spec.args,
|
||||
cwd: spec.cwd,
|
||||
env: spec.env,
|
||||
detached: spec.detached,
|
||||
stdio: ['pipe', 'pipe', 'pipe']
|
||||
})
|
||||
const listeners = new Set<(exit: ProviderProcessExit) => void>()
|
||||
let observed: ProviderProcessExit | null = null
|
||||
let spawnFailed = false
|
||||
let lastCloseResult: ProviderProcessCloseResult | null = null
|
||||
const exitProof = new RetryableProcessExitProof(options.acceptClose ?? acceptProviderRootExit)
|
||||
let stderrTail = ''
|
||||
// An undrained stderr pipe blocks the child once it fills.
|
||||
child.stderr.setEncoding('utf8').on('data', (chunk: string) => {
|
||||
stderrTail = (stderrTail + chunk).slice(-STDERR_TAIL_MAX_CHARS)
|
||||
})
|
||||
let resolveExit = (): void => {}
|
||||
const exitPromise = new Promise<void>((resolve) => {
|
||||
resolveExit = resolve
|
||||
})
|
||||
const observeExit = (exit: ProviderProcessExit): void => {
|
||||
if (observed) {
|
||||
return
|
||||
}
|
||||
observed = exit
|
||||
resolveExit()
|
||||
for (const listener of listeners) {
|
||||
listener(exit)
|
||||
}
|
||||
listeners.clear()
|
||||
}
|
||||
child.on('exit', (code, signal) => observeExit({ code, signal, processless: false }))
|
||||
child.on('error', () => {
|
||||
spawnFailed ||= child.pid === undefined
|
||||
})
|
||||
child.on('close', (code, signal) => {
|
||||
const processless = spawnFailed && child.pid === undefined
|
||||
if (processless) {
|
||||
observeExit({ code, signal, processless })
|
||||
}
|
||||
})
|
||||
const rootVerdict = (): DescendantTreeVerdict =>
|
||||
observed ? 'exited' : child.pid === undefined ? 'unverifiable' : 'live'
|
||||
const terminateTree = (): Promise<ProviderProcessTeardownVerdict> =>
|
||||
terminateProviderProcessTree(child, { site: options.site, platform })
|
||||
|
||||
return {
|
||||
child,
|
||||
supervised: spec.supervised,
|
||||
exitPromise,
|
||||
get processless() {
|
||||
return observed?.processless ?? false
|
||||
},
|
||||
get rootVerdict() {
|
||||
return rootVerdict()
|
||||
},
|
||||
get rootExitObserved() {
|
||||
return observed !== null && !observed.processless
|
||||
},
|
||||
stderrTail: () => stderrTail,
|
||||
get lastCloseResult() {
|
||||
return lastCloseResult
|
||||
},
|
||||
onExit(listener) {
|
||||
if (observed) {
|
||||
listener(observed)
|
||||
} else {
|
||||
listeners.add(listener)
|
||||
}
|
||||
},
|
||||
terminateTree,
|
||||
close(tree) {
|
||||
return exitProof.run(async () => {
|
||||
// The one already-exited guard. It observed nothing, so an earlier close's findings stay.
|
||||
if (observed && !tree) {
|
||||
const result: ProviderProcessCloseResult = { root: 'exited', tree: null }
|
||||
lastCloseResult ??= result
|
||||
return result
|
||||
}
|
||||
const result = await closeProviderProcess({
|
||||
child,
|
||||
exitPromise,
|
||||
rootVerdict,
|
||||
supervised: spec.supervised,
|
||||
policy,
|
||||
tree,
|
||||
terminateTree
|
||||
})
|
||||
lastCloseResult = result
|
||||
return result
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,87 @@
|
||||
import type { SpawnedProcess } from '../../shared/child-process/run-process'
|
||||
import type { DescendantTreeVerdict } from '../pty-descendant-exit-verification'
|
||||
import { waitForProcessExitUntil } from './provider-process-exit-deadline'
|
||||
import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from './provider-process-supervisor'
|
||||
import type { ProviderProcessTeardownVerdict } from './provider-process-teardown'
|
||||
|
||||
export type ProviderProcessTree = {
|
||||
capture(): Promise<void>
|
||||
refresh?: () => Promise<void>
|
||||
reap(): Promise<DescendantTreeVerdict>
|
||||
readonly treeVerdict: DescendantTreeVerdict
|
||||
}
|
||||
|
||||
export type ProviderProcessClosePolicy = {
|
||||
gracefulExitMs: number
|
||||
forcedExitMs: number
|
||||
signalSupervisorOnClose?: boolean
|
||||
}
|
||||
|
||||
export type ProviderProcessCloseInput = {
|
||||
child: Pick<SpawnedProcess, 'pid' | 'kill' | 'stdin'>
|
||||
exitPromise: Promise<void>
|
||||
rootVerdict: () => DescendantTreeVerdict
|
||||
supervised?: boolean
|
||||
policy: ProviderProcessClosePolicy
|
||||
tree?: ProviderProcessTree
|
||||
terminateTree: () => Promise<ProviderProcessTeardownVerdict>
|
||||
}
|
||||
|
||||
export type ProviderProcessCloseResult = {
|
||||
root: DescendantTreeVerdict
|
||||
/** Null when this close made no observation of the descendants. */
|
||||
tree: DescendantTreeVerdict | null
|
||||
}
|
||||
|
||||
export const ROOT_ONLY_GRACEFUL_EXIT_MS = 1_500
|
||||
const ROOT_ONLY_FORCED_EXIT_MS = 1_000
|
||||
|
||||
/** Default for providers without a descendant reaper: end stdin, wait, then the fallback teardown. */
|
||||
export function rootOnlyProviderClosePolicy(supervised: boolean): ProviderProcessClosePolicy {
|
||||
return {
|
||||
gracefulExitMs: supervised ? PROVIDER_SUPERVISOR_MAX_STOP_MS : ROOT_ONLY_GRACEFUL_EXIT_MS,
|
||||
forcedExitMs: ROOT_ONLY_FORCED_EXIT_MS
|
||||
}
|
||||
}
|
||||
|
||||
/** A root-only close is done once the root is gone; an unproven tree is reported, not retried. */
|
||||
export function acceptProviderRootExit(result: ProviderProcessCloseResult): boolean {
|
||||
return result.root === 'exited'
|
||||
}
|
||||
|
||||
/** The supervisor owns the POSIX signal ladder; its wrapper must outlive that ladder. */
|
||||
export async function closeProviderProcess(
|
||||
input: ProviderProcessCloseInput
|
||||
): Promise<ProviderProcessCloseResult> {
|
||||
const { child, policy, tree } = input
|
||||
if (tree) {
|
||||
await tree.capture()
|
||||
}
|
||||
try {
|
||||
child.stdin?.end()
|
||||
} catch {
|
||||
// A broken pipe still owes the reap.
|
||||
}
|
||||
if (input.supervised && policy.signalSupervisorOnClose && input.rootVerdict() !== 'exited') {
|
||||
child.kill('SIGTERM')
|
||||
}
|
||||
let reaped = false
|
||||
let fallbackTree: DescendantTreeVerdict | null = null
|
||||
if (input.rootVerdict() !== 'exited') {
|
||||
await waitForProcessExitUntil(input.exitPromise, policy.gracefulExitMs)
|
||||
if (input.rootVerdict() !== 'exited') {
|
||||
reaped = true
|
||||
await tree?.refresh?.()
|
||||
if (tree) {
|
||||
await tree.reap()
|
||||
} else {
|
||||
fallbackTree = await input.terminateTree()
|
||||
}
|
||||
await waitForProcessExitUntil(input.exitPromise, policy.forcedExitMs)
|
||||
}
|
||||
}
|
||||
if (!reaped && input.rootVerdict() === 'exited' && tree && tree.treeVerdict !== 'exited') {
|
||||
await tree.reap()
|
||||
}
|
||||
return { root: input.rootVerdict(), tree: tree ? tree.treeVerdict : fallbackTree }
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import {
|
||||
findSelfInitiatedTreeKills,
|
||||
resetSelfInitiatedTreeKillLogForTest
|
||||
} from '../crash-reporting/self-initiated-tree-kill-log'
|
||||
import type { DescendantTreeVerdict } from '../pty-descendant-exit-verification'
|
||||
import { terminateProviderProcessTree } from './provider-process-teardown'
|
||||
|
||||
/** Above pid_max on every supported POSIX host, so the group signal is a real ESRCH. */
|
||||
@@ -21,7 +22,7 @@ describe('terminateProviderProcessTree', () => {
|
||||
resetSelfInitiatedTreeKillLogForTest()
|
||||
})
|
||||
|
||||
it('waits for the Windows tree kill before releasing the wrapper', async () => {
|
||||
it('waits for the Windows tree kill before releasing the wrapper, and claims no observation', async () => {
|
||||
const target = child()
|
||||
const release = Promise.withResolvers<void>()
|
||||
const terminateWindowsTree = vi.fn(() => release.promise)
|
||||
@@ -33,7 +34,8 @@ describe('terminateProviderProcessTree', () => {
|
||||
})
|
||||
expect(target.kill).not.toHaveBeenCalled()
|
||||
release.resolve()
|
||||
await teardown
|
||||
// taskkill resolves alike on success, failure and timeout.
|
||||
await expect(teardown).resolves.toBeNull()
|
||||
|
||||
expect(terminateWindowsTree).toHaveBeenCalledWith(1234, { site: 'codex-app-server-teardown' })
|
||||
expect(target.kill).toHaveBeenCalledWith('SIGKILL')
|
||||
@@ -54,7 +56,7 @@ describe('terminateProviderProcessTree', () => {
|
||||
it('waits for an owned POSIX snapshot before killing the wrapper', async () => {
|
||||
const target = child()
|
||||
const snapshot = { rootPgid: 1234, descendants: [], capturedAtMs: 1 }
|
||||
const release = Promise.withResolvers<boolean>()
|
||||
const release = Promise.withResolvers<DescendantTreeVerdict>()
|
||||
|
||||
const teardown = terminateProviderProcessTree(target, {
|
||||
site: 'codex-app-server-teardown',
|
||||
@@ -64,8 +66,8 @@ describe('terminateProviderProcessTree', () => {
|
||||
})
|
||||
await vi.waitFor(() => expect(target.kill).toHaveBeenCalledWith('SIGSTOP'))
|
||||
expect(target.kill).not.toHaveBeenCalledWith('SIGKILL')
|
||||
release.resolve(true)
|
||||
await teardown
|
||||
release.resolve('exited')
|
||||
await expect(teardown).resolves.toBe('exited')
|
||||
|
||||
expect(target.kill).toHaveBeenLastCalledWith('SIGKILL')
|
||||
})
|
||||
@@ -83,7 +85,7 @@ describe('terminateProviderProcessTree', () => {
|
||||
captureDescendants,
|
||||
signalProcessGroup
|
||||
})
|
||||
).resolves.toBe(true)
|
||||
).resolves.toBeNull()
|
||||
|
||||
expect(signalProcessGroup).toHaveBeenCalledWith(1234, 'SIGKILL')
|
||||
expect(captureDescendants).not.toHaveBeenCalled()
|
||||
@@ -105,7 +107,7 @@ describe('terminateProviderProcessTree', () => {
|
||||
throw Object.assign(new Error('denied'), { code: 'EPERM' })
|
||||
}
|
||||
})
|
||||
).resolves.toBe(false)
|
||||
).resolves.toBe('unverifiable')
|
||||
|
||||
expect(target.kill).not.toHaveBeenCalled()
|
||||
})
|
||||
@@ -129,9 +131,9 @@ describe('terminateProviderProcessTree', () => {
|
||||
descendants: [],
|
||||
capturedAtMs: 1
|
||||
}),
|
||||
terminateDescendants: async () => true
|
||||
terminateDescendants: async () => 'exited'
|
||||
})
|
||||
).resolves.toBe(true)
|
||||
).resolves.toBe('exited')
|
||||
|
||||
expect(target.kill).toHaveBeenLastCalledWith('SIGKILL')
|
||||
expect(findSelfInitiatedTreeKills(Date.now())).toEqual([])
|
||||
@@ -146,10 +148,10 @@ describe('terminateProviderProcessTree', () => {
|
||||
site: 'codex-app-server-teardown',
|
||||
platform: 'darwin',
|
||||
captureDescendants: async () => ({ rootPgid: 1234, descendants: [], capturedAtMs: 1 }),
|
||||
terminateDescendants: async () => true,
|
||||
terminateDescendants: async () => 'exited',
|
||||
signalProcessGroup
|
||||
})
|
||||
).resolves.toBe(true)
|
||||
).resolves.toBe('exited')
|
||||
|
||||
expect(signalProcessGroup).toHaveBeenCalledWith(1234, 'SIGKILL')
|
||||
expect(findSelfInitiatedTreeKills(Date.now())).toEqual([
|
||||
@@ -161,6 +163,54 @@ describe('terminateProviderProcessTree', () => {
|
||||
])
|
||||
})
|
||||
|
||||
it.each([
|
||||
['live', 'live'],
|
||||
['unverifiable', 'unverifiable']
|
||||
] as const)(
|
||||
'reports a %s descendant snapshot as-is and leaves the stopped root resumable',
|
||||
async (observed, verdict) => {
|
||||
const target = child()
|
||||
await expect(
|
||||
terminateProviderProcessTree(target, {
|
||||
site: 'codex-app-server-teardown',
|
||||
platform: 'darwin',
|
||||
captureDescendants: async () => ({ rootPgid: 1234, descendants: [], capturedAtMs: 1 }),
|
||||
terminateDescendants: async () => observed
|
||||
})
|
||||
).resolves.toBe(verdict)
|
||||
expect(target.kill).toHaveBeenLastCalledWith('SIGCONT')
|
||||
}
|
||||
)
|
||||
|
||||
it('claims no observation when the POSIX process table cannot be read', async () => {
|
||||
const target = child()
|
||||
await expect(
|
||||
terminateProviderProcessTree(target, {
|
||||
site: 'codex-app-server-teardown',
|
||||
platform: 'darwin',
|
||||
captureDescendants: async () => null
|
||||
})
|
||||
).resolves.toBeNull()
|
||||
expect(target.kill).toHaveBeenLastCalledWith('SIGKILL')
|
||||
})
|
||||
|
||||
it('claims no observation from an ESRCH dedicated group and reports a spawnless child as unverifiable', async () => {
|
||||
await expect(
|
||||
terminateProviderProcessTree(child(), {
|
||||
site: 'provider-test-teardown',
|
||||
platform: 'linux',
|
||||
dedicatedProcessGroup: true,
|
||||
signalProcessGroup: () => {
|
||||
throw Object.assign(new Error('gone'), { code: 'ESRCH' })
|
||||
}
|
||||
})
|
||||
).resolves.toBeNull()
|
||||
const spawnless = { pid: undefined, kill: vi.fn<ChildProcess['kill']>(() => true) }
|
||||
await expect(
|
||||
terminateProviderProcessTree(spawnless, { site: 'provider-test-teardown', platform: 'linux' })
|
||||
).resolves.toBe('unverifiable')
|
||||
})
|
||||
|
||||
it('tears down 40 dedicated groups without process-table scans or cross-group fanout', async () => {
|
||||
const killMocks = Array.from({ length: 40 }, () => vi.fn<ChildProcess['kill']>(() => true))
|
||||
const targets = killMocks.map((kill, index) => ({
|
||||
@@ -182,7 +232,7 @@ describe('terminateProviderProcessTree', () => {
|
||||
)
|
||||
)
|
||||
|
||||
expect(results).toEqual(Array.from({ length: targets.length }, () => true))
|
||||
expect(results).toEqual(Array.from({ length: targets.length }, () => null))
|
||||
expect(signalProcessGroup.mock.calls).toEqual(targets.map((target) => [target.pid, 'SIGKILL']))
|
||||
expect(captureDescendants).not.toHaveBeenCalled()
|
||||
expect(killMocks.every((kill) => kill.mock.calls.length === 0)).toBe(true)
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
import type { ChildProcessHandle } from '../../shared/child-process/run-process'
|
||||
import { captureDescendantSnapshot, type DescendantSnapshot } from '../pty-descendant-termination'
|
||||
import { terminateDescendantSnapshotAndWait } from '../pty-descendant-exit-verification'
|
||||
import {
|
||||
terminateDescendantSnapshotWithVerdict,
|
||||
type DescendantTreeVerdict
|
||||
} from '../pty-descendant-exit-verification'
|
||||
import { terminateWindowsProcessTree } from '../windows-process-tree-kill'
|
||||
import { recordSelfInitiatedTreeKill } from '../crash-reporting/self-initiated-tree-kill-log'
|
||||
|
||||
const activeTeardowns = new WeakMap<object, Promise<boolean>>()
|
||||
/** What the teardown observed of the descendants; null when it signalled but observed nothing. */
|
||||
export type ProviderProcessTeardownVerdict = DescendantTreeVerdict | null
|
||||
|
||||
const activeTeardowns = new WeakMap<object, Promise<ProviderProcessTeardownVerdict>>()
|
||||
|
||||
type TeardownChild = Pick<ChildProcessHandle, 'pid' | 'kill'>
|
||||
|
||||
@@ -13,20 +19,24 @@ export type ProviderProcessTeardownDeps = {
|
||||
platform?: NodeJS.Platform
|
||||
dedicatedProcessGroup?: boolean
|
||||
captureDescendants?: (rootPid: number) => Promise<DescendantSnapshot | null>
|
||||
terminateDescendants?: (snapshot: DescendantSnapshot) => Promise<boolean>
|
||||
terminateDescendants?: (snapshot: DescendantSnapshot) => Promise<DescendantTreeVerdict>
|
||||
terminateWindowsTree?: (rootPid: number, deps?: { site?: string }) => Promise<void>
|
||||
signalProcessGroup?: (pgid: number, signal: NodeJS.Signals) => void
|
||||
}
|
||||
|
||||
function terminateDedicatedPosixGroup(rootPid: number, deps: ProviderProcessTeardownDeps): boolean {
|
||||
function terminateDedicatedPosixGroup(
|
||||
rootPid: number,
|
||||
deps: ProviderProcessTeardownDeps
|
||||
): ProviderProcessTeardownVerdict {
|
||||
const signalGroup =
|
||||
deps.signalProcessGroup ??
|
||||
((pgid: number, signal: NodeJS.Signals) => process.kill(-pgid, signal))
|
||||
try {
|
||||
signalGroup(rootPid, 'SIGKILL')
|
||||
} catch (error) {
|
||||
// ESRCH says only that the group is empty; a descendant that left it, or a root that never led it, may live.
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Node process.kill errors expose an optional errno code; only that field is read.
|
||||
return (error as NodeJS.ErrnoException).code === 'ESRCH'
|
||||
return (error as NodeJS.ErrnoException).code === 'ESRCH' ? null : 'unverifiable'
|
||||
}
|
||||
// Outside the try: that catch is the ESRCH contract, not a breadcrumb handler.
|
||||
recordSelfInitiatedTreeKill({
|
||||
@@ -34,23 +44,29 @@ function terminateDedicatedPosixGroup(rootPid: number, deps: ProviderProcessTear
|
||||
site: deps.site,
|
||||
scope: 'posix-process-group'
|
||||
})
|
||||
return true
|
||||
// A delivered signal is not an observed exit.
|
||||
return null
|
||||
}
|
||||
|
||||
async function terminatePosixTree(
|
||||
child: TeardownChild,
|
||||
rootPid: number,
|
||||
deps: ProviderProcessTeardownDeps
|
||||
): Promise<boolean> {
|
||||
): Promise<ProviderProcessTeardownVerdict> {
|
||||
child.kill('SIGSTOP')
|
||||
const capture = deps.captureDescendants ?? captureDescendantSnapshot
|
||||
const snapshot = await capture(rootPid).catch(() => null)
|
||||
if (!snapshot) {
|
||||
// No observation rather than the reaper's `unverifiable`: Codex's diagnostic treated this as accepted.
|
||||
// The reaper-move follow-up maps it to `unverifiable` and takes that Codex change deliberately.
|
||||
child.kill('SIGKILL')
|
||||
return true
|
||||
return null
|
||||
}
|
||||
const terminate = deps.terminateDescendants ?? terminateDescendantSnapshotAndWait
|
||||
const descendantsExited = await terminate(snapshot)
|
||||
const terminate =
|
||||
deps.terminateDescendants ??
|
||||
((captured: DescendantSnapshot) => terminateDescendantSnapshotWithVerdict(captured))
|
||||
const verdict = await terminate(snapshot)
|
||||
const descendantsExited = verdict === 'exited'
|
||||
// A detached POSIX launch is the leader of its own process group. Group
|
||||
// signalling reaches grandchildren even after they daemonise/reparent,
|
||||
// while the stopped root and captured pgid make the ownership proof exact.
|
||||
@@ -80,28 +96,29 @@ async function terminatePosixTree(
|
||||
}
|
||||
if (!descendantsExited) {
|
||||
child.kill('SIGCONT')
|
||||
return false
|
||||
return verdict
|
||||
}
|
||||
child.kill('SIGKILL')
|
||||
return true
|
||||
return 'exited'
|
||||
}
|
||||
|
||||
/** Stops every process owned by one provider launch before releasing its wrapper. */
|
||||
async function terminateOnce(
|
||||
child: TeardownChild,
|
||||
deps: ProviderProcessTeardownDeps
|
||||
): Promise<boolean> {
|
||||
): Promise<ProviderProcessTeardownVerdict> {
|
||||
const rootPid = child.pid
|
||||
if (!rootPid) {
|
||||
child.kill('SIGKILL')
|
||||
return false
|
||||
return 'unverifiable'
|
||||
}
|
||||
if ((deps.platform ?? process.platform) === 'win32') {
|
||||
const terminate = deps.terminateWindowsTree ?? terminateWindowsProcessTree
|
||||
await terminate(rootPid, { site: deps.site })
|
||||
// taskkill owns the tree; this preserves the prior direct-child fallback when it fails.
|
||||
child.kill('SIGKILL')
|
||||
return true
|
||||
// taskkill resolves alike on success, failure and timeout, so nothing was observed.
|
||||
return null
|
||||
}
|
||||
if (deps.dedicatedProcessGroup) {
|
||||
return terminateDedicatedPosixGroup(rootPid, deps)
|
||||
@@ -112,13 +129,15 @@ async function terminateOnce(
|
||||
export function terminateProviderProcessTree(
|
||||
child: TeardownChild,
|
||||
deps: ProviderProcessTeardownDeps
|
||||
): Promise<boolean> {
|
||||
): Promise<ProviderProcessTeardownVerdict> {
|
||||
const key = child
|
||||
const active = activeTeardowns.get(key)
|
||||
if (active) {
|
||||
return active
|
||||
}
|
||||
const attempt = terminateOnce(child, deps).catch(() => false)
|
||||
const attempt = terminateOnce(child, deps).catch(
|
||||
(): ProviderProcessTeardownVerdict => 'unverifiable'
|
||||
)
|
||||
activeTeardowns.set(key, attempt)
|
||||
void attempt.then(() => {
|
||||
if (activeTeardowns.get(key) === attempt) {
|
||||
|
||||
@@ -4,7 +4,7 @@ import { RetryableProcessExitProof } from './retryable-process-exit-proof'
|
||||
|
||||
describe('RetryableProcessExitProof', () => {
|
||||
it('shares a concurrent attempt and retains proven exit', async () => {
|
||||
const proof = new RetryableProcessExitProof()
|
||||
const proof = new RetryableProcessExitProof((result: boolean) => result)
|
||||
const proveExit = vi.fn(async () => true)
|
||||
|
||||
const first = proof.run(proveExit)
|
||||
@@ -16,7 +16,7 @@ describe('RetryableProcessExitProof', () => {
|
||||
})
|
||||
|
||||
it('permits another attempt after exit was not proven', async () => {
|
||||
const proof = new RetryableProcessExitProof()
|
||||
const proof = new RetryableProcessExitProof((result: boolean) => result)
|
||||
const proveExit = vi.fn().mockResolvedValueOnce(false).mockResolvedValueOnce(true)
|
||||
|
||||
await expect(proof.run(proveExit)).resolves.toBe(false)
|
||||
@@ -25,7 +25,7 @@ describe('RetryableProcessExitProof', () => {
|
||||
})
|
||||
|
||||
it('permits another attempt after proof rejects', async () => {
|
||||
const proof = new RetryableProcessExitProof()
|
||||
const proof = new RetryableProcessExitProof((result: boolean) => result)
|
||||
const proveExit = vi
|
||||
.fn()
|
||||
.mockRejectedValueOnce(new Error('probe failed'))
|
||||
|
||||
@@ -1,15 +1,17 @@
|
||||
export class RetryableProcessExitProof {
|
||||
private inFlight: Promise<boolean> | null = null
|
||||
export class RetryableProcessExitProof<Result> {
|
||||
private inFlight: Promise<Result> | null = null
|
||||
|
||||
run(proveExit: () => Promise<boolean>): Promise<boolean> {
|
||||
constructor(private readonly isProven: (result: Result) => boolean) {}
|
||||
|
||||
run(proveExit: () => Promise<Result>): Promise<Result> {
|
||||
if (this.inFlight) {
|
||||
return this.inFlight
|
||||
}
|
||||
const attempt = proveExit()
|
||||
this.inFlight = attempt
|
||||
void attempt.then(
|
||||
(proven) => {
|
||||
if (!proven) {
|
||||
(result) => {
|
||||
if (!this.isProven(result)) {
|
||||
this.clear(attempt)
|
||||
}
|
||||
},
|
||||
@@ -18,7 +20,7 @@ export class RetryableProcessExitProof {
|
||||
return attempt
|
||||
}
|
||||
|
||||
private clear(attempt: Promise<boolean>): void {
|
||||
private clear(attempt: Promise<Result>): void {
|
||||
if (this.inFlight === attempt) {
|
||||
this.inFlight = null
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user