Files
orca/src/main/codex/codex-app-server-connection.test.ts
T
Brennan Benson 57fedeed79 fix(native-chat): the agent's exit ends its record, and an unconfirmed stop is joined instead of held (#24862)
* 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."
2026-10-05 12:14:47 -07:00

893 lines
32 KiB
TypeScript

import { EventEmitter } from 'node:events'
import { realpathSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { PassThrough } from 'node:stream'
import { providerDiagnosticOf } from '../../shared/agent-session-failure'
import { afterEach, describe, expect, it, vi } from 'vitest'
import type { spawnProcess } from '../../shared/child-process/run-process'
import {
isCodexAppServerRequestError,
openCodexAppServerConnection,
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(() => {
vi.useRealTimers()
vi.restoreAllMocks()
if (originalCodexHome === undefined) {
delete process.env.CODEX_HOME
} else {
process.env.CODEX_HOME = originalCodexHome
}
})
/**
* A real `node -e` child speaking the same JSONL framing Codex does. Slower than
* a stub, but it is the only thing that proves the spawn, the environment, and
* both traffic directions actually work end to end.
*/
const FAKE_APP_SERVER = String.raw`
const readline = require('node:readline')
const send = (payload) => process.stdout.write(JSON.stringify(payload) + '\n')
readline.createInterface({ input: process.stdin }).on('line', (line) => {
const message = JSON.parse(line)
if (message.method === 'initialize') return send({ id: message.id, result: {} })
if (message.method === 'test/env') {
return send({ id: message.id, result: { codexHome: process.env.CODEX_HOME ?? null } })
}
if (message.method === 'test/cwd') {
return send({ id: message.id, result: { cwd: process.cwd() } })
}
if (message.method === 'test/notify') {
send({ method: 'turn/started', params: { threadId: 'thread-1', turn: { id: 'turn-7' } } })
return send({ id: message.id, result: {} })
}
if (message.method === 'test/ask') {
return send({ id: 99, method: 'item/fileChange/requestApproval', params: { itemId: 'i1' } })
}
if (message.method === 'test/refuse') {
return send({ id: message.id, error: { code: -32602, message: 'bad params' } })
}
if (message.method === 'test/missing') {
return send({ id: message.id, error: { code: -32601, message: 'method not found' } })
}
if (message.id === 99) {
return send({ method: 'test/answered', params: message })
}
})
`
async function openFakeServer(
handlers: CodexAppServerConnectionHandlers = {},
env?: Record<string, string>,
envToDelete?: string[],
cwd?: string
): Promise<CodexAppServerConnection> {
return openCodexAppServerConnection(
{ command: process.execPath, args: ['-e', FAKE_APP_SERVER], env, envToDelete, cwd },
handlers
)
}
type StubChild = EventEmitter & {
stdout: PassThrough
stderr: PassThrough
stdin: PassThrough
pid: number
kill: ReturnType<typeof vi.fn>
}
/** Full control over framing and death, which a real child cannot give. */
function stubChild(options: { exitOnStdinEnd?: boolean } = {}): {
child: StubChild
spawnImpl: typeof spawnProcess
written: Record<string, unknown>[]
} {
const child = new EventEmitter() as StubChild
child.stdout = new PassThrough()
child.stderr = new PassThrough()
child.stdin = new PassThrough()
// Keep the synthetic pid outside any real process table so teardown never
// mistakes an unrelated process for this stub.
child.pid = 9_999_999
child.kill = vi.fn()
const written: Record<string, unknown>[] = []
child.stdin.on('data', (chunk: Buffer) => {
for (const line of chunk.toString('utf8').split('\n')) {
if (line.trim()) {
written.push(JSON.parse(line) as Record<string, unknown>)
}
}
})
if (options.exitOnStdinEnd !== false) {
child.stdin.on('finish', () => child.emit('exit', 0, null))
}
return { child, spawnImpl: (() => child) as unknown as typeof spawnProcess, written }
}
/** Answers the handshake so `openCodexAppServerConnection` can resolve. */
function answerInitialize(child: StubChild): void {
child.stdin.once('data', () => {
child.stdout.write(`${JSON.stringify({ id: 1, result: {} })}\n`)
})
}
/** Stream writes land a tick later, so the stderr tail is only complete here. */
async function flushStreams(): Promise<void> {
await new Promise((resolve) => setImmediate(resolve))
}
function rejection(promise: Promise<unknown>): Promise<Error> {
return promise.then(
() => {
throw new Error('expected the call to reject')
},
(error: Error) => error
)
}
function commandCompletionFixture(
targetBytes: number,
itemId = 'item-large'
): { line: string; output: string } {
const frame = {
method: 'item/completed',
params: {
turnId: 'turn-large',
item: { id: itemId, type: 'commandExecution', aggregated_output: '' }
}
}
const emptyBytes = Buffer.byteLength(JSON.stringify(frame), 'utf8')
const remaining = targetBytes - emptyBytes
if (remaining < 0) {
throw new Error(`target ${targetBytes} is smaller than fixture envelope ${emptyBytes}`)
}
const output = `${'\n'.repeat(Math.floor(remaining / 2))}${remaining % 2 ? 'x' : ''}`
frame.params.item.aggregated_output = output
const line = JSON.stringify(frame)
expect(Buffer.byteLength(line, 'utf8')).toBe(targetBytes)
return { line: `${line}\n`, output }
}
function commandCompletionLine(targetBytes: number): string {
return commandCompletionFixture(targetBytes).line
}
function responseLine(targetBytes: number, id: number): string {
const frame = { id, result: { data: '' } }
const emptyBytes = Buffer.byteLength(JSON.stringify(frame), 'utf8')
frame.result.data = 'x'.repeat(targetBytes - emptyBytes)
const line = JSON.stringify(frame)
expect(Buffer.byteLength(line, 'utf8')).toBe(targetBytes)
return `${line}\n`
}
describe('openCodexAppServerConnection', () => {
it('advertises the experimental API required for rollout-path resume', async () => {
const { child, spawnImpl, written } = stubChild()
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
expect(written[0]).toMatchObject({
method: 'initialize',
params: { capabilities: { experimentalApi: true } }
})
await connection.close()
})
it('reports the spawned pid before it sends the handshake', async () => {
const { child, spawnImpl, written } = stubChild()
answerInitialize(child)
const writtenAtSpawn: number[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onSpawned: async (pid) => {
writtenAtSpawn.push(written.length)
expect(pid).toBe(child.pid)
}
},
spawnImpl
)
// The owner is durable before initialize, so a crash mid-handshake leaves it stoppable.
expect(writtenAtSpawn).toEqual([0])
expect(written[0]).toMatchObject({ method: 'initialize' })
await connection.close()
})
it('reaps the child and never handshakes when its spawn cannot be recorded', async () => {
const { spawnImpl, written } = stubChild()
await expect(
openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onSpawned: async () => {
throw new Error('agent_session_checkpoint_stale')
}
},
spawnImpl
)
).rejects.toThrow('agent_session_checkpoint_stale')
expect(written).toEqual([])
})
it('completes the handshake and keeps the child alive across calls', async () => {
const notifications: { method: string; params: unknown }[] = []
const connection = await openFakeServer({
onNotification: (method, params) => notifications.push({ method, params })
})
await connection.request('test/notify')
await connection.request('test/notify')
expect(connection.pid).toBeGreaterThan(0)
expect(connection.closed).toBe(false)
expect(notifications).toHaveLength(2)
expect(notifications[0]).toEqual({
method: 'turn/started',
params: { threadId: 'thread-1', turn: { id: 'turn-7' } }
})
await connection.close()
expect(connection.closed).toBe(true)
})
it('applies the environment overlay after stripping inherited keys', async () => {
process.env.CODEX_HOME = '/tmp/inherited-home'
const pinned = await openFakeServer({}, { CODEX_HOME: '/tmp/pinned-home' })
expect(await pinned.request('test/env')).toEqual({ codexHome: '/tmp/pinned-home' })
await pinned.close()
const stripped = await openFakeServer({}, undefined, ['CODEX_HOME'])
expect(await stripped.request('test/env')).toEqual({ codexHome: null })
await stripped.close()
})
it('starts the provider in the resolved workspace directory', async () => {
const workspace = realpathSync(tmpdir())
const connection = await openFakeServer({}, undefined, undefined, workspace)
await expect(connection.request('test/cwd')).resolves.toEqual({ cwd: workspace })
await connection.close()
})
it('routes a server request to the handler and writes the reply back', async () => {
const requests: { id: number | string; method: string }[] = []
let resolveAnswered: (params: unknown) => void = () => {}
const answered = new Promise<unknown>((resolve) => {
resolveAnswered = resolve
})
const connection = await openFakeServer({
onServerRequest: (request) => {
requests.push({ id: request.id, method: request.method })
connection.respond(request.id, { decision: 'accept' })
},
onNotification: (method, params) => {
if (method === 'test/answered') {
resolveAnswered(params)
}
}
})
connection.notify('test/ask')
expect(await answered).toEqual({ id: 99, result: { decision: 'accept' } })
expect(requests).toEqual([{ id: 99, method: 'item/fileChange/requestApproval' }])
await connection.close()
})
it('classifies a refusal apart from a missing method', async () => {
const connection = await openFakeServer()
const refusal = await connection.request('test/refuse').catch((error: unknown) => error)
const missing = await connection.request('test/missing').catch((error: unknown) => error)
expect(isCodexAppServerRequestError(refusal)).toBe(true)
expect((refusal as Error).message).toContain('bad params')
// Codex's own words, apart from Orca's prefix, for a person to read.
expect(providerDiagnosticOf(refusal)).toEqual({ text: 'bad params', audience: 'person' })
expect(providerDiagnosticOf(missing)).toBeUndefined()
expect(isCodexAppServerUnsupportedError(missing)).toBe(true)
expect(isCodexAppServerRequestError(missing)).toBe(false)
await connection.close()
})
it('reassembles a message split mid-character across chunks', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const notifications: unknown[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onNotification: (_method, params) => notifications.push(params) },
spawnImpl
)
const payload = Buffer.from(
`${JSON.stringify({ method: 'item/agentMessage/delta', params: { delta: '日本語' } })}\n`,
'utf8'
)
const split = payload.indexOf(Buffer.from('日', 'utf8')) + 1
child.stdout.write(payload.subarray(0, split))
child.stdout.write(payload.subarray(split))
await vi.waitFor(() => expect(notifications).toHaveLength(1))
expect(notifications[0]).toEqual({ delta: '日本語' })
await connection.close()
})
it('surfaces valid but unclassified frames instead of dropping them', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const frames: { kind: string; payload: unknown }[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onUnhandledFrame: (kind, payload) => frames.push({ kind, payload }) },
spawnImpl
)
child.stdout.write(`${JSON.stringify({ id: 'late-string-id', result: { value: 1 } })}\n`)
child.stdout.write(`${JSON.stringify({ id: null, error: { message: 'parse error' } })}\n`)
await vi.waitFor(() => expect(frames).toHaveLength(2))
expect(frames.map((frame) => frame.kind)).toEqual(['frame:unclassified', 'frame:unclassified'])
await connection.close()
})
it('logs a reply to a timed-out request instead of surfacing it as a frame', async () => {
vi.useFakeTimers()
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
const { child, spawnImpl, written } = stubChild()
answerInitialize(child)
const frames: string[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onUnhandledFrame: (kind) => frames.push(kind) },
spawnImpl
)
const slow = rejection(connection.request('turn/interrupt', undefined, { timeoutMs: 50 }))
await vi.advanceTimersByTimeAsync(60)
expect((await slow).name).toBe('CodexAppServerTimeoutError')
const id = Number(written.find((frame) => frame.method === 'turn/interrupt')?.id)
child.stdout.write(`${JSON.stringify({ id, result: {} })}\n`)
child.stdout.write(`${JSON.stringify({ id: 999, error: { message: 'no such request' } })}\n`)
await vi.waitFor(() => expect(warn).toHaveBeenCalledTimes(2))
expect(frames).toEqual([])
expect(warn.mock.calls.map((call) => call[0])).toEqual([
`[codex-app-server] late reply to turn/interrupt after timeout (id ${id})`,
'[codex-app-server] reply with no waiting request (id 999)'
])
expect(warn.mock.calls[1][1]).toBe('no such request')
await vi.advanceTimersByTimeAsync(0)
})
it('fails in-flight requests and reports an unexpected exit once', async () => {
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const exits: string[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onExit: (error) => exits.push(error.message) },
spawnImpl
)
const inFlight = rejection(connection.request('turn/start'))
child.stderr.write('codex crashed\n')
await flushStreams()
child.emit('exit', 1, null)
child.emit('close', 1, null)
expect((await inFlight).message).toContain('codex crashed')
expect(exits).toHaveLength(1)
await connection.close()
})
it('classifies a CLI without the app-server subcommand as unsupported', async () => {
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
const opening = openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
).catch((error: unknown) => error)
child.stderr.write("error: unrecognized subcommand 'app-server'\n")
await flushStreams()
child.emit('exit', 2, null)
child.emit('close', 2, null)
expect(isCodexAppServerUnsupportedError(await opening)).toBe(true)
})
it('exposes an unproven handshake child for later cleanup', async () => {
vi.useFakeTimers()
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
child.stdin.once('data', () => {
child.stdout.write(
`${JSON.stringify({ id: 1, error: { code: -32602, message: 'initialize failed' } })}\n`
)
})
const opening = rejection(
openCodexAppServerConnection({ command: 'codex', args: ['app-server'] }, {}, spawnImpl)
)
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 3_500)
const error = (await opening) as Error & { connection?: CodexAppServerConnection }
expect(error.name).toBe('CodexAppServerHandshakeExitUnprovenError')
expect(error.connection).toBeDefined()
child.emit('close', 1, null)
await expect(error.connection?.close()).resolves.toBe(true)
})
it('times out one request without ending the connection', async () => {
vi.useFakeTimers()
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
const slow = rejection(connection.request('turn/start', undefined, { timeoutMs: 50 }))
await vi.advanceTimersByTimeAsync(60)
expect((await slow).name).toBe('CodexAppServerTimeoutError')
expect(connection.closed).toBe(false)
await vi.advanceTimersByTimeAsync(0)
})
it('kills a child that ignores stdin EOF', async () => {
vi.useFakeTimers()
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
child.kill.mockImplementation(() => {
child.emit('exit', null, 'SIGKILL')
return true
})
const closing = connection.close()
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 500)
await closing
await vi.waitFor(() => expect(child.kill).toHaveBeenCalledWith('SIGKILL'))
})
it('reports unproven close when forced termination did not produce an exit event', async () => {
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
const forcedKill = new Promise<void>((resolve) => {
child.kill.mockImplementation((signal) => {
if (signal === 'SIGKILL') {
resolve()
}
})
})
const closing = connection.close()
await flushStreams()
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS)
await forcedKill
await vi.advanceTimersByTimeAsync(0)
await vi.advanceTimersByTimeAsync(1_000)
await expect(closing).resolves.toBe(false)
expect(vi.getTimerCount()).toBe(0)
}, 10_000)
it('shares one eventual exit proof across concurrent close callers', async () => {
vi.useFakeTimers()
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
child.kill.mockImplementation(() => {
setTimeout(() => child.emit('exit', null, 'SIGKILL'), 10)
return true
})
const first = connection.close()
const second = connection.close()
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'])
})
it('allows a later close to observe exit after an unproven attempt', async () => {
vi.useFakeTimers()
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
const first = connection.close()
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 3_500)
await expect(first).resolves.toBe(false)
child.emit('exit', 0, null)
await expect(connection.close()).resolves.toBe(true)
})
// A root that outlived one kill is killed again by the next ask, which then proves it gone.
it('kills the root again on a close after an unproven attempt', async () => {
vi.useFakeTimers()
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
const first = connection.close()
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 3_500)
await expect(first).resolves.toBe(false)
child.kill.mockImplementation((signal) => {
if (signal === 'SIGKILL') {
setTimeout(() => child.emit('exit', null, 'SIGKILL'), 10)
}
return true
})
const second = connection.close()
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS + 3_500)
await expect(second).resolves.toBe(true)
expect(child.kill.mock.calls.filter(([signal]) => signal === 'SIGKILL')).toHaveLength(2)
})
it.each([1_090_188, 2_900_090])(
'accepts a realistic %i-byte escaped command completion and keeps processing',
async (frameBytes) => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const completed: unknown[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onNotification: (method, params) => {
if (method === 'item/completed') {
completed.push(params)
}
}
},
spawnImpl
)
const line = Buffer.from(commandCompletionLine(frameBytes), 'utf8')
const split = Math.floor(line.length / 3)
child.stdout.write(line.subarray(0, split))
child.stdout.write(line.subarray(split, split * 2))
child.stdout.write(line.subarray(split * 2))
child.stdout.write('{"method":"turn/completed","params":{"turn":{"id":"turn-large"}}}\n')
await vi.waitFor(() => expect(completed).toHaveLength(1))
expect(
(completed[0] as { item: { aggregated_output: string } }).item.aggregated_output.length
).toBeGreaterThan(500_000)
expect(connection.closed).toBe(false)
await connection.close()
}
)
it('accepts two realistic large command completions without losing either payload', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const completed: { item: { id: string; aggregated_output: string } }[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onNotification: (method, params) => {
if (method === 'item/completed') {
completed.push(params as { item: { id: string; aggregated_output: string } })
}
}
},
spawnImpl
)
const fixtures = [
commandCompletionFixture(1_090_188, 'item-large-a'),
commandCompletionFixture(2_900_090, 'item-large-b')
]
child.stdout.write(fixtures[0]!.line)
child.stdout.write(fixtures[1]!.line)
await vi.waitFor(() => expect(completed).toHaveLength(2))
expect(completed.map((entry) => entry.item.id)).toEqual(['item-large-a', 'item-large-b'])
expect(
completed.map((entry) => Buffer.byteLength(entry.item.aggregated_output, 'utf8'))
).toEqual(fixtures.map((fixture) => Buffer.byteLength(fixture.output, 'utf8')))
expect(connection.closed).toBe(false)
await connection.close()
})
it('accepts a response beyond the daemon wire limit and keeps the provider alive', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{},
spawnImpl
)
const large = connection.request('thread/resume')
child.stdout.write(responseLine(16 * 1024 * 1024 + 1, 2))
await expect(large).resolves.toMatchObject({ data: expect.any(String) })
expect(child.kill).not.toHaveBeenCalled()
expect(connection.closed).toBe(false)
const followup = connection.request('turn/start')
child.stdout.write('{"id":3,"result":{"turn":{"id":"turn-next"}}}\n')
await expect(followup).resolves.toEqual({ turn: { id: 'turn-next' } })
await connection.close()
})
it('delivers a notification beyond the daemon wire limit whole, never as an oversized frame', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const frames: string[] = []
const deltas: unknown[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onUnhandledFrame: (kind) => frames.push(kind),
onNotification: (method, params) => {
if (method === 'item/commandExecution/outputDelta') {
deltas.push(params)
}
}
},
spawnImpl
)
const params = { threadId: 'thread-1', itemId: 'exec-1', delta: '' }
params.delta = 'x'.repeat(16 * 1024 * 1024 + 1)
child.stdout.write(
`${JSON.stringify({ method: 'item/commandExecution/outputDelta', params })}\n`
)
await vi.waitFor(() => expect(deltas).toEqual([params]))
expect(frames).toEqual([])
expect(connection.closed).toBe(false)
await connection.close()
})
it('keeps malformed and non-object JSON non-fatal and processes the next record', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const frames: { kind: string; payload: unknown }[] = []
const notifications: string[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onUnhandledFrame: (kind, payload) => frames.push({ kind, payload }),
onNotification: (method) => notifications.push(method)
},
spawnImpl
)
child.stdout.write('not json\n[]\n{"method":"turn/completed","params":{}}\n')
await vi.waitFor(() => expect(notifications).toEqual(['turn/completed']))
expect(frames).toEqual([
{ kind: 'frame:invalid-json', payload: 'not json' },
{ kind: 'frame:invalid-json', payload: '[]' }
])
expect(connection.closed).toBe(false)
await connection.close()
})
it('pauses between coalesced records and resumes the retained remainder', async () => {
const { child, spawnImpl } = stubChild()
answerInitialize(child)
const notifications: string[] = []
let connection: CodexAppServerConnection
connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onNotification: (method) => {
notifications.push(method)
if (notifications.length === 1) {
connection.pauseReading?.()
}
}
},
spawnImpl
)
child.stdout.write(
'{"method":"item/started","params":{}}\n{"method":"item/completed","params":{}}\n'
)
await vi.waitFor(() => expect(notifications).toEqual(['item/started']))
connection.resumeReading?.()
await vi.waitFor(() => expect(notifications).toEqual(['item/started', 'item/completed']))
await connection.close()
})
it.each([
{
kind: 'notification',
frame: { method: 'turn/started', params: { turn: { id: 'turn-1' } } }
},
{
kind: 'server request',
frame: { id: 41, method: 'item/fileChange/requestApproval', params: { itemId: 'item-1' } }
}
])('surfaces a synchronous $kind handler failure as a terminal exit', async ({ frame }) => {
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const exits: string[] = []
const fail = (): never => {
throw new Error('structured sink failed')
}
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onNotification: fail,
onServerRequest: fail,
onExit: (error) => exits.push(error.message)
},
spawnImpl
)
child.kill.mockImplementation(() => {
child.emit('exit', null, 'SIGKILL')
return true
})
const inFlight = rejection(connection.request('turn/start'))
child.stdout.write(`${JSON.stringify(frame)}\n`)
expect((await inFlight).message).toContain('structured sink failed')
expect(exits).toEqual([expect.stringContaining('structured sink failed')])
expect(connection.closed).toBe(true)
await vi.waitFor(() => expect(child.kill).toHaveBeenCalledWith('SIGKILL'))
await connection.close()
})
it('reports one exit for a death that arrives through two listeners', async () => {
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const exits: string[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onExit: (error) => exits.push(error.message) },
spawnImpl
)
child.emit('error', new Error('provider transport failed'))
child.stderr.write('provider died\n')
await flushStreams()
child.emit('exit', null, 'SIGKILL')
child.emit('close', null, 'SIGKILL')
expect(exits).toHaveLength(1)
// The first cause survives; the generic exit that follows does not overwrite it.
expect(exits[0]).toContain('provider transport failed')
await connection.close()
})
it('does not report recovery for a handler failure until child exit is observed', async () => {
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const exits: string[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{
onNotification: () => {
throw new Error('structured sink failed')
},
onExit: (error) => exits.push(error.message)
},
spawnImpl
)
const inFlight = rejection(connection.request('turn/start'))
child.stdout.write('{"method":"turn/started","params":{}}\n')
await flushStreams()
expect(exits).toHaveLength(0)
expect((await inFlight).message).toContain('structured sink failed')
child.emit('exit', null, 'SIGKILL')
expect(exits).toHaveLength(1)
await connection.close()
})
it('treats a broken stdin pipe as the end of the transport', async () => {
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const exits: string[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onExit: (error) => exits.push(error.message) },
spawnImpl
)
child.kill.mockImplementation(() => {
child.emit('exit', null, 'SIGKILL')
return true
})
const inFlight = rejection(connection.request('turn/start'))
child.stdin.emit('error', new Error('write EPIPE'))
expect((await inFlight).message).toContain('EPIPE')
expect(exits).toHaveLength(1)
// A child nobody can write to is not a live session: the owner must see the
// connection as gone rather than keep issuing calls that can only time out.
expect(connection.closed).toBe(true)
await vi.waitFor(() => expect(child.kill).toHaveBeenCalledWith('SIGKILL'))
expect((await rejection(connection.request('turn/start'))).message).toContain('EPIPE')
await connection.close()
})
it('reports a graceful close as expected when stdin breaks during the reap', async () => {
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
const { child, spawnImpl } = stubChild({ exitOnStdinEnd: false })
answerInitialize(child)
const exits: (boolean | undefined)[] = []
const connection = await openCodexAppServerConnection(
{ command: 'codex', args: ['app-server'] },
{ onExit: (_error, exit) => exits.push(exit?.expected) },
spawnImpl
)
child.stdin.on('finish', () => child.stdin.emit('error', new Error('write EPIPE')))
const forcedKill = new Promise<void>((resolve) => {
child.kill.mockImplementation((signal) => {
child.emit('exit', null, 'SIGKILL')
if (signal === 'SIGKILL') {
resolve()
}
return true
})
})
const inFlight = rejection(connection.request('turn/start'))
const closing = connection.close()
await flushStreams()
await vi.advanceTimersByTimeAsync(GRACEFUL_EXIT_MS)
await forcedKill
await vi.advanceTimersByTimeAsync(0)
await expect(closing).resolves.toBe(true)
expect((await inFlight).message).toContain('EPIPE')
// The root's exit is the close's own end, never an unexpected death.
expect(exits).toEqual([true])
expect(vi.getTimerCount()).toBe(0)
})
})