fix(native-chat): a Stop Codex took settles the message whose turn never opened (follow-up to #25217) (#26105)

* fix(native-chat): a Stop the agent took settles the message whose turn never opened

When Codex takes a person's Stop on a turn it never opened, no turn-ended
event follows, so the message stayed pending: the chat read Working with
Stop shown, and the next turn's clock counted from that message.

The Stop's settle now handles that case from its own answer: when the
provider names the turn it took and that turn has no record, the sends it
was for are withdrawn by the same rule a Codex child's end after a Stop
already uses. The send's own row then says it was stopped before the agent
started, and Working, the clock and the opening-send hold follow from it
being settled.

Deletes the opening-send hold's special case for a taken Stop's note, which
this makes dead, and moves the Stop note key back next to its only users.

* test(native-chat): a Claude Stop before the echo settles the send through the CLI's end, and the next message goes out

* fix(native-chat): an interrupted Codex turn that never started records nothing, and the withdrawal reads from the handover

Codex can abort a turn before it starts: it answers the interrupt, then
sends turn/completed for a turn that never sent turn/started. The
translator wrote an empty interrupted turn for it, so the person saw that
turn beside the message's own "Stopped before the agent started" row, and
when the end was read before the interrupt's answer the Stop's settle
found a record and withdrew nothing. Codex records a turn's prompt only
once the turn starts, so an interrupted turn this child never started,
with no item or prompt read, now gets no record. Failed ends keep theirs.

The withdrawal now asks whether a turn opened since the send's handover
row, the same point the opening-send hold reads, rather than since its
acceptance: a turn record written in between is not the send's turn.

* test(native-chat): say why the item-read case pins only that the send is not taken back
This commit is contained in:
Brennan Benson
2026-10-07 01:25:50 -07:00
committed by GitHub
parent 47f4d3f527
commit e6c298e8a9
13 changed files with 418 additions and 147 deletions
@@ -174,14 +174,25 @@ export class CodexJournalTurnBoundaries {
error: readString(readRecord(readRecord(readRecord(event.params).turn).error), 'message'),
completedAt
})
const turnLifecycle = this.ownsRecord(event.threadId, turnId)
? this.settled(event.threadId, turnId, {
state: codexTurnLifecycleState(status),
outcome: codexTurnOutcome(status),
completedAt,
durationMs: readCodexTurnDurationMs(event.params)
})
: null
// Codex records a turn's prompt only once the turn starts, so an interrupted turn this child
// never started, with no item or prompt read, ran nothing and gets no record: the Stop that
// took it settles its send.
const ranNothing =
status === 'interrupted' &&
!this.deps.activeTurns.has(event.threadId, turnId) &&
!this.deps.items.ordinals.hasItems(event.threadId, turnId) &&
![...this.deps.pendingPrompts.values()].some(
(prompt) => prompt.threadId === event.threadId && prompt.turnId === turnId
)
const turnLifecycle =
this.ownsRecord(event.threadId, turnId) && !ranNothing
? this.settled(event.threadId, turnId, {
state: codexTurnLifecycleState(status),
outcome: codexTurnOutcome(status),
completedAt,
durationMs: readCodexTurnDurationMs(event.params)
})
: null
const requestOrigin = this.deps.activeTurns.requestOrigin(event.threadId, turnId)
const latestDispatchSequence = this.deps.activeTurns.latestDispatchSequence(
event.threadId,
@@ -57,6 +57,11 @@ export class CodexJournalActiveTurns {
return [...(this.byThread.get(threadId) ?? [])].at(-1) ?? null
}
/** Whether this child started the turn: its running record is this translator's. */
has(threadId: string, turnId: string): boolean {
return this.byThread.get(threadId)?.has(turnId) === true
}
startedAt(threadId: string, turnId: string): number | undefined {
return this.startedAtByTurn.get(this.turnKey(threadId, turnId))
}
+5
View File
@@ -100,6 +100,11 @@ export class CodexTurnOrdinals {
}
}
/** Whether any of the turn's items was read. */
hasItems(threadId: string, turnId: string): boolean {
return this.turns.has(this.turnKey(threadId, turnId))
}
ordinalFor(threadId: string, turnId: string, codexItemId: string): number {
const turnKey = this.turnKey(threadId, turnId)
let turn = this.turns.get(turnKey)
@@ -35,6 +35,7 @@ import {
} from './structured-agent-session-host-test-data'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
import { claudeAndCodexDeclared } from './structured-agent-session-adapter-router-test-support'
import { isStructuredAgentSessionMainAgentWorking } from '../../../shared/structured-agent-session-main-agent-working'
import { DISPATCH_DOUBT_PROVIDER_IDLE } from '../agent-session-journal/journal-dispatch-doubt-reasons'
import { claudeUnwrittenUserMessageError } from '../../claude/claude-agent-sdk-user-message-queue'
@@ -290,6 +291,52 @@ it('releases the follow-up when the CLI goes idle on a started send it never ech
expect(submissions.find((entry) => entry.clientMessageId === steer)?.handedOverAt).toBeDefined()
})
/** A person's Stop naming no turn. */
function stop() {
return host.cancel(CALLER, {
envelope: {
sessionId: SESSION,
clientOperationId: hostTestOperationId(),
expectedRuntimeFence: store.getRecord(SESSION)!.lease.runtimeFence,
payloadFingerprint: computeAgentSessionPayloadFingerprint({
method: 'agentSession.cancel',
sessionId: SESSION,
fields: {}
})
}
})
}
// Claude's interrupt names no turn, but its Stop ends the CLI, and that end settles the send.
it('settles a send a Stop took before its echo, and the next message goes out', async () => {
const connection = claude.connections[0]!
const first = await send('FIRST prompt')
await eventually(() => expect(written(connection, 'FIRST prompt')).toBeDefined())
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
await eventually(() => expect(connection.closeCount).toBe(1))
await eventually(async () =>
expect(
(await snapshot()).submissions.find((entry) => entry.clientMessageId === first)?.dispatchState
).not.toBe('pending')
)
const after = await snapshot()
expect(
isStructuredAgentSessionMainAgentWorking(
null,
after.submissions,
store.getRecord(SESSION)!.lease.runtimeFence
)
).toBe(false)
await send('NEXT prompt')
await eventually(() => {
const resumed = claude.connections.at(-1)!
expect(resumed).not.toBe(connection)
expect(written(resumed, 'NEXT prompt')).toBeDefined()
})
})
/** As the real connection: once a close begins it refuses every write; the first `failures`
* closes come back unproven. */
function closeUnprovenFor(connection: FakeConnection, failures: number): void {
@@ -339,19 +386,7 @@ async function stoppedWithUnprovenClose(): Promise<FakeConnection> {
).toBe('accepted')
)
closeUnprovenFor(connection, 1)
const stopped = await host.cancel(CALLER, {
envelope: {
sessionId: SESSION,
clientOperationId: hostTestOperationId(),
expectedRuntimeFence: store.getRecord(SESSION)!.lease.runtimeFence,
payloadFingerprint: computeAgentSessionPayloadFingerprint({
method: 'agentSession.cancel',
sessionId: SESSION,
fields: {}
})
}
})
expect(stopped).toMatchObject({ ok: true })
expect(await stop()).toMatchObject({ ok: true })
frame(connection, {
type: 'result',
subtype: 'error_during_execution',
@@ -18,7 +18,8 @@ import {
} from '../../../shared/agent-session-failure-words'
import {
agentJournalItemKey,
agentJournalSubmissionKey
agentJournalSubmissionKey,
parseAgentJournalItemKey
} from '../../../shared/agent-session-journal-item-key'
import {
AGENT_JOURNAL_THREAD_SCOPE,
@@ -35,10 +36,6 @@ import {
readAgentJournalTurn
} from '../../../shared/agent-session-turn-record'
import type { JournalLifecycleMutationInput } from '../agent-session-journal/journal-row-builders'
import {
isStructuredAgentSessionStopNote,
structuredAgentSessionStopNoteIdentity
} from '../../../shared/structured-agent-session-stop-note-key'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type {
AgentSessionCommandAdmission,
@@ -122,7 +119,18 @@ export function isStructuredAgentSessionCommandTurnId(turnId: string): boolean {
return turnId.startsWith('compact:')
}
export { isStructuredAgentSessionStopNote, structuredAgentSessionStopNoteIdentity }
const STOP_NOTE_PREFIX = 'stop:'
/** A Stop's note, keyed by the turn it stopped (or, with no turn, by its operation). */
export function structuredAgentSessionStopNoteIdentity(stopKey: string): AgentJournalItemIdentity {
return { provider: 'orca', clientMessageId: `${STOP_NOTE_PREFIX}${stopKey}` }
}
/** Whether a journal row is a Stop's note. */
export function isStructuredAgentSessionStopNote(itemId: string): boolean {
const identity = parseAgentJournalItemKey(itemId)
return identity?.provider === 'orca' && identity.clientMessageId.startsWith(STOP_NOTE_PREFIX)
}
/** Whether an earlier Stop already asked the running command `turnId` names to end. Read from the
* journal, so nothing is held that could outlive the command. */
@@ -252,6 +252,123 @@ describe('a Stop that names no turn', () => {
pending.resolve({ data: [], nextCursor: null })
})
// Codex can abort a turn before it starts: it answers the interrupt, then ends the turn it never
// started. That turn ran nothing, so it draws no turn of its own, whichever of the two is read first.
it.each([
['its end follows the answer', 'after'],
['its end is read before the answer', 'before'],
['no end comes', 'none']
])('withdraws a send whose turn Codex never started when it takes a Stop: %s', async (_, end) => {
const codex = fakeCodex()
codex.routes['turn/start'] = () => ({ turn: { id: 'turn-1' } })
const ended = () =>
codex.connections[0].handlers.onNotification?.('turn/completed', {
threadId: THREAD_ID,
turn: { id: 'turn-1', status: 'interrupted' }
})
codex.routes['turn/interrupt'] = () => {
if (end === 'before') {
ended()
} else if (end === 'after') {
setTimeout(ended)
}
return {}
}
const adapter = adapterFor(codex, { codexHome: '/codex/home' })
await adapter.acquire({
identity: identityFor(SESSION),
fence: 1,
spawnToken: 'spawn-unstarted',
events: hostEvents()
})
dispatch.mockImplementation((input) => adapter.dispatch(input))
cancelTurn.mockImplementation((input) => adapter.cancelTurn(input))
const first = send('first send')
await first.result
await eventually(async () => expect((await submission(first.id))?.handedOverAt).toBeDefined())
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
await new Promise((resolve) => setTimeout(resolve, 20))
await host.flushStreamedEvents(SESSION)
expect(await submission(first.id)).toMatchObject({
dispatchState: 'rejected',
reason: DISPATCH_REJECTED_CANCELLED
})
const { items } = await host.journalSnapshot(SESSION)
expect(items.filter((item) => item.body.kind === 'turn')).toEqual([])
expect(await statusRows()).toEqual([])
})
it('keeps the record of a turn Codex never started when something of it was read', async () => {
const codex = fakeCodex()
codex.routes['turn/start'] = () => ({ turn: { id: 'turn-1' } })
codex.routes['turn/interrupt'] = () => {
const notify = codex.connections[0].handlers.onNotification
notify?.('item/completed', {
threadId: THREAD_ID,
turnId: 'turn-1',
item: { type: 'agentMessage', id: 'message-1', text: 'partial' }
})
notify?.('turn/completed', {
threadId: THREAD_ID,
turn: { id: 'turn-1', status: 'interrupted' }
})
return {}
}
const adapter = adapterFor(codex, { codexHome: '/codex/home' })
await adapter.acquire({
identity: identityFor(SESSION),
fence: 1,
spawnToken: 'spawn-unstarted-item',
events: hostEvents()
})
dispatch.mockImplementation((input) => adapter.dispatch(input))
cancelTurn.mockImplementation((input) => adapter.cancelTurn(input))
const first = send('first send')
await first.result
await eventually(async () => expect((await submission(first.id))?.handedOverAt).toBeDefined())
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
await host.flushStreamedEvents(SESSION)
const { items } = await host.journalSnapshot(SESSION)
expect(items.filter((item) => item.body.kind === 'turn').map((item) => item.body)).toEqual([
expect.objectContaining({ turnId: 'turn-1', state: 'interrupted' })
])
// Something of it ran, so it is never taken back as unrun. Codex sends no item before a turn's
// start, so this order only guards the record; the send's own end is not this test's.
expect((await submission(first.id))?.dispatchState).not.toBe('rejected')
})
// Only an interrupted end ran nothing for certain: a failure keeps its record, which carries why.
it('keeps the record of a turn Codex failed without starting it', async () => {
const codex = fakeCodex()
codex.routes['turn/start'] = () => ({ turn: { id: 'turn-1' } })
const adapter = adapterFor(codex, { codexHome: '/codex/home' })
await adapter.acquire({
identity: identityFor(SESSION),
fence: 1,
spawnToken: 'spawn-unstarted-failed',
events: hostEvents()
})
dispatch.mockImplementation((input) => adapter.dispatch(input))
const first = send('first send')
await first.result
await eventually(async () => expect((await submission(first.id))?.handedOverAt).toBeDefined())
codex.connections[0].handlers.onNotification?.('turn/completed', {
threadId: THREAD_ID,
turn: { id: 'turn-1', status: 'failed', error: { message: 'model overloaded' } }
})
await host.flushStreamedEvents(SESSION)
const { items } = await host.journalSnapshot(SESSION)
expect(items.filter((item) => item.body.kind === 'turn').map((item) => item.body)).toEqual([
expect.objectContaining({ turnId: 'turn-1' })
])
})
it('fails a turn Codex refuses for a model picked before the list, then sends again', async () => {
const pending = Promise.withResolvers<unknown>()
const codex = fakeCodex({ 'model/list': () => pending.promise })
@@ -301,12 +418,20 @@ describe('a Stop that names no turn', () => {
pending.resolve({ data: [], nextCursor: null })
})
// The provider took the Stop on the turn it named and never opened it: no turn-ended event will
// come, so the Stop's answer settles the send, and its own row says it never started.
it('delivers a newly accepted message after Stop withdraws the prior send', async () => {
const first = send('first send')
await first.result
await eventually(() => expect(dispatch).toHaveBeenCalledOnce())
cancelTurn.mockResolvedValueOnce({ cancelled: true, turnId: 'turn-never-opened' })
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
expect(await submission(first.id)).toMatchObject({
dispatchState: 'rejected',
reason: DISPATCH_REJECTED_CANCELLED
})
expect(await statusRows()).toEqual([])
const replacement = send('replacement send')
await replacement.result
@@ -338,6 +463,69 @@ describe('a Stop that names no turn', () => {
expect(note?.turnScope).toEqual({ kind: 'turn', turnItemId: record?.itemId })
})
// Codex's own shape: the turn's start is read before the interrupt's answer, so the turn has a
// record, which the Stop ends; the send opened it and stays as it is.
it('ends the turn the provider opened before answering, and withdraws nothing', async () => {
const first = send('first send')
await first.result
await eventually(() => expect(dispatch).toHaveBeenCalledOnce())
cancelTurn.mockImplementationOnce(async () => {
hostEvents().appendItem(
{ provider: 'legacy', agent: 'codex', sessionId: SESSION, recordId: 'turn:turn-1' },
{ kind: 'turn', turnId: 'turn-1', state: 'running' },
{ turnScope: AGENT_JOURNAL_THREAD_SCOPE }
)
return { cancelled: true, turnId: 'turn-1' }
})
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
expect((await submission(first.id))?.dispatchState).toBe('pending')
const turn = (await host.journalSnapshot(SESSION)).items.find(
(item) => item.body.kind === 'turn'
)
expect(turn?.body).toMatchObject({ state: 'interrupted' })
expect(await statusRows()).toEqual(['Cancellation requested.'])
})
it('withdraws nothing on a refusal, even one that names the turn', async () => {
const first = send('first send')
await first.result
await eventually(async () => expect((await submission(first.id))?.handedOverAt).toBeDefined())
cancelTurn.mockResolvedValueOnce({
cancelled: false,
turnId: 'turn-never-opened',
refusal: {
detail: { text: 'no active turn to interrupt', audience: 'person' },
turnNotRunning: true
}
})
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: false } })
expect((await submission(first.id))?.dispatchState).toBe('pending')
})
it('answers the Stop when settling the send it took fails, and reports the failure', async () => {
const first = send('first send')
await first.result
await eventually(() => expect(dispatch).toHaveBeenCalledOnce())
const journal = host.collaboratorsForTests().sessions.get(SESSION)?.journal
if (!journal) {
throw new Error('expected the conversation open')
}
vi.spyOn(journal, 'resolveDispatch').mockRejectedValueOnce(new Error('the journal threw'))
cancelTurn.mockResolvedValueOnce({ cancelled: true, turnId: 'turn-never-opened' })
expect(await stop()).toMatchObject({ ok: true, value: { cancelled: true } })
expect(log.entries).toContainEqual(
expect.objectContaining({
fields: expect.objectContaining({ scope: 'stop-settle', sessionId: SESSION })
})
)
})
it('is a valid cancel, and only a plain Stop may omit the turn', () => {
const envelope = {
sessionId: SESSION,
@@ -0,0 +1,48 @@
// What a Stop the provider took settles from its own answer, with no turn-ended event to wait on.
import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record'
import { runningTurnLifecycleRevisions } from './structured-agent-session-stale-turn-verdict'
import type { AgentSessionTurnContext } from './structured-agent-session-turns'
import { withdrawCodexSendsNoTurnOpenedFor } from './structured-agent-session-unopened-send-withdrawal'
/**
* Settles the turn the provider says the Stop took: a turn with records that still reads running
* ends interrupted, once, while the Stop still binds it; the provider's stream lands as it arrives,
* so the fold already holds any end it sent. A turn with no record never opened, so the send it was
* for is withdrawn as its child's end would withdraw it. Returns whether the turn had opened.
* Bookkeeping: a failure is reported, never the Stop's.
*/
export async function settleTakenStop(
ctx: AgentSessionTurnContext,
turnId: string
): Promise<boolean> {
let opened = false
try {
const records = ctx.journal
.snapshot()
.items.filter((item) => readAgentJournalTurn(item.body)?.turnId === turnId)
opened = records.length > 0
if (!opened) {
await withdrawCodexSendsNoTurnOpenedFor(ctx.journal, ctx.fence)
return false
}
const mutations = runningTurnLifecycleRevisions(records, {
state: 'interrupted',
completedAt: ctx.now()
})
if (mutations.length > 0) {
await ctx.journal.appendLifecycleBatch({
settlementId: `stop-settled:${turnId}`,
mutations,
fence: ctx.fence
})
}
} catch (error) {
ctx.logger.warn('settling what a Stop took at its settle failed', {
scope: 'stop-settle',
sessionId: ctx.sessionId,
error
})
}
return opened
}
@@ -6,7 +6,6 @@ import type {
AgentJournalStatusItem,
AgentJournalSubmission
} from '../../../shared/agent-session-journal-types'
import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record'
import type { AgentSessionCancelResult } from '../../../shared/agent-session-wire'
import { latestJournalDispatchObservation } from '../agent-session-journal/journal-dispatch-observation'
import type { AgentSessionCancelOutcome } from './structured-agent-session-adapter'
@@ -26,11 +25,11 @@ import {
structuredAgentSessionStoppedTurnId
} from './structured-agent-session-turn-stop-notes'
import { isStructuredAgentSessionMainAgentWorking } from '../../../shared/structured-agent-session-main-agent-working'
import { runningTurnLifecycleRevisions } from './structured-agent-session-stale-turn-verdict'
import type { StructuredAgentSessionStopWindDown } from './structured-agent-session-stop-wind-down'
import type { JournalStopFailedOn } from '../agent-session-journal/queued-message-pause'
import { structuredAgentSessionFailedStopMark } from './structured-agent-session-stopping'
import { sendStopCanTakeBack } from './structured-agent-session-unopened-send-withdrawal'
import { settleTakenStop } from './structured-agent-session-taken-stop-settle'
import type { AgentSessionTurnContext, TurnOutcome } from './structured-agent-session-turns'
/** Whether the fold reads working. Every write has landed by its call's return, so a Stop reads it
@@ -61,7 +60,7 @@ function stillRunsStoppedTurn(
return stoppedTurnId === null || ctx.journal.activeTurnId() === stoppedTurnId
}
/** Whether the child's end took back every send the Stop found in flight, none of them having run
/** Whether the Stop took back every send it found in flight, none of them having run
* (`withdrawCodexSendsNoTurnOpenedFor`). */
function tookBackEverySend(
ctx: Pick<AgentSessionTurnContext, 'journal'>,
@@ -91,34 +90,6 @@ function stopRefusedNote(
}
}
/** A turn the Stop's interrupt took that still reads running: the Stop ends it, once, while it
* still binds it. The provider's stream lands as it arrives, so the fold already holds any end it
* sent. Bookkeeping: a failure is reported, never the Stop's. */
async function endStoppedTurnAtSettle(ctx: AgentSessionTurnContext, turnId: string): Promise<void> {
try {
const running = ctx.journal
.snapshot()
.items.filter((item) => readAgentJournalTurn(item.body)?.turnId === turnId)
const mutations = runningTurnLifecycleRevisions(running, {
state: 'interrupted',
completedAt: ctx.now()
})
if (mutations.length > 0) {
await ctx.journal.appendLifecycleBatch({
settlementId: `stop-settled:${turnId}`,
mutations,
fence: ctx.fence
})
}
} catch (error) {
ctx.logger.warn("ending a stopped turn at its Stop's settle failed", {
scope: 'stop-settle',
sessionId: ctx.sessionId,
error
})
}
}
/** What the Stop's settle binds: the turn it stopped, and whether its wind-down closes it. */
type StopSettleBinding = {
turnId?: string
@@ -341,7 +312,11 @@ async function cancelAndNote(
binding.failedOn = structuredAgentSessionFailedStopMark(ctx.journal)
}
if (taken === true && stoppedTurn !== undefined && !binding.closedByWindDown && !input.scope) {
await endStoppedTurnAtSettle(ctx, stoppedTurn)
const opened = await settleTakenStop(ctx, stoppedTurn)
// As after a child's end: the send's own row says it never started.
if (!opened && stoppedTurnId === null && tookBackEverySend(ctx, sentBeforeStop)) {
note = null
}
}
const value = { ...(input.turnId !== undefined ? { turnId: input.turnId } : {}), cancelled }
if (input.scope || note === null) {
@@ -1,8 +1,11 @@
// Which sends a Stop's child end takes back: the same rule decides the withdrawal and whether the
// Which sends a Stop takes back, at its child's end or once Codex takes it: the same rule decides the withdrawal and whether the
// Stop writes a row, so a send it must leave alone is left alone by both.
import { describe, expect, it, vi } from 'vitest'
import type { AgentJournalSubmission } from '../../../shared/agent-session-journal-types'
import type {
AgentJournalSubmission,
AgentJournalTurnScope
} from '../../../shared/agent-session-journal-types'
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
import {
sendStopCanTakeBack,
@@ -32,7 +35,20 @@ function sendBeforeStop(overrides: Partial<Send> = {}): Send {
}
}
function journalWith(send: Send) {
/** A turn that ended at sequence 4, before a handover at 5. */
const turnEndedBeforeHandover = {
itemId: 'turn-0',
revision: 1,
sequence: 4,
observedAt: 4,
body: { kind: 'turn' as const, turnId: 'turn-0', state: 'completed' as const }
}
function journalWith(
send: Send,
turnScope: AgentJournalTurnScope = { kind: 'thread' },
earlier: (typeof turnEndedBeforeHandover)[] = []
) {
const resolveDispatch = vi.fn(async () => ({ epoch: 'epoch-1', sequence: 11 }))
const journal: UnopenedSendJournal = {
agent: 'codex',
@@ -41,13 +57,14 @@ function journalWith(send: Send) {
},
snapshot: () => ({
items: [
...earlier,
{
itemId: agentJournalSubmissionKey(send.clientMessageId),
revision: 0,
sequence: 5,
observedAt: 5,
body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'look around' }] },
turnScope: { kind: 'thread' }
turnScope
}
]
}),
@@ -76,7 +93,7 @@ describe('a send a Stop can take back', () => {
})
})
describe("a Codex child's end under a person's Stop", () => {
describe("a Codex child's end, or a Stop Codex took, under a person's Stop", () => {
it('withdraws a send whose answer this process lost', async () => {
const { journal, resolveDispatch } = journalWith(sendBeforeStop({ dispatchState: 'unknown' }))
@@ -97,6 +114,33 @@ describe("a Codex child's end under a person's Stop", () => {
expect(resolveDispatch).not.toHaveBeenCalled()
})
// The send could start a turn only once handed over, so a turn that ended before then is not one.
it('withdraws a send when a turn record was written between its acceptance and its handover', async () => {
const { journal, resolveDispatch } = journalWith(
sendBeforeStop({ acceptedSequence: 3 }),
{ kind: 'thread' },
[turnEndedBeforeHandover]
)
await withdrawCodexSendsNoTurnOpenedFor(journal, 1)
expect(resolveDispatch).toHaveBeenCalledWith(
expect.objectContaining({ clientMessageId: 'send-1', state: 'rejected' })
)
})
// A send that joined a running turn may be in it.
it('leaves a send steered into a running turn as it is', async () => {
const { journal, resolveDispatch } = journalWith(sendBeforeStop(), {
kind: 'turn',
turnItemId: 'turn-item-1'
})
await withdrawCodexSendsNoTurnOpenedFor(journal, 1)
expect(resolveDispatch).not.toHaveBeenCalled()
})
it("leaves a queued card's send as it is", async () => {
const { journal, resolveDispatch } = journalWith(sendBeforeStop({ handedOverAt: undefined }))
@@ -1,5 +1,5 @@
// Which unanswered sends a dying Codex child takes with it as never sent, read from the journal at
// the settlement that lands, so a retried wind-down reads the same rows.
// Which unanswered sends a person's Stop takes back as never sent, when Codex takes it or the child
// ends, read from the journal at the settlement that lands, so a retried settle reads the same rows.
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
@@ -44,7 +44,7 @@ export function sendStopCanTakeBack(
/**
* Withdraws the sends a Codex child left unanswered when a person's Stop, in force since they were
* sent, ends it: each one that started its own turn (handed over with no turn running, so its
* sent, ends the child or is taken by Codex: each one that started its own turn (handed over with no turn running, so its
* message belongs to no turn) when no turn has opened since. Codex records a prompt only once its
* turn starts (core tasks/regular.rs:50, session/turn.rs:886-902), so they never ran. A send that
* joined a running turn may be in it, so it, a send with no recorded place, and any other end stay
@@ -63,22 +63,26 @@ export async function withdrawCodexSendsNoTurnOpenedFor(
const turn = readAgentJournalTurn(item.body)
return turn ? [{ running: turn.state === 'running', sequence: item.sequence }] : []
})
// Its message's place, recorded at handover: the turn it joined, or the conversation.
const startedOwnTurn = (clientMessageId: string): boolean =>
items.find((item) => item.itemId === agentJournalSubmissionKey(clientMessageId))?.turnScope
?.kind === 'thread'
const opened = (acceptedSequence: number): boolean =>
turns.some((turn) => turn.running || turn.sequence > acceptedSequence)
const unopened = journal
.submissions()
.filter(
(entry) =>
sendStopCanTakeBack(entry) &&
entry.acceptedSequence !== undefined &&
entry.acceptedSequence < stop.sequence &&
startedOwnTurn(entry.clientMessageId) &&
!opened(entry.acceptedSequence)
// Its message's place, recorded at handover: the turn it joined, or the conversation. A send
// that started its own turn has the conversation's; one with no turn recorded since then never ran.
const ownTurnHandover = (clientMessageId: string): number | undefined => {
const handover = items.find(
(item) => item.itemId === agentJournalSubmissionKey(clientMessageId)
)
return handover?.turnScope?.kind === 'thread' ? handover.sequence : undefined
}
const openedSince = (handoverSequence: number): boolean =>
turns.some((turn) => turn.running || turn.sequence > handoverSequence)
const unopened = journal.submissions().filter((entry) => {
const handover = ownTurnHandover(entry.clientMessageId)
return (
sendStopCanTakeBack(entry) &&
entry.acceptedSequence !== undefined &&
entry.acceptedSequence < stop.sequence &&
handover !== undefined &&
!openedSince(handover)
)
})
const withdrawn = agentSessionFailureWords(agentSessionFailureFact('cancelled'), {
surface: 'rejection'
})
@@ -1,12 +1,9 @@
import { describe, expect, it } from 'vitest'
import { agentSessionFailureFact } from './agent-session-failure'
import { agentSessionFailureWords } from './agent-session-failure-words'
import { agentJournalItemKey, agentJournalSubmissionKey } from './agent-session-journal-item-key'
import { agentJournalSubmissionKey } from './agent-session-journal-item-key'
import {
structuredAgentSessionSendOpeningTurn,
type StructuredAgentSessionOpeningSendItem
} from './structured-agent-session-opening-send'
import { structuredAgentSessionStopNoteIdentity } from './structured-agent-session-stop-note-key'
const send = {
clientMessageId: 'send-1',
@@ -25,12 +22,6 @@ const turnRecord = (sequence: number): StructuredAgentSessionOpeningSendItem =>
sequence,
body: { kind: 'turn', turnId: 'turn-1', state: 'running' }
})
const stopNote = (sequence: number): StructuredAgentSessionOpeningSendItem => ({
itemId: agentJournalItemKey(structuredAgentSessionStopNoteIdentity(`op-${sequence}`)),
sequence,
body: { kind: 'status', text: 'Cancellation requested.' },
turnScope: { kind: 'thread' }
})
function opening(items: StructuredAgentSessionOpeningSendItem[]): boolean {
return structuredAgentSessionSendOpeningTurn([send], (visit) => items.forEach(visit), 1)
@@ -42,31 +33,11 @@ describe('a send still opening its turn', () => {
expect(opening([handover, turnRecord(6)])).toBe(false)
})
// A Stop the provider took before the turn opened may write no turn record at all.
it("ends at a taken Stop's note written after the handover, so a later message is not held", () => {
expect(opening([handover, stopNote(6)])).toBe(false)
})
// A refused or unconfirmed Stop leaves the turn opening: the next message still joins it.
it("keeps holding past a Stop's note that carries a failure", () => {
const refused = stopNote(6)
expect(
opening([
handover,
{
...refused,
body: {
kind: 'status',
...agentSessionFailureWords(agentSessionFailureFact('cancelUnconfirmed'), {
surface: 'row'
})
}
}
])
).toBe(true)
})
it("keeps holding past a Stop's note written before the handover", () => {
expect(opening([stopNote(4), handover])).toBe(true)
// A Stop the provider took before the turn opened withdraws the send: nothing is opening.
it('ends once the send is settled, with no turn record', () => {
const withdrawn = { ...send, dispatchState: 'rejected' } as const
expect(structuredAgentSessionSendOpeningTurn([withdrawn], (visit) => visit(handover), 1)).toBe(
false
)
})
})
@@ -10,7 +10,6 @@ import type {
AgentJournalTurnScope
} from './agent-session-journal-types'
import { readAgentJournalTurn } from './agent-session-turn-record'
import { isStructuredAgentSessionStopNote } from './structured-agent-session-stop-note-key'
/** What the read needs of each journal item; the host visits its fold without a snapshot. */
export type StructuredAgentSessionOpeningSendItem = {
@@ -22,15 +21,15 @@ export type StructuredAgentSessionOpeningSendItem = {
}
/**
* A send handed over while no turn ran, still unsettled, with no turn record and no taken Stop's
* note written since its handover (`fence`, when given: handed over by that child). A message handed over now would join
* A send handed over while no turn ran, still unsettled, with no turn record written since its
* handover (`fence`, when given: handed over by that child). A message handed over now would join
* the turn that send is opening (Codex steers it in once the turn opens, Claude folds it into the
* running cycle), yet its handover row would be written before that turn exists, and so read as
* belonging to none. So the host holds it until the turn opens, and a client draws it after the
* live turn meanwhile. Every way the send stops opening — its turn record, its echo, its refusal, a
* lost answer's doubt, its child's end, a Stop that took — is a journal commit. After a Stop that
* took, the next message is a new instruction, not one for that turn; a Stop that failed or was
* refused (its note carries the failure) leaves the turn opening, so the next message still waits.
* lost answer's doubt, its child's end, a Stop that took — is a journal commit that settles it or
* records its turn. A Stop that failed or was refused leaves the turn opening, so the next message
* still waits.
*/
export function structuredAgentSessionSendOpeningTurn(
submissions: readonly Pick<
@@ -79,12 +78,7 @@ export function structuredAgentSessionOpeningSendItemId(
newest.handover = item.itemId
newest.handoverAt = item.sequence
}
if (
(readAgentJournalTurn(item.body) && isRootAgentJournalItem(item)) ||
(isStructuredAgentSessionStopNote(item.itemId) &&
item.body.kind === 'status' &&
item.body.failure === undefined)
) {
if (readAgentJournalTurn(item.body) && isRootAgentJournalItem(item)) {
newest.settledAt = Math.max(newest.settledAt, item.sequence)
}
})
@@ -1,17 +0,0 @@
// A Stop's note row is keyed so the host and every client can tell it from other Orca rows.
import { parseAgentJournalItemKey } from './agent-session-journal-item-key'
import type { AgentJournalItemIdentity } from './agent-session-journal-types'
const STOP_NOTE_PREFIX = 'stop:'
/** A Stop's note, keyed by the turn it stopped (or, with no turn, by its operation). */
export function structuredAgentSessionStopNoteIdentity(stopKey: string): AgentJournalItemIdentity {
return { provider: 'orca', clientMessageId: `${STOP_NOTE_PREFIX}${stopKey}` }
}
/** Whether a journal row is a Stop's note. */
export function isStructuredAgentSessionStopNote(itemId: string): boolean {
const identity = parseAgentJournalItemKey(itemId)
return identity?.provider === 'orca' && identity.clientMessageId.startsWith(STOP_NOTE_PREFIX)
}