mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
fix: settle dead TUI handoffs before reacquire
This commit is contained in:
+1
@@ -70,6 +70,7 @@ async function failingFlowRunner(
|
||||
throw new Error('unused')
|
||||
},
|
||||
importTuiHistory: async () => {},
|
||||
retryPendingSettlement: async () => true,
|
||||
publish: () => {},
|
||||
schedule: async () => {
|
||||
throw new Error('scheduling failed')
|
||||
|
||||
@@ -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,
|
||||
|
||||
+1
@@ -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),
|
||||
|
||||
+2
@@ -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
|
||||
|
||||
+24
@@ -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)
|
||||
})
|
||||
})
|
||||
|
||||
+21
-10
@@ -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
|
||||
@@ -85,7 +97,6 @@ export async function retryPendingStructuredAgentSessionSettlement(input: {
|
||||
...latest,
|
||||
lease: {
|
||||
...latest.lease,
|
||||
handoffStage: null,
|
||||
settlementRetryRequired: undefined,
|
||||
settlementRetryId: undefined,
|
||||
lastRenewedAt: input.now()
|
||||
|
||||
@@ -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[] = []
|
||||
|
||||
Reference in New Issue
Block a user