mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 16:02:03 +00:00
fix(codex): a message whose turn was stopped before Codex took it is withdrawn, not stuck (#23618)
* fix(native-chat): land a late settlement from a streamed turn's end after that turn's rows A settlement that says a streamed turn ended waits for the session's event sink to drain before writing its dispatch row. The journal reducer still refuses to overwrite an accepted or rejected send. * fix(codex): settle a send from the end of the turn Codex answered it into The turn/start answer names the turn that holds a send. The send's echo entry now keeps that binding, in memory only. If the bound turn is interrupted without echoing the send, the send is withdrawn: Codex clears a turn's pending input on interrupt, so the model never saw it. If the turn fails first, the send is rejected in Codex's words. A completed turn settles nothing, since Codex records pending input when it finishes and the echo is still due. An answer read after its turn already ended is settled by that end. The echo is still the acceptance and carries the item key. * test(codex): a send settles from the end of the turn Codex answered it into The fake Codex keeps 0.157's turn bookkeeping, and can deliver the turn/start answer after turn/started or after turn/completed. The tests cover: - a Stop before any echo withdraws the send, and the working rule reads idle; - a steered follow-up is withdrawn when the turn is interrupted; - a failed turn rejects the send in Codex's words; - a completed turn leaves the send to its echo; - a normal echo and a late echo; - two steered sends in one turn; - an answer read after the turn ended; - a timed-out answer; - child-thread turns; - how a binding dies. * refactor(native-chat): drop the stream flush before a late turn-end settlement Nothing reads the order of a dispatch row against the turn's terminal row: the reducer keeps a settled send terminal and the working state is derived from both. The echo acceptance on the same path never waited either, and the wait could drop the settlement on a failed sink barrier. * fix(codex): settle a failed turn's sends at its end, not at its error Codex keeps a failed turn's pending input and records it after the error frame, before turn/completed. Settling at the error rejected a steered follow-up the model had in fact received, so a Retry would send it twice. * docs(codex): say a completed turn echoes what it took before it ends Codex records a completed turn's pending input before `turn/completed`, so a bound send that turn never echoed is left for recovery, not awaiting an echo. The comments and one test title said the echo was still due. * test(codex): settle a send whose answer is read after its turn failed or completed A failed turn that ended before the answer rejects the send in Codex's words, once; a completed one leaves it admitted and still armed for its echo. * refactor(codex): read a failed turn's reason with the typed thread-fact reader
This commit is contained in:
@@ -1,8 +1,17 @@
|
||||
import type { ProviderDiagnostic } from '../../shared/agent-session-failure'
|
||||
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
|
||||
|
||||
/** Sends awaiting their echo. A send whose echo never arrives is
|
||||
* retired by the journal's pending-submission recovery on exit, not from here. */
|
||||
/** Sends awaiting their echo. One bound to a turn that ended without taking it settles from that
|
||||
* end; any other whose echo never arrives is retired by the journal's recovery on exit. */
|
||||
export const MAX_CODEX_PENDING_DISPATCH_ECHOES = 256
|
||||
/** Turn ends kept for an answer read after the turn it names had already ended. */
|
||||
export const MAX_CODEX_RECORDED_TURN_ENDS = 64
|
||||
|
||||
/** How a primary-thread turn ended, as Codex reported it. */
|
||||
export type CodexTurnEnd =
|
||||
| { status: 'completed' }
|
||||
| { status: 'interrupted' }
|
||||
| { status: 'failed'; detail?: ProviderDiagnostic }
|
||||
|
||||
export type CodexDispatchRequestOrigin = {
|
||||
requestedAt: number
|
||||
@@ -24,6 +33,16 @@ export type CodexDispatchEchoes = {
|
||||
settle: (clientMessageId: string) => boolean
|
||||
/** Drops an armed send whose write never reached the provider. */
|
||||
disarm: (clientMessageId: string) => void
|
||||
/**
|
||||
* Binds a send to the turn Codex answered it into. Returns that turn's end when the answer is
|
||||
* read after it; a send that end settles is no longer armed.
|
||||
*/
|
||||
bindTurn: (clientMessageId: string, threadId: string, turnId: string) => CodexTurnEnd | null
|
||||
/**
|
||||
* Records a turn's end and returns the sends bound to it that it settles: all of them unless it
|
||||
* completed, which echoes its pending input first, so one it never echoed waits for recovery.
|
||||
*/
|
||||
endTurn: (threadId: string, turnId: string, end: CodexTurnEnd) => string[]
|
||||
/** Submission origin for this exact send, retained until its echo settles it. */
|
||||
requestOrigin: (clientMessageId: string) => CodexDispatchRequestOrigin | null
|
||||
/** Highest causal sequence assigned to a dispatch in this session. */
|
||||
@@ -33,8 +52,11 @@ export type CodexDispatchEchoes = {
|
||||
}
|
||||
|
||||
export function createCodexDispatchEchoes(): CodexDispatchEchoes {
|
||||
const armed = new Map<string, { requestedAt: number | null; sequence: number }>()
|
||||
const armed = new Map<string, { requestedAt: number | null; sequence: number; turn?: string }>()
|
||||
const endedTurns = new Map<string, CodexTurnEnd>()
|
||||
let nextSequence = 0
|
||||
const turnKey = (threadId: string, turnId: string): string => JSON.stringify([threadId, turnId])
|
||||
const settles = (end: CodexTurnEnd): boolean => end.status !== 'completed'
|
||||
return {
|
||||
arm(clientMessageId, requestedAt) {
|
||||
const existing = armed.get(clientMessageId)
|
||||
@@ -52,6 +74,40 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes {
|
||||
},
|
||||
settle: (clientMessageId) => armed.delete(clientMessageId),
|
||||
disarm: (clientMessageId) => void armed.delete(clientMessageId),
|
||||
bindTurn: (clientMessageId, threadId, turnId) => {
|
||||
const entry = armed.get(clientMessageId)
|
||||
if (!entry) {
|
||||
return null
|
||||
}
|
||||
const turn = turnKey(threadId, turnId)
|
||||
entry.turn = turn
|
||||
const end = endedTurns.get(turn) ?? null
|
||||
if (end && settles(end)) {
|
||||
armed.delete(clientMessageId)
|
||||
}
|
||||
return end
|
||||
},
|
||||
endTurn: (threadId, turnId, end) => {
|
||||
const turn = turnKey(threadId, turnId)
|
||||
endedTurns.delete(turn)
|
||||
endedTurns.set(turn, end)
|
||||
for (const oldest of endedTurns.keys()) {
|
||||
if (endedTurns.size <= MAX_CODEX_RECORDED_TURN_ENDS) {
|
||||
break
|
||||
}
|
||||
endedTurns.delete(oldest)
|
||||
}
|
||||
if (!settles(end)) {
|
||||
return []
|
||||
}
|
||||
const settled = [...armed].flatMap(([clientMessageId, entry]) =>
|
||||
entry.turn === turn ? [clientMessageId] : []
|
||||
)
|
||||
for (const clientMessageId of settled) {
|
||||
armed.delete(clientMessageId)
|
||||
}
|
||||
return settled
|
||||
},
|
||||
requestOrigin: (clientMessageId) => {
|
||||
const origin = armed.get(clientMessageId)
|
||||
return origin?.requestedAt === null || origin === undefined
|
||||
@@ -61,6 +117,7 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes {
|
||||
latestSequence: () => nextSequence - 1,
|
||||
clear: () => {
|
||||
armed.clear()
|
||||
endedTurns.clear()
|
||||
nextSequence = 0
|
||||
},
|
||||
get size() {
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
import type {
|
||||
AgentJournalItemIdentity,
|
||||
AgentJournalMessageItem,
|
||||
AgentSessionJournalIdentity
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
@@ -30,11 +29,9 @@ type FakeConnection = Omit<CodexAppServerConnection, 'closed'> & {
|
||||
calls: { method: string; params?: Record<string, unknown> }[]
|
||||
}
|
||||
|
||||
export type LateSettlement = {
|
||||
sessionId: string
|
||||
clientMessageId: string
|
||||
providerIdentity: AgentJournalItemIdentity
|
||||
}
|
||||
export type LateSettlement = Parameters<
|
||||
NonNullable<CodexStructuredSessionAdapterDeps['onDispatchSettledLate']>
|
||||
>[0]
|
||||
|
||||
/** A `codex app-server` whose turn traffic the test drives by hand. */
|
||||
export function fakeCodexAppServer(routes: Record<string, CodexTestRoute> = {}): {
|
||||
|
||||
@@ -34,6 +34,7 @@ import {
|
||||
translateCodexNotification
|
||||
} from './codex-structured-provider-events'
|
||||
import { CodexStructuredTurnCancellation } from './codex-structured-turn-cancellation'
|
||||
import { settleCodexSendsInEndedTurn } from './codex-structured-turn-end-settlement'
|
||||
import { createCodexStructuredNotificationRetry } from './codex-structured-notification-retry'
|
||||
import { acquireCodexStructuredSession } from './codex-structured-session-acquire'
|
||||
import { changeCodexThreadGoal } from './codex-structured-thread-goal'
|
||||
@@ -154,6 +155,10 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
}
|
||||
if (event.type === 'notification') {
|
||||
this.compactions.codex(event.sessionId, event.method, event.params)
|
||||
// Only an admitted turn end settles; a refused one settles on the retry that lands.
|
||||
settleCodexSendsInEndedTurn(session, event.method, event.params, (settlement) =>
|
||||
this.deps.onDispatchSettledLate?.({ sessionId: event.sessionId, ...settlement })
|
||||
)
|
||||
// After the admission check, so a refused frame is observed by the strip
|
||||
// only on the retry that also reaches the journal.
|
||||
if (session.backgroundTasks.observe(event)) {
|
||||
|
||||
@@ -3,6 +3,7 @@ import type {
|
||||
AgentSessionJournalIdentity
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import type { AgentJournalDispatchRejection } from '../../shared/agent-session-failure-words'
|
||||
import { cancelProcessAcquisition } from '../../shared/child-process/cancel-process-acquisition'
|
||||
import type {
|
||||
CodexAppServerConnection,
|
||||
@@ -77,12 +78,14 @@ export type CodexStructuredSessionAdapterDeps = {
|
||||
sessionId: string,
|
||||
state: AgentSessionBackgroundTaskState | null
|
||||
) => void
|
||||
/** Identity for a send admitted earlier, once Codex echoes the user message. */
|
||||
onDispatchSettledLate?: (input: {
|
||||
sessionId: string
|
||||
clientMessageId: string
|
||||
providerIdentity: AgentJournalItemIdentity
|
||||
}) => void
|
||||
/** A send admitted earlier: its identity once Codex echoes it, or its rejection when the turn
|
||||
* Codex answered it into ended without taking it. */
|
||||
onDispatchSettledLate?: (
|
||||
input: { sessionId: string; clientMessageId: string } & (
|
||||
| { providerIdentity: AgentJournalItemIdentity }
|
||||
| ({ state: 'rejected' } & AgentJournalDispatchRejection)
|
||||
)
|
||||
) => void
|
||||
/** Codex reported its thread not running with no turn open: a send whose
|
||||
* dispatch was never answered is owed nothing after this. */
|
||||
onPrimaryThreadStoppedRunning?: (input: { sessionId: string }) => void
|
||||
|
||||
@@ -47,6 +47,11 @@ export function readCodexTurnStatus(payload: unknown): string | null {
|
||||
return nonEmptyString(record(root.turn)?.status) ?? nonEmptyString(root.status)
|
||||
}
|
||||
|
||||
/** A failed `turn/completed` carries Codex's reason as `turn.error.message`. */
|
||||
export function readCodexTurnErrorMessage(payload: unknown): string | null {
|
||||
return nonEmptyString(record(record(record(payload)?.turn)?.error)?.message)
|
||||
}
|
||||
|
||||
/** Codex's own turn duration, already in milliseconds; absent or malformed reads as null. */
|
||||
export function readCodexTurnDurationMs(payload: unknown): number | null {
|
||||
const root = record(payload)
|
||||
|
||||
@@ -0,0 +1,322 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentJournalItemBody } from '../../shared/agent-session-journal-types'
|
||||
import { classifyDispatchRejection } from '../../shared/structured-agent-session-dispatch-rejection'
|
||||
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
|
||||
import {
|
||||
acquiredCodexAdapter,
|
||||
CODEX_TEST_THREAD_ID,
|
||||
CODEX_TEST_USER_MESSAGE,
|
||||
fakeCodexAppServer,
|
||||
type LateSettlement
|
||||
} from './codex-structured-dispatch-test-support'
|
||||
import { codexTurnLifecycleFake } from './codex-turn-lifecycle-fake'
|
||||
|
||||
async function turnEndRig() {
|
||||
const codex = fakeCodexAppServer()
|
||||
const turns = codexTurnLifecycleFake(CODEX_TEST_THREAD_ID, () => {
|
||||
const handlers = codex.connections.at(-1)?.handlers
|
||||
return (method, params) => handlers?.onNotification?.(method, params)
|
||||
})
|
||||
Object.assign(codex.routes, turns.routes)
|
||||
const settlements: LateSettlement[] = []
|
||||
const bodies: AgentJournalItemBody[] = []
|
||||
const adapter = await acquiredCodexAdapter({
|
||||
codex,
|
||||
settlements,
|
||||
sink: {
|
||||
appendItem: (_identity, body) => bodies.push(body),
|
||||
appendTombstone: () => {},
|
||||
publish: () => {}
|
||||
}
|
||||
})
|
||||
const send = (clientMessageId: string) =>
|
||||
adapter.dispatch({
|
||||
sessionId: 'session-1',
|
||||
clientMessageId,
|
||||
body: CODEX_TEST_USER_MESSAGE,
|
||||
fence: 7
|
||||
})
|
||||
const settledIds = () => settlements.map(({ clientMessageId }) => clientMessageId)
|
||||
const categoryOf = (settlement: LateSettlement | undefined) =>
|
||||
settlement && 'state' in settlement ? classifyDispatchRejection(settlement).category : null
|
||||
return { codex, turns, adapter, send, settlements, settledIds, categoryOf, bodies }
|
||||
}
|
||||
|
||||
describe('a Codex send its turn ended without echoing', () => {
|
||||
it('is withdrawn when the turn is interrupted, once', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await expect(rig.send('client-1')).resolves.toEqual({ state: 'admitted' })
|
||||
rig.turns.start()
|
||||
|
||||
rig.turns.end('interrupted')
|
||||
|
||||
expect(rig.settlements).toEqual([
|
||||
expect.objectContaining({
|
||||
sessionId: 'session-1',
|
||||
clientMessageId: 'client-1',
|
||||
state: 'rejected'
|
||||
})
|
||||
])
|
||||
expect(rig.categoryOf(rig.settlements[0])).toBe('withdrawn')
|
||||
})
|
||||
|
||||
it('is rejected in Codex words when its failed turn ends without echoing it', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
|
||||
rig.codex.connections[0]!.handlers.onNotification?.('error', {
|
||||
threadId: CODEX_TEST_THREAD_ID,
|
||||
turnId: 'turn-1',
|
||||
willRetry: false,
|
||||
error: { message: 'usage limit reached' }
|
||||
})
|
||||
rig.turns.end('failed', 'usage limit reached')
|
||||
|
||||
expect(rig.settlements).toEqual([
|
||||
expect.objectContaining({
|
||||
clientMessageId: 'client-1',
|
||||
state: 'rejected',
|
||||
rejection: {
|
||||
kind: 'providerRejected',
|
||||
detail: { text: 'usage limit reached', audience: 'person' }
|
||||
}
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
it('is accepted when Codex records it after the error that fails its turn', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
rig.turns.echo('client-1')
|
||||
await rig.send('client-2')
|
||||
|
||||
// A failed turn keeps its steered input: Codex records it after the error, before the end.
|
||||
rig.codex.connections[0]!.handlers.onNotification?.('error', {
|
||||
threadId: CODEX_TEST_THREAD_ID,
|
||||
turnId: 'turn-1',
|
||||
willRetry: false,
|
||||
error: { message: 'usage limit reached' }
|
||||
})
|
||||
rig.turns.echo('client-2')
|
||||
rig.turns.end('failed', 'usage limit reached')
|
||||
|
||||
expect(rig.settlements).toEqual([
|
||||
expect.objectContaining({ clientMessageId: 'client-1', providerIdentity: expect.anything() }),
|
||||
expect.objectContaining({
|
||||
clientMessageId: 'client-2',
|
||||
providerIdentity: expect.objectContaining({ turnId: 'turn-1' })
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
it('stays pending, still armed, when its turn completes without echoing it', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
|
||||
rig.turns.end('completed')
|
||||
expect(rig.settlements).toEqual([])
|
||||
|
||||
// Codex echoes before a completed end; this late one only proves the send is still armed.
|
||||
rig.turns.echo('client-1')
|
||||
expect(rig.settlements).toEqual([
|
||||
expect.objectContaining({
|
||||
clientMessageId: 'client-1',
|
||||
providerIdentity: expect.objectContaining({ turnId: 'turn-1' })
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
it('accepts a send its turn echoed, with its key, exactly once', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
rig.turns.echo('client-1')
|
||||
|
||||
rig.turns.end('interrupted')
|
||||
|
||||
expect(rig.settlements).toEqual([
|
||||
{
|
||||
sessionId: 'session-1',
|
||||
clientMessageId: 'client-1',
|
||||
providerIdentity: {
|
||||
provider: 'codex',
|
||||
threadId: CODEX_TEST_THREAD_ID,
|
||||
turnId: 'turn-1',
|
||||
ordinal: 0
|
||||
}
|
||||
}
|
||||
])
|
||||
})
|
||||
|
||||
it('ignores an echo that arrives after the withdrawal', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
rig.turns.end('interrupted')
|
||||
|
||||
rig.turns.echo('client-1')
|
||||
|
||||
expect(rig.settledIds()).toEqual(['client-1'])
|
||||
expect(rig.categoryOf(rig.settlements[0])).toBe('withdrawn')
|
||||
})
|
||||
|
||||
it('settles each of two sends steered into one interrupted turn once', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
await rig.send('client-2')
|
||||
expect(rig.turns.turnId).toBe('turn-1')
|
||||
|
||||
rig.turns.end('interrupted')
|
||||
|
||||
expect(rig.settledIds().sort()).toEqual(['client-1', 'client-2'])
|
||||
})
|
||||
|
||||
it('settles a send whose answer arrived after turn/started when the turn is interrupted', async () => {
|
||||
const rig = await turnEndRig()
|
||||
const release = rig.turns.holdNextAnswer()
|
||||
const sending = rig.send('client-1')
|
||||
await vi.waitFor(() => expect(rig.turns.turnId).toBe('turn-1'))
|
||||
rig.turns.start()
|
||||
release()
|
||||
await expect(sending).resolves.toEqual({ state: 'admitted' })
|
||||
|
||||
rig.turns.end('interrupted')
|
||||
|
||||
expect(rig.settledIds()).toEqual(['client-1'])
|
||||
})
|
||||
|
||||
it('settles a send by the recorded end of a turn that finished before its answer arrived', async () => {
|
||||
const rig = await turnEndRig()
|
||||
const release = rig.turns.holdNextAnswer()
|
||||
const sending = rig.send('client-1')
|
||||
await vi.waitFor(() => expect(rig.turns.turnId).toBe('turn-1'))
|
||||
rig.turns.start()
|
||||
rig.turns.end('interrupted')
|
||||
const rowsAtEnd = rig.bodies.filter((body) => body.kind === 'turn').length
|
||||
release()
|
||||
|
||||
const outcome = await sending
|
||||
expect(outcome.state === 'rejected' && classifyDispatchRejection(outcome).category).toBe(
|
||||
'withdrawn'
|
||||
)
|
||||
expect(rig.settlements).toEqual([])
|
||||
// The answer opens nothing: the turn keeps its terminal row.
|
||||
expect(rig.bodies.filter((body) => body.kind === 'turn')).toHaveLength(rowsAtEnd)
|
||||
expect(rig.bodies.findLast((body) => body.kind === 'turn')).toMatchObject({
|
||||
state: 'interrupted'
|
||||
})
|
||||
})
|
||||
|
||||
it('rejects a send in Codex words by the recorded end of a failed turn that finished before its answer', async () => {
|
||||
const rig = await turnEndRig()
|
||||
const release = rig.turns.holdNextAnswer()
|
||||
const sending = rig.send('client-1')
|
||||
await vi.waitFor(() => expect(rig.turns.turnId).toBe('turn-1'))
|
||||
rig.turns.start()
|
||||
rig.turns.end('failed', 'usage limit reached')
|
||||
release()
|
||||
|
||||
await expect(sending).resolves.toMatchObject({
|
||||
state: 'rejected',
|
||||
rejection: {
|
||||
kind: 'providerRejected',
|
||||
detail: { text: 'usage limit reached', audience: 'person' }
|
||||
}
|
||||
})
|
||||
// Settled once, by the answer: a late echo is no longer owed anything.
|
||||
rig.turns.echo('client-1')
|
||||
expect(rig.settlements).toEqual([])
|
||||
})
|
||||
|
||||
it('leaves a send pending, still armed, when its answer is read after its turn completed', async () => {
|
||||
const rig = await turnEndRig()
|
||||
const release = rig.turns.holdNextAnswer()
|
||||
const sending = rig.send('client-1')
|
||||
await vi.waitFor(() => expect(rig.turns.turnId).toBe('turn-1'))
|
||||
rig.turns.start()
|
||||
rig.turns.end('completed')
|
||||
release()
|
||||
|
||||
await expect(sending).resolves.toEqual({ state: 'admitted' })
|
||||
expect(rig.settlements).toEqual([])
|
||||
rig.turns.echo('client-1')
|
||||
expect(rig.settledIds()).toEqual(['client-1'])
|
||||
expect(rig.settlements[0]).toHaveProperty('providerIdentity')
|
||||
})
|
||||
|
||||
it('leaves a send whose answer timed out to its echo', async () => {
|
||||
const rig = await turnEndRig()
|
||||
rig.codex.routes['turn/start'] = () => {
|
||||
throw new Error('codex app-server turn/start exceeded 30000ms')
|
||||
}
|
||||
await expect(rig.send('client-1')).rejects.toThrow('exceeded')
|
||||
|
||||
rig.codex.connections[0]!.handlers.onNotification?.('turn/started', {
|
||||
threadId: CODEX_TEST_THREAD_ID,
|
||||
turn: { id: 'turn-9' }
|
||||
})
|
||||
rig.codex.connections[0]!.handlers.onNotification?.('item/completed', {
|
||||
threadId: CODEX_TEST_THREAD_ID,
|
||||
turn: { id: 'turn-9' },
|
||||
item: { type: 'userMessage', id: 'item-9', clientId: 'client-1', content: [] }
|
||||
})
|
||||
|
||||
expect(rig.settlements).toEqual([
|
||||
expect.objectContaining({
|
||||
clientMessageId: 'client-1',
|
||||
providerIdentity: expect.objectContaining({ turnId: 'turn-9' })
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
it('ignores a turn end on a child thread', async () => {
|
||||
const rig = await turnEndRig()
|
||||
await rig.send('client-1')
|
||||
rig.turns.start()
|
||||
|
||||
rig.codex.connections[0]!.handlers.onNotification?.('turn/completed', {
|
||||
threadId: 'thread-child',
|
||||
turn: { id: 'turn-1', status: 'interrupted' }
|
||||
})
|
||||
|
||||
expect(rig.settlements).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
describe('a send bound to a turn', () => {
|
||||
it('dies with the settlement its turn end makes', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
|
||||
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual(['client-1'])
|
||||
expect(echoes.size).toBe(0)
|
||||
expect(echoes.settle('client-1')).toBe(false)
|
||||
})
|
||||
|
||||
it('dies with its child, which forgets recorded turn ends too', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.endTurn('thread-2', 'turn-2', { status: 'interrupted' })
|
||||
|
||||
echoes.clear()
|
||||
|
||||
expect(echoes.size).toBe(0)
|
||||
echoes.arm('client-2')
|
||||
expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2')).toBeNull()
|
||||
})
|
||||
|
||||
it('is matched by thread as well as turn id', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
|
||||
expect(echoes.endTurn('thread-2', 'turn-1', { status: 'interrupted' })).toEqual([])
|
||||
expect(echoes.size).toBe(1)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,96 @@
|
||||
// A send Codex answered into a turn that then ended without echoing it. Codex clears
|
||||
// a turn's pending input when it is interrupted, so that send never reached the model
|
||||
// and is withdrawn, as a Stop's host-side withdrawal is. Any other end records pending
|
||||
// input before `turn/completed`, a failed turn after its `error` frame, so only that
|
||||
// frame settles: a failed turn that never echoed the send refused it, in Codex's words,
|
||||
// and a completed one leaves it pending for the journal's recovery on exit.
|
||||
|
||||
import {
|
||||
agentSessionFailureFact,
|
||||
providerDiagnostic,
|
||||
type ProviderDiagnostic,
|
||||
type SubmissionRejectionFact
|
||||
} from '../../shared/agent-session-failure'
|
||||
import {
|
||||
agentSessionFailureWords,
|
||||
type AgentJournalDispatchRejection
|
||||
} from '../../shared/agent-session-failure-words'
|
||||
import type { CodexTurnEnd } from './codex-structured-dispatch-echo'
|
||||
import type { CodexSession } from './codex-structured-session-state'
|
||||
import {
|
||||
readCodexThreadId,
|
||||
readCodexTurnErrorMessage,
|
||||
readCodexTurnId,
|
||||
readCodexTurnStatus
|
||||
} from './codex-structured-thread-facts'
|
||||
import { TUI_AGENT_DISPLAY_NAMES } from '../../shared/tui-agent-display-names'
|
||||
|
||||
/** A message Codex rejected, in the words that name Codex and its legacy markers. */
|
||||
export function codexDispatchRejection(
|
||||
failure: SubmissionRejectionFact
|
||||
): AgentJournalDispatchRejection {
|
||||
return agentSessionFailureWords(failure, {
|
||||
surface: 'rejection',
|
||||
agentName: TUI_AGENT_DISPLAY_NAMES.codex,
|
||||
provider: 'codex'
|
||||
})
|
||||
}
|
||||
|
||||
export type CodexTurnEndSettlement = {
|
||||
clientMessageId: string
|
||||
state: 'rejected'
|
||||
} & AgentJournalDispatchRejection
|
||||
|
||||
function errorDetail(params: unknown): ProviderDiagnostic | undefined {
|
||||
const message = readCodexTurnErrorMessage(params)
|
||||
return message ? providerDiagnostic(message, 'person') : undefined
|
||||
}
|
||||
|
||||
/** The end a primary-thread notification reports for its turn, or null for any other frame. */
|
||||
export function readCodexTurnEnd(method: string, params: unknown): CodexTurnEnd | null {
|
||||
if (method !== 'turn/completed') {
|
||||
return null
|
||||
}
|
||||
const status = readCodexTurnStatus(params)
|
||||
if (status === 'interrupted') {
|
||||
return { status: 'interrupted' }
|
||||
}
|
||||
if (status === 'failed') {
|
||||
const detail = errorDetail(params)
|
||||
return { status: 'failed', ...(detail ? { detail } : {}) }
|
||||
}
|
||||
return { status: 'completed' }
|
||||
}
|
||||
|
||||
/** How an ended turn settles a send it never echoed; null leaves the send to its echo. */
|
||||
export function codexTurnEndRejection(end: CodexTurnEnd): AgentJournalDispatchRejection | null {
|
||||
if (end.status === 'interrupted') {
|
||||
return agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' })
|
||||
}
|
||||
if (end.status === 'failed') {
|
||||
return codexDispatchRejection(
|
||||
agentSessionFailureFact('providerRejected', end.detail ? { detail: end.detail } : {})
|
||||
)
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
/** Settles the sends bound to the turn this admitted notification ended. */
|
||||
export function settleCodexSendsInEndedTurn(
|
||||
session: Pick<CodexSession, 'threadId' | 'dispatchEchoes'>,
|
||||
method: string,
|
||||
params: unknown,
|
||||
settle: (settlement: CodexTurnEndSettlement) => void
|
||||
): void {
|
||||
const turnId = readCodexTurnId(params)
|
||||
const end = readCodexTurnEnd(method, params)
|
||||
if (!turnId || !end || (readCodexThreadId(params) ?? session.threadId) !== session.threadId) {
|
||||
return
|
||||
}
|
||||
const rejection = codexTurnEndRejection(end)
|
||||
for (const clientMessageId of session.dispatchEchoes.endTurn(session.threadId, turnId, end)) {
|
||||
if (rejection) {
|
||||
settle({ clientMessageId, state: 'rejected', ...rejection })
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,10 +1,4 @@
|
||||
import { agentSessionFailureFact, providerDiagnosticOf } from '../../shared/agent-session-failure'
|
||||
import {
|
||||
agentSessionFailureWords,
|
||||
type AgentJournalDispatchRejection
|
||||
} from '../../shared/agent-session-failure-words'
|
||||
import type { SubmissionRejectionFact } from '../../shared/agent-session-failure'
|
||||
import { TUI_AGENT_DISPLAY_NAMES } from '../../shared/tui-agent-display-names'
|
||||
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
|
||||
import type { NativeChatBlock } from '../../shared/native-chat-types'
|
||||
import type { AgentSessionDispatchOutcome } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
|
||||
@@ -14,6 +8,11 @@ import {
|
||||
} from './codex-app-server-connection'
|
||||
import { isCodexAppServerUnsupportedError } from './codex-app-server-session'
|
||||
import type { CodexDispatchEchoes } from './codex-structured-dispatch-echo'
|
||||
import { readCodexTurnId } from './codex-structured-thread-facts'
|
||||
import {
|
||||
codexDispatchRejection,
|
||||
codexTurnEndRejection
|
||||
} from './codex-structured-turn-end-settlement'
|
||||
import { decodeStructuredAgentSessionOptionValue } from '../../shared/structured-agent-session-option-codec'
|
||||
|
||||
// Writing a Codex turn and learning which message landed where, which are not
|
||||
@@ -21,7 +20,8 @@ import { decodeStructuredAgentSessionOptionValue } from '../../shared/structured
|
||||
// message issued while a turn is running is COALESCED into that turn: the same
|
||||
// turn id comes back, no second `turn/started` fires, and the user message is
|
||||
// echoed only when the running turn reaches it. So the response proves
|
||||
// admission and nothing about identity, which the echo settles later.
|
||||
// admission and nothing about identity, which the echo settles later. The turn
|
||||
// it names is kept with the send, so that turn's end can settle it.
|
||||
|
||||
/** Keys Codex accepts as per-turn overrides. An unlisted key would otherwise
|
||||
* become an arbitrary client-controlled `turn/start` parameter. Permission posture is owned by
|
||||
@@ -92,7 +92,8 @@ function codexTurnOptions(host: CodexTurnHost): Record<string, string> {
|
||||
|
||||
/**
|
||||
* Hands one submission to Codex. False means the bounded correlation window
|
||||
* refused it before the write; otherwise resolves when Codex has taken it.
|
||||
* refused it before the write; otherwise resolves with the turn Codex answered
|
||||
* it into, or null when the answer named none.
|
||||
*/
|
||||
export async function startCodexTurn(
|
||||
host: CodexTurnHost,
|
||||
@@ -102,13 +103,13 @@ export async function startCodexTurn(
|
||||
requestedAt?: number
|
||||
timeoutMs?: number
|
||||
}
|
||||
): Promise<boolean> {
|
||||
): Promise<{ turnId: string | null } | false> {
|
||||
// Armed before the write: the echo and `turn/started` can both land while the
|
||||
// response is in flight, and the start must snapshot this send in its frontier.
|
||||
if (!host.dispatchEchoes.arm(input.clientMessageId, input.requestedAt)) {
|
||||
return false
|
||||
}
|
||||
await host.connection.request(
|
||||
const answer = await host.connection.request(
|
||||
'turn/start',
|
||||
{
|
||||
threadId: host.threadId,
|
||||
@@ -118,16 +119,7 @@ export async function startCodexTurn(
|
||||
},
|
||||
{ timeoutMs: input.timeoutMs }
|
||||
)
|
||||
return true
|
||||
}
|
||||
|
||||
/** A message Codex rejected, in the words that name Codex and its legacy markers. */
|
||||
function codexDispatchRejection(failure: SubmissionRejectionFact): AgentJournalDispatchRejection {
|
||||
return agentSessionFailureWords(failure, {
|
||||
surface: 'rejection',
|
||||
agentName: TUI_AGENT_DISPLAY_NAMES.codex,
|
||||
provider: 'codex'
|
||||
})
|
||||
return { turnId: readCodexTurnId(answer) }
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -141,10 +133,9 @@ export async function dispatchCodexTurn(
|
||||
input: { clientMessageId: string; body: AgentJournalMessageItem; requestedAt?: number },
|
||||
timeoutMs: number | undefined
|
||||
): Promise<AgentSessionDispatchOutcome> {
|
||||
let answer: { turnId: string | null } | false
|
||||
try {
|
||||
if (!(await startCodexTurn(session, { ...input, timeoutMs }))) {
|
||||
return { state: 'rejected', ...codexDispatchRejection(agentSessionFailureFact('queueFull')) }
|
||||
}
|
||||
answer = await startCodexTurn(session, { ...input, timeoutMs })
|
||||
} catch (error) {
|
||||
if (isCodexAppServerRequestError(error) || isCodexAppServerUnsupportedError(error)) {
|
||||
// Codex answered and declined, so no echo for this write can arrive.
|
||||
@@ -161,5 +152,13 @@ export async function dispatchCodexTurn(
|
||||
// Keep the correlation armed so a later echo can prove delivery.
|
||||
throw error
|
||||
}
|
||||
return { state: 'admitted' }
|
||||
if (!answer) {
|
||||
return { state: 'rejected', ...codexDispatchRejection(agentSessionFailureFact('queueFull')) }
|
||||
}
|
||||
// An answer read after the turn it names already ended is settled by that end.
|
||||
const endedFirst = answer.turnId
|
||||
? session.dispatchEchoes.bindTurn(input.clientMessageId, session.threadId, answer.turnId)
|
||||
: null
|
||||
const rejection = endedFirst ? codexTurnEndRejection(endedFirst) : null
|
||||
return rejection ? { state: 'rejected', ...rejection } : { state: 'admitted' }
|
||||
}
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
// Codex 0.157's turn bookkeeping as `turn/start` and `turn/interrupt` see it, for
|
||||
// tests. From app-server `turn_processor.rs`: `turn/start` picks the turn before it
|
||||
// answers, a send while a turn is open is steered into it under the same id with no
|
||||
// second `turn/started`, and `turn_interrupt_inner` refuses with -32600 until the
|
||||
// turn has started. The answer can be held, so a test can deliver it after the
|
||||
// turn's own frames, as the wire allows.
|
||||
|
||||
import { CodexAppServerRequestError } from './codex-app-server-connection'
|
||||
|
||||
type Notify = (method: string, params: unknown) => void
|
||||
|
||||
export type CodexTurnLifecycleFake = {
|
||||
routes: {
|
||||
'turn/start': () => unknown
|
||||
'turn/interrupt': (params: Record<string, unknown> | undefined) => unknown
|
||||
}
|
||||
/** The next `turn/start` answer waits until the returned release runs. */
|
||||
holdNextAnswer: () => () => void
|
||||
/** Codex emits `turn/started` for the turn it picked last. */
|
||||
start: () => void
|
||||
/** Codex ends the picked or running turn on its own. */
|
||||
end: (status: 'completed' | 'interrupted' | 'failed', errorMessage?: string) => void
|
||||
/** Codex echoes a user message it recorded in the current turn. */
|
||||
echo: (clientId: string) => void
|
||||
readonly turnId: string | null
|
||||
}
|
||||
|
||||
function refusal(message: string): CodexAppServerRequestError {
|
||||
return new CodexAppServerRequestError(
|
||||
'turn/interrupt',
|
||||
-32600,
|
||||
`codex app-server turn/interrupt failed: ${message}`,
|
||||
message
|
||||
)
|
||||
}
|
||||
|
||||
export function codexTurnLifecycleFake(
|
||||
threadId: string,
|
||||
notify: () => Notify
|
||||
): CodexTurnLifecycleFake {
|
||||
let minted = 0
|
||||
let echoes = 0
|
||||
let picked: string | null = null
|
||||
let active: string | null = null
|
||||
let lastTurn: string | null = null
|
||||
let held: Promise<void> | null = null
|
||||
const finish = (turnId: string, status: string, errorMessage?: string): void => {
|
||||
picked = null
|
||||
active = null
|
||||
notify()('turn/completed', {
|
||||
threadId,
|
||||
turn: { id: turnId, status, ...(errorMessage ? { error: { message: errorMessage } } : {}) }
|
||||
})
|
||||
}
|
||||
return {
|
||||
routes: {
|
||||
'turn/start': () => {
|
||||
const turnId = active ?? picked ?? `turn-${++minted}`
|
||||
picked ??= active ? null : turnId
|
||||
lastTurn = turnId
|
||||
const answer = { turn: { id: turnId, status: 'inProgress' } }
|
||||
const wait = held
|
||||
held = null
|
||||
return wait ? wait.then(() => answer) : answer
|
||||
},
|
||||
'turn/interrupt': (params) => {
|
||||
const turnId = params?.turnId
|
||||
if (!active) {
|
||||
throw refusal('no active turn to interrupt')
|
||||
}
|
||||
if (active !== turnId) {
|
||||
throw refusal(`expected active turn id ${String(turnId)} but found ${active}`)
|
||||
}
|
||||
// Codex answers the interrupt once the turn has aborted.
|
||||
finish(active, 'interrupted')
|
||||
return {}
|
||||
}
|
||||
},
|
||||
holdNextAnswer: () => {
|
||||
let release!: () => void
|
||||
held = new Promise<void>((resolve) => {
|
||||
release = resolve
|
||||
})
|
||||
return () => release()
|
||||
},
|
||||
start: () => {
|
||||
if (!picked) {
|
||||
throw new Error('no picked turn to start')
|
||||
}
|
||||
active = picked
|
||||
notify()('turn/started', { threadId, turn: { id: active, status: 'inProgress' } })
|
||||
},
|
||||
end: (status, errorMessage) => {
|
||||
const turnId = active ?? picked
|
||||
if (!turnId) {
|
||||
throw new Error('no turn to end')
|
||||
}
|
||||
finish(turnId, status, errorMessage)
|
||||
},
|
||||
echo: (clientId) => {
|
||||
const turnId = active ?? picked ?? lastTurn
|
||||
notify()('item/completed', {
|
||||
threadId,
|
||||
turn: { id: turnId },
|
||||
item: { type: 'userMessage', id: `item-user-${++echoes}`, clientId, content: [] }
|
||||
})
|
||||
},
|
||||
get turnId() {
|
||||
return active ?? picked
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,217 @@
|
||||
// A Codex send the turn it went into never took: the Stop that interrupts that turn
|
||||
// withdraws it, and nothing reads as working after. Driven through the shipped host,
|
||||
// journal and Codex adapter; only the Codex child is fake, keeping Codex 0.157's
|
||||
// turn bookkeeping.
|
||||
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type {
|
||||
CodexAppServerConnection,
|
||||
CodexAppServerConnectionHandlers,
|
||||
openCodexAppServerConnection
|
||||
} from '../codex/codex-app-server-connection'
|
||||
import { codexTurnLifecycleFake } from '../codex/codex-turn-lifecycle-fake'
|
||||
import { computeAgentSessionPayloadFingerprint } from '../../shared/agent-session-mutation-envelope'
|
||||
import type { AgentJournalSubmission } from '../../shared/agent-session-journal-types'
|
||||
import { classifyDispatchRejection } from '../../shared/structured-agent-session-dispatch-rejection'
|
||||
import { owesStructuredAgentSessionWork } from '../../shared/structured-agent-session-owed-work'
|
||||
import {
|
||||
HOST_TEST_SESSION as SESSION,
|
||||
HOST_TEST_THREAD as THREAD,
|
||||
hostTestAttachParams,
|
||||
hostTestMessage
|
||||
} from '../native-chat/agent-session-wire/structured-agent-session-host-test-data'
|
||||
import type { StructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-host'
|
||||
import {
|
||||
ensureStructuredAgentSessionHost,
|
||||
stopStructuredAgentSessionRuntime
|
||||
} from './structured-agent-session-runtime'
|
||||
|
||||
// A Stop must never reach for real processes on this machine under a made-up pid.
|
||||
vi.mock('../codex/codex-structured-turn-processes', () => ({
|
||||
captureCodexTurnProcesses: async () => null,
|
||||
terminateCodexTurnProcesses: async () => true
|
||||
}))
|
||||
|
||||
const CALLER = { callerKey: 'codex-turn-end-test' }
|
||||
const MODEL = {
|
||||
model: 'gpt-test',
|
||||
displayName: 'GPT Test',
|
||||
hidden: false,
|
||||
supportedReasoningEfforts: [],
|
||||
defaultReasoningEffort: null,
|
||||
isDefault: true
|
||||
}
|
||||
|
||||
let root: string
|
||||
let host: StructuredAgentSessionHost
|
||||
let fence: number
|
||||
let handlers: CodexAppServerConnectionHandlers | undefined
|
||||
let answers: number
|
||||
let turns: ReturnType<typeof codexTurnLifecycleFake>
|
||||
let operations = 0
|
||||
|
||||
/** The durable ledger stamps its own clock and refuses an id far from it. */
|
||||
const operationId = (): string => `${Date.now()}-${(++operations).toString(16).padStart(32, '0')}`
|
||||
|
||||
function envelope(method: string, fields: Record<string, unknown>) {
|
||||
return {
|
||||
sessionId: SESSION,
|
||||
clientOperationId: operationId(),
|
||||
expectedRuntimeFence: fence,
|
||||
payloadFingerprint: computeAgentSessionPayloadFingerprint({
|
||||
method,
|
||||
sessionId: SESSION,
|
||||
fields
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
async function send(text: string): Promise<string> {
|
||||
const body = hostTestMessage(text)
|
||||
const sent = await host.send(CALLER, { envelope: envelope('agentSession.send', { body }), body })
|
||||
if (!sent.ok) {
|
||||
throw new Error(JSON.stringify(sent.refusal))
|
||||
}
|
||||
return sent.value.clientMessageId
|
||||
}
|
||||
|
||||
async function stop(turnId: string): Promise<void> {
|
||||
const stopped = await host.cancel(CALLER, {
|
||||
envelope: envelope('agentSession.cancel', { turnId }),
|
||||
turnId
|
||||
})
|
||||
if (!stopped.ok) {
|
||||
throw new Error(JSON.stringify(stopped.refusal))
|
||||
}
|
||||
}
|
||||
|
||||
async function settled(): Promise<{
|
||||
submissions: readonly AgentJournalSubmission[]
|
||||
owesWork: boolean
|
||||
}> {
|
||||
await host.flushStreamedEvents(SESSION)
|
||||
const snapshot = await host.journalSnapshot(SESSION)
|
||||
return {
|
||||
submissions: snapshot.submissions,
|
||||
// The shared rule every surface reads "working" from.
|
||||
owesWork: owesStructuredAgentSessionWork(snapshot.items, snapshot.submissions, fence)
|
||||
}
|
||||
}
|
||||
|
||||
function verdictOf(submissions: readonly AgentJournalSubmission[], clientMessageId: string) {
|
||||
const submission = submissions.find((entry) => entry.clientMessageId === clientMessageId)
|
||||
return submission?.dispatchState === 'rejected'
|
||||
? classifyDispatchRejection(submission).category
|
||||
: submission?.dispatchState
|
||||
}
|
||||
|
||||
beforeEach(async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-codex-turn-end-'))
|
||||
answers = 0
|
||||
turns = codexTurnLifecycleFake(
|
||||
THREAD,
|
||||
() => (method, params) => handlers?.onNotification?.(method, params)
|
||||
)
|
||||
const openConnection: typeof openCodexAppServerConnection = async (
|
||||
_launch,
|
||||
connectionHandlers = {}
|
||||
) => {
|
||||
handlers = connectionHandlers
|
||||
const connection: CodexAppServerConnection = {
|
||||
pid: 4321,
|
||||
closed: false,
|
||||
request: async (method, params) => {
|
||||
if (method === 'thread/start' || method === 'thread/resume') {
|
||||
return { thread: { id: THREAD } }
|
||||
}
|
||||
if (method === 'model/list') {
|
||||
return { data: [MODEL], nextCursor: null }
|
||||
}
|
||||
if (method === 'turn/start') {
|
||||
answers += 1
|
||||
return turns.routes['turn/start']()
|
||||
}
|
||||
if (method === 'turn/interrupt') {
|
||||
return turns.routes['turn/interrupt'](params)
|
||||
}
|
||||
return {}
|
||||
},
|
||||
notify: () => {},
|
||||
respond: () => {},
|
||||
respondWithError: () => {},
|
||||
close: async () => true
|
||||
}
|
||||
return connection
|
||||
}
|
||||
host = await ensureStructuredAgentSessionHost({
|
||||
stateDirectory: root,
|
||||
hostId: 'local',
|
||||
claimKeyId: 'key-1',
|
||||
resolveWorkspacePath: async () => root,
|
||||
resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }),
|
||||
resolveCodexCommand: () => 'codex',
|
||||
resolveEnvironment: async () => ({ PATH: process.env.PATH }),
|
||||
openCodexConnection: openConnection,
|
||||
readProcessStartTime: async () => 1_700_000_000_000
|
||||
})
|
||||
const attachParams = hostTestAttachParams(null, { providerHandle: undefined })
|
||||
attachParams.envelope.clientOperationId = operationId()
|
||||
const attached = await host.attach(CALLER, attachParams)
|
||||
if (!attached.ok) {
|
||||
throw new Error(JSON.stringify(attached.refusal))
|
||||
}
|
||||
fence = attached.value.fence
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await stopStructuredAgentSessionRuntime()
|
||||
await rm(root, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('a Codex send its turn ended without taking it', () => {
|
||||
it('is withdrawn by a Stop before any echo, and nothing reads as working after', async () => {
|
||||
const sent = await send('look around')
|
||||
await vi.waitFor(() => expect(answers).toBe(1))
|
||||
turns.start()
|
||||
|
||||
await stop('turn-1')
|
||||
|
||||
await vi.waitFor(async () =>
|
||||
expect(verdictOf((await settled()).submissions, sent)).not.toBe('pending')
|
||||
)
|
||||
const after = await settled()
|
||||
expect(verdictOf(after.submissions, sent)).toBe('withdrawn')
|
||||
expect(after.owesWork).toBe(false)
|
||||
|
||||
// Codex echoing it late changes nothing: a settled answer stands.
|
||||
await host.settleLateDispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: sent,
|
||||
providerIdentity: { provider: 'codex', threadId: THREAD, turnId: 'turn-1', ordinal: 0 }
|
||||
})
|
||||
expect(verdictOf((await settled()).submissions, sent)).toBe('withdrawn')
|
||||
})
|
||||
|
||||
it('withdraws a follow-up Codex steered into the turn a Stop ends, with no working latch', async () => {
|
||||
const opening = await send('look around')
|
||||
await vi.waitFor(() => expect(answers).toBe(1))
|
||||
turns.start()
|
||||
turns.echo(opening)
|
||||
const followUp = await send('and check the tests')
|
||||
// Steered: Codex answers with the running turn, and fires no second turn/started.
|
||||
await vi.waitFor(() => expect(answers).toBe(2))
|
||||
|
||||
await stop('turn-1')
|
||||
|
||||
await vi.waitFor(async () =>
|
||||
expect(verdictOf((await settled()).submissions, followUp)).not.toBe('pending')
|
||||
)
|
||||
const after = await settled()
|
||||
expect(verdictOf(after.submissions, opening)).toBe('accepted')
|
||||
expect(verdictOf(after.submissions, followUp)).toBe('withdrawn')
|
||||
expect(after.owesWork).toBe(false)
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user