fix(native-chat): recover unexpected Claude exits

This commit is contained in:
Merge Sim
2026-08-31 17:44:32 -07:00
parent 077fae0c8f
commit 8cf9bc92c0
9 changed files with 244 additions and 66 deletions
@@ -9,6 +9,8 @@ import type { ClaudeSession } from './claude-structured-session-state'
function sessionFor(send = vi.fn().mockResolvedValue(undefined)): ClaudeSession {
return {
connection: { send } as unknown as ClaudeSession['connection'],
ended: false,
requestedClose: false,
providerSessionId: 'provider-session',
leafUuid: null,
fence: 1,
@@ -17,7 +19,8 @@ function sessionFor(send = vi.fn().mockResolvedValue(undefined)): ClaudeSession
options: new Map(),
reportedOptions: {},
events: undefined,
translator: null
translator: null,
acquisitionGeneration: 'generation-test'
}
}
@@ -23,12 +23,36 @@ export class ClaudeStructuredProviderEvents {
}
handleExit(sessionId: string, attempt: ClaudeAcquisitionAttempt, error: Error): void {
if (attempt.connection) {
attempt.exitProven = true
}
const session = this.sessions.get(sessionId)
if (!session || session.connection !== attempt.connection) {
if (!session || session.connection !== attempt.connection || session.ended) {
return
}
this.sessions.delete(sessionId)
this.emit(session, session.events, { type: 'ended', sessionId, reason: error.message })
this.finishExit(sessionId, session, error)
}
handleClosed(sessionId: string, error: Error): boolean {
const session = this.sessions.get(sessionId)
if (!session || session.ended) {
return false
}
this.finishExit(sessionId, session, error)
return true
}
private finishExit(sessionId: string, session: ClaudeSession, error: Error): void {
const event = {
type: 'ended' as const,
sessionId,
reason: error.message,
cause: session.requestedClose ? ('requested-close' as const) : ('unexpected-exit' as const),
fence: session.fence,
acquisitionGeneration: session.acquisitionGeneration
}
session.ended = true
this.emit(session, session.events, event)
settleClaudeExitedSession(session)
}
@@ -226,6 +226,72 @@ describe('ClaudeStructuredSessionAdapter.acquire', () => {
expect(events[0]).toMatchObject({ type: 'message', message: { subtype: 'init' } })
})
it('reports an unexpected child exit with the shared lifecycle identity', async () => {
const claude = fakeClaude()
const events: ClaudeStructuredSessionEvent[] = []
const adapter = await acquired(claude, {}, events)
claude.connections[0].handlers.onExit?.(new Error('provider exited'))
expect(events.at(-1)).toMatchObject({
type: 'ended',
sessionId: 'session-1',
reason: 'provider exited',
cause: 'unexpected-exit',
fence: 7,
acquisitionGeneration: expect.any(String)
})
await expect(
Promise.resolve().then(() =>
adapter.dispatch({
sessionId: 'session-1',
clientMessageId: 'client-1',
body: USER_MESSAGE,
fence: 7
})
)
).rejects.toThrow('no live claude stream-json')
})
it('drops a late exit from a superseded child after reacquisition', async () => {
const claude = fakeClaude()
const events: ClaudeStructuredSessionEvent[] = []
const adapter = await acquired(claude, {}, events)
const first = claude.connections[0]
await adapter.acquire({ identity: identityFor(), fence: 8, spawnToken: 'spawn-10' })
const endedBeforeLateExit = events.filter((event) => event.type === 'ended').length
first.handlers.onExit?.(new Error('stale child exited'))
expect(events.filter((event) => event.type === 'ended')).toHaveLength(endedBeforeLateExit)
})
it('retries teardown when a child cannot yet prove exit', async () => {
const claude = fakeClaude()
const adapter = await acquired(claude)
const connection = claude.connections[0]
connection.close = async () => {
connection.closeCount += 1
return false
}
await expect(adapter.closeAll()).rejects.toThrow(
'claude structured session teardown could not prove provider-child exit'
)
expect(connection.closeCount).toBe(3)
})
it('marks a sink-triggered kill as an unexpected exit', async () => {
const claude = fakeClaude()
const events: ClaudeStructuredSessionEvent[] = []
const adapter = await acquired(claude, {}, events)
await expect(adapter.forceCloseSession('session-1')).resolves.toBe(true)
expect(events.filter((event) => event.type === 'ended')).toMatchObject([
{ cause: 'unexpected-exit', fence: 7, acquisitionGeneration: expect.any(String) }
])
})
it('restores persisted model and effort before publishing a reacquired session', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude, { resumed: true })
@@ -32,10 +32,12 @@ import { ClaudeStructuredProviderEvents } from './claude-structured-provider-eve
import {
cancelClaudeAcquisitionAttempt,
ClaudeAcquisitionRegistry,
mintClaudeAcquisitionGeneration,
type ClaudeSession,
type ClaudeStructuredSessionAdapterDeps
} from './claude-structured-session-state'
import { closeClaudePublishedSession } from './claude-structured-session-close'
import { closeProcessRegistry } from '../../shared/child-process/close-process-registry'
export type { ClaudeStructuredLaunch } from './claude-structured-launch-resolution'
export type {
@@ -162,13 +164,17 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
process,
options: restoredClaudeStructuredSessionOptions(input.options),
...(this.deps.mintLinkId ? { linkId: this.deps.mintLinkId() } : {}),
observedAt: this.deps.now?.() ?? Date.now()
observedAt: this.deps.now?.() ?? Date.now(),
acquisitionGeneration: mintClaudeAcquisitionGeneration(this.deps)
})
const acquired: AgentSessionAcquisition = publication.acquisition
const liveSession = publication.session
acquisitionEvents.publish(liveSession)
await restoreClaudeStructuredSessionOptions(liveSession, this.deps.requestTimeoutMs)
this.acquisitions.assertCurrent(sessionId, attempt)
if (attempt.exitProven || connection.closed) {
throw new Error(`claude stream-json for session ${sessionId} exited while being acquired`)
}
this.acquisitions.deleteIfCurrent(sessionId, attempt)
this.sessions.set(sessionId, liveSession)
attempt.published = true
@@ -248,23 +254,35 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
return closeClaudePublishedSession({
sessions: this.sessions,
sessionId,
providerEvents: this.events,
...(this.deps.persistHandle ? { persistHandle: this.deps.persistHandle } : {}),
...(this.deps.onEvent ? { onEvent: this.deps.onEvent } : {})
})
}
forceCloseSession = (sessionId: string): Promise<boolean> =>
closeClaudePublishedSession({
sessions: this.sessions,
sessionId,
providerEvents: this.events,
requestedClose: false,
allowFailedSettlement: true
})
async closeAll(): Promise<void> {
this.acquisitions.close()
const ids = new Set([...this.sessions.keys(), ...this.acquisitions.sessionIds()])
const outcomes = await Promise.all([...ids].map((sessionId) => this.closeSession(sessionId)))
if (outcomes.some((exited) => !exited)) {
throw new Error('claude structured session teardown could not prove provider-child exit')
}
await closeProcessRegistry({
attempts: 3,
hasEntries: () => this.sessions.size > 0 || this.acquisitions.size > 0,
entryIds: () => new Set([...this.sessions.keys(), ...this.acquisitions.sessionIds()]),
closeEntry: (sessionId) => this.closeSession(sessionId),
failureMessage: 'claude structured session teardown could not prove provider-child exit'
})
}
private session(sessionId: string): ClaudeSession {
const session = this.sessions.get(sessionId)
if (!session) {
if (!session || session.ended) {
throw new Error(`no live claude stream-json session for ${sessionId}`)
}
return session
@@ -1,4 +1,5 @@
import type { ClaudeSession, ClaudeStructuredSessionEvent } from './claude-structured-session-state'
import type { ClaudeStructuredProviderEvents } from './claude-structured-provider-events'
export function settleClaudeDispatchWaiters(session: ClaudeSession): void {
for (const waiter of session.dispatchWaiters.splice(0)) {
@@ -23,11 +24,25 @@ export async function closeClaudePublishedSession(input: {
fence: number
}) => Promise<void>
onEvent?: (event: ClaudeStructuredSessionEvent) => void
providerEvents?: ClaudeStructuredProviderEvents
requestedClose?: boolean
expectedFence?: number
expectedAcquisitionGeneration?: string
unexpectedReason?: Error
allowFailedSettlement?: boolean
}): Promise<boolean> {
const session = input.sessions.get(input.sessionId)
if (!session) {
return true
}
if (
(input.expectedFence !== undefined && session.fence !== input.expectedFence) ||
(input.expectedAcquisitionGeneration !== undefined &&
session.acquisitionGeneration !== input.expectedAcquisitionGeneration)
) {
return false
}
session.requestedClose = input.requestedClose ?? true
settleClaudeDispatchWaiters(session)
const pending = session.prompts.clear()
await Promise.allSettled(
@@ -44,34 +59,36 @@ export async function closeClaudePublishedSession(input: {
if (!(await session.connection.close())) {
return false
}
input.sessions.delete(input.sessionId)
let persistenceError: unknown
try {
await input.persistHandle?.({
sessionId: input.sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
input.onEvent?.({
type: 'handle',
sessionId: input.sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
} catch (error) {
persistenceError = error
} finally {
const ended = {
type: 'ended',
sessionId: input.sessionId,
reason: 'claude session closed'
} as const
session.translator?.handle(ended)
input.onEvent?.(ended)
session.translator?.dispose()
if (!session.ended && session.requestedClose) {
try {
await input.persistHandle?.({
sessionId: input.sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
input.onEvent?.({
type: 'handle',
sessionId: input.sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
} catch (error) {
persistenceError = error
}
}
if (!session.ended) {
input.providerEvents?.handleClosed(
input.sessionId,
input.unexpectedReason ?? new Error('claude session closed')
)
}
if (!session.ended) {
return false
}
input.sessions.delete(input.sessionId)
if (persistenceError) {
throw persistenceError
}
@@ -17,6 +17,7 @@ export function createClaudeSessionPublication(input: {
process: AgentSessionAcquisition['process']
linkId?: string
observedAt: number
acquisitionGeneration: string
options?: ReadonlyMap<string, string>
}): { acquisition: AgentSessionAcquisition; session: ClaudeSession } {
const model = readClaudeFrameString(input.init.message, 'model')
@@ -24,6 +25,7 @@ export function createClaudeSessionPublication(input: {
return {
acquisition: {
process: input.process,
acquisitionGeneration: input.acquisitionGeneration,
link: claudeProviderHandleLink({
sessionId: input.init.providerSessionId,
leafUuid: input.leafUuid,
@@ -35,6 +37,8 @@ export function createClaudeSessionPublication(input: {
},
session: {
connection: input.connection,
ended: false,
requestedClose: false,
providerSessionId: input.init.providerSessionId,
leafUuid: input.leafUuid,
fence: input.fence,
@@ -46,7 +50,8 @@ export function createClaudeSessionPublication(input: {
...(effort ? { effort } : {})
},
translator: input.translator,
events: input.events
events: input.events,
acquisitionGeneration: input.acquisitionGeneration
}
}
}
@@ -1,4 +1,7 @@
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
import { randomUUID } from 'node:crypto'
import { cancelProcessAcquisition } from '../../shared/child-process/cancel-process-acquisition'
import type { StructuredAgentSessionLifecycleEvent } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type {
ClaudeStreamJsonConnection,
@@ -30,7 +33,7 @@ export type ClaudeStructuredSessionEvent =
fence: number
}
| { type: 'auth-diagnostic'; sessionId: string; diagnostic: ClaudeAuthDiagnostic }
| { type: 'ended'; sessionId: string; reason: string }
| StructuredAgentSessionLifecycleEvent
export type ClaudeStructuredSessionAdapterDeps = {
resolveLaunch: (input: {
@@ -41,6 +44,7 @@ export type ClaudeStructuredSessionAdapterDeps = {
readProcessStartTime?: (pid: number) => Promise<number | null>
mintLinkId?: () => string
now?: () => number
mintAcquisitionGeneration?: () => string
requestTimeoutMs?: number
initTimeoutMs?: number
dispatchAckTimeoutMs?: number
@@ -76,6 +80,8 @@ export type ClaudeDispatchWaiter = {
export type ClaudeSession = {
connection: ClaudeStreamJsonConnection
ended: boolean
requestedClose: boolean
providerSessionId: string
leafUuid: string | null
fence: number
@@ -85,6 +91,7 @@ export type ClaudeSession = {
reportedOptions: { model?: string; effort?: string }
translator: ClaudeJournalTranslator | null
events: StructuredAgentSessionEventSink | undefined
acquisitionGeneration: string
}
export type ClaudeAcquisitionAttempt = {
@@ -93,6 +100,7 @@ export type ClaudeAcquisitionAttempt = {
buffered: (() => void)[]
published: boolean
cancelled: boolean
exitProven: boolean
finished: Promise<void>
finish: () => void
}
@@ -110,11 +118,18 @@ export function createClaudeAcquisitionAttempt(
buffered: [],
published: false,
cancelled: false,
exitProven: false,
finished,
finish
}
}
export function mintClaudeAcquisitionGeneration(
deps: Pick<ClaudeStructuredSessionAdapterDeps, 'mintAcquisitionGeneration'>
): string {
return deps.mintAcquisitionGeneration?.() ?? randomUUID()
}
export class ClaudeAcquisitionRegistry {
private readonly attempts = new Map<string, ClaudeAcquisitionAttempt>()
private closing = false
@@ -170,8 +185,12 @@ export async function cancelClaudeAcquisitionAttempt(
if (!attempt) {
return true
}
attempt.cancelled = true
const exited = attempt.connection ? await attempt.connection.close() : true
await attempt.finished
return exited
return cancelProcessAcquisition({
cancel: () => {
attempt.cancelled = true
},
connection: () => attempt.connection,
exitProven: () => attempt.exitProven,
finished: attempt.finished
})
}
@@ -0,0 +1,37 @@
export type StructuredAgentSessionExitEvent = {
type: string
sessionId: string
cause?: string
}
export function createStructuredAgentSessionExitRecovery(
handle: (event: StructuredAgentSessionExitEvent) => Promise<unknown> | undefined,
reportError: (sessionId: string, error: unknown) => void
): {
queue: (event: StructuredAgentSessionExitEvent) => void
wait: () => Promise<void>
} {
let chain: Promise<void> = Promise.resolve()
const queue = (event: StructuredAgentSessionExitEvent): void => {
if (event.type !== 'ended' || event.cause !== 'unexpected-exit') {
return
}
chain = chain
.then(async () => {
await handle(event)
})
.catch((error) => reportError(event.sessionId, error))
}
return {
queue,
wait: async () => {
for (;;) {
const observed = chain
await observed
if (observed === chain) {
return
}
}
}
}
}
@@ -39,6 +39,7 @@ import { readEchoedAgentSessionSpawnToken } from './agent-session-spawn-token-re
import { agentSessionPtyWriteGate } from './agent-session-pty-write-gate'
import { resolveLoginShellEnvironment } from '../startup/login-shell-environment'
import { recordAgentSessionProviderHandle } from './agent-session-provider-handle-transition'
import { createStructuredAgentSessionExitRecovery } from './structured-agent-session-exit-recovery'
/** Sibling of the journal tree rather than inside it: one file adjudicates every
* session's lease, while a journal is per session. */
@@ -161,7 +162,14 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
})
try {
let host: StructuredAgentSessionHost | null = null
let recoveryChain = Promise.resolve()
const recovery = createStructuredAgentSessionExitRecovery(
(event) =>
host?.handleAdapterEvent(
event as Parameters<StructuredAgentSessionHost['handleAdapterEvent']>[0]
),
(sessionId, error) =>
deps.onError?.({ scope: `structured-agent-session-exit:${sessionId}`, error })
)
const codex = new CodexStructuredSessionAdapter({
resolveLaunch: createCodexStructuredLaunchResolver({
store,
@@ -171,21 +179,7 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
}),
...(deps.openCodexConnection ? { openConnection: deps.openCodexConnection } : {}),
...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {}),
onEvent: (event) => {
if (event.type !== 'ended' || !('cause' in event) || event.cause !== 'unexpected-exit') {
return
}
// Serialize recovery with teardown. Exit callbacks arrive from child
// process tasks, so a fire-and-forget callback can otherwise append
// after the host has flushed and its journal directory is removed.
recoveryChain = recoveryChain.then(async () => {
try {
await host?.handleAdapterEvent(event)
} catch (error) {
deps.onError?.({ scope: `structured-agent-session-exit:${event.sessionId}`, error })
}
})
}
onEvent: recovery.queue
})
const claude = new ClaudeStructuredSessionAdapter({
resolveLaunch: createClaudeStructuredLaunchResolver({
@@ -196,6 +190,7 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
}),
...(deps.openClaudeConnection ? { openConnection: deps.openClaudeConnection } : {}),
readProcessStartTime: deps.readClaudeProcessStartTime ?? deps.readProcessStartTime,
onEvent: recovery.queue,
persistHandle: async (observed) => {
const now = Date.now()
await store.transitionHandoff(observed.sessionId, (record) =>
@@ -246,13 +241,7 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
waitForRecovery: async () => {
// A recovery may synchronously trigger another exit while it is
// reacquiring. Observe until the chain stops growing.
for (;;) {
const observed = recoveryChain
await observed
if (observed === recoveryChain) {
return
}
}
await recovery.wait()
}
}
} catch (error) {