fix: preserve unconfirmed turn cancellation state

This commit is contained in:
Merge Sim
2026-09-06 14:14:55 -07:00
parent 8e7c47ffd0
commit f0809ad96c
3 changed files with 9 additions and 95 deletions
@@ -27,7 +27,7 @@ afterEach(async () => {
})
describe('performCancel', () => {
it('acknowledges a confirmed cancel without settling the foreground turn', async () => {
it('acknowledges only the request and leaves the running lifecycle row intact', async () => {
root = await mkdtemp(join(tmpdir(), 'orca-turn-cancel-'))
const journal = await journals.open({ identity: IDENTITY, journalDir: root })
const lifecycleIdentity = {
@@ -74,54 +74,6 @@ describe('performCancel', () => {
])
})
it('settles an already-finished turn and keeps repeated stop notes bounded', async () => {
root = await mkdtemp(join(tmpdir(), 'orca-turn-cancel-finished-'))
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 cancelTurn = vi.fn(async () => ({ cancelled: false }))
const ctx: AgentSessionTurnContext = {
sessionId: 'session-1',
journal,
fence: 1,
adapter: { cancelTurn } as unknown as StructuredAgentSessionAdapter,
persistOptions: async () => undefined,
resolvedBy: 'client-1',
publish: vi.fn(),
now: () => 1
}
const first = await performCancel(ctx, {
clientOperationId: 'cancel-1',
turnId: 'turn-1'
})
const itemCountAfterFirst = journal.snapshot().items.length
const second = await performCancel(ctx, {
clientOperationId: 'cancel-2',
turnId: 'turn-1'
})
expect(first).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: false } })
expect(second).toEqual({ ok: true, value: { turnId: 'turn-1', cancelled: false } })
expect(cancelTurn).toHaveBeenCalledTimes(2)
expect(journal.snapshot().items).toHaveLength(itemCountAfterFirst)
expect(journal.snapshot().items.map((item) => item.body)).toEqual([
{ 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 })
@@ -17,8 +17,6 @@ import type {
AgentSessionDispatchOutcome,
StructuredAgentSessionAdapter
} from './structured-agent-session-adapter'
import type { JournalLifecycleMutationInput } from '../agent-session-journal/journal-row-builders'
import { runningTurnLifecycleTombstoneMutations } from './structured-agent-session-unexpected-exit'
export { performSetOption } from './structured-agent-session-turns-options'
export { performPrompt } from './structured-agent-session-turns-prompt'
@@ -156,7 +154,6 @@ export async function performCancel(
}
): Promise<TurnOutcome<AgentSessionCancelResult>> {
let cancelled = false
let alreadyFinished = false
let note = 'Cancellation requested.'
try {
cancelled = input.scope
@@ -175,7 +172,6 @@ export async function performCancel(
})
).cancelled
if (!cancelled) {
alreadyFinished = true
note = 'The provider had already finished this turn.'
}
} catch (error) {
@@ -186,23 +182,7 @@ export async function performCancel(
if (input.scope) {
return { ok: true, value: { turnId: input.turnId, cancelled } }
}
if (alreadyFinished) {
const mutations: JournalLifecycleMutationInput[] = [
{
kind: 'item',
identity: { provider: 'orca', clientMessageId: input.turnId },
body: { kind: 'status', text: note }
},
...runningTurnLifecycleTombstoneMutations(ctx.journal.snapshot().items, input.turnId)
]
await ctx.journal.appendLifecycleBatch({
settlementId: `cancel-finished:${ctx.sessionId}:${input.turnId}`,
fence: ctx.fence,
mutations
})
ctx.publish()
} else {
await appendStatus(ctx, input.turnId, note)
}
// Keyed by the operation id so a replayed cancel upserts one item, not two.
await appendStatus(ctx, input.clientOperationId, note)
return { ok: true, value: { turnId: input.turnId, cancelled } }
}
@@ -193,8 +193,8 @@ function unexpectedExitFallbackMutations(
stableSettlementId: string
): JournalLifecycleMutationInput[] {
const mutations: JournalLifecycleMutationInput[] = []
const items = session.journal.snapshot().items
for (const item of items) {
const tombstones: JournalLifecycleMutationInput[] = []
for (const item of session.journal.snapshot().items) {
const identity = parseAgentJournalItemKey(item.itemId)
if (!identity) {
continue
@@ -203,37 +203,19 @@ function unexpectedExitFallbackMutations(
if (terminal) {
mutations.push({ kind: 'item', identity, body: terminal })
}
if (item.body.kind === 'status' && item.body.turnLifecycle?.state === 'running') {
tombstones.push({ kind: 'tombstone', identity })
}
}
mutations.push({
kind: 'item',
identity: { provider: 'orca', clientMessageId: stableSettlementId },
body: { kind: 'status', text: boundJournalStatusText(`Provider exited: ${event.reason}`) }
})
mutations.push(...runningTurnLifecycleTombstoneMutations(items))
mutations.push(...tombstones)
return mutations
}
export function runningTurnLifecycleTombstoneMutations(
items: readonly AgentJournalRenderItem[],
turnId?: string
): JournalLifecycleMutationInput[] {
const tombstones: JournalLifecycleMutationInput[] = []
for (const item of items) {
if (
item.body.kind !== 'status' ||
item.body.turnLifecycle?.state !== 'running' ||
(turnId !== undefined && item.body.turnLifecycle.turnId !== turnId)
) {
continue
}
const identity = parseAgentJournalItemKey(item.itemId)
if (identity) {
tombstones.push({ kind: 'tombstone', identity })
}
}
return tombstones
}
function terminalExitBody(item: AgentJournalRenderItem): AgentJournalItemBody | null {
if (item.body.kind === 'tool-call' && item.body.state === 'running') {
return { ...item.body, state: 'failed' }