Merge remote-tracking branch 'origin/main' into brennanb2025/structured-chat-liveness-gate

This commit is contained in:
Merge Sim
2026-09-06 23:37:13 -07:00
28 changed files with 586 additions and 27 deletions
@@ -70,6 +70,7 @@ async function failingFlowRunner(
throw new Error('unused')
},
importTuiHistory: async () => {},
retryPendingSettlement: async () => true,
publish: () => {},
schedule: async () => {
throw new Error('scheduling failed')
@@ -0,0 +1,52 @@
import { describe, expect, it, vi } from 'vitest'
import {
agentSessionLeaseFixture,
agentSessionRecordFixture
} from '../../../shared/agent-session-record.test-fixture'
import type { StructuredAgentSessionHandoffDeps } from './structured-agent-session-handoff-types'
import { closeRetainedTuiOwner } from './structured-agent-session-handoff-owner-close'
const NOW = 1_800_000_000_000
describe('closeRetainedTuiOwner', () => {
it('does not latch unexpected-exit settlement after an intentional close', async () => {
let record = agentSessionRecordFixture(agentSessionLeaseFixture())
const closeTuiOwner = vi.fn(async () => ({}))
const releaseOwner = vi.fn()
const owner = {
terminal: { handle: 'terminal-1', tabId: 'tab-1', paneKey: 'pane-1', ptyId: 'pty-1' },
process: record.lease.ownerProcess!,
link: record.providerHandleChain[0]!
}
const deps = {
store: {
transitionHandoff: async (
_sessionId: string,
transition: (current: typeof record) => typeof record
) => {
record = transition(record)
return record
}
},
transport: { closeTuiOwner },
now: () => NOW
} as unknown as StructuredAgentSessionHandoffDeps
await closeRetainedTuiOwner({
sessionId: record.sessionId,
deps,
owner: () => owner,
requireRecord: () => record,
releaseOwner
})
expect(closeTuiOwner).toHaveBeenCalledWith(owner)
expect(releaseOwner).toHaveBeenCalledWith(record.sessionId)
expect(record.lease).toMatchObject({
claimStatus: 'released',
deathEvidence: { kind: 'exit-observed' }
})
expect(record.lease.settlementRetryRequired).toBeUndefined()
expect(record.lease.settlementRetryId).toBeUndefined()
})
})
@@ -26,7 +26,8 @@ export async function closeRetainedTuiOwner(input: {
record: current,
expectedFence: record.lease.runtimeFence,
probe: { outcome: 'exit-observed' },
now: input.deps.now()
now: input.deps.now(),
journalSettlement: 'not-required'
})
)
input.releaseOwner(input.sessionId)
@@ -95,6 +95,13 @@ export async function handoffStructuredSessionToNative(
...(transcriptPath ? { transcriptPath } : {})
})
}
if (record.lease.settlementRetryRequired) {
const settled = await deps.retryPendingSettlement(sessionId)
if (!settled) {
throw new Error('The provider-exit terminal journal settlement is still pending.')
}
record = context.requireRecord(sessionId)
}
const spawnToken = randomUUID()
record = await reserveStoredAgentSessionHandoffOwner(deps.store, {
sessionId,
@@ -73,6 +73,7 @@ export function createStructuredAgentSessionHandoffTestCoordinator(
{ fence, recovered: true }
)
},
retryPendingSettlement: async () => true,
publish: (_sessionId, status) => input.statuses.push(status),
schedule: async (_sessionId, task) => task(),
now: () => input.now
@@ -81,6 +81,7 @@ export type StructuredAgentSessionHandoffDeps = {
fence: number
transcriptPath?: string
}) => Promise<void>
retryPendingSettlement: (sessionId: string) => Promise<boolean>
prepareTuiHistoryCatchup?: (sessionId: string, fence: number) => Promise<void>
recoverTuiHistoryCatchup?: (sessionId: string, fence: number) => Promise<void>
activateTuiHistoryCatchup?: (sessionId: string) => Promise<void>
@@ -179,6 +179,7 @@ function createCoordinator(): StructuredAgentSessionHandoffCoordinator {
{ fence, recovered: true }
)
},
retryPendingSettlement: async () => true,
prepareTuiHistoryCatchup,
recoverTuiHistoryCatchup,
activateTuiHistoryCatchup,
@@ -274,6 +275,7 @@ describe('structured session handoff failure handling', () => {
}),
acquireNativeStop: (_sessionId, turnId) => acquireNativeStop(turnId),
importTuiHistory: vi.fn(async () => undefined),
retryPendingSettlement: vi.fn(async () => true),
prepareTuiHistoryCatchup,
recoverTuiHistoryCatchup,
activateTuiHistoryCatchup,
@@ -14,6 +14,7 @@ import { recoverDeadTuiHandoffStatus } from './structured-agent-session-dead-tui
import { readNativeSessionOptions } from './structured-agent-session-option-restoration'
import type { AgentSessionSubscribers } from './structured-agent-session-subscribers'
import { StructuredTuiTranscriptCatchup } from './structured-tui-transcript-catchup'
import { retryLoadedStructuredAgentSessionSettlement } from './structured-agent-session-settlement-retry'
type HostHandoffAccess = {
session: (sessionId: string) => StructuredAgentSessionHostSession
@@ -92,6 +93,13 @@ export function createStructuredAgentSessionHostHandoff(
acquireNativeStop: async (sessionId, turnId, fence) =>
(await deps.adapter.cancelTurn({ sessionId, turnId, fence })).cancelled,
importTuiHistory: (input) => importTuiHistory(deps, host, input),
retryPendingSettlement: (sessionId) =>
retryLoadedStructuredAgentSessionSettlement({
deps,
sessionId,
session: host.session(sessionId),
now: host.now
}),
prepareTuiHistoryCatchup: (sessionId, fence) => tuiHistoryCatchup.prepare(sessionId, fence),
recoverTuiHistoryCatchup: (sessionId, fence) => tuiHistoryCatchup.recover(sessionId, fence),
activateTuiHistoryCatchup: (sessionId) => tuiHistoryCatchup.activate(sessionId),
@@ -127,7 +127,7 @@ export class StructuredAgentSessionHost {
),
evict: (sessionId) => this.close(sessionId)
})
this.restore = createStructuredAgentSessionHostRestore(deps, {
this.restore = createStructuredAgentSessionHostRestore(deps, this.sessions, () => this.now(), {
reconcile: this.reconcileLeases,
resolveRecovery: (sessionId) => this.runtimeState.resolveRecovery(sessionId),
serialize: (sessionId, task) => this.serialize(sessionId, task),
@@ -104,6 +104,7 @@ describe('structured session live TUI restart survival', () => {
suspendNative: vi.fn(),
acquireNative: vi.fn(),
importTuiHistory: vi.fn(),
retryPendingSettlement: vi.fn(async () => true),
publish: vi.fn(),
schedule: async (_sessionId, task) => task(),
now: () => NOW
@@ -224,6 +225,7 @@ describe('structured session live TUI restart survival', () => {
suspendNative: vi.fn(),
acquireNative: vi.fn(),
importTuiHistory: vi.fn(),
retryPendingSettlement: vi.fn(async () => true),
publish: vi.fn(),
schedule: async (_sessionId, task) => task(),
now: () => NOW
@@ -3,12 +3,14 @@ import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
import { activeStructuredAgentSessionTurnId } from '../../../shared/structured-agent-session-projection'
import type { AgentSessionHandoffRequest } from '../../../shared/agent-session-wire'
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import { recoverStoredDeadTuiOwnerForHandoff } from '../../runtime/agent-session-handoff-record-transitions'
import { openAgentSessionJournal } from '../agent-session-journal/journal-store-factory'
import { StructuredAgentSessionHandoffCoordinator } from './structured-agent-session-handoff'
import type { StructuredAgentSessionHandoffTransport } from './structured-agent-session-handoff-types'
import { retryLoadedStructuredAgentSessionSettlement } from './structured-agent-session-settlement-retry'
const NOW = 1_800_000_000_000
const SESSION = 'session-proven-dead-retry'
@@ -89,6 +91,15 @@ describe('structured session proven-dead TUI retry', () => {
},
journalDir: join(root, 'journal')
})
await journal.appendItem(
{ provider: 'orca', clientMessageId: 'running-turn' },
{
kind: 'status',
text: 'Working',
turnLifecycle: { turnId: 'turn-1', state: 'running' }
},
{ fence: store.getRecord(SESSION)?.lease.runtimeFence ?? tuiFence }
)
const closeTuiOwner =
vi.fn<NonNullable<StructuredAgentSessionHandoffTransport['closeTuiOwner']>>()
const coordinator = new StructuredAgentSessionHandoffCoordinator({
@@ -134,6 +145,17 @@ describe('structured session proven-dead TUI retry', () => {
},
acquireNativeStop: vi.fn(async () => true),
importTuiHistory: vi.fn(),
retryPendingSettlement: (sessionId) =>
retryLoadedStructuredAgentSessionSettlement({
deps: { store },
sessionId,
session: {
journal,
fence: store.getRecord(sessionId)?.lease.runtimeFence ?? 1,
acquisitionGeneration: null
},
now: () => NOW
}),
publish: vi.fn(),
schedule: async (_sessionId, task) => task(),
now: () => NOW
@@ -174,5 +196,7 @@ describe('structured session proven-dead TUI retry', () => {
claimStatus: 'live',
handoffStage: null
})
expect(store.getRecord(SESSION)?.lease.settlementRetryRequired).toBeUndefined()
expect(activeStructuredAgentSessionTurnId(journal.snapshot().items)).toBe(null)
})
})
@@ -27,6 +27,7 @@ describe('StructuredAgentSessionReadableRestorer', () => {
serialize: async (_sessionId, task) => task(),
hasSession: () => false,
onReadable: () => undefined,
retrySettlement: async () => true,
restoreHandoff: async () => undefined
})
@@ -20,6 +20,10 @@ export class StructuredAgentSessionReadableRestorer {
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
hasSession: (sessionId: string) => boolean
onReadable: (sessionId: string, restored: RestoredStructuredAgentSessionRead) => void
retrySettlement: (
sessionId: string,
params: RestoredStructuredAgentSessionRead['params']
) => Promise<boolean>
restoreHandoff: (sessionId: string) => Promise<void>
}
) {}
@@ -1,5 +1,6 @@
import { beforeEach, describe, expect, it, vi } from 'vitest'
import type { AgentSessionRecord } from '../../../shared/agent-session-record'
import type { AgentSessionAttachParams } from './structured-agent-session-attach'
const { restoreRead } = vi.hoisted(() => ({
restoreRead: vi.fn()
@@ -9,7 +10,10 @@ vi.mock('./structured-agent-session-read-restore', () => ({
restoreStructuredAgentSessionRead: restoreRead
}))
import { restoreStructuredAgentSessionsOnRestart } from './structured-agent-session-restart-restore'
import {
restoreOneStructuredAgentSessionRead,
restoreStructuredAgentSessionsOnRestart
} from './structured-agent-session-restart-restore'
describe('restart journal restoration', () => {
beforeEach(() => restoreRead.mockReset())
@@ -45,6 +49,7 @@ describe('restart journal restoration', () => {
serialize: async (_sessionId, task) => task(),
hasSession: () => false,
onReadable: () => undefined,
retrySettlement: async () => true,
restoreHandoff: async () => undefined
})
@@ -56,4 +61,96 @@ describe('restart journal restoration', () => {
expect(restoreRead).toHaveBeenCalledTimes(records.length)
expect(peak).toBe(4)
})
it('runs pending settlement retry after recovery resolution and before handoff', async () => {
const calls: string[] = []
const params: AgentSessionAttachParams = {
envelope: {
sessionId: 'session-1',
clientOperationId: 'read-restore:session-1',
expectedRuntimeFence: 4,
payloadFingerprint: 'fingerprint'
},
location: {
executionHostId: 'local',
wslDistro: null,
workspaceId: 'workspace-1',
workspaceKind: 'folder'
},
provider: 'codex',
agent: 'codex',
accountHome: { variable: 'CODEX_HOME', path: '/tmp/codex' },
runtimeKind: 'native'
}
restoreRead.mockResolvedValue({
journal: {},
params,
fence: 4,
hasProviderChild: false,
acquisitionGeneration: null
})
await restoreOneStructuredAgentSessionRead(
{
store: {} as never,
journalRoot: '/tmp/journals',
reconcile: async () => null,
resolveRecovery: async () => {
calls.push('resolveRecovery')
},
serialize: async (_sessionId, task) => task(),
hasSession: () => false,
onReadable: () => {
calls.push('onReadable')
},
retrySettlement: async (_sessionId, restoredParams) => {
calls.push(
restoredParams === params ? 'retrySettlement:restored-params' : 'retrySettlement'
)
return true
},
restoreHandoff: async () => {
calls.push('restoreHandoff')
}
},
'session-1'
)
expect(calls).toEqual([
'resolveRecovery',
'onReadable',
'retrySettlement:restored-params',
'restoreHandoff'
])
})
it('does not rerun settlement retry when a second restore finds the session already open', async () => {
const retrySettlement = vi.fn(async () => true)
const restoreHandoff = vi.fn(async () => undefined)
restoreRead.mockResolvedValue({
journal: {},
params: {},
fence: 4,
hasProviderChild: false,
acquisitionGeneration: null
})
await restoreOneStructuredAgentSessionRead(
{
store: {} as never,
journalRoot: '/tmp/journals',
reconcile: async () => null,
resolveRecovery: async () => undefined,
serialize: async (_sessionId, task) => task(),
hasSession: () => true,
onReadable: () => undefined,
retrySettlement,
restoreHandoff
},
'session-1'
)
expect(retrySettlement).not.toHaveBeenCalled()
expect(restoreHandoff).toHaveBeenCalledOnce()
})
})
@@ -30,6 +30,10 @@ export type StructuredAgentSessionReadRestoreDeps = {
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
hasSession: (sessionId: string) => boolean
onReadable: (sessionId: string, restored: RestoredStructuredAgentSessionRead) => void
retrySettlement: (
sessionId: string,
params: RestoredStructuredAgentSessionRead['params']
) => Promise<boolean>
restoreHandoff: (sessionId: string) => Promise<void>
}
@@ -64,6 +68,7 @@ export async function restoreOneStructuredAgentSessionRead(
return
}
input.onReadable(sessionId, restored)
await input.retrySettlement(sessionId, restored.params)
await input.restoreHandoff(sessionId)
})
}
@@ -58,6 +58,7 @@ function harness(
serialize,
hasSession: (sessionId) => live.has(sessionId),
onReadable: (sessionId, restored) => live.set(sessionId, restored),
retrySettlement: async () => true,
restoreHandoff
})
return { restorer, live, restoreHandoff, serializedIds }
@@ -16,8 +16,10 @@ import { StructuredAgentSessionReadableRestorer } from './structured-agent-sessi
import { StructuredAgentSessionRestartRestoreGate } from './structured-agent-session-restart-restore-gate'
import type {
StructuredAgentSessionHostDeps,
StructuredAgentSessionHostSession,
StructuredAgentSessionReveal
} from './structured-agent-session-host-types'
import { retryPendingStructuredAgentSessionSettlement } from './structured-agent-session-settlement-retry'
/** Throws its refusal as the code itself, matching `resumeHeldStructuredAgentSession`. */
export async function revealStructuredAgentSession(
@@ -55,9 +57,11 @@ export async function revealStructuredAgentSession(
*/
export function createStructuredAgentSessionHostRestore(
deps: StructuredAgentSessionHostDeps,
sessions: Map<string, StructuredAgentSessionHostSession>,
now: () => number,
wiring: Omit<
ConstructorParameters<typeof StructuredAgentSessionReadableRestorer>[0],
'store' | 'journalRoot' | 'supportsRecord'
'store' | 'journalRoot' | 'supportsRecord' | 'retrySettlement'
>
): {
restoreReadableSessions: (sessionIds?: readonly string[]) => Promise<void>
@@ -67,6 +71,8 @@ export function createStructuredAgentSessionHostRestore(
store: deps.store,
journalRoot: deps.journalRoot,
supportsRecord: (record) => adapterSupportsRecord(deps.adapter, record),
retrySettlement: (sessionId, params) =>
retryPendingStructuredAgentSessionSettlement({ deps, sessions, sessionId, params, now }),
...wiring
})
const gate = new StructuredAgentSessionRestartRestoreGate()
@@ -46,15 +46,27 @@ export async function retryPendingStructuredAgentSessionSettlement(input: {
hasProviderChild: false,
acquisitionGeneration: null
} as StructuredAgentSessionHostSession)
return retryLoadedStructuredAgentSessionSettlement({
deps: input.deps,
sessionId: input.sessionId,
session: retrySession,
now: input.now
})
}
export async function retryLoadedStructuredAgentSessionSettlement(input: {
deps: Pick<StructuredAgentSessionHostDeps, 'store' | 'onEventSinkError'>
sessionId: string
session: Pick<StructuredAgentSessionHostSession, 'journal' | 'fence' | 'acquisitionGeneration'>
now: () => number
}): Promise<boolean> {
const record = input.deps.store.getRecord(input.sessionId)
if (!record?.lease.settlementRetryRequired || !record.lease.settlementRetryId) {
return true
}
const retrySession = input.session
retrySession.fence = record.lease.runtimeFence
const context: StructuredAgentSessionUnexpectedExitContext = {
store: input.deps.store,
sessions: input.sessions,
flushLifecycle: async () => ({ ok: true as const }),
publishFence: () => undefined,
hasResumeCapableHolder: () => false,
serialize: async <T>(_id: string, task: () => Promise<T>) => task(),
now: input.now,
const context: Pick<StructuredAgentSessionUnexpectedExitContext, 'onBarrierError'> = {
onBarrierError: (id, error) => input.deps.onEventSinkError?.({ sessionId: id, error })
}
const ok = await retryUnexpectedExitSettlement({
@@ -65,7 +77,7 @@ export async function retryPendingStructuredAgentSessionSettlement(input: {
reason: record.lease.deathEvidence?.detail ?? 'provider exited',
cause: 'unexpected-exit',
fence: record.lease.runtimeFence,
acquisitionGeneration: current?.acquisitionGeneration ?? 'recovery'
acquisitionGeneration: retrySession.acquisitionGeneration ?? 'recovery'
},
session: retrySession,
stableSettlementId: record.lease.settlementRetryId
@@ -81,11 +93,14 @@ export async function retryPendingStructuredAgentSessionSettlement(input: {
) {
throw new Error('agent_session_checkpoint_stale')
}
// A dead-TUI retry still needs its stopped-owner stage; recovery-only stages end here.
const preserveHandoff = latest.lease.handoffStage === 'old-owner-stopped'
return {
...latest,
lease: {
...latest.lease,
handoffStage: null,
handoffStage: preserveHandoff ? latest.lease.handoffStage : null,
handoffOperationId: preserveHandoff ? latest.lease.handoffOperationId : null,
settlementRetryRequired: undefined,
settlementRetryId: undefined,
lastRenewedAt: input.now()
@@ -74,6 +74,52 @@ describe('performCancel', () => {
])
})
it('keeps the running lifecycle when cancellation cannot be confirmed', async () => {
root = await mkdtemp(join(tmpdir(), 'orca-turn-cancel-unconfirmed-'))
const journal = await journals.open({ identity: IDENTITY, journalDir: root })
await journal.appendItem(
{
provider: 'legacy',
agent: 'codex',
sessionId: 'session-1',
recordId: 'turn-lifecycle:turn-1'
},
{
kind: 'status',
text: 'Agent is working…',
turnLifecycle: { turnId: 'turn-1', state: 'running' }
},
{ fence: 1 }
)
const ctx: AgentSessionTurnContext = {
sessionId: 'session-1',
journal,
fence: 1,
adapter: {
cancelTurn: vi.fn(async () => ({ cancelled: false }))
} as unknown as StructuredAgentSessionAdapter,
persistOptions: async () => undefined,
resolvedBy: 'client-1',
publish: vi.fn(),
now: () => 1
}
const result = await performCancel(ctx, {
clientOperationId: 'cancel-unconfirmed-1',
turnId: 'turn-1'
})
expect(result).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: false } })
expect(journal.snapshot().items.map((item) => item.body)).toEqual([
{
kind: 'status',
text: 'Agent is working…',
turnLifecycle: { turnId: 'turn-1', state: 'running' }
},
{ kind: 'status', text: 'The provider had already finished this turn.' }
])
})
it('stops background tasks without interrupting the foreground turn or writing a row', async () => {
root = await mkdtemp(join(tmpdir(), 'orca-background-task-cancel-'))
const journal = await journals.open({ identity: IDENTITY, journalDir: root })
@@ -161,9 +161,9 @@ export function isStructuredAgentSessionRecoveryTicketCurrent(
}
export async function retryUnexpectedExitSettlement(input: {
context: StructuredAgentSessionUnexpectedExitContext
context: Pick<StructuredAgentSessionUnexpectedExitContext, 'onBarrierError'>
event: UnexpectedExitLifecycleEvent
session: StructuredAgentSessionHostSession
session: Pick<StructuredAgentSessionHostSession, 'journal' | 'fence'>
stableSettlementId: string
}): Promise<boolean> {
try {
@@ -189,7 +189,7 @@ export async function retryUnexpectedExitSettlement(input: {
function unexpectedExitFallbackMutations(
event: UnexpectedExitLifecycleEvent,
session: StructuredAgentSessionHostSession,
session: Pick<StructuredAgentSessionHostSession, 'journal'>,
stableSettlementId: string
): JournalLifecycleMutationInput[] {
const mutations: JournalLifecycleMutationInput[] = []
@@ -15,6 +15,7 @@ import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
import { evaluateAgentSessionAcquisition } from '../../../shared/agent-session-lease-adjudication'
import { activeStructuredAgentSessionTurnId } from '../../../shared/structured-agent-session-projection'
import type {
AgentSessionClaimStatus,
AgentSessionHandoffStage,
@@ -25,6 +26,9 @@ import type {
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import { AGENT_SESSION_STORE_FILE_NAME } from '../../runtime/agent-session-record-store-file'
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
import { openAgentSessionJournal } from '../agent-session-journal/journal-store-factory'
import { journalDirectoryFor } from '../agent-session-journal/journal-paths'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import { StructuredAgentSessionHost } from './structured-agent-session-host'
import type { StructuredAgentSessionHostDeps } from './structured-agent-session-host-types'
import {
@@ -126,7 +130,8 @@ function openHost(overrides: Partial<StructuredAgentSessionHostDeps> = {}): void
dispatch: vi.fn(),
cancelTurn: vi.fn(),
answerPrompt: vi.fn(),
setOption: vi.fn()
setOption: vi.fn(),
supportsCreate: () => true
} as unknown as StructuredAgentSessionAdapter,
journalRoot: root,
claimKeyId: 'key-1',
@@ -174,7 +179,149 @@ function isAcquirable(lease: NonNullable<ReturnType<typeof store.getRecord>>['le
)
}
async function seedRunningTurn(provider: 'codex' | 'claude' = 'codex'): Promise<void> {
const journal = await openAgentSessionJournal({
identity: {
sessionId: SESSION,
workspaceId: LOCATION.workspaceId,
hostId: LOCATION.executionHostId,
agent: provider,
providerHandle:
provider === 'codex'
? { kind: 'codex', threadId: THREAD }
: { kind: 'claude', sessionId: 'provider-session-alpha-1', leafUuid: null }
},
journalDir: journalDirectoryFor(root, { workspaceId: LOCATION.workspaceId, sessionId: SESSION })
})
await journal.appendItem(
provider === 'codex'
? { provider: 'codex', threadId: THREAD, turnId: 'turn-1', ordinal: 0 }
: { provider: 'claude', sessionId: 'provider-session-alpha-1', uuid: 'uuid-running' },
{
kind: 'status',
text: 'Agent is working...',
turnLifecycle: { turnId: 'turn-1', state: 'running' }
},
{ fence: 13 }
)
await journal.close()
}
function restoredJournal(): AgentSessionJournal {
const restored = (
host as unknown as { sessions: Map<string, { journal: AgentSessionJournal }> }
).sessions.get(SESSION)
if (!restored) {
throw new Error('expected a restored session journal')
}
return restored.journal
}
describe('already-wedged profiles become usable on load', () => {
it.each(['codex', 'claude'] as const)(
'settles a wedged %s journal on boot without opening a provider child',
async (provider) => {
const record = wedgedRecord({
claimStatus: 'live',
handoffStage: null,
ownerProcess: DEAD_OWNER
})
const providerRecord: AgentSessionRecord =
provider === 'codex'
? record
: {
...record,
provider: 'claude',
accountHome: { variable: 'CLAUDE_CONFIG_DIR', path: '/home/dev/.claude' },
lease: { ...record.lease, provenHandleLinkId: 'claude-13-link' },
providerHandleChain: [
{
linkId: 'claude-13-link',
handle: {
provider: 'claude',
sessionId: 'provider-session-alpha-1',
leafUuid: null
},
origin: 'created',
mintedAtFence: 13,
observedAt: NOW - 10_000
}
]
}
await seedStore(providerRecord)
await seedRunningTurn(provider)
openHost()
await host.restoreReadableSessions()
expect(host.hasSession(SESSION)).toBe(true)
const firstCursor = restoredJournal().cursor()
expect(activeStructuredAgentSessionTurnId(restoredJournal().snapshot().items)).toBe(null)
expect(store.getRecord(SESSION)?.lease).toMatchObject({
claimStatus: 'released',
handoffStage: null,
settlementRetryRequired: undefined,
settlementRetryId: undefined
})
expect(acquire).not.toHaveBeenCalled()
await host.flushAllStreamedEvents()
store = await AgentSessionRecordStore.open({
directory: join(root, 'store'),
hostId: 'local'
})
openHost()
await host.restoreReadableSessions()
expect(restoredJournal().cursor()).toEqual(firstCursor)
expect(activeStructuredAgentSessionTurnId(restoredJournal().snapshot().items)).toBe(null)
}
)
it('settles restart eviction through attach when a hold arrives before the boot sweep', async () => {
await seedStore(
wedgedRecord({ claimStatus: 'live', handoffStage: null, ownerProcess: DEAD_OWNER })
)
await seedRunningTurn()
openHost()
await host.hold(SESSION, 'desktop-chat:restart')
expect(acquire).toHaveBeenCalledOnce()
expect(activeStructuredAgentSessionTurnId(restoredJournal().snapshot().items)).toBe(null)
expect(store.getRecord(SESSION)?.lease).toMatchObject({
claimStatus: 'live',
handoffStage: null,
settlementRetryRequired: undefined,
settlementRetryId: undefined
})
})
it('settles an observed-exit latch through attach before the boot sweep', async () => {
const record = wedgedRecord({ claimStatus: 'released', handoffStage: 'recovering' })
record.lease.settlementRetryRequired = true
record.lease.settlementRetryId = `provider-exit:${SESSION}:12:generation-1`
record.lease.deathEvidence = {
kind: 'exit-observed',
detail: 'provider exited: transport closed',
observedAt: NOW - 1_000
}
await seedStore(record)
await seedRunningTurn()
openHost()
expect(await host.attach(CALLER, hostTestAttachParams(13))).toMatchObject({ ok: true })
expect(acquire).toHaveBeenCalledOnce()
expect(activeStructuredAgentSessionTurnId(restoredJournal().snapshot().items)).toBe(null)
expect(store.getRecord(SESSION)?.lease).toMatchObject({
claimStatus: 'live',
handoffStage: null,
settlementRetryRequired: undefined,
settlementRetryId: undefined
})
})
it('re-adjudicates a conflicted manual-recovery record whose owner is provably gone', async () => {
// A crash can leave a conflicted current-schema row in manual recovery; positive death proof
// must make it acquirable again without discarding the provider handle.
@@ -0,0 +1,115 @@
import { describe, expect, it } from 'vitest'
import {
agentSessionLeaseFixture,
agentSessionRecordFixture
} from '../../shared/agent-session-record.test-fixture'
import { evictAgentSessionOwner } from './agent-session-lease-transitions'
import { applyAgentSessionRestartAdjudication } from './agent-session-restart-lease-transitions'
const NOW = 1_800_000_000_000
describe('proven-dead agent session eviction settlement', () => {
it('latches restart eviction with a stable id while keeping the lease resumable', () => {
const record = agentSessionRecordFixture(
agentSessionLeaseFixture({ runtimeKind: 'native', unreconciled: true })
)
const evicted = applyAgentSessionRestartAdjudication({
record,
probe: { outcome: 'pid-absent' },
now: NOW
})
expect(evicted.lease).toMatchObject({
claimStatus: 'released',
runtimeFence: 8,
handoffStage: null,
settlementRetryRequired: true,
settlementRetryId: 'restart-eviction:session-alpha-1:8',
deathEvidence: { kind: 'pid-absent', detail: 'recorded pid absent on host' }
})
})
it('latches recovery eviction from the same evicted disposition', () => {
const record = agentSessionRecordFixture(
agentSessionLeaseFixture({ runtimeKind: 'native', handoffStage: 'recovering' })
)
const evicted = evictAgentSessionOwner({
record,
expectedFence: 7,
probe: { outcome: 'identity-mismatch', field: 'process-start-time' },
now: NOW,
journalSettlement: 'required'
})
expect(evicted.lease).toMatchObject({
claimStatus: 'released',
runtimeFence: 8,
handoffStage: null,
settlementRetryRequired: true,
settlementRetryId: 'restart-eviction:session-alpha-1:8',
deathEvidence: { kind: 'identity-mismatch', detail: 'mismatched process-start-time' }
})
})
it('never latches an indeterminate owner', () => {
const restartRecord = agentSessionRecordFixture(
agentSessionLeaseFixture({ runtimeKind: 'native', unreconciled: true })
)
const recovered = applyAgentSessionRestartAdjudication({
record: restartRecord,
probe: { outcome: 'indeterminate', reason: 'remote host unavailable' },
now: NOW
})
const recoveryRecord = agentSessionRecordFixture(
agentSessionLeaseFixture({ runtimeKind: 'native', handoffStage: 'recovering' })
)
expect(recovered.lease).toMatchObject({
handoffStage: 'recovering',
ownerProcess: { pid: 4242 }
})
expect(recovered.lease).not.toHaveProperty('settlementRetryRequired')
expect(recovered.lease).not.toHaveProperty('settlementRetryId')
expect(() =>
evictAgentSessionOwner({
record: recoveryRecord,
expectedFence: 7,
probe: { outcome: 'indeterminate', reason: 'remote host unavailable' },
now: NOW,
journalSettlement: 'required'
})
).toThrow('agent_session_ownership_unknown')
expect(recoveryRecord.lease).not.toHaveProperty('settlementRetryRequired')
expect(recoveryRecord.lease).not.toHaveProperty('settlementRetryId')
})
it('preserves a null handoff stage when the latch survives another restart', () => {
const record = agentSessionRecordFixture(
agentSessionLeaseFixture({
runtimeKind: 'native',
ownerProcess: null,
reservedSpawnToken: null,
claimStatus: 'released',
handoffStage: null,
settlementRetryRequired: true,
settlementRetryId: 'restart-eviction:session-alpha-1:8',
unreconciled: true
})
)
const restored = applyAgentSessionRestartAdjudication({
record,
probe: { outcome: 'indeterminate', reason: 'remote host unavailable' },
now: NOW
})
expect(restored.lease).toMatchObject({
handoffStage: null,
settlementRetryRequired: true,
settlementRetryId: 'restart-eviction:session-alpha-1:8',
unreconciled: false
})
})
})
@@ -28,7 +28,7 @@ export function recoverDeadTuiOwnerForHandoff(args: {
) {
throw new Error('agent_session_ownership_unknown')
}
const evicted = evictAgentSessionOwner(args)
const evicted = evictAgentSessionOwner({ ...args, journalSettlement: 'required' })
return withLease(evicted, {
...evicted.lease,
handoffStage: 'old-owner-stopped',
@@ -8,6 +8,7 @@
import {
adjudicateAgentSessionRestart,
agentSessionRestartEvictionSettlementId,
evaluateAgentSessionAcquisition,
type AgentSessionOwnerProbe
} from '../../shared/agent-session-lease-adjudication'
@@ -207,6 +208,7 @@ export function evictAgentSessionOwner(args: {
expectedFence: number
probe: AgentSessionOwnerProbe
now: number
journalSettlement: 'required' | 'not-required'
}): AgentSessionRecord {
const { record } = args
assertFence(record.lease, args.expectedFence)
@@ -232,6 +234,7 @@ export function evictAgentSessionOwner(args: {
if (adjudication.disposition !== 'evicted') {
throw new Error('agent_session_ownership_unknown')
}
const settlementRequired = args.journalSettlement === 'required'
return withLease(record, {
...record.lease,
runtimeFence: adjudication.nextFence,
@@ -242,7 +245,11 @@ export function evictAgentSessionOwner(args: {
claimStatus: 'released',
lastRenewedAt: args.now,
handoffOperationId: null,
deathEvidence: adjudication.evidence
deathEvidence: adjudication.evidence,
settlementRetryRequired: settlementRequired ? true : undefined,
settlementRetryId: settlementRequired
? agentSessionRestartEvictionSettlementId(record.lease, adjudication)
: undefined
})
}
@@ -254,7 +254,9 @@ export class AgentSessionRecordStore {
probe: AgentSessionOwnerProbe
now: number
}): Promise<AgentSessionRecord> {
return this.mutate(args.sessionId, (record) => evictAgentSessionOwner({ ...args, record }))
return this.mutate(args.sessionId, (record) =>
evictAgentSessionOwner({ ...args, record, journalSettlement: 'required' })
)
}
async transitionHandoff(
@@ -17,6 +17,9 @@ export function adjudicateRestartedAgentSessionHandoff(
if (adjudication.disposition === 'readopt') {
return updateLease(record, { ...record.lease, unreconciled: false, lastRenewedAt: now })
}
if (adjudication.disposition === 'settlement-pending') {
return updateLease(record, { ...record.lease, unreconciled: false, lastRenewedAt: now })
}
if (adjudication.disposition === 'free') {
return updateLease(record, {
...record.lease,
@@ -8,6 +8,7 @@
import {
adjudicateAgentSessionRestart,
agentSessionRestartEvictionSettlementId,
type AgentSessionOwnerProbe
} from '../../shared/agent-session-lease-adjudication'
import type {
@@ -50,6 +51,9 @@ export function applyAgentSessionRestartAdjudication(args: {
// Why: re-adoption is not a new generation, so the fence does not move.
return withLease(record, { ...record.lease, unreconciled: false, lastRenewedAt: args.now })
}
if (adjudication.disposition === 'settlement-pending') {
return withLease(record, { ...record.lease, unreconciled: false, lastRenewedAt: args.now })
}
if (adjudication.disposition === 'free') {
// Why: an already-free lease that reloads into `recovering` is unopenable forever; clearing
// the stage restores it without moving the fence or touching the recorded death evidence.
@@ -74,7 +78,9 @@ export function applyAgentSessionRestartAdjudication(args: {
unreconciled: false,
lastRenewedAt: args.now,
handoffOperationId: null,
deathEvidence: adjudication.evidence
deathEvidence: adjudication.evidence,
settlementRetryRequired: true,
settlementRetryId: agentSessionRestartEvictionSettlementId(record.lease, adjudication)
})
}
const stage: AgentSessionHandoffStage =
+10 -5
View File
@@ -47,12 +47,21 @@ export type AgentSessionAcquisitionDecision =
export type AgentSessionRestartAdjudication =
| { disposition: 'readopt' }
/** A journal settlement latch survives restart without changing its handoff stage. */
| { disposition: 'settlement-pending' }
/** Nothing is outstanding — no owner, no reservation. Clear any latched stage; the fence stays. */
| { disposition: 'free'; reason: string }
| { disposition: 'evicted'; nextFence: number; evidence: AgentSessionDeathEvidence }
| { disposition: 'recovering'; stage: AgentSessionHandoffStage; reason: string }
| { disposition: 'conflicted'; reason: string }
export function agentSessionRestartEvictionSettlementId(
lease: Pick<AgentSessionLease, 'sessionId'>,
eviction: Extract<AgentSessionRestartAdjudication, { disposition: 'evicted' }>
): string {
return `restart-eviction:${lease.sessionId}:${eviction.nextFence}`
}
/** Stages that can legally admit a new owner at all; the rest have an owner or no evidence. */
const STAGES_ADMITTING_NEW_OWNER: ReadonlySet<AgentSessionHandoffStage> = new Set([
'old-owner-stopped',
@@ -201,11 +210,7 @@ export function adjudicateAgentSessionRestart(args: {
if (lease.settlementRetryRequired) {
// A watched provider death can leave terminal rows unsettled. This latch is not owner
// uncertainty and must survive restart until the journal settlement is durably accepted.
return {
disposition: 'recovering',
stage: 'recovering',
reason: 'provider-exit settlement requires retry'
}
return { disposition: 'settlement-pending' }
}
if (lease.reservedSpawnToken === null && lease.claimStatus !== 'reserved') {
// Why: the spawn token is minted before the child and is the only thing a child could be