mirror of
https://github.com/stablyai/orca.git
synced 2026-10-02 16:02:15 +00:00
fix(native-chat): a Stop of a start that never landed binds no later turn, whatever sent it
A person's Stop pressed while the agent starts names no turn, and the send it stopped is cancelled before it opens one. The Stop then bound the next turn anything opened (orchestration mail, a restart continuation, the queue's drain, all of which send as the host), so a host eviction of that turn wrote nothing and its crash or close read as the person's cancellation. A Stop that named no turn now binds only a turn no send journaled after it opened: any send since, of any origin and not refused, opens its own. The E2 test's mail send is accepted as Codex accepts it, instead of opening the stopped send's own turn first.
This commit is contained in:
@@ -26,10 +26,7 @@ function stopIsAPersons(reason: JournalStopEvent['reason']): boolean {
|
||||
}
|
||||
}
|
||||
|
||||
type TurnEndState = Pick<
|
||||
JournalReducerState,
|
||||
'items' | 'queuePauseMarks' | 'latestPersonTurnSequence'
|
||||
>
|
||||
type TurnEndState = Pick<JournalReducerState, 'items' | 'queuePauseMarks' | 'submissions'>
|
||||
|
||||
/** Whether `stop`, a person's, makes the end of turn `turnId` theirs: it named that turn, or named
|
||||
* none and stopped the turn item `itemId` opened. */
|
||||
@@ -48,20 +45,26 @@ function stopIsTurnCancellation(
|
||||
}
|
||||
|
||||
/** Pressed before any turn showed, a Stop stopped the first turn opened after it, and no later
|
||||
* one: unless a send a person made since was accepted, whose turn that is. `itemId` null: a turn
|
||||
* not yet opened. */
|
||||
* one: unless a send journaled since, of any origin, was not refused, whose turn that is. A Stop
|
||||
* whose stopped send never opens a turn so binds nothing. `itemId` null: a turn not yet opened. */
|
||||
function turnlessStopStopped(
|
||||
state: TurnEndState,
|
||||
stop: JournalLatestStop,
|
||||
itemId: string | null
|
||||
): boolean {
|
||||
const createdAt = itemId === null ? null : (state.items.get(itemId)?.sequence ?? null)
|
||||
if (
|
||||
(createdAt !== null && createdAt <= stop.sequence) ||
|
||||
state.latestPersonTurnSequence >= stop.sequence
|
||||
) {
|
||||
if (createdAt !== null && createdAt <= stop.sequence) {
|
||||
return false
|
||||
}
|
||||
for (const submission of state.submissions.values()) {
|
||||
if (
|
||||
submission.dispatchState !== 'rejected' &&
|
||||
submission.acceptedSequence !== undefined &&
|
||||
submission.acceptedSequence > stop.sequence
|
||||
) {
|
||||
return false
|
||||
}
|
||||
}
|
||||
for (const [otherId, item] of state.items) {
|
||||
if (
|
||||
otherId !== itemId &&
|
||||
@@ -103,19 +106,17 @@ export function personStopDecidesTurn(
|
||||
turnId: string | null,
|
||||
endedAt?: number
|
||||
): boolean {
|
||||
if (turnId !== null) {
|
||||
const itemId = [...state.items].find(
|
||||
([, item]) => readAgentJournalTurn(item.body)?.turnId === turnId
|
||||
)?.[0]
|
||||
return stopEndsTurnAsCancellation(state, turnId, itemId ?? null, endedAt)
|
||||
}
|
||||
const stop = state.queuePauseMarks.latestStop
|
||||
return (
|
||||
stop !== null &&
|
||||
stop.event.turnId === undefined &&
|
||||
stopIsAPersons(stop.event.reason) &&
|
||||
turnlessStopStopped(state, stop, null)
|
||||
)
|
||||
if (stop === null || !stopIsAPersons(stop.event.reason)) {
|
||||
return false
|
||||
}
|
||||
if (turnId === null) {
|
||||
return stop.event.turnId === undefined && turnlessStopStopped(state, stop, null)
|
||||
}
|
||||
const itemId = [...state.items].find(
|
||||
([, item]) => readAgentJournalTurn(item.body)?.turnId === turnId
|
||||
)?.[0]
|
||||
return stopEndsTurnAsCancellation(state, turnId, itemId ?? null, endedAt)
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+128
@@ -0,0 +1,128 @@
|
||||
// Which turn a person's Stop that named no turn binds: only the one its stopped send opens. A Stop of
|
||||
// a start that never landed stopped a send that opens no turn, and a send journaled after the Stop
|
||||
// opens its own; neither is the Stop's, whatever sent it.
|
||||
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import {
|
||||
AGENT_JOURNAL_THREAD_SCOPE,
|
||||
type AgentJournalItemIdentity
|
||||
} from '../../../shared/agent-session-journal-types'
|
||||
import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record'
|
||||
import type { JournalStopEvent } from '../agent-session-journal/journal-row-schema'
|
||||
import { settleStaleStructuredAgentSessionState } from './structured-agent-session-dead-generation-settlement'
|
||||
import { HOST_TEST_SESSION } from './structured-agent-session-host-test-data'
|
||||
import {
|
||||
createQueuedMessageTestRig,
|
||||
eventually,
|
||||
type QueuedMessageTestRig
|
||||
} from './structured-agent-session-queued-message-rig.test-fixture'
|
||||
|
||||
let rig: QueuedMessageTestRig
|
||||
|
||||
afterEach(() => rig.dispose())
|
||||
|
||||
const MAIL_TURN: AgentJournalItemIdentity = {
|
||||
provider: 'codex',
|
||||
threadId: 'thread-1',
|
||||
turnId: 'turn-mail',
|
||||
ordinal: 999
|
||||
}
|
||||
|
||||
function journal() {
|
||||
const open = rig.host.collaboratorsForTests().sessions.get(HOST_TEST_SESSION)?.journal
|
||||
if (!open) {
|
||||
throw new Error('expected the conversation open')
|
||||
}
|
||||
return open
|
||||
}
|
||||
|
||||
function stopEvents(): JournalStopEvent[] {
|
||||
const since = journal().readSince({ epoch: journal().epoch, sequence: 0 })
|
||||
if (!since.ok) {
|
||||
throw new Error(`expected rows, got reset ${since.reset}`)
|
||||
}
|
||||
return since.rows.flatMap((row) =>
|
||||
row.kind === 'tombstone' && row.stopEvent ? [row.stopEvent] : []
|
||||
)
|
||||
}
|
||||
|
||||
function childPhase() {
|
||||
return rig.host.collaboratorsForTests().sessions.get(HOST_TEST_SESSION)?.child?.phase
|
||||
}
|
||||
|
||||
function fence(): number {
|
||||
return rig.store.getRecord(HOST_TEST_SESSION)?.lease.runtimeFence ?? 1
|
||||
}
|
||||
|
||||
function mailTurn() {
|
||||
return journal()
|
||||
.snapshot()
|
||||
.items.map((item) => readAgentJournalTurn(item.body))
|
||||
.find((turn) => turn?.turnId === 'turn-mail')
|
||||
}
|
||||
|
||||
/** A person's Stop of a start that never landed, whose send opens no turn; then orchestration mail
|
||||
* starts a new child, which lands and runs the mail's turn. */
|
||||
async function mailTurnAfterStopOfStart(): Promise<void> {
|
||||
rig = await createQueuedMessageTestRig({ starting: true, restartable: true })
|
||||
let release: () => void = () => undefined
|
||||
rig.awaitStarted.mockImplementation(
|
||||
() => new Promise<undefined>((resolve) => (release = () => resolve(undefined)))
|
||||
)
|
||||
rig.send('work on this')
|
||||
await eventually(() => expect(childPhase()).toBe('starting'))
|
||||
expect(await rig.stop()).toMatchObject({ ok: true })
|
||||
release()
|
||||
expect(stopEvents()).toEqual([expect.objectContaining({ reason: 'user-stop' })])
|
||||
expect(stopEvents()[0]).not.toHaveProperty('turnId')
|
||||
await eventually(() => expect(childPhase()).toBeUndefined())
|
||||
rig.awaitStarted.mockImplementation(async () => undefined)
|
||||
await rig.send('mail for the worker', undefined, { internal: true }).result
|
||||
await eventually(() => expect(rig.dispatch).toHaveBeenCalled())
|
||||
await journal().appendItem(
|
||||
MAIL_TURN,
|
||||
{ kind: 'turn', turnId: 'turn-mail', state: 'running', startedAt: Date.now() },
|
||||
{ fence: fence(), turnScope: AGENT_JOURNAL_THREAD_SCOPE }
|
||||
)
|
||||
}
|
||||
|
||||
describe('a Stop of a start that never landed binds no later turn', () => {
|
||||
it("writes the host's event when it evicts the mail turn, which reads as news", async () => {
|
||||
await mailTurnAfterStopOfStart()
|
||||
let atClose: JournalStopEvent[] = []
|
||||
rig.closeSession.mockImplementationOnce(async () => {
|
||||
atClose = stopEvents()
|
||||
return true
|
||||
})
|
||||
|
||||
await rig.host.close(HOST_TEST_SESSION, 'evict')
|
||||
|
||||
expect(atClose.map((event) => event.reason)).toEqual(['user-stop', 'evict'])
|
||||
await rig.host.journalSnapshot(HOST_TEST_SESSION)
|
||||
expect(mailTurn()).toMatchObject({ state: 'interrupted' })
|
||||
expect(mailTurn()).not.toHaveProperty('outcome')
|
||||
})
|
||||
|
||||
it('settles a crash of the mail turn on relaunch as news', async () => {
|
||||
await mailTurnAfterStopOfStart()
|
||||
const owner = fence()
|
||||
rig.crashRestartHostProcess()
|
||||
await rig.host.journalSnapshot(HOST_TEST_SESSION)
|
||||
|
||||
await settleStaleStructuredAgentSessionState({
|
||||
journal: journal(),
|
||||
sessionId: HOST_TEST_SESSION,
|
||||
fence: owner + 1,
|
||||
acquisitionGeneration: 'generation-2',
|
||||
deathEvidence: {
|
||||
kind: 'exit-observed',
|
||||
detail: 'the relaunch proved the old child gone',
|
||||
observedAt: Date.now() + 60_000,
|
||||
ownerFence: owner
|
||||
}
|
||||
})
|
||||
|
||||
expect(mailTurn()).toMatchObject({ state: 'interrupted' })
|
||||
expect(mailTurn()).not.toHaveProperty('outcome')
|
||||
})
|
||||
})
|
||||
+3
-25
@@ -60,28 +60,6 @@ async function runningTurn(turnId = 'turn-1'): Promise<string> {
|
||||
return working
|
||||
}
|
||||
|
||||
/** The stopped send's turn opens after its turnless Stop and ends cut: the turn that Stop bound. */
|
||||
async function stoppedTurnEnded(): Promise<void> {
|
||||
const fence = rig.store.getRecord(HOST_TEST_SESSION)?.lease.runtimeFence ?? 1
|
||||
const scope = { fence, turnScope: AGENT_JOURNAL_THREAD_SCOPE }
|
||||
const identity = {
|
||||
provider: 'codex' as const,
|
||||
threadId: 'thread-1',
|
||||
turnId: 'turn-stopped',
|
||||
ordinal: 998
|
||||
}
|
||||
await journal().appendItem(
|
||||
identity,
|
||||
{ kind: 'turn', turnId: 'turn-stopped', state: 'running', startedAt: 1 },
|
||||
scope
|
||||
)
|
||||
await journal().appendItem(
|
||||
identity,
|
||||
{ kind: 'turn', turnId: 'turn-stopped', state: 'interrupted', completedAt: Date.now() },
|
||||
scope
|
||||
)
|
||||
}
|
||||
|
||||
async function queuedDraft(text: string): Promise<string> {
|
||||
const queued = await rig.send(text, 'queue-if-active').result
|
||||
if (!queued.ok || !('queued' in queued.value)) {
|
||||
@@ -257,10 +235,11 @@ describe("a person's Stop pause and the Stop events after it", () => {
|
||||
await queuedDraft('queued behind the turn')
|
||||
expect(await rig.stop()).toMatchObject({ ok: true })
|
||||
await rig.settleAccepted(working, 'stopped')
|
||||
await stoppedTurnEnded()
|
||||
expect(await rig.queuePause()).toEqual({ reason: 'stopped' })
|
||||
// Orchestration mail starts a turn the host sent, which lifts nothing.
|
||||
await rig.send('mail for the lead', undefined, { internal: true }).result
|
||||
const mail = rig.send('mail for the lead', undefined, { internal: true })
|
||||
await mail.result
|
||||
await rig.settleAccepted(mail.id, 'mail')
|
||||
await journal().appendItem(
|
||||
{ provider: 'codex', threadId: 'thread-1', turnId: 'turn-mail', ordinal: 999 },
|
||||
{ kind: 'turn', turnId: 'turn-mail', state: 'running', startedAt: 1 },
|
||||
@@ -285,7 +264,6 @@ describe("a person's Stop pause and the Stop events after it", () => {
|
||||
const held = await queuedDraft('queued behind the turn')
|
||||
expect(await rig.stop()).toMatchObject({ ok: true })
|
||||
await rig.settleAccepted(working, 'stopped')
|
||||
await stoppedTurnEnded()
|
||||
// The agent at rest goes, writing nothing; mail then starts a new child, which never lands.
|
||||
await idleSweep().tick()
|
||||
expect(await rig.queuePause()).toEqual({ reason: 'stopped' })
|
||||
|
||||
Reference in New Issue
Block a user