mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 00:02:29 +00:00
fix(native-chat): close stale turns and retry rejected sends
This commit is contained in:
@@ -4,7 +4,10 @@ import type {
|
||||
AgentJournalItemIdentity
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
|
||||
import { projectStructuredItemsToNativeChat } from '../../shared/structured-agent-session-projection'
|
||||
import {
|
||||
projectStructuredAgentSessionStatus,
|
||||
projectStructuredItemsToNativeChat
|
||||
} from '../../shared/structured-agent-session-projection'
|
||||
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
|
||||
import { CodexTurnOrdinals } from './codex-structured-item-translation'
|
||||
import {
|
||||
@@ -138,6 +141,74 @@ describe('codex journal translation', () => {
|
||||
expect(tap.tombstones).toEqual(['legacy:codex:session-1:turn-lifecycle%3Aturn-1'])
|
||||
})
|
||||
|
||||
it('closes every active turn when the provider session ends after a later turn starts', () => {
|
||||
const tap = recorder()
|
||||
const translator = createCodexJournalTranslator({
|
||||
sink: tap.sink,
|
||||
primaryThreadId: () => THREAD_ID
|
||||
})
|
||||
|
||||
translator.handle(notification('turn/started', { turn: { id: 'turn-stale' } }))
|
||||
translator.handle(notification('turn/started', { turn: { id: 'turn-later' } }))
|
||||
translator.handle({ type: 'ended', sessionId: SESSION_ID, reason: 'app-server exited' })
|
||||
|
||||
expect(tap.rows.filter((row) => row.body.kind === 'status')).toHaveLength(2)
|
||||
expect(tap.rows.map((row) => row.body)).toEqual([
|
||||
expect.objectContaining({ turnLifecycle: { turnId: 'turn-stale', state: 'running' } }),
|
||||
expect.objectContaining({ turnLifecycle: { turnId: 'turn-later', state: 'running' } })
|
||||
])
|
||||
expect(tap.tombstones).toEqual([
|
||||
'legacy:codex:session-1:turn-lifecycle%3Aturn-stale',
|
||||
'legacy:codex:session-1:turn-lifecycle%3Aturn-later'
|
||||
])
|
||||
// The tombstones remove both running rows from the reduced journal; no
|
||||
// lifecycle identity remains live after a session end.
|
||||
expect(
|
||||
projectStructuredAgentSessionStatus(
|
||||
tap.rows
|
||||
.filter((row) => !tap.tombstones.includes(row.key))
|
||||
.map((row, sequence) => ({
|
||||
itemId: row.key,
|
||||
revision: 1,
|
||||
sequence: sequence + 1,
|
||||
observedAt: sequence + 1,
|
||||
body: row.body
|
||||
}))
|
||||
)
|
||||
).toBe('idle')
|
||||
})
|
||||
|
||||
it('matches out-of-order completions to each turn identity', () => {
|
||||
const tap = recorder()
|
||||
const translator = createCodexJournalTranslator({
|
||||
sink: tap.sink,
|
||||
primaryThreadId: () => THREAD_ID
|
||||
})
|
||||
|
||||
translator.handle(notification('turn/started', { turn: { id: 'turn-stale' } }))
|
||||
translator.handle(notification('turn/started', { turn: { id: 'turn-later' } }))
|
||||
translator.handle(notification('turn/completed', { turn: { id: 'turn-stale' } }))
|
||||
translator.handle(notification('turn/completed', { turn: { id: 'turn-later' } }))
|
||||
|
||||
expect(tap.tombstones).toEqual([
|
||||
'legacy:codex:session-1:turn-lifecycle%3Aturn-stale',
|
||||
'legacy:codex:session-1:turn-lifecycle%3Aturn-later'
|
||||
])
|
||||
expect(
|
||||
projectStructuredAgentSessionStatus(
|
||||
tap.rows
|
||||
.filter((row) => !tap.tombstones.includes(row.key))
|
||||
.map((row, sequence) => ({
|
||||
itemId: row.key,
|
||||
revision: 1,
|
||||
sequence: sequence + 1,
|
||||
observedAt: sequence + 1,
|
||||
body: row.body
|
||||
}))
|
||||
)
|
||||
).toBe('idle')
|
||||
})
|
||||
|
||||
it('journals a user turn and the assistant answer under durable codex keys', () => {
|
||||
const { translator, tap } = translatorWith()
|
||||
|
||||
|
||||
@@ -65,17 +65,33 @@ export function createCodexJournalTranslator(
|
||||
const identities = new Map<string, AgentJournalItemIdentity>()
|
||||
/** What each announced item is, so an approval can name what it approves. */
|
||||
const details = new Map<string, string>()
|
||||
const currentTurnIds = new Map<string, string>()
|
||||
/** Turns announced by the provider and not yet closed. */
|
||||
const currentTurnIds = new Map<string, Set<string>>()
|
||||
const genericRowsByTurn = new Map<string, number>()
|
||||
const suppressedRowsByTurn = new Map<string, number>()
|
||||
let fallbackSequence = 0
|
||||
|
||||
const currentTurnIdFor = (threadId: string): string | null =>
|
||||
[...(currentTurnIds.get(threadId) ?? [])].at(-1) ?? null
|
||||
|
||||
const rememberTurn = (threadId: string, turnId: string): void => {
|
||||
currentTurnIds.set(threadId, new Set([...(currentTurnIds.get(threadId) ?? []), turnId]))
|
||||
}
|
||||
|
||||
const forgetTurn = (threadId: string, turnId: string): void => {
|
||||
const active = currentTurnIds.get(threadId)
|
||||
active?.delete(turnId)
|
||||
if (!active?.size) {
|
||||
currentTurnIds.delete(threadId)
|
||||
}
|
||||
}
|
||||
|
||||
const appendUnhandled = (kind: string, payload: unknown, threadId = 'session'): void => {
|
||||
const translated = unhandledProviderFrameJournalItem('codex', kind, payload)
|
||||
if (!translated) {
|
||||
return
|
||||
}
|
||||
const turnId = readCodexTurnId(payload) ?? currentTurnIds.get(threadId) ?? 'outside-turn'
|
||||
const turnId = readCodexTurnId(payload) ?? currentTurnIdFor(threadId) ?? 'outside-turn'
|
||||
const bucket = `${encodeURIComponent(threadId)}:${encodeURIComponent(turnId)}`
|
||||
const rowCount = genericRowsByTurn.get(bucket) ?? 0
|
||||
// The cap bounds noise, never evidence: an error frame is always journaled,
|
||||
@@ -152,7 +168,7 @@ export function createCodexJournalTranslator(
|
||||
coalesceMs: deps.coalesceMs,
|
||||
schedule: deps.schedule,
|
||||
identityFor: (threadId, params, item) => {
|
||||
const turnId = readCodexTurnId(params) ?? currentTurnIds.get(threadId) ?? null
|
||||
const turnId = readCodexTurnId(params) ?? currentTurnIdFor(threadId)
|
||||
return identityFor(threadId, turnId, item)
|
||||
}
|
||||
})
|
||||
@@ -167,7 +183,7 @@ export function createCodexJournalTranslator(
|
||||
if (!item) {
|
||||
return false
|
||||
}
|
||||
const turnId = readCodexTurnId(event.params) ?? currentTurnIds.get(event.threadId) ?? null
|
||||
const turnId = readCodexTurnId(event.params) ?? currentTurnIdFor(event.threadId)
|
||||
const identity = identityFor(event.threadId, turnId, item)
|
||||
const translated = codexJournalItem(item)
|
||||
const command = readString(item, 'command')
|
||||
@@ -238,7 +254,7 @@ export function createCodexJournalTranslator(
|
||||
if (!turnId) {
|
||||
continue
|
||||
}
|
||||
currentTurnIds.set(threadId, turnId)
|
||||
currentTurnIds.set(threadId, new Set([turnId]))
|
||||
for (const item of Array.isArray(turn.items) ? turn.items : []) {
|
||||
handleItemEvent({ threadId, method: 'item/completed', params: { turnId, item } })
|
||||
}
|
||||
@@ -250,8 +266,11 @@ export function createCodexJournalTranslator(
|
||||
handle: (event) => {
|
||||
if (event.type === 'ended') {
|
||||
streams.flush()
|
||||
for (const [threadId, turnId] of currentTurnIds) {
|
||||
publishTurnLifecycle(event.sessionId, threadId, turnId, 'completed')
|
||||
for (const [threadId, turnIds] of currentTurnIds) {
|
||||
for (const turnId of turnIds) {
|
||||
publishTurnLifecycle(event.sessionId, threadId, turnId, 'completed')
|
||||
ordinals.forgetTurn(threadId, turnId)
|
||||
}
|
||||
}
|
||||
currentTurnIds.clear()
|
||||
return
|
||||
@@ -279,20 +298,20 @@ export function createCodexJournalTranslator(
|
||||
if (event.method === 'turn/started') {
|
||||
const turnId = readCodexTurnId(event.params)
|
||||
if (turnId) {
|
||||
currentTurnIds.set(event.threadId, turnId)
|
||||
rememberTurn(event.threadId, turnId)
|
||||
publishTurnLifecycle(event.sessionId, event.threadId, turnId, 'running')
|
||||
}
|
||||
return
|
||||
}
|
||||
if (event.method === 'turn/completed') {
|
||||
const turnId = readCodexTurnId(event.params) ?? currentTurnIds.get(event.threadId)
|
||||
const turnId = readCodexTurnId(event.params) ?? currentTurnIdFor(event.threadId)
|
||||
if (turnId) {
|
||||
publishTurnLifecycle(event.sessionId, event.threadId, turnId, 'completed')
|
||||
ordinals.forgetTurn(event.threadId, turnId)
|
||||
forgetTurn(event.threadId, turnId)
|
||||
}
|
||||
// A later item with no turn of its own belongs to no turn, not to the
|
||||
// one that just ended.
|
||||
currentTurnIds.delete(event.threadId)
|
||||
// A later item without its own turn id falls back to another active
|
||||
// turn, if one exists; completed turns are never adopted again.
|
||||
return
|
||||
}
|
||||
if (event.method === 'item/started' || event.method === 'item/completed') {
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
import { act, renderHook, waitFor } from '@testing-library/react'
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentJournalSubmission } from '../../../../shared/agent-session-journal-types'
|
||||
import type { AgentSessionWireRefusalCode } from '../../../../shared/agent-session-wire'
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
@@ -46,6 +47,50 @@ function acceptedResult(fence: number) {
|
||||
}
|
||||
}
|
||||
|
||||
function acceptedResultFor(clientMessageId: string, fence: number) {
|
||||
return {
|
||||
ok: true,
|
||||
replayed: false,
|
||||
fence,
|
||||
cursor: { epoch: 'epoch-1', sequence: fence },
|
||||
value: {
|
||||
clientMessageId,
|
||||
submission: {
|
||||
clientMessageId,
|
||||
fence,
|
||||
payloadFingerprint: 'fingerprint',
|
||||
dispatchState: 'accepted',
|
||||
providerItemId: `provider-${clientMessageId}`,
|
||||
reason: null,
|
||||
submittedAt: fence,
|
||||
resolvedAt: fence
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function unknownResultFor(clientMessageId: string, submittedAt: number) {
|
||||
return {
|
||||
ok: true,
|
||||
replayed: false,
|
||||
fence: 1,
|
||||
cursor: { epoch: 'epoch-1', sequence: submittedAt },
|
||||
value: {
|
||||
clientMessageId,
|
||||
submission: {
|
||||
clientMessageId,
|
||||
fence: 1,
|
||||
payloadFingerprint: 'fingerprint',
|
||||
dispatchState: 'unknown' as const,
|
||||
providerItemId: null,
|
||||
reason: 'socket closed',
|
||||
submittedAt,
|
||||
resolvedAt: submittedAt
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function refusedResult(code: AgentSessionWireRefusalCode) {
|
||||
return { ok: false, refusal: { code, message: code } }
|
||||
}
|
||||
@@ -175,4 +220,126 @@ describe('useStructuredAgentSessionOutbox', () => {
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
it('retries an unknown head and advances a queued tail', async () => {
|
||||
vi.mocked(globalThis.crypto.randomUUID)
|
||||
.mockReturnValueOnce('aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa')
|
||||
.mockReturnValueOnce('bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb')
|
||||
mocks.call
|
||||
.mockImplementationOnce(async (_target, _method, params) => {
|
||||
const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope
|
||||
.clientOperationId
|
||||
return unknownResultFor(clientMessageId, 10)
|
||||
})
|
||||
.mockImplementationOnce(async (_target, _method, params) => {
|
||||
const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope
|
||||
.clientOperationId
|
||||
return acceptedResultFor(clientMessageId, 11)
|
||||
})
|
||||
.mockImplementationOnce(async (_target, _method, params) => {
|
||||
const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope
|
||||
.clientOperationId
|
||||
return acceptedResultFor(clientMessageId, 12)
|
||||
})
|
||||
const { result, rerender } = renderHook(
|
||||
({ submissions }: { submissions: readonly AgentJournalSubmission[] }) =>
|
||||
useStructuredAgentSessionOutbox({
|
||||
sessionId: 'session-1',
|
||||
target: LOCAL_TARGET,
|
||||
fence: 1,
|
||||
submissions
|
||||
}),
|
||||
{ initialProps: { submissions: [] as readonly AgentJournalSubmission[] } }
|
||||
)
|
||||
|
||||
act(() => {
|
||||
expect(result.current.send('first')).toBe(true)
|
||||
})
|
||||
await waitFor(() => expect(result.current.outbox[0]?.state).toBe('unconfirmed'))
|
||||
const firstId = result.current.outbox[0]!.clientMessageId
|
||||
rerender({
|
||||
submissions: [
|
||||
{
|
||||
clientMessageId: firstId,
|
||||
fence: 1,
|
||||
payloadFingerprint: 'fingerprint',
|
||||
dispatchState: 'unknown',
|
||||
providerItemId: null,
|
||||
reason: 'socket closed',
|
||||
submittedAt: 10,
|
||||
resolvedAt: 10
|
||||
}
|
||||
]
|
||||
})
|
||||
act(() => {
|
||||
expect(result.current.send('second')).toBe(true)
|
||||
})
|
||||
expect(result.current.outbox).toHaveLength(2)
|
||||
|
||||
act(() => result.current.retry(firstId))
|
||||
await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(3))
|
||||
await waitFor(() => expect(result.current.outbox).toHaveLength(0))
|
||||
const retryParams = mocks.call.mock.calls[1]?.[2] as { retryUnknown?: true } | undefined
|
||||
expect(retryParams?.retryUnknown).toBe(true)
|
||||
})
|
||||
|
||||
it('rotates a history-rejected unknown head so the queued tail can advance', async () => {
|
||||
vi.mocked(globalThis.crypto.randomUUID)
|
||||
.mockReturnValueOnce('aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa')
|
||||
.mockReturnValueOnce('bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb')
|
||||
.mockReturnValueOnce('cccccccc-cccc-4ccc-8ccc-cccccccccccc')
|
||||
mocks.call
|
||||
.mockImplementationOnce(async (_target, _method, params) => {
|
||||
const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope
|
||||
.clientOperationId
|
||||
return unknownResultFor(clientMessageId, 10)
|
||||
})
|
||||
.mockImplementationOnce(async (_target, _method, params) => {
|
||||
const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope
|
||||
.clientOperationId
|
||||
return acceptedResultFor(clientMessageId, 11)
|
||||
})
|
||||
.mockImplementationOnce(async (_target, _method, params) => {
|
||||
const clientMessageId = (params as { envelope: { clientOperationId: string } }).envelope
|
||||
.clientOperationId
|
||||
return acceptedResultFor(clientMessageId, 12)
|
||||
})
|
||||
const { result, rerender } = renderHook(
|
||||
({ submissions }: { submissions: readonly AgentJournalSubmission[] }) =>
|
||||
useStructuredAgentSessionOutbox({
|
||||
sessionId: 'session-1',
|
||||
target: LOCAL_TARGET,
|
||||
fence: 1,
|
||||
submissions
|
||||
}),
|
||||
{ initialProps: { submissions: [] as readonly AgentJournalSubmission[] } }
|
||||
)
|
||||
|
||||
act(() => expect(result.current.send('first')).toBe(true))
|
||||
await waitFor(() => expect(result.current.outbox[0]?.state).toBe('unconfirmed'))
|
||||
const firstId = result.current.outbox[0]!.clientMessageId
|
||||
act(() => expect(result.current.send('second')).toBe(true))
|
||||
rerender({
|
||||
submissions: [
|
||||
{
|
||||
clientMessageId: firstId,
|
||||
fence: 1,
|
||||
payloadFingerprint: 'fingerprint',
|
||||
dispatchState: 'rejected',
|
||||
providerItemId: null,
|
||||
reason: 'not_delivered',
|
||||
submittedAt: 10,
|
||||
resolvedAt: 10
|
||||
}
|
||||
]
|
||||
})
|
||||
|
||||
act(() => result.current.retry(firstId))
|
||||
await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(3))
|
||||
await waitFor(() => expect(result.current.outbox).toHaveLength(0))
|
||||
const retryParams = mocks.call.mock.calls[1]?.[2] as
|
||||
| { envelope: { clientOperationId: string } }
|
||||
| undefined
|
||||
expect(retryParams?.envelope.clientOperationId).not.toBe(firstId)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -247,13 +247,38 @@ export function useStructuredAgentSessionOutbox(args: {
|
||||
const retry = (clientMessageId: string): void => {
|
||||
blockedIdRef.current = null
|
||||
setError(null)
|
||||
const unknown = submissions.find(
|
||||
(submission) =>
|
||||
submission.clientMessageId === clientMessageId && submission.dispatchState === 'unknown'
|
||||
const submission = submissions.find(
|
||||
(candidate) => candidate.clientMessageId === clientMessageId
|
||||
)
|
||||
const current = outboxRef.current.find((entry) => entry.clientMessageId === clientMessageId)
|
||||
// A provider-history reconciliation can settle an earlier unknown as
|
||||
// rejected before the user presses Retry. Reusing that operation id only
|
||||
// replays the settled rejection forever, so rotate the id for a safe resend.
|
||||
if (current && submission?.dispatchState === 'rejected') {
|
||||
const rotated = outboxRef.current.map((entry) =>
|
||||
entry.clientMessageId === clientMessageId
|
||||
? {
|
||||
...entry,
|
||||
clientMessageId: structuredSessionOperationId(),
|
||||
state: 'queued' as const,
|
||||
retryAfterUnknownSubmittedAt: null
|
||||
}
|
||||
: entry
|
||||
)
|
||||
if (!writeOutbox(sessionId, rotated)) {
|
||||
setError('Message could not be saved to the outbox')
|
||||
return
|
||||
}
|
||||
outboxRef.current = rotated
|
||||
setOutbox(rotated)
|
||||
return
|
||||
}
|
||||
const retryAfterUnknownSubmittedAt =
|
||||
unknown?.submittedAt ?? (current?.state === 'unconfirmed' ? -1 : null)
|
||||
submission?.dispatchState === 'unknown'
|
||||
? submission.submittedAt
|
||||
: current?.state === 'unconfirmed'
|
||||
? -1
|
||||
: null
|
||||
const next = outboxRef.current.map((entry) =>
|
||||
entry.clientMessageId === clientMessageId
|
||||
? {
|
||||
|
||||
Reference in New Issue
Block a user