diff --git a/src/main/codex/codex-app-server-connection.test.ts b/src/main/codex/codex-app-server-connection.test.ts index 4e7cd016bd6..7fe63c6e418 100644 --- a/src/main/codex/codex-app-server-connection.test.ts +++ b/src/main/codex/codex-app-server-connection.test.ts @@ -11,8 +11,12 @@ import { type CodexAppServerConnection, type CodexAppServerConnectionHandlers } from './codex-app-server-connection' +import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from './codex-app-server-posix-supervisor' import { isCodexAppServerUnsupportedError } from './codex-app-server-session' +// close() waits out the supervisor's own stop before forcing the tree. +const GRACEFUL_EXIT_MS = process.platform === 'win32' ? 1_500 : PROVIDER_SUPERVISOR_MAX_STOP_MS + const originalCodexHome = process.env.CODEX_HOME afterEach(() => { @@ -392,7 +396,7 @@ describe('openCodexAppServerConnection', () => { openCodexAppServerConnection({ command: 'codex', args: ['app-server'] }, {}, spawnImpl) ) - await vi.advanceTimersByTimeAsync(5_000) + await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 3_500) const error = (await opening) as Error & { connection?: CodexAppServerConnection } expect(error.name).toBe('CodexAppServerHandshakeExitUnprovenError') @@ -434,7 +438,7 @@ describe('openCodexAppServerConnection', () => { }) const closing = connection.close() - await vi.advanceTimersByTimeAsync(2_000) + await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 500) await closing await vi.waitFor(() => expect(child.kill).toHaveBeenCalledWith('SIGKILL')) @@ -468,7 +472,7 @@ describe('openCodexAppServerConnection', () => { const first = connection.close() const second = connection.close() - await vi.advanceTimersByTimeAsync(4_100) + await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 2_600) await expect(Promise.all([first, second])).resolves.toEqual([true, true]) expect(child.kill.mock.calls.map(([signal]) => signal)).toEqual(['SIGSTOP', 'SIGKILL']) @@ -485,7 +489,7 @@ describe('openCodexAppServerConnection', () => { ) const first = connection.close() - await vi.advanceTimersByTimeAsync(5_000) + await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 3_500) await expect(first).resolves.toBe(false) child.emit('exit', 0, null) diff --git a/src/main/codex/codex-app-server-connection.ts b/src/main/codex/codex-app-server-connection.ts index b0cf133e1e2..ef566bca336 100644 --- a/src/main/codex/codex-app-server-connection.ts +++ b/src/main/codex/codex-app-server-connection.ts @@ -1,6 +1,9 @@ import { spawnProcess } from '../../shared/child-process/run-process' import { RetryableProcessExitProof } from '../../shared/child-process/retryable-process-exit-proof' -import { createProviderSpawnSpec } from './codex-app-server-posix-supervisor' +import { + createProviderSpawnSpec, + PROVIDER_SUPERVISOR_MAX_STOP_MS +} from './codex-app-server-posix-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' @@ -247,7 +250,11 @@ export async function openCodexAppServerConnection( // Already destroyed; the reap below still runs. } if (!exited) { - await waitForProcessExitUntil(exitPromise, GRACEFUL_EXIT_MS) + // 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() if (!treeExited) { diff --git a/src/main/codex/codex-app-server-posix-supervisor.integration.test.ts b/src/main/codex/codex-app-server-posix-supervisor.integration.test.ts new file mode 100644 index 00000000000..e23cdbc88c3 --- /dev/null +++ b/src/main/codex/codex-app-server-posix-supervisor.integration.test.ts @@ -0,0 +1,333 @@ +import { spawn, type ChildProcess } from 'node:child_process' +import { existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import { + POSIX_PROVIDER_SUPERVISOR_SCRIPT, + PROVIDER_SIGTERM_GRACE_MS, + PROVIDER_STDIN_END_GRACE_MS, + PROVIDER_SUPERVISOR_MAX_STOP_MS, + supervisedPosixLaunch, + type ProviderSupervisorOptions +} from './codex-app-server-posix-supervisor' + +// The provider leads its own group; its grandchild shares that group and ignores SIGTERM. +const PROVIDER = String.raw` + const { spawn } = require('node:child_process') + if (process.env.ORCA_TEST_PROVIDER_IGNORES_SIGTERM) process.on('SIGTERM', () => {}) + if (process.env.ORCA_TEST_PROVIDER_SIGNAL_FILE) { + process.on('SIGTERM', () => { + require('node:fs').writeFileSync(process.env.ORCA_TEST_PROVIDER_SIGNAL_FILE, 'SIGTERM') + process.exit(0) + }) + } + process.stdout.on('error', () => {}) + const grandchild = spawn( + process.execPath, + ['-e', "process.on('SIGTERM', () => {}); process.stdout.write('armed'); setInterval(() => {}, 60000)"], + { stdio: ['ignore', 'pipe', 'ignore'] } + ) + grandchild.stdout.once('data', () => { + process.stdout.write(JSON.stringify({ provider: process.pid, grandchild: grandchild.pid }) + '\n') + if (process.env.ORCA_TEST_PROVIDER_STREAMS_OUTPUT) setInterval(() => process.stdout.write('.'), 2) + }) + setInterval(() => {}, 60000) +` + +// Exits the moment its stdin ends, as Codex does on a normal close. +const EXITS_ON_STDIN_END_PROVIDER = String.raw` + process.stdin.on('end', () => process.exit(0)).resume() + process.stdout.write(JSON.stringify({ provider: process.pid }) + '\n') +` + +// Ignores stdin end and SIGTERM, recording when SIGTERM arrived, so only SIGKILL ends it. +const RECORDS_SIGTERM_PROVIDER = String.raw` + process.on('SIGTERM', () => { + require('node:fs').writeFileSync(process.env.ORCA_TEST_PROVIDER_SIGNAL_FILE, String(Date.now())) + }) + process.stdout.write(JSON.stringify({ provider: process.pid }) + '\n') + setInterval(() => {}, 60000) +` + +// Stands in for Orca: launches the supervisor as its own child, then can be killed outright. A +// second child holds the supervisor's stdin open, so only the parent-death watch can notice. +const OWNER = String.raw` + const { spawn } = require('node:child_process') + const spec = JSON.parse(Buffer.from(process.env.ORCA_PROVIDER_SUPERVISOR_SPEC, 'base64').toString()) + spec.ownerPid = process.pid + const supervisor = spawn(process.execPath, ['-e', process.env.ORCA_TEST_SUPERVISOR_SCRIPT], { + env: { ...process.env, ORCA_PROVIDER_SUPERVISOR_SPEC: Buffer.from(JSON.stringify(spec)).toString('base64') }, + stdio: ['pipe', 'pipe', 'ignore'], + detached: true + }) + const holder = spawn(process.execPath, ['-e', 'setInterval(() => {}, 60000)'], { + stdio: ['ignore', supervisor.stdin, 'ignore'] + }) + process.stdout.write(JSON.stringify({ supervisor: supervisor.pid, holder: holder.pid }) + '\n') + supervisor.stdout.pipe(process.stdout) + setInterval(() => {}, 60000) +` + +// Preloaded into the supervisor: signals it the instant its provider exists, the spawn window. +const SIGNAL_AFTER_SPAWN_PRELOAD = String.raw` + const childProcess = require('node:child_process') + const spawn = childProcess.spawn + childProcess.spawn = (...args) => { + const child = spawn(...args) + require('node:fs').writeFileSync(process.env.ORCA_TEST_PROVIDER_PID_FILE, String(child.pid)) + process.kill(process.pid, 'SIGTERM') + return child + } +` + +const recordedPids = new Set() +const tempDirs: string[] = [] + +function tempDir(): string { + const dir = mkdtempSync(join(tmpdir(), 'orca-supervisor-')) + tempDirs.push(dir) + return dir +} + +function alive(pid: number): boolean { + try { + process.kill(pid, 0) + return true + } catch (error) { + return !(error instanceof Error && 'code' in error && error.code === 'ESRCH') + } +} + +async function waitFor(predicate: () => boolean, timeoutMs: number): Promise { + const deadline = Date.now() + timeoutMs + while (!predicate()) { + if (Date.now() >= deadline) { + return false + } + await new Promise((resolve) => setTimeout(resolve, 25)) + } + return true +} + +function readPids(child: ChildProcess, keys: readonly string[]): Promise> { + return new Promise((resolve, reject) => { + const pids: Record = {} + let buffered = '' + const timeout = setTimeout(() => reject(new Error(`no ${keys.join('/')} pids`)), 10_000) + const onData = (chunk: Buffer): void => { + buffered += chunk.toString() + const lines = buffered.split('\n') + buffered = lines.pop() ?? '' + for (const line of lines) { + const parsed: unknown = JSON.parse(line) + for (const [key, pid] of Object.entries(parsed ?? {})) { + if (typeof pid === 'number') { + pids[key] = pid + recordedPids.add(pid) + } + } + } + if (keys.every((key) => key in pids)) { + clearTimeout(timeout) + // Later output is not pids; the stream keeps flowing without a listener. + child.stdout!.off('data', onData) + resolve(pids) + } + } + child.stdout!.on('data', onData) + }) +} + +function launchSupervisor( + options: ProviderSupervisorOptions, + env: Record = {}, + provider: { command: string; args: string[] } = { + command: process.execPath, + args: ['-e', PROVIDER] + }, + nodeArgs: string[] = [] +): { supervisor: ChildProcess; exit: Promise<{ code: number | null; signal: string | null }> } { + const launch = supervisedPosixLaunch(provider, { ...process.env, ...env }, options) + const supervisor = spawn(launch.command, [...nodeArgs, ...launch.args], { + env: launch.env, + stdio: ['pipe', 'pipe', 'ignore'], + detached: true + }) + recordedPids.add(supervisor.pid!) + const exit = new Promise<{ code: number | null; signal: string | null }>((resolve) => + supervisor.once('exit', (code, signal) => resolve({ code, signal })) + ) + return { supervisor, exit } +} + +async function launchUnderOwner( + options: ProviderSupervisorOptions, + env: Record = {} +): Promise<{ owner: ChildProcess; pids: Record }> { + const launch = supervisedPosixLaunch( + { command: process.execPath, args: ['-e', PROVIDER] }, + { ...process.env, ...env }, + options + ) + const owner = spawn(process.execPath, ['-e', OWNER], { + env: { ...launch.env, ORCA_TEST_SUPERVISOR_SCRIPT: POSIX_PROVIDER_SUPERVISOR_SCRIPT }, + stdio: ['ignore', 'pipe', 'ignore'] + }) + recordedPids.add(owner.pid!) + const pids = await readPids(owner, ['supervisor', 'holder', 'provider', 'grandchild']) + return { owner, pids } +} + +afterEach(() => { + for (const pid of recordedPids) { + if (alive(pid)) { + process.kill(pid, 'SIGKILL') + } + } + recordedPids.clear() + for (const dir of tempDirs.splice(0)) { + rmSync(dir, { recursive: true, force: true }) + } +}) + +describe.runIf(process.platform !== 'win32')('POSIX provider supervisor processes', () => { + it('reaps the provider group on SIGTERM and exits only after the group is gone', async () => { + const { supervisor, exit } = launchSupervisor({ sigtermGraceMs: 300 }) + const { provider, grandchild } = await readPids(supervisor, ['provider', 'grandchild']) + + let groupAliveAtExit: boolean | null = null + void exit.then(() => { + groupAliveAtExit = alive(-provider) + }) + supervisor.kill('SIGTERM') + + await expect(exit).resolves.toEqual({ code: null, signal: 'SIGTERM' }) + expect(groupAliveAtExit).toBe(false) + expect(alive(provider)).toBe(false) + expect(alive(grandchild)).toBe(false) + }) + + it('escalates a SIGTERM-ignoring provider to SIGKILL after the grace from the spec', async () => { + const graceMs = 200 + const { supervisor, exit } = launchSupervisor( + { sigtermGraceMs: graceMs }, + { ORCA_TEST_PROVIDER_IGNORES_SIGTERM: '1' } + ) + const { provider, grandchild } = await readPids(supervisor, ['provider', 'grandchild']) + + const signalledAt = Date.now() + supervisor.kill('SIGTERM') + const exited = await Promise.race([exit, new Promise((resolve) => setTimeout(resolve, 5_000))]) + + expect(exited).toEqual({ code: null, signal: 'SIGTERM' }) + expect(Date.now() - signalledAt).toBeGreaterThanOrEqual(graceMs) + expect(Date.now() - signalledAt).toBeLessThan(PROVIDER_SIGTERM_GRACE_MS) + expect(alive(provider)).toBe(false) + expect(alive(grandchild)).toBe(false) + }) + + it('reaps a provider spawned in the instant before a stop arrives', async () => { + const dir = tempDir() + const preload = join(dir, 'signal-after-spawn.js') + const pidFile = join(dir, 'provider-pid') + writeFileSync(preload, SIGNAL_AFTER_SPAWN_PRELOAD) + const { exit } = launchSupervisor( + { sigtermGraceMs: 300 }, + { ORCA_TEST_PROVIDER_PID_FILE: pidFile }, + { command: process.execPath, args: ['-e', 'setInterval(() => {}, 60000)'] }, + ['--require', preload] + ) + + await expect(exit).resolves.toEqual({ code: null, signal: 'SIGTERM' }) + const provider = Number(readFileSync(pidFile, 'utf8')) + recordedPids.add(provider) + expect(await waitFor(() => !alive(provider), 3_000)).toBe(true) + }) + + it('never spawns the provider when its owner is not its parent at start', async () => { + const marker = join(tempDir(), 'provider-started') + const { exit } = launchSupervisor( + { ownerPid: process.pid === 1 ? 2 : 1 }, + {}, + { command: 'touch', args: [marker] } + ) + + await expect(exit).resolves.toEqual({ code: 1, signal: null }) + await new Promise((resolve) => setTimeout(resolve, 200)) + expect(existsSync(marker)).toBe(false) + }) + + it.each([ + ['', {}], + // Output after the owner's death meets a closed pipe, which must not end the supervisor first. + [' while the provider is writing output', { ORCA_TEST_PROVIDER_STREAMS_OUTPUT: '1' }] + ])('reaps the provider group when its owner dies%s', async (_, env) => { + const graceMs = 300 + const { owner, pids } = await launchUnderOwner({ sigtermGraceMs: graceMs }, env) + + const killedAt = Date.now() + owner.kill('SIGKILL') + + expect(await waitFor(() => !alive(-pids.provider), 3_000)).toBe(true) + // The grandchild ignores SIGTERM, so the group lasts until the grace ends in SIGKILL. + expect(Date.now() - killedAt).toBeGreaterThanOrEqual(graceMs) + expect(await waitFor(() => !alive(pids.supervisor), 3_000)).toBe(true) + expect(alive(pids.grandchild)).toBe(false) + }) + + it('closes a provider that exits on stdin end without waiting out any grace', async () => { + const { supervisor, exit } = launchSupervisor( + {}, + {}, + { + command: process.execPath, + args: ['-e', EXITS_ON_STDIN_END_PROVIDER] + } + ) + const { provider } = await readPids(supervisor, ['provider']) + + const endedAt = Date.now() + supervisor.stdin!.end() + + await expect(exit).resolves.toEqual({ code: 0, signal: null }) + expect(Date.now() - endedAt).toBeLessThan(PROVIDER_STDIN_END_GRACE_MS) + expect(alive(provider)).toBe(false) + }) + + it('gives a provider 1 s after stdin end, then 3 s after SIGTERM before SIGKILL', async () => { + const signalFile = join(tempDir(), 'provider-sigterm-at') + const { supervisor, exit } = launchSupervisor( + {}, + { ORCA_TEST_PROVIDER_SIGNAL_FILE: signalFile }, + { command: process.execPath, args: ['-e', RECORDS_SIGTERM_PROVIDER] } + ) + const { provider } = await readPids(supervisor, ['provider']) + + const endedAt = Date.now() + supervisor.stdin!.end() + const exited = await exit + const exitedAt = Date.now() + const signalledAt = Number(readFileSync(signalFile, 'utf8')) + + expect(exited).toEqual({ code: 137, signal: null }) + // Timers may fire a tick early against another process's clock. + expect(signalledAt - endedAt).toBeGreaterThanOrEqual(1_000 - 20) + expect(exitedAt - signalledAt).toBeGreaterThanOrEqual(3_000 - 20) + expect(exitedAt - endedAt).toBeLessThan(PROVIDER_SUPERVISOR_MAX_STOP_MS + 1_000) + expect(alive(provider)).toBe(false) + }) + + it('asks the provider to stop with SIGTERM when its owner dies', async () => { + const signalFile = join(tempDir(), 'provider-signal') + const { owner, pids } = await launchUnderOwner( + { sigtermGraceMs: 300 }, + { ORCA_TEST_PROVIDER_SIGNAL_FILE: signalFile } + ) + + owner.kill('SIGKILL') + + expect(await waitFor(() => !alive(-pids.provider), 3_000)).toBe(true) + expect(existsSync(signalFile) && readFileSync(signalFile, 'utf8')).toBe('SIGTERM') + }) +}) diff --git a/src/main/codex/codex-app-server-posix-supervisor.test.ts b/src/main/codex/codex-app-server-posix-supervisor.test.ts index f51157a443c..28a088716f6 100644 --- a/src/main/codex/codex-app-server-posix-supervisor.test.ts +++ b/src/main/codex/codex-app-server-posix-supervisor.test.ts @@ -3,6 +3,8 @@ import type { CodexAppServerLaunch } from './codex-app-server-connection' import { createProviderSpawnSpec, POSIX_PROVIDER_SUPERVISOR_SCRIPT, + PROVIDER_SIGTERM_GRACE_MS, + PROVIDER_STDIN_END_GRACE_MS, supervisedPosixLaunch } from './codex-app-server-posix-supervisor' @@ -27,7 +29,10 @@ describe('structured provider supervision', () => { expect.objectContaining({ command: '/opt/codex', args: ['app-server', '--flag'], - cwd: '/work/repo' + cwd: '/work/repo', + ownerPid: process.pid, + stdinEndGraceMs: PROVIDER_STDIN_END_GRACE_MS, + sigtermGraceMs: PROVIDER_SIGTERM_GRACE_MS }) ) expect( @@ -38,7 +43,7 @@ describe('structured provider supervision', () => { ) expect(POSIX_PROVIDER_SUPERVISOR_SCRIPT).toContain('delete childEnv.ELECTRON_RUN_AS_NODE') expect(spec.env.ELECTRON_RUN_AS_NODE).toBe('1') - expect(POSIX_PROVIDER_SUPERVISOR_SCRIPT).toContain('process.ppid !== originalParent') + expect(POSIX_PROVIDER_SUPERVISOR_SCRIPT).toContain('process.ppid !== spec.ownerPid') expect(POSIX_PROVIDER_SUPERVISOR_SCRIPT).toContain( "process.stdin.once('close', scheduleOwnerShutdown)" ) @@ -49,6 +54,18 @@ describe('structured provider supervision', () => { expect(POSIX_PROVIDER_SUPERVISOR_SCRIPT).not.toContain('process.ppid === 1') }) + it('refuses a grace longer than recovery waits before SIGKILL', () => { + const stdinEnd = (stdinEndGraceMs: number) => () => + supervisedPosixLaunch(launch, {}, { stdinEndGraceMs }) + const sigterm = (sigtermGraceMs: number) => () => + supervisedPosixLaunch(launch, {}, { sigtermGraceMs }) + + expect(stdinEnd(PROVIDER_STDIN_END_GRACE_MS)).not.toThrow() + expect(stdinEnd(PROVIDER_STDIN_END_GRACE_MS + 1)).toThrow(RangeError) + expect(sigterm(PROVIDER_SIGTERM_GRACE_MS)).not.toThrow() + expect(sigterm(PROVIDER_SIGTERM_GRACE_MS + 1)).toThrow(RangeError) + }) + it('uses direct provider spawning on Windows because the job owns the tree', () => { expect(createProviderSpawnSpec(launch, { PATH: '/bin' }, 'win32')).toEqual({ program: '/opt/codex', diff --git a/src/main/codex/codex-app-server-posix-supervisor.ts b/src/main/codex/codex-app-server-posix-supervisor.ts index f83ba6f374b..30714ee7911 100644 --- a/src/main/codex/codex-app-server-posix-supervisor.ts +++ b/src/main/codex/codex-app-server-posix-supervisor.ts @@ -1,9 +1,33 @@ import type { CodexAppServerLaunch } from './codex-app-server-connection' +/** Time the provider gets to exit on its own after its stdin ends, before SIGTERM. */ +export const PROVIDER_STDIN_END_GRACE_MS = 1_000 +/** Time the provider group gets to flush and exit after SIGTERM, before SIGKILL. */ +export const PROVIDER_SIGTERM_GRACE_MS = 3_000 +/** + * How long the supervisor waits for a SIGKILLed group to disappear. A killed process never runs + * again, so this only covers the kernel finishing the kill; waiting forever could hang close or + * recovery on a process stuck in the kernel, such as one blocked on a hung network drive. + */ +export const PROVIDER_GROUP_REAP_TIMEOUT_MS = 1_500 +/** + * Longest a supervisor can take to stop once asked (stdin end, grace, SIGTERM, grace, SIGKILL, + * reap); a SIGKILL sooner can orphan its group. The graces are also the largest a spec may carry. + */ +export const PROVIDER_SUPERVISOR_MAX_STOP_MS = + PROVIDER_STDIN_END_GRACE_MS + PROVIDER_SIGTERM_GRACE_MS + PROVIDER_GROUP_REAP_TIMEOUT_MS + /** Inline supervisor source kept dependency-free for the spawned Node child. */ export const POSIX_PROVIDER_SUPERVISOR_SCRIPT = ` const { spawn } = require('node:child_process') const spec = JSON.parse(Buffer.from(process.env.ORCA_PROVIDER_SUPERVISOR_SPEC, 'base64').toString()) +// A detached supervisor is reparented when its owner exits. The new parent may +// be PID 1 or a platform subreaper, so any other parent means no live owner. +const ownerGone = () => process.ppid !== spec.ownerPid +// Registered before the spawn, so a stop that lands while the provider starts still reaps it. +for (const signal of ['SIGTERM', 'SIGINT', 'SIGHUP']) process.on(signal, () => stopProviderGroup(signal)) +// Orca can die before this runs; spawning then would start a provider nothing watches. +if (ownerGone()) process.exit(1) const childEnv = { ...process.env } delete childEnv.ORCA_PROVIDER_SUPERVISOR_SPEC delete childEnv.ELECTRON_RUN_AS_NODE @@ -13,7 +37,6 @@ const child = spawn(spec.command, spec.args, { stdio: ['pipe', 'pipe', 'pipe'], detached: true }) -const originalParent = process.ppid let timer let ownerShutdownTimer let settling = false @@ -26,29 +49,47 @@ const providerGroupExists = () => { return Boolean(error && error.code !== 'ESRCH') } } -const reapOwnedProviderGroup = async () => { - if (!child.pid) return false - try { process.kill(-child.pid, 'SIGKILL') } catch (error) { - if (error && error.code !== 'ESRCH') return false - } - const deadline = Date.now() + 1500 +const waitForProviderGroupExit = async (timeoutMs) => { + const deadline = Date.now() + timeoutMs while (providerGroupExists()) { if (Date.now() >= deadline) return false await new Promise((resolve) => setTimeout(resolve, 25)) } return true } -const terminateOwnedGroup = () => { +const reapOwnedProviderGroup = async () => { + if (!child.pid) return false + try { process.kill(-child.pid, 'SIGKILL') } catch (error) { + if (error && error.code !== 'ESRCH') return false + } + return waitForProviderGroupExit(${PROVIDER_GROUP_REAP_TIMEOUT_MS}) +} +const finishWithProviderOutcome = (code, signal) => { + if (!signal) return process.exit(code ?? 1) + // Re-raise with the default action; this supervisor's own handler would swallow it. + process.removeAllListeners(signal) + process.kill(process.pid, signal) +} +// Every stop is the same: SIGTERM the group, SIGKILL it after its grace, and exit only once it is +// gone. Whoever stops this pid judges the provider by it, so a dead supervisor means a dead group. +const stopProviderGroup = (receivedSignal) => { if (settling) return settling = true clearInterval(timer) - void reapOwnedProviderGroup().then((reaped) => process.exit(reaped ? 137 : 1)) + if (ownerShutdownTimer) clearTimeout(ownerShutdownTimer) + try { process.kill(-child.pid, 'SIGTERM') } catch {} + void waitForProviderGroupExit(spec.sigtermGraceMs) + .then((exited) => exited || reapOwnedProviderGroup()) + .then((reaped) => { + if (!reaped) return process.exit(1) + finishWithProviderOutcome(137, receivedSignal) + }) } const scheduleOwnerShutdown = () => { if (settling || ownerShutdownTimer) return // A normal close ends the provider's stdin first; allow it to flush and // exit before forcing the group, while still bounding an orphaned child. - ownerShutdownTimer = setTimeout(terminateOwnedGroup, 1250) + ownerShutdownTimer = setTimeout(() => stopProviderGroup(null), spec.stdinEndGraceMs) ownerShutdownTimer.unref() } process.stdin.once('end', scheduleOwnerShutdown) @@ -56,10 +97,9 @@ process.stdin.once('close', scheduleOwnerShutdown) process.stdin.pipe(child.stdin) child.stdout.pipe(process.stdout) child.stderr.pipe(process.stderr) -for (const stream of [process.stdin, child.stdin, child.stdout, child.stderr]) stream.on('error', () => {}) -const finishWithProviderOutcome = (code, signal) => { - if (!signal) return process.exit(code ?? 1) - process.kill(process.pid, signal) +// A dead owner's stdout pipe raises EPIPE; unhandled, it would end this pid before the group. +for (const stream of [process.stdin, process.stdout, process.stderr, child.stdin, child.stdout, child.stderr]) { + stream.on('error', () => {}) } const reapProviderExit = async (code, signal) => { if (settling) return @@ -70,12 +110,7 @@ const reapProviderExit = async (code, signal) => { finishWithProviderOutcome(code, signal) } timer = setInterval(() => { - // A detached supervisor is reparented when its owner exits. The new parent - // may be PID 1 or a platform subreaper, so any parent change is proof that - // this process group no longer has a live Orca owner. - if (process.ppid !== originalParent) { - terminateOwnedGroup() - } + if (ownerGone()) stopProviderGroup(null) }, 100) timer.unref() child.once('error', () => { @@ -87,16 +122,41 @@ child.once('exit', (code, signal) => { }) ` +export type ProviderSupervisorOptions = { + cwd?: string + /** The process the supervisor serves; it must be the supervisor's parent. */ + ownerPid?: number + stdinEndGraceMs?: number + sigtermGraceMs?: number +} + +// A longer grace than the max stop allows would let recovery or close SIGKILL mid-stop. +function assertGraceWithin(name: string, graceMs: number, maxMs: number): void { + if (!(graceMs >= 0 && graceMs <= maxMs)) { + throw new RangeError(`Provider supervisor ${name} grace ${graceMs} ms is outside 0-${maxMs} ms`) + } +} + export function supervisedPosixLaunch( launch: CodexAppServerLaunch, childEnv: NodeJS.ProcessEnv, - cwd = launch.cwd ?? process.cwd() + { + cwd = launch.cwd ?? process.cwd(), + ownerPid = process.pid, + stdinEndGraceMs = PROVIDER_STDIN_END_GRACE_MS, + sigtermGraceMs = PROVIDER_SIGTERM_GRACE_MS + }: ProviderSupervisorOptions = {} ): { command: string; args: string[]; env: NodeJS.ProcessEnv } { + assertGraceWithin('stdin-end', stdinEndGraceMs, PROVIDER_STDIN_END_GRACE_MS) + assertGraceWithin('SIGTERM', sigtermGraceMs, PROVIDER_SIGTERM_GRACE_MS) const supervisorSpec = Buffer.from( JSON.stringify({ command: launch.command, args: launch.args, - cwd + cwd, + ownerPid, + stdinEndGraceMs, + sigtermGraceMs }) ).toString('base64') return { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-eviction.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-eviction.ts index 2c03143ff41..226781117bd 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-eviction.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-eviction.ts @@ -52,7 +52,7 @@ export type StructuredAgentSessionEvictionContext = { } /** The resume offer is advisory; a stalled sink must not hold the child's stop behind it. */ -const SNAPSHOT_DRAIN_TIMEOUT_MS = 1_000 +export const SNAPSHOT_DRAIN_TIMEOUT_MS = 1_000 export type StructuredAgentSessionEvictionStep = { name: string diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.test.ts index 5fd0440a1fe..760b2603d99 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.test.ts @@ -1,5 +1,13 @@ import { describe, expect, it, vi } from 'vitest' -import { structuredAgentSessionHostTeardownPhases } from './structured-agent-session-host-teardown' +import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from '../../codex/codex-app-server-posix-supervisor' +import { SNAPSHOT_DRAIN_TIMEOUT_MS } from './structured-agent-session-eviction' +import { + CHILD_EVICTION_TIMEOUT_MS, + structuredAgentSessionHostTeardownPhases +} from './structured-agent-session-host-teardown' + +// Codex close observes the supervisor's exit after its timers fire, which run late on a loaded host. +const SUPERVISOR_EXIT_OBSERVATION_HEADROOM_MS = 1_000 describe('structured agent-session host teardown', () => { it('names every phase, so the quit-path order is pinned rather than incidental', () => { @@ -23,6 +31,16 @@ describe('structured agent-session host teardown', () => { ]) }) + it("fits the Codex supervisor's longest stop inside quit's child-eviction bound", () => { + // A quit that times out first leaves the provider running and its lease unreleased. Eviction + // drains the sink for the resume offer before it stops the child, inside the same bound. + expect( + SNAPSHOT_DRAIN_TIMEOUT_MS + + PROVIDER_SUPERVISOR_MAX_STOP_MS + + SUPERVISOR_EXIT_OBSERVATION_HEADROOM_MS + ).toBeLessThan(CHILD_EVICTION_TIMEOUT_MS) + }) + it('bounds stalled recovery publication without preventing later cleanup', async () => { vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) const warning = vi.spyOn(console, 'warn').mockImplementation(() => {}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.ts index 97a27d5c246..eade95a7dd2 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-teardown.ts @@ -28,7 +28,7 @@ const RESUME_MARKER_RECORD_TIMEOUT_MS = 2_000 /** Eight steps at ten seconds each would outlast the global quit deadline, and a quit that dies * mid-eviction leaves the lease unreleased — the exact state restart has to clean up. Bounded * well below that deadline so the phases after this one still get to run. */ -const CHILD_EVICTION_TIMEOUT_MS = 8_000 +export const CHILD_EVICTION_TIMEOUT_MS = 8_000 /** Bounds a phase without swallowing its failure, which `withTimeout` alone would. */ async function withPhaseTimeout(run: () => Promise, timeoutMs: number): Promise { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.test.ts index 7fbd785894d..063e8c64af7 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.test.ts @@ -1,3 +1,4 @@ +import { spawn } from 'node:child_process' import { mkdtemp, rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' @@ -8,6 +9,7 @@ import { writeOlderBuildLease } from '../../runtime/agent-session-older-build-lease.test-fixture' import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' +import { supervisedPosixLaunch } from '../../codex/codex-app-server-posix-supervisor' import { resolveStructuredSessionRecovery, type StructuredSessionRecoveryResolutionDeps @@ -62,7 +64,7 @@ async function reserve(store: AgentSessionRecordStore) { }) } -async function liveOwner(store: AgentSessionRecordStore) { +async function liveOwner(store: AgentSessionRecordStore, pid = 4242) { const reserved = await reserve(store) const fence = reserved.record.lease.runtimeFence await store.commitProcessIdentity({ @@ -70,7 +72,7 @@ async function liveOwner(store: AgentSessionRecordStore) { fence, process: { hostId: 'local', - pid: 4242, + pid, processStartTimeMs: NOW - 1_000, spawnToken: 'spawn-recovery' }, @@ -302,3 +304,72 @@ describe('structured session recovery resolution', () => { ).toBe('not-applicable') }) }) + +describe.runIf(process.platform !== 'win32')( + 'structured session recovery of a supervised owner', + () => { + const recordedPids: number[] = [] + const alive = (pid: number): boolean => { + try { + process.kill(pid, 0) + return true + } catch (error) { + return !(error instanceof Error && 'code' in error && error.code === 'ESRCH') + } + } + + afterEach(() => { + for (const pid of recordedPids.splice(0)) { + if (alive(pid)) { + process.kill(pid, 'SIGKILL') + } + } + }) + + it('evicts only after the supervisor has reaped a provider that ignores SIGTERM', async () => { + const launch = supervisedPosixLaunch( + { + command: process.execPath, + args: [ + '-e', + "process.on('SIGTERM', () => {}); process.stdout.write(process.pid + '\\n'); setInterval(() => {}, 60000)" + ] + }, + process.env + ) + const supervisor = spawn(launch.command, launch.args, { + env: launch.env, + stdio: ['pipe', 'pipe', 'ignore'], + detached: true + }) + recordedPids.push(supervisor.pid!) + const provider = await new Promise((resolve, reject) => { + const timeout = setTimeout(() => reject(new Error('provider never started')), 10_000) + supervisor.stdout.once('data', (chunk: Buffer) => { + clearTimeout(timeout) + resolve(Number(chunk.toString().trim())) + }) + }) + recordedPids.push(provider) + const store = await openStore() + await liveOwner(store, supervisor.pid!) + await latch(store) + + const result = await resolveStructuredSessionRecovery( + { + store, + // The recorded pid is the supervisor's; its absence is the proof recovery evicts on. + probeRecord: async () => + alive(supervisor.pid!) ? MATCHED : { outcome: 'pid-absent' as const }, + now: () => NOW + 10_000 + }, + SESSION + ) + + expect(result).toBe('resolved') + // Evicted on proof of death, which must hold for the provider too, not only the supervisor. + expect(store.getRecord(SESSION)?.lease.deathEvidence).not.toBeNull() + expect(alive(provider)).toBe(false) + }) + } +) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.ts index 5e04cbdf90e..36d950a09a4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-recovery-resolution.ts @@ -15,6 +15,7 @@ import { type AgentSessionOwnerProbe } from '../../../shared/agent-session-lease-adjudication' import type { AgentSessionRecord } from '../../../shared/agent-session-record' +import { PROVIDER_SUPERVISOR_MAX_STOP_MS } from '../../codex/codex-app-server-posix-supervisor' import { releaseUnprovenAgentSessionOwner } from '../../runtime/agent-session-lease-transitions' import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' @@ -28,8 +29,13 @@ export type StructuredSessionRecoveryResolutionDeps = { delay?: (ms: number) => Promise } -const STOP_PROBES_PER_SIGNAL = 4 const STOP_PROBE_INTERVAL_MS = 250 +// A POSIX Codex owner is its provider supervisor, which exits only after its provider group. A +// SIGKILL that lands first leaves the group running, so SIGTERM outlasts the supervisor's stop. +const STOP_PROBES: Record = { + SIGTERM: Math.ceil(PROVIDER_SUPERVISOR_MAX_STOP_MS / STOP_PROBE_INTERVAL_MS) + 1, + SIGKILL: 4 +} const UNRESOLVED_REFUSALS: ReadonlySet = new Set([ 'agent_session_ownership_unknown', @@ -94,7 +100,7 @@ async function stopOwnerAndReprobe( let probe: AgentSessionOwnerProbe = { outcome: 'indeterminate', reason: 'owner stop requested' } for (const signal of ['SIGTERM', 'SIGKILL'] as const) { stop(pid, signal) - for (let attempt = 0; attempt < STOP_PROBES_PER_SIGNAL; attempt += 1) { + for (let attempt = 0; attempt < STOP_PROBES[signal]; attempt += 1) { probe = await deps.probeRecord(record) if (isProvenDeadProbe(probe)) { return probe