mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 08:03:12 +00:00
Preserve Codex exit receipt across close retries
This commit is contained in:
@@ -0,0 +1,97 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types'
|
||||
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
|
||||
import { CodexBackgroundTaskTracker } from './codex-background-task-tracker'
|
||||
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
|
||||
import { closeCodexPublishedSession } from './codex-structured-session-close'
|
||||
import type { CodexSession } from './codex-structured-session-state'
|
||||
|
||||
afterEach(() => vi.useRealTimers())
|
||||
|
||||
describe('requested-close durable turn timing', () => {
|
||||
it.each([true, false])(
|
||||
'keeps the first exit receipt when retry requestedClose=%s',
|
||||
async (requestedClose) => {
|
||||
vi.useFakeTimers()
|
||||
vi.setSystemTime(1_000)
|
||||
const terminalBodies: AgentJournalItemBody[] = []
|
||||
let refuseSettlement = true
|
||||
const sink: StructuredAgentSessionEventSink = {
|
||||
appendItem: () => {},
|
||||
appendTombstone: () => {},
|
||||
publish: () => {},
|
||||
tryAppendLifecycleBatch: (_id, mutations) => {
|
||||
if (refuseSettlement) {
|
||||
return { accepted: false, reason: 'backpressure' }
|
||||
}
|
||||
for (const mutation of mutations) {
|
||||
if (mutation.kind === 'item') {
|
||||
terminalBodies.push(mutation.body)
|
||||
}
|
||||
}
|
||||
return { accepted: true }
|
||||
}
|
||||
}
|
||||
const translator = createCodexJournalTranslator({
|
||||
sink,
|
||||
sessionId: 'session-1',
|
||||
primaryThreadId: () => 'thread-1',
|
||||
now: () => Date.now()
|
||||
})
|
||||
expect(
|
||||
translator.handle({
|
||||
type: 'notification',
|
||||
sessionId: 'session-1',
|
||||
threadId: 'thread-1',
|
||||
method: 'turn/started',
|
||||
params: { turn: { id: 'turn-1' } },
|
||||
observedAt: 1_000
|
||||
})
|
||||
).toEqual({ accepted: true })
|
||||
const session = {
|
||||
connection: { close: vi.fn(async () => true) },
|
||||
backgroundTasks: new CodexBackgroundTaskTracker('thread-1'),
|
||||
ended: false,
|
||||
requestedClose: false,
|
||||
fence: 7,
|
||||
acquisitionGeneration: 'generation-1',
|
||||
threadId: 'thread-1',
|
||||
prompts: { clear: vi.fn() },
|
||||
translator
|
||||
} as unknown as CodexSession
|
||||
const sessions = new Map([['session-1', session]])
|
||||
const onEvent = vi.fn()
|
||||
|
||||
vi.setSystemTime(2_000)
|
||||
await expect(closeCodexPublishedSession(sessions, 'session-1')).resolves.toBe(false)
|
||||
expect(sessions.get('session-1')).toBe(session)
|
||||
expect(session.ended).toBe(false)
|
||||
|
||||
refuseSettlement = false
|
||||
vi.setSystemTime(60_000)
|
||||
await expect(
|
||||
closeCodexPublishedSession(sessions, 'session-1', onEvent, {
|
||||
expectedAcquisitionGeneration: 'replacement-generation'
|
||||
})
|
||||
).resolves.toBe(false)
|
||||
expect(onEvent).not.toHaveBeenCalled()
|
||||
await expect(
|
||||
closeCodexPublishedSession(sessions, 'session-1', onEvent, { requestedClose })
|
||||
).resolves.toBe(true)
|
||||
expect(sessions.has('session-1')).toBe(false)
|
||||
expect(onEvent).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
cause: requestedClose ? 'requested-close' : 'unexpected-exit',
|
||||
acquisitionGeneration: 'generation-1',
|
||||
observedAt: 2_000
|
||||
})
|
||||
)
|
||||
expect(terminalBodies.find((body) => body.kind === 'turn')).toMatchObject({
|
||||
kind: 'turn',
|
||||
state: 'interrupted',
|
||||
startedAt: 1_000,
|
||||
completedAt: 2_000
|
||||
})
|
||||
}
|
||||
)
|
||||
})
|
||||
@@ -157,7 +157,12 @@ export function createCodexJournalTranslator(
|
||||
primaryThreadId: deps.primaryThreadId?.() ?? null,
|
||||
ordinals: items.ordinals,
|
||||
settledTurnLifecycle: (threadId, turnId) =>
|
||||
turnBoundaries.settled(threadId, turnId, 'interrupted', deps.now?.() ?? Date.now())
|
||||
turnBoundaries.settled(
|
||||
threadId,
|
||||
turnId,
|
||||
'interrupted',
|
||||
event.observedAt ?? deps.now?.() ?? Date.now()
|
||||
)
|
||||
})
|
||||
if (!admission.accepted) {
|
||||
return admission
|
||||
|
||||
@@ -182,7 +182,8 @@ describe('CodexStructuredSessionAdapter lifecycle', () => {
|
||||
reason: 'codex app-server connection ended',
|
||||
cause: 'unexpected-exit',
|
||||
fence: 7,
|
||||
acquisitionGeneration: 'generation-1'
|
||||
acquisitionGeneration: 'generation-1',
|
||||
observedAt: expect.any(Number)
|
||||
})
|
||||
await expect(
|
||||
adapter.dispatch({
|
||||
|
||||
@@ -24,13 +24,15 @@ export function handleCodexSessionExit(input: {
|
||||
input.prompts?.clear()
|
||||
return false
|
||||
}
|
||||
session.exitObservedAt ??= Date.now()
|
||||
const event: StructuredAgentSessionLifecycleEvent = {
|
||||
type: 'ended',
|
||||
sessionId: input.sessionId,
|
||||
reason: input.error.message,
|
||||
cause: session.requestedClose ? 'requested-close' : 'unexpected-exit',
|
||||
fence: session.fence,
|
||||
acquisitionGeneration: session.acquisitionGeneration
|
||||
acquisitionGeneration: session.acquisitionGeneration,
|
||||
observedAt: session.exitObservedAt
|
||||
} as const
|
||||
// A synchronous sink rejection (usually backpressure) is handed to host
|
||||
// recovery, which appends the bounded fallback before reacquisition.
|
||||
|
||||
@@ -45,7 +45,7 @@ export type CodexStructuredSessionEvent =
|
||||
}
|
||||
| StructuredAgentSessionLifecycleEvent
|
||||
/** Translator-only compatibility for callers that do not participate in host recovery. */
|
||||
| { type: 'ended'; sessionId: string; reason: string }
|
||||
| { type: 'ended'; sessionId: string; reason: string; observedAt?: number }
|
||||
|
||||
export type CodexStructuredSessionAdapterDeps = {
|
||||
resolveLaunch: (input: {
|
||||
@@ -74,6 +74,8 @@ export type CodexStructuredSessionAdapterDeps = {
|
||||
export type CodexSession = {
|
||||
connection: CodexAppServerConnection
|
||||
ended: boolean
|
||||
/** First observed child exit survives rejected settlement admission. */
|
||||
exitObservedAt?: number
|
||||
requestedClose: boolean
|
||||
fence: number
|
||||
acquisitionGeneration: string
|
||||
|
||||
@@ -102,6 +102,8 @@ export type StructuredAgentSessionLifecycleEvent = {
|
||||
cause: 'unexpected-exit' | 'requested-close'
|
||||
fence: number
|
||||
acquisitionGeneration: string
|
||||
/** Host receipt of the child exit, retained across settlement retries. */
|
||||
observedAt?: number
|
||||
/** Translator could not admit terminal rows; host recovery must append its bounded fallback. */
|
||||
settlementRetryRequired?: boolean
|
||||
}
|
||||
|
||||
+4
-3
@@ -60,8 +60,8 @@ function lifecycleItem(
|
||||
}
|
||||
|
||||
describe('provider-exit recovery tickets', () => {
|
||||
it('retains the exit receipt through failed settlement and a later loaded retry', async () => {
|
||||
let now = 2_000
|
||||
it.each([undefined, 2_000])('keeps exit receipt %s on retry', async (observedAt) => {
|
||||
let now = observedAt === undefined ? 2_000 : 30_000
|
||||
let record = {
|
||||
lease: {
|
||||
handoffStage: null,
|
||||
@@ -115,7 +115,8 @@ describe('provider-exit recovery tickets', () => {
|
||||
reason: 'provider exited',
|
||||
cause: 'unexpected-exit',
|
||||
fence: 7,
|
||||
acquisitionGeneration: GENERATION
|
||||
acquisitionGeneration: GENERATION,
|
||||
observedAt
|
||||
}
|
||||
)
|
||||
expect(record.lease.settlementRetryRequired).toBe(true)
|
||||
|
||||
@@ -52,7 +52,7 @@ export async function settleUnexpectedStructuredAgentSessionExit(
|
||||
}
|
||||
const unexpectedEvent = event as UnexpectedExitLifecycleEvent
|
||||
// Receipt of the exit is the one end time the host may record for a running turn.
|
||||
const observedAt = context.now()
|
||||
const observedAt = event.observedAt ?? context.now()
|
||||
return context.serialize(unexpectedEvent.sessionId, async () => {
|
||||
const session = context.sessions.get(unexpectedEvent.sessionId)
|
||||
if (
|
||||
|
||||
Reference in New Issue
Block a user