mirror of
https://github.com/stablyai/orca.git
synced 2026-10-04 16:02:08 +00:00
* fix(codex): settle a structured send on admission, and stop minting a colliding identity Two sends could be written into the journal under one durable identity. Codex coalesces a mid-turn `turn/start` into the running turn rather than refusing it -- measured against real `codex app-server` builds 0.147.0, 0.150.1 and 0.153.4, none of which refuse and none of which fire a second `turn/started`. The dispatch path read the turn id from the turn/start response and stamped every accepted send `ordinal: 0`. Since a coalesced send gets the running turn's id back, two submissions persisted the same `providerItemId`. That string is durable, and it is the key a restore uses to match a submission against provider history, so the second message's real history row matched nothing and rendered as an extra bubble on replay. On 0.147.0 it is worse than a collision: the coalesced response returns a turn id that never starts and never completes, so the persisted key named a turn absent from history and NEITHER message could match. Identity is now minted from the echoed user message at `identityFor` -- the single point that mints the journal row's own identity -- so the settled key is by construction the one replay computes, rather than a parallel calculation that can drift. Dispatch returns `admitted` when the transport write completes; identity settles on the echo through a channel that did not previously exist for Codex. Waiters are keyed by client message id instead of being shifted off the front of an array by arrival order, and they are cleared on session close and child exit -- previously a timeout was the only thing that ever ended one. `TURN_ID_WAIT_MS` is deleted. It was never reachable on any build measured: `readCodexTurnId` returns non-null on all three, so the 10s wait never fired. The comment justifying it claimed older builds acknowledge before the id exists, which no tested build does. Three comments asserting Codex answers a mid-turn send with `turn already running` are corrected. Their only backing was a test fixture inventing that error string. The correction is factual only -- every changed line in `src/main/runtime/orchestration/` is a comment, and mid-turn delivery is still refused for both providers. Whether that policy is right is a separate question; it was resting on a false premise. Known gap, stated rather than implied: this prevents new collisions and does not repair journals already written with a colliding or phantom key. Those conversations keep duplicating on restore. Repairing them means re-matching persisted submissions against provider history and rewriting `providerItemId` -- which is what `journal-submission-reconciler.ts` is written for, and it still has no production caller. * test(codex): drop the synchronous-accept contract and the colliding `:0` from the integration fakes Three tests in the structured-session integration suites encoded the dispatch contract this branch replaces, and two of them pinned the defect it fixes. They asserted `agentSession.send` answers `dispatchState: 'accepted'` carrying `providerItemId: codex:<thread>:<turn>:0` at send time. That ordinal was never observed; it was stamped on every accepted send, which is exactly the collision this branch removes -- a send coalesced into a running turn is answered with the running turn's id, so two submissions persisted one durable key. The visible failure was a 30s timeout rather than a failed assertion. The fake client advertised no `agent-session.pending-send-result.v1`, and without it the host holds the reply until the send settles: a shim for clients too old to render a pending bubble. The fake provider then echoed the user message with no `clientId`, so nothing could correlate that echo back to the submission, and the wait ran to its own 30s ceiling. Real Codex sends `clientId` on that echo, and the fake now does too, which is what makes it a model of the provider rather than a sketch of one. The identity assertion is kept rather than dropped. Each send now asserts `pending` with no identity at admission, then asserts the submission settles `accepted` at `codex:<thread>:<turn>:0` once the echo lands. Same ordinal, but earned from `identityFor` on the echo -- the key a replay recomputes -- instead of guessed from the turn/start response. Ablated: removing `clientId` from the two echoes leaves both submissions `pending` and fails both assertions, so the assertion is load-bearing and not satisfied by something incidental. Both suites' client fixtures now advertise the capability set the desktop renderer sends in `src/main/ipc/runtime.ts`, which is what these suites mean by a client. The older-client settlement wait keeps its own coverage in `src/main/runtime/rpc/methods/structured-agent-session.test.ts`. `structured-agent-session-runtime-exit.test.ts` asserts `pending` for the same reason; it drives the host directly, so it never took the compatibility path, and what proves delivery there is still the turn the reacquired provider starts. The replay suite's "without dispatching it twice" property is untouched: one `turn/start` call, one replayed ledger row. * fix(codex): preserve unsettled dispatch correlations * test(codex): type the dispatch fixtures instead of asserting over them main's new casting gate (#20367 base) flags type assertions on changed lines. Replace them with checked types: the recording sink already satisfies its interface, both CodexSession fixtures are now annotated and carry real collaborators, the settlement assertion compares whole identities, and the integration helper reads submissions through the host's public journalSnapshot instead of its private session map. * fix(test): merge the duplicate doubt-reasons import the merge left behind Both sides added an import from journal-dispatch-doubt-reasons and the merge kept both statements, which the whole-repo native plugin gate refuses under --deny-warnings. * test(codex): a Fast mode turn is admitted, not accepted #20506 landed its Fast mode tests against the dispatch contract this branch replaces: a Codex send now returns admitted and settles its identity on the provider echo. The tier assertions the test exists for are untouched. --------- Co-authored-by: Merge Sim <sim@local>
308 lines
10 KiB
TypeScript
308 lines
10 KiB
TypeScript
import { describe, expect, it, vi } from 'vitest'
|
|
import type {
|
|
AgentJournalMessageItem,
|
|
AgentSessionJournalIdentity
|
|
} from '../../shared/agent-session-journal-types'
|
|
import {
|
|
CodexAppServerRequestError,
|
|
type CodexAppServerConnection,
|
|
type CodexAppServerConnectionHandlers,
|
|
type CodexAppServerLaunch,
|
|
type openCodexAppServerConnection
|
|
} from './codex-app-server-connection'
|
|
import { CodexAppServerUnsupportedError } from './codex-app-server-session'
|
|
import {
|
|
CodexStructuredSessionAdapter,
|
|
type CodexStructuredSessionAdapterDeps,
|
|
type CodexStructuredSessionEvent
|
|
} from './codex-structured-session-adapter'
|
|
|
|
const THREAD_ID = 'thread-abc'
|
|
const USER_MESSAGE: AgentJournalMessageItem = {
|
|
kind: 'message',
|
|
role: 'user',
|
|
blocks: [{ type: 'text', text: 'ship it' }]
|
|
}
|
|
|
|
type Route = (params: Record<string, unknown> | undefined) => unknown
|
|
type FakeConnection = Omit<CodexAppServerConnection, 'closed'> & {
|
|
closed: boolean
|
|
launch: CodexAppServerLaunch
|
|
handlers: CodexAppServerConnectionHandlers
|
|
calls: { method: string; params?: Record<string, unknown> }[]
|
|
}
|
|
|
|
function identity(): AgentSessionJournalIdentity {
|
|
return {
|
|
sessionId: 'session-1',
|
|
workspaceId: 'ws-1',
|
|
hostId: 'host-1',
|
|
agent: 'codex',
|
|
providerHandle: { kind: 'codex', threadId: THREAD_ID }
|
|
}
|
|
}
|
|
|
|
function fakeCodex(): {
|
|
connections: FakeConnection[]
|
|
openConnection: typeof openCodexAppServerConnection
|
|
routes: Record<string, Route>
|
|
} {
|
|
const connections: FakeConnection[] = []
|
|
const routes: Record<string, Route> = {
|
|
'thread/resume': () => ({ thread: { id: THREAD_ID } })
|
|
}
|
|
const openConnection = (async (launch, handlers = {}) => {
|
|
const connection: FakeConnection = {
|
|
launch,
|
|
handlers,
|
|
calls: [],
|
|
pid: 4321,
|
|
closed: false,
|
|
request: async (method, params) => {
|
|
connection.calls.push({ method, params })
|
|
return routes[method]?.(params) ?? {}
|
|
},
|
|
notify: () => {},
|
|
respond: () => {},
|
|
respondWithError: () => {},
|
|
close: async () => {
|
|
connection.closed = true
|
|
return true
|
|
}
|
|
}
|
|
connections.push(connection)
|
|
return connection
|
|
}) as typeof openCodexAppServerConnection
|
|
return { connections, openConnection, routes }
|
|
}
|
|
|
|
async function acquired(
|
|
codex: ReturnType<typeof fakeCodex>,
|
|
events: CodexStructuredSessionEvent[] = [],
|
|
processControl: Partial<
|
|
Pick<
|
|
CodexStructuredSessionAdapterDeps,
|
|
'captureTurnProcesses' | 'terminateTurnProcesses' | 'now'
|
|
>
|
|
> = {}
|
|
): Promise<CodexStructuredSessionAdapter> {
|
|
const adapter = new CodexStructuredSessionAdapter({
|
|
resolveLaunch: async () => ({
|
|
command: 'codex',
|
|
args: ['app-server'],
|
|
cwd: '/work/repo',
|
|
codexHome: null,
|
|
resumeThreadId: THREAD_ID
|
|
}),
|
|
onEvent: (event) => events.push(event),
|
|
openConnection: codex.openConnection,
|
|
readProcessStartTime: async () => 1_700_000_000_000,
|
|
captureTurnProcesses: async () => ({ platform: 'win32', identities: new Map() }),
|
|
terminateTurnProcesses: async () => true,
|
|
...processControl
|
|
})
|
|
await adapter.acquire({ identity: identity(), fence: 7, spawnToken: 'spawn-9' })
|
|
return adapter
|
|
}
|
|
|
|
function completeTurn(codex: ReturnType<typeof fakeCodex>, turnId = 'turn-1'): void {
|
|
codex.connections[0].handlers.onNotification?.('turn/completed', {
|
|
threadId: THREAD_ID,
|
|
turn: { id: turnId, status: 'interrupted' }
|
|
})
|
|
}
|
|
|
|
describe('CodexStructuredSessionAdapter.cancelTurn', () => {
|
|
it('confirms an interrupt Codex acknowledged', async () => {
|
|
const codex = fakeCodex()
|
|
const adapter = await acquired(codex)
|
|
|
|
await expect(
|
|
adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
).resolves.toEqual({ cancelled: true })
|
|
expect(codex.connections[0].calls.at(-1)).toEqual({
|
|
method: 'turn/interrupt',
|
|
params: { threadId: THREAD_ID, turnId: 'turn-1' }
|
|
})
|
|
})
|
|
|
|
it('reports not-cancelled when Codex declines or lacks the method', async () => {
|
|
const declined = fakeCodex()
|
|
declined.routes['turn/interrupt'] = () => {
|
|
throw new CodexAppServerRequestError('turn/interrupt', -32602, 'no such turn')
|
|
}
|
|
const absent = fakeCodex()
|
|
absent.routes['turn/interrupt'] = () => {
|
|
throw new CodexAppServerUnsupportedError('no turn/interrupt')
|
|
}
|
|
|
|
await expect(
|
|
(await acquired(declined)).cancelTurn({
|
|
sessionId: 'session-1',
|
|
turnId: 'turn-1',
|
|
fence: 7
|
|
})
|
|
).resolves.toEqual({ cancelled: false })
|
|
await expect(
|
|
(await acquired(absent)).cancelTurn({
|
|
sessionId: 'session-1',
|
|
turnId: 'turn-1',
|
|
fence: 7
|
|
})
|
|
).resolves.toEqual({ cancelled: false })
|
|
})
|
|
|
|
it('rethrows an unsettled interrupt so the turn is not shown as cancelled', async () => {
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/interrupt'] = () => {
|
|
throw new Error('codex app-server turn/interrupt exceeded 30000ms')
|
|
}
|
|
|
|
await expect(
|
|
(await acquired(codex)).cancelTurn({
|
|
sessionId: 'session-1',
|
|
turnId: 'turn-1',
|
|
fence: 7
|
|
})
|
|
).rejects.toThrow('exceeded 30000ms')
|
|
})
|
|
|
|
it('publishes terminal state only after streaming interruption is physically settled', async () => {
|
|
const events: CodexStructuredSessionEvent[] = []
|
|
let finishTermination!: (terminated: boolean) => void
|
|
const termination = new Promise<boolean>((resolve) => {
|
|
finishTermination = resolve
|
|
})
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/interrupt'] = () => {
|
|
completeTurn(codex)
|
|
return {}
|
|
}
|
|
const adapter = await acquired(codex, events, {
|
|
terminateTurnProcesses: async () => termination
|
|
})
|
|
codex.connections[0].handlers.onNotification?.('item/agentMessage/delta', {
|
|
threadId: THREAD_ID,
|
|
turnId: 'turn-1',
|
|
itemId: 'item-1',
|
|
delta: 'still streaming'
|
|
})
|
|
|
|
const pending = adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
await vi.waitFor(() => expect(codex.connections[0].calls.at(-1)?.method).toBe('turn/interrupt'))
|
|
expect(events).toContainEqual(expect.objectContaining({ method: 'item/agentMessage/delta' }))
|
|
expect(events).not.toContainEqual(expect.objectContaining({ method: 'turn/completed' }))
|
|
|
|
finishTermination(true)
|
|
await expect(pending).resolves.toEqual({ cancelled: true })
|
|
expect(events.at(-1)).toMatchObject({ method: 'turn/completed' })
|
|
})
|
|
|
|
it('starts physical termination without waiting for the interrupt receipt', async () => {
|
|
let finishInterrupt!: () => void
|
|
const interruptReceipt = new Promise<void>((resolve) => {
|
|
finishInterrupt = resolve
|
|
})
|
|
const terminateTurnProcesses = vi.fn(async () => true)
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/interrupt'] = () => interruptReceipt
|
|
const adapter = await acquired(codex, [], { terminateTurnProcesses })
|
|
|
|
const pending = adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
await vi.waitFor(() => expect(terminateTurnProcesses).toHaveBeenCalledOnce())
|
|
finishInterrupt()
|
|
|
|
await expect(pending).resolves.toEqual({ cancelled: true })
|
|
})
|
|
|
|
it('keeps the turn live when process termination cannot be verified', async () => {
|
|
const events: CodexStructuredSessionEvent[] = []
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/interrupt'] = () => {
|
|
completeTurn(codex)
|
|
return {}
|
|
}
|
|
const adapter = await acquired(codex, events, {
|
|
terminateTurnProcesses: async () => false
|
|
})
|
|
|
|
await expect(
|
|
adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
).resolves.toEqual({ cancelled: false })
|
|
expect(events).toContainEqual(expect.objectContaining({ method: 'turn/completed' }))
|
|
})
|
|
|
|
it('accepts an immediate resend after verified interruption', async () => {
|
|
let nextTurn = 0
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/start'] = () => ({ turn: { id: `turn-${++nextTurn}` } })
|
|
codex.routes['turn/interrupt'] = () => {
|
|
completeTurn(codex)
|
|
return {}
|
|
}
|
|
const adapter = await acquired(codex)
|
|
|
|
await adapter.dispatch({
|
|
sessionId: 'session-1',
|
|
clientMessageId: 'client-1',
|
|
body: USER_MESSAGE,
|
|
fence: 7
|
|
})
|
|
await adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
|
|
await expect(
|
|
adapter.dispatch({
|
|
sessionId: 'session-1',
|
|
clientMessageId: 'client-2',
|
|
body: USER_MESSAGE,
|
|
fence: 7
|
|
})
|
|
).resolves.toEqual({
|
|
state: 'admitted'
|
|
})
|
|
})
|
|
|
|
it('keeps the receipt time of a completion deferred behind physical termination', async () => {
|
|
const events: CodexStructuredSessionEvent[] = []
|
|
let clock = 5_000
|
|
let finishTermination!: (terminated: boolean) => void
|
|
const termination = new Promise<boolean>((resolve) => {
|
|
finishTermination = resolve
|
|
})
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/interrupt'] = () => {
|
|
completeTurn(codex)
|
|
return {}
|
|
}
|
|
const adapter = await acquired(codex, events, {
|
|
terminateTurnProcesses: async () => termination,
|
|
now: () => clock
|
|
})
|
|
|
|
const pending = adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
await vi.waitFor(() => expect(codex.connections[0].calls.at(-1)?.method).toBe('turn/interrupt'))
|
|
clock = 9_000
|
|
finishTermination(true)
|
|
await expect(pending).resolves.toEqual({ cancelled: true })
|
|
|
|
expect(events.at(-1)).toMatchObject({ method: 'turn/completed', observedAt: 5_000 })
|
|
})
|
|
|
|
it('does not strand a deferred completion when the interrupt receipt fails', async () => {
|
|
const events: CodexStructuredSessionEvent[] = []
|
|
const codex = fakeCodex()
|
|
codex.routes['turn/interrupt'] = () => {
|
|
completeTurn(codex)
|
|
throw new Error('interrupt receipt lost')
|
|
}
|
|
const adapter = await acquired(codex, events, {
|
|
terminateTurnProcesses: async () => true
|
|
})
|
|
|
|
await expect(
|
|
adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
|
|
).rejects.toThrow('interrupt receipt lost')
|
|
expect(events).toContainEqual(expect.objectContaining({ method: 'turn/completed' }))
|
|
})
|
|
})
|