fix(native-chat): a message the chat said was not sent is never sent later on its own (#24232)

* fix(native-chat): keep a message the chat said was not sent held until its Retry

A native-chat send the host refused (for example "Chats were saved by a newer
Orca. Your message was not sent.") or that never reached the host showed "not
sent" with a Retry button, but the hold that stopped it lived only in the
outbox hook's memory. The message itself was saved in the outbox, so the next
launch lifted the hold and sent it with no Retry; a copy the user retyped in
the meantime was held behind it and went out as well.

The hold is now read from the failure the message already saves: a queued
entry that carries its last failure waits for the user's Retry, on this launch
and every later one, and an entry saved by an earlier build in that shape is
held too. The drain passes over a held message instead of stopping behind it,
so what the user sends next goes out as they send it. Retry clears the saved
failure. On a host from before accepted-send, a new agent owner observed while
the chat is open still sends the refused message again, as before; a relaunch
does not.

A held message keeps its operation id, and a host refuses an id older than a
day as expired for good, so its Retry could never go through; a new id could
deliver a message an earlier attempt already delivered. Such a message now
goes back to the composer with a notice to check the chat before sending it
again, and leaves the outbox.

* fix(native-chat): keep refused messages as rows until Retry, and only release them for an older host's new owner

- A message the host refuses as expired under an id it kept stays a saved row
  reading "Orca couldn't confirm what happened. Check the chat.", and its Retry
  sends it under a new id. It no longer moves into the message box, where an
  automatic resend after a relaunch could put text the user never asked for,
  held only in memory.
- An owner change releases a refused message only on a host known to predate
  accepted sends, and only for the refusals such a host gives while it restarts
  the chat's agent. Those rows say Orca will send it again when the agent
  restarts, beside their Retry. A host whose capability check has not answered,
  or failed, no longer releases anything.
- Every failed message ahead of the one the queue stopped on keeps its Retry,
  since that Retry sends it at once.
- A journal row saying the host cannot tell whether a message landed, and the
  unconfirmed probe's resend, replace an earlier attempt's saved failure, so
  the message is probed rather than held.
- The drain stages from the hook's own outbox, so a hold kept only in memory
  after a failed save survives the next send; the hold is written once more
  after that failed save.

* fix(native-chat): a refused message waits for its Retry on every host, and a send is staged from the latest outbox

A message the chat showed as not sent no longer goes out on its own when an
older host's chat gets a new agent owner. Resending it on the owner change
sent it after messages typed later, still delivered a retyped copy twice,
and its "Orca will send it again when the agent restarts" row promised a
resend that often never came. It now waits for the user's Retry, as it does
on every current host. A send still in flight when the owner changes is
still sent again under its id; it was never shown as failed.

The drain admitted and staged the next send from the render's outbox. An
owner change requeues the send it interrupted in an effect earlier in the
same commit, and staging from the render's list wrote the old list back,
leaving that send stuck as sending. The drain now reads the latest list,
which still carries a hold kept only in memory after a failed save.

Tests pass the view's target as one stable object, as the view does: a new
object each render re-ran the owner-change requeue, which hid the drain bug.

* fix(native-chat): keep a not-sent message out of newer turns, and word its saved cause only when seen

A message shown as not sent stays in the outbox and draws below every
newer turn. It was an ordinary user row there, so it counted as the newest
user row: while a new send waited for its turn to open, that turn's
"Working for" clock drew under the old message, and once the turn ended an
empty "Worked for" divider was left under it. The projection now marks
such a bubble (held for its Retry, or rejected) as unsent; turn membership
gives it no turn and never makes it the live one, on hosts that state turn
scopes and on those that do not; and the transcript draws it after the
live activity, as it draws a message waiting behind /compact. Mobile has no
outbox, so its rows never carry the mark and its grouping is unchanged.

A held message read back from storage repeated the cause it was saved with,
which may no longer hold: "Update Orca to keep using them" after the user
updated Orca. The outbox hook now remembers, in memory only, which messages
failed while the chat was open; only those word their cause. Any other held
message reads "Your message was not sent." with its Retry, and a Retry the
cause still stops brings the full words back. An expired id keeps its words,
since that cause cannot clear.

* fix(native-chat): follow the bottom and light a tick for a chat whose only rows are not sent

A message shown as not sent draws after the windowed transcript. When it was
the only row, the windowed list was empty, and following the bottom or "Jump
to latest" asked the virtualizer for an end it computes from its own rows:
the top. It now scrolls to the container's own bottom when no row is
windowed.

A rejected send the journal recorded keeps its place but opens no turn, so
the rail lit no tick when it was the row being read. A user row in no turn
now lights its own tick.

Also pins that a refusal seen while the chat is open reaches the rendered
notice in full, and reads only "not sent" after the chat is reopened until a
Retry is refused again.

* test(native-chat): name the relaunch test parameter for how the refusal arrives

The low-evidence lint rejects "shape" as a symbol name.
This commit is contained in:
Brennan Benson
2026-09-30 22:49:26 -07:00
committed by GitHub
parent be575800c2
commit 07e9fdfd13
55 changed files with 1816 additions and 379 deletions
@@ -0,0 +1,216 @@
// @vitest-environment happy-dom
// A message shown as not sent stays in the outbox and draws below newer turns, after their live
// activity. It is in no turn: the newer turn's "Working for" clock and "Worked for" bar stay with
// that turn, never under it.
import '@testing-library/jest-dom/vitest'
import { cleanup, render } from '@testing-library/react'
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'
import type {
AgentJournalItemBody,
AgentJournalRenderItem,
AgentJournalSubmission,
AgentJournalTurnScope
} from '../../../../shared/agent-session-journal-types'
import { agentJournalSubmissionKey } from '../../../../shared/agent-session-journal-item-key'
import type { NativeChatSettledTurns } from '../../../../shared/native-chat-turn-status'
import {
createStructuredAgentSessionOutboxEntry,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
import { NativeChatMessageList } from './NativeChatMessageList'
import { installNativeChatMessageListTestViewport } from './native-chat-message-list-test-viewport'
import { projectStructuredAgentSessionMessages } from './structured-agent-session-message-projection'
let restoreViewport = (): void => {}
beforeAll(() => {
restoreViewport = installNativeChatMessageListTestViewport()
})
afterAll(() => restoreViewport())
afterEach(cleanup)
type Phase = 'in flight' | 'running' | 'done'
const SEED = agentJournalSubmissionKey('seed')
const NEW = agentJournalSubmissionKey('new')
function said(role: 'user' | 'assistant', text: string): AgentJournalItemBody {
return { kind: 'message', role, blocks: [{ type: 'text', text }] }
}
function journal(phase: Phase, scoped: boolean): AgentJournalRenderItem[] {
const thread: AgentJournalTurnScope = { kind: 'thread' }
const row = (
itemId: string,
body: AgentJournalItemBody,
turnScope: AgentJournalTurnScope
): AgentJournalRenderItem => ({
itemId,
revision: 0,
sequence: rows.length + 1,
observedAt: 1000 + rows.length,
body,
...(scoped ? { turnScope } : {})
})
const rows: AgentJournalRenderItem[] = []
rows.push(row(SEED, said('user', 'SEED PROMPT'), thread))
rows.push(
row(
'turn-seed',
{ kind: 'turn', turnId: 'turn-seed', state: 'completed', userItemId: SEED },
thread
)
)
rows.push(
row('seed-answer', said('assistant', 'SEED OK'), { kind: 'turn', turnItemId: 'turn-seed' })
)
rows.push(row(NEW, said('user', 'NEW PROMPT'), thread))
if (phase !== 'in flight') {
rows.push(
row(
'turn-new',
{
kind: 'turn',
turnId: 'turn-new',
state: phase === 'done' ? 'completed' : 'running',
userItemId: NEW
},
thread
)
)
}
if (phase === 'done') {
rows.push(
row('new-answer', said('assistant', 'NEW OK'), { kind: 'turn', turnItemId: 'turn-new' })
)
}
return rows
}
function submission(
id: string,
dispatchState: AgentJournalSubmission['dispatchState']
): AgentJournalSubmission {
return {
clientMessageId: id,
fence: 1,
payloadFingerprint: id,
dispatchState,
providerItemId: null,
reason: null,
submittedAt: 1,
resolvedAt: null
}
}
function unsentEntry(kind: 'held' | 'rejected'): StructuredAgentSessionOutboxEntry {
const entry = createStructuredAgentSessionOutboxEntry({
clientMessageId: 'held',
sessionId: 'session-1',
text: 'HELD PROMPT',
attachments: [],
queuedAt: 500
})
return kind === 'held'
? {
...entry,
lastAttemptAt: 500,
lastFailure: { kind: 'refused', code: 'agent_session_journal_unreadable' }
}
: { ...entry, state: 'rejected', lastFailure: { kind: 'rejected', reason: null } }
}
function list(phase: Phase, scoped: boolean, outbox: StructuredAgentSessionOutboxEntry[]) {
const items = journal(phase, scoped)
const submissions = [
submission('seed', 'accepted'),
submission('new', phase === 'done' ? 'accepted' : 'pending')
]
const settledTurns: NativeChatSettledTurns = new Map([
[SEED, { startedAt: 1, workedSeconds: 3 }],
...(phase === 'done' ? [[NEW, { startedAt: 2, workedSeconds: 5 }] as const] : [])
])
return (
<NativeChatMessageList
session={{
messages: projectStructuredAgentSessionMessages(items, outbox, submissions),
status: phase === 'done' ? 'ready' : 'working',
sessionId: 'session-1',
agent: 'claude',
hasMore: false,
loadingEarlier: false,
olderHistoryGeneration: 0,
loadEarlier: vi.fn(),
readPhase: 'ready'
}}
journalItems={items}
journalSubmissions={submissions}
isWorking={phase !== 'done'}
workingStartedAt={phase === 'done' ? null : Date.now() - 1500}
settledTurns={settledTurns}
expandSignal={false}
fontScale={1}
/>
)
}
/** The drawn sequence of the prompts, bars and live activity line, top to bottom. */
function drawn(container: HTMLElement): string[] {
const out: string[] = []
for (const element of container.querySelectorAll('*')) {
if (element.hasAttribute('data-native-chat-turn-activity')) {
out.push('ACTIVITY')
} else if (element.hasAttribute('data-native-chat-turn-status')) {
out.push(
element.getAttribute('data-native-chat-turn-status') === 'active' ? 'WORKING' : 'WORKED'
)
} else if (element.childElementCount === 0 && (element.textContent ?? '').endsWith('PROMPT')) {
out.push(element.textContent ?? '')
}
}
return out
}
describe.each([
['held for its Retry', 'held'],
['rejected', 'rejected']
] as const)('a message %s, below a newer turn', (_label, kind) => {
it.each([
['states each row turn', true],
['states no turn scope', false]
])(
'keeps the newer turn bar with that turn while it runs and once done (host %s)',
(_host, scoped) => {
const outbox = [unsentEntry(kind)]
const { container, rerender } = render(list('in flight', scoped, outbox))
expect(drawn(container)).toEqual([
'SEED PROMPT',
'WORKED',
'NEW PROMPT',
'WORKING',
'ACTIVITY',
'HELD PROMPT'
])
rerender(list('running', scoped, outbox))
expect(drawn(container)).toEqual([
'SEED PROMPT',
'WORKED',
'NEW PROMPT',
'WORKING',
'ACTIVITY',
'HELD PROMPT'
])
rerender(list('done', scoped, outbox))
expect(drawn(container)).toEqual([
'SEED PROMPT',
'WORKED',
'NEW PROMPT',
'WORKED',
'HELD PROMPT'
])
expect(container.textContent).toContain('Worked for 5s')
}
)
})
@@ -174,8 +174,8 @@ export function createStructuredSessionMocks() {
loadOlder: mocks.loadOlder,
prompts: mocks.promptItems,
outbox: outbox.outbox,
failedHere: outbox.failedHere,
submissions: mocks.submissions,
blockedClientMessageId: outbox.blockedClientMessageId,
send: outbox.send,
retry: outbox.retry,
isWorking: mocks.isWorking,
@@ -131,19 +131,19 @@ export function NativeChatStructuredSession(
() =>
structuredAgentSessionDeliveryNotices(
controller.outbox,
controller.blockedClientMessageId,
agentLabel,
retryDelivery,
rejectionRows,
startFailures
startFailures,
controller.failedHere
),
[
controller.outbox,
controller.blockedClientMessageId,
agentLabel,
retryDelivery,
rejectionRows,
startFailures
startFailures,
controller.failedHere
]
)
const viewState = selectNativeChatViewState(session, { readRetries: true })
@@ -79,8 +79,8 @@ vi.mock('./use-structured-agent-session', async () => {
loadOlder: vi.fn(),
prompts: mocks.promptItems,
outbox: outbox.outbox,
failedHere: outbox.failedHere,
submissions: mocks.submissions,
blockedClientMessageId: outbox.blockedClientMessageId,
send: outbox.send,
retry: outbox.retry,
isWorking: false,
@@ -183,6 +183,7 @@ vi.mock('./NativeChatQuestionCard', () => ({
}))
import { NativeChatStructuredSession } from './NativeChatStructuredSession'
import { appendStructuredAgentSessionOutboxMessage } from './structured-agent-session-outbox-storage'
describe('NativeChatStructuredSession delivery', () => {
afterEach(() => {
@@ -768,4 +769,44 @@ describe('NativeChatStructuredSession delivery', () => {
vi.useRealTimers()
}
}, 30000)
// A cause seen while the chat was open is worded in full; one read back after the chat is
// reopened may have cleared, until a Retry it still stops brings it back.
it('words a refusal seen here in full, and after a reopen only once its Retry is refused', async () => {
const newerOrca =
'Chats were saved by a newer Orca. Your message was not sent. Update Orca to keep using them.'
mocks.mode = 'outbox'
mocks.call.mockResolvedValue({
ok: false,
refusal: {
code: 'agent_session_journal_unreadable',
message: 'newer',
details: { reason: 'journalWrittenByNewerOrca' }
}
})
const view = (): React.JSX.Element => (
<NativeChatStructuredSession
isVisible
isFocusedGroup
tabId="tab-held"
sessionId="session-held"
target={{ kind: 'local' }}
agent="codex"
/>
)
const first = render(view())
// The composer's own enqueue.
act(() => {
expect(appendStructuredAgentSessionOutboxMessage('session-held', 'hello')).not.toBeNull()
})
await waitFor(() => expect(screen.getByText(newerOrca)).toBeTruthy())
first.unmount()
render(view())
await waitFor(() => expect(screen.getByText('Your message was not sent.')).toBeTruthy())
expect(mocks.call).toHaveBeenCalledOnce()
fireEvent.click(screen.getByRole('button', { name: /Retry/ }))
await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(2))
await waitFor(() => expect(screen.getByText(newerOrca)).toBeTruthy())
})
})
@@ -122,4 +122,31 @@ describe('active rail item', () => {
})
).toBe('u2')
})
// A rejected send the journal recorded keeps its place but opens no turn.
it('lights a prompt in no turn by its own id', () => {
const slots: NativeChatRailSlot[] = [
...TURNS,
{ turnKey: undefined, message: { id: 'rejected', role: 'user' } }
]
expect(
findActiveNativeChatRailItem({
slots,
virtualItems: rows(11),
scrollTop: 800,
clientHeight: VIEWPORT,
scrollHeight: 1100,
previousActiveId: null
})
).toBe('rejected')
expect(
findActiveNativeChatRailItem({
slots,
virtualItems: rows(11),
scrollTop: 1010,
...MID_SCROLL,
scrollHeight: 2000
})
).toBe('rejected')
})
})
@@ -24,9 +24,15 @@ export type NativeChatRailVirtualItem = {
end: number
}
/** Only the field the rail reads, so a test needs no slot builder. */
/** Only the fields the rail reads, so a test needs no slot builder. */
export type NativeChatRailSlot = {
turnKey: string | undefined
/** A user row in no turn (one shown as not sent) lights its own tick. */
message?: { id: string; role: string }
}
function railTickOf(slot: NativeChatRailSlot | undefined): string | null {
return slot?.turnKey ?? (slot?.message?.role === 'user' ? slot.message.id : null)
}
export function findActiveNativeChatRailItem({
@@ -53,7 +59,7 @@ export function findActiveNativeChatRailItem({
const atBottom = scrollHeight - clientHeight - scrollTop <= NATIVE_CHAT_BOTTOM_THRESHOLD_PX
if (atBottom) {
const last = virtualItems.at(-1)
return last === undefined ? previousActiveId : (slots[last.index]?.turnKey ?? null)
return last === undefined ? previousActiveId : railTickOf(slots[last.index])
}
let fold: NativeChatRailVirtualItem | undefined
@@ -65,12 +71,12 @@ export function findActiveNativeChatRailItem({
// Scrolled above everything the window holds: the first windowed row is the
// nearest thing to the fold.
if (fold === undefined) {
return slots[virtualItems[0]?.index ?? -1]?.turnKey ?? null
return railTickOf(slots[virtualItems[0]?.index ?? -1])
}
// The window lags the scroll by a commit, so a fold past every row it holds is
// a stale read, not an answer. Holding the previous tick beats blanking one.
if (fold.end <= scrollTop) {
return previousActiveId
}
return slots[fold.index]?.turnKey ?? null
return railTickOf(slots[fold.index])
}
@@ -242,8 +242,8 @@ export function nativeChatSlotIndexOf(
return slots.findIndex((slot) => slot.kind === 'message' && slot.message.id === messageId)
}
/** Splits off the slots of messages waiting behind the live turn: they draw after its live
* activity, not inside it. */
/** Splits off the slots of messages waiting behind the live turn, and of ones shown as not sent
* that the journal holds no place for: they draw after the live activity, not inside it. */
export function splitNativeChatSlotsWaitingBehindLiveTurn(
slots: readonly NativeChatTranscriptSlot[],
journalItems: readonly AgentJournalRenderItem[] | undefined
@@ -253,7 +253,9 @@ export function splitNativeChatSlotsWaitingBehindLiveTurn(
journalItems
)
const isWaiting = (slot: NativeChatTranscriptSlot): boolean =>
slot.kind === 'message' && waiting.has(slot.message.id)
slot.kind === 'message' &&
(waiting.has(slot.message.id) ||
(slot.message.unsent === true && slot.message.journalPosition === undefined))
return {
slots: slots.filter((slot) => !isWaiting(slot)),
waitingSlots: slots.filter(isWaiting)
@@ -36,19 +36,21 @@ function entry(
}
}
const NOT_FAILED_HERE: ReadonlySet<string> = new Set()
function texts(
outbox: StructuredAgentSessionOutboxEntry[],
blocked: string | null = null,
submissions: readonly AgentJournalSubmission[] = [],
startFailures: readonly AgentSessionFailureFact[] = []
): Record<string, string> {
// Every failure seen while the chat was open, so each words its whole cause.
const notices = structuredAgentSessionDeliveryNotices(
outbox,
blocked,
'Claude',
() => {},
submissions,
startFailures
startFailures,
new Set(outbox.map((candidate) => candidate.clientMessageId))
)
return Object.fromEntries([...notices].map(([id, notice]) => [id, notice.text]))
}
@@ -71,11 +73,11 @@ describe('the notice on each message that did not go through', () => {
}
})
],
null,
'Claude',
retry,
[],
[]
[],
NOT_FAILED_HERE
)
expect([...notices.keys()]).toEqual([
@@ -92,25 +94,22 @@ describe('the notice on each message that did not go through', () => {
expect(retry).toHaveBeenCalledExactlyOnceWith('second')
})
it('chooses the words from the saved refusal on the message the queue stopped on', () => {
it('chooses the words from the saved refusal on a refused message', () => {
expect(
texts(
[
entry('held', {
lastFailure: { kind: 'refused', code: 'agent_session_owner_restart_failed' }
})
],
'held'
)
texts([
entry('held', {
lastFailure: { kind: 'refused', code: 'agent_session_owner_restart_failed' }
})
])
).toEqual({
[agentJournalSubmissionKey('held')]: "The agent couldn't restart. Your message was not sent."
})
})
// The same rule as a rejected row's: its own Retry is the resend step, and any other step stays.
it('leaves a retry step to the Retry beside the message the queue stopped on', () => {
it('leaves a retry step to the Retry beside a refused message', () => {
const held = (lastFailure: StructuredAgentSessionOutboxEntry['lastFailure']) =>
texts([entry('held', { lastFailure })], 'held')
texts([entry('held', { lastFailure })])
expect(
held({
kind: 'refused',
@@ -137,53 +136,55 @@ describe('the notice on each message that did not go through', () => {
expect(texts([entry('doubt', { state: 'unconfirmed' })])).toEqual({
[agentJournalSubmissionKey('doubt')]: 'Message delivery is unconfirmed.'
})
expect(texts([entry('bare')], 'bare')).toEqual({
expect(texts([entry('bare', { state: 'rejected' })])).toEqual({
[agentJournalSubmissionKey('bare')]: 'Message was not sent.'
})
})
it('never says a send attempted before a Stop was not sent: the host may hold it', () => {
const interrupted = entry('stopped', { state: 'queued', lastAttemptAt: 5, outlivedStop: true })
expect(texts([interrupted], 'stopped')).toEqual({
expect(texts([interrupted])).toEqual({
[agentJournalSubmissionKey('stopped')]: 'Message delivery is unconfirmed.'
})
const neverSent = entry('unsent', { state: 'queued', outlivedStop: true })
expect(texts([neverSent], 'unsent')).toEqual({
expect(texts([neverSent])).toEqual({
[agentJournalSubmissionKey('unsent')]: 'Message was not sent.'
})
})
// The drain's own rule: a message behind the one the queue stopped on is only waiting, so it says
// nothing. A rejected message holds nothing up and keeps its words.
it('says why on the message the queue stopped on and on every rejected one', () => {
// nothing. A rejected or refused message holds nothing up and keeps its words.
it('says why on the message the queue stopped on and on every rejected or refused one', () => {
expect(
texts([
entry('sent', { state: 'dispatching' }),
entry('rejected', { state: 'rejected' }),
entry('failed', { lastFailure: { kind: 'failed' } }),
entry('stuck', { state: 'unconfirmed' }),
entry('behind', { state: 'unconfirmed' }),
entry('queued')
])
).toEqual({
[agentJournalSubmissionKey('rejected')]: 'Message was not sent.',
[agentJournalSubmissionKey('failed')]: 'Your message was not sent.',
[agentJournalSubmissionKey('stuck')]: 'Message delivery is unconfirmed.'
})
})
// Its Retry would put it back in the queue to wait unseen behind the stopped message.
it('keeps a rejected message its words but not its Retry while the queue is stopped', () => {
it('keeps a rejected message behind the stopped one its words but not its Retry', () => {
const retry = vi.fn()
for (const [outbox, blocked] of [
[[entry('stuck', { state: 'unconfirmed' }), entry('rejected', { state: 'rejected' })], null],
[[entry('rejected', { state: 'rejected' }), entry('held')], 'held']
] as const) {
for (const outbox of [
[entry('stuck', { state: 'unconfirmed' }), entry('rejected', { state: 'rejected' })],
[entry('held', { outlivedStop: true }), entry('rejected', { state: 'rejected' })]
]) {
const notices = structuredAgentSessionDeliveryNotices(
[...outbox],
blocked,
outbox,
'Claude',
retry,
[],
[]
[],
NOT_FAILED_HERE
)
expect(notices.get(agentJournalSubmissionKey('rejected'))).toEqual({
text: 'Message was not sent.'
@@ -191,6 +192,37 @@ describe('the notice on each message that did not go through', () => {
}
})
// Ahead of the stopped message, its Retry sends it at once, so the row offers it.
it.each([
['in doubt', { state: 'unconfirmed' as const }],
['outlived by a Stop', { outlivedStop: true as const }]
])('gives a failed message ahead of one %s its Retry', (_label, patch) => {
const retry = vi.fn()
const notices = structuredAgentSessionDeliveryNotices(
[
entry('refused', {
lastAttemptAt: 1,
lastFailure: {
kind: 'refused',
code: 'agent_session_journal_unreadable',
details: { reason: 'journalWrittenByNewerOrca' }
}
}),
entry('rejected', { state: 'rejected' }),
entry('stuck', { lastAttemptAt: 2, ...patch })
],
'Claude',
retry,
[],
[],
NOT_FAILED_HERE
)
for (const id of ['refused', 'rejected', 'stuck']) {
notices.get(agentJournalSubmissionKey(id))?.onRetry?.()
}
expect(retry.mock.calls).toEqual([['refused'], ['rejected'], ['stuck']])
})
// Beside its own Retry the resend step is the button; without one the words keep it.
it('leaves out sending again only where the message has its own Retry', () => {
const startFailed = (clientMessageId: string): StructuredAgentSessionOutboxEntry =>
@@ -206,7 +238,7 @@ describe('the notice on each message that did not go through', () => {
[agentJournalSubmissionKey('first')]: 'Claude stopped before it finished starting.',
[agentJournalSubmissionKey('second')]: 'Claude stopped before it finished starting.'
})
expect(texts([startFailed('rejected'), entry('held')], 'held')).toMatchObject({
expect(texts([entry('held', { outlivedStop: true }), startFailed('rejected')])).toMatchObject({
[agentJournalSubmissionKey('rejected')]:
'Claude stopped before it finished starting. Send your message to try again.'
})
@@ -296,7 +328,6 @@ describe('the notice on each message that did not go through', () => {
expect(
texts(
facts.map(([id]) => rejected(id)),
null,
facts.map(([id, fact]) => recorded(id, fact))
)
).toEqual(
@@ -357,7 +388,6 @@ describe('the notice on each message that did not go through', () => {
})
})
],
null,
[
{
clientMessageId: 'recorded',
@@ -378,6 +408,54 @@ describe('the notice on each message that did not go through', () => {
})
})
// An earlier attempt under the id may have landed, so the row never says it was not sent.
it('words a kept message whose id expired as an outcome Orca cannot confirm', () => {
const notices = structuredAgentSessionDeliveryNotices(
[
entry('expired', {
lastAttemptAt: 1,
lastFailure: { kind: 'refused', code: 'agent_session_operation_expired' }
}),
entry('fresh', {
state: 'rejected',
lastFailure: { kind: 'refused', code: 'agent_session_operation_expired' }
})
],
'Claude',
vi.fn(),
[],
[],
NOT_FAILED_HERE
)
expect(notices.get(agentJournalSubmissionKey('expired'))).toMatchObject({
text: "Orca couldn't confirm what happened. Check the chat."
})
expect(notices.get(agentJournalSubmissionKey('expired'))?.onRetry).toBeDefined()
// A first attempt's id was replaced when it was refused, so nothing can have landed.
expect(notices.get(agentJournalSubmissionKey('fresh'))?.text).toBe('Your message was not sent.')
})
// A cause read back from storage may have cleared; a Retry it still stops brings it back.
it('words a held cause only where its failure was seen while the chat was open', () => {
const held = entry('held', {
lastAttemptAt: 1,
lastFailure: {
kind: 'refused',
code: 'agent_session_journal_unreadable',
details: { reason: 'journalWrittenByNewerOrca' }
}
})
const words = (failedHere: ReadonlySet<string>) =>
structuredAgentSessionDeliveryNotices([held], 'Claude', vi.fn(), [], [], failedHere).get(
agentJournalSubmissionKey('held')
)
expect(words(NOT_FAILED_HERE)).toMatchObject({ text: 'Your message was not sent.' })
expect(words(NOT_FAILED_HERE)?.onRetry).toBeDefined()
expect(words(new Set(['held']))?.text).toBe(
'Chats were saved by a newer Orca. Your message was not sent. Update Orca to keep using them.'
)
})
it('says nothing on a message that is only waiting its turn or on its way', () => {
expect(texts([entry('queued'), entry('sending', { state: 'dispatching' })])).toEqual({})
})
@@ -443,7 +521,6 @@ describe('the notice on each message that did not go through', () => {
rejected('second', startFailed),
rejected('other', otherRefusal)
],
null,
[
recorded('first', startFailed),
recorded('second', startFailed),
@@ -461,12 +538,12 @@ describe('the notice on each message that did not go through', () => {
it('keeps the full notice when the rejection is not loaded, or no start row states it', () => {
const shown = "Claude couldn't start. Start a new chat to continue."
expect(texts([rejected('first', startFailed)], null, [], [startFailed])).toEqual({
expect(texts([rejected('first', startFailed)], [], [startFailed])).toEqual({
[agentJournalSubmissionKey('first')]: 'Written by the host.'
})
expect(
texts([rejected('first', startFailed)], null, [recorded('first', startFailed)], [])
).toEqual({ [agentJournalSubmissionKey('first')]: shown })
expect(texts([rejected('first', startFailed)], [recorded('first', startFailed)], [])).toEqual(
{ [agentJournalSubmissionKey('first')]: shown }
)
})
})
})
@@ -2,9 +2,10 @@
//
// Derived from the outbox on every render and never stored: each failed or held message carries
// its own typed failure, so each row words its own reason. Read through the drain's own rule: while
// the queue is stopped, only the message it stopped on has a Retry; another's would wait unseen
// behind it. One waiting behind says nothing; a rejected message holds nothing up, so it keeps its
// words and gets its Retry once the queue moves.
// the queue is stopped, the message it stopped on and any failed one ahead of it have a Retry, as
// each would go out at once; one behind it would wait unseen. One waiting behind says nothing; a
// rejected or refused message holds nothing up, so it keeps its words and gets its Retry once the
// queue moves.
//
// A message the host recorded and then rejected is worded from the journal's own fact, found by id;
// the message keeps only a smaller copy, read when its submission is not loaded. A rejection that
@@ -24,9 +25,13 @@ import type {
import { agentSessionWriteNotDoneParts } from '../../../../shared/agent-session-refusal-notice'
import { isStructuredAgentSessionStartFailureRow } from '../../../../shared/structured-agent-session-start-failure-row-key'
import {
admitStructuredAgentSessionOutboxEntry,
structuredAgentSessionEntryIdExpired,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
import {
admitStructuredAgentSessionOutboxEntry,
structuredAgentSessionEntryHeldForRetry
} from '../../../../shared/structured-agent-session-outbox-admission'
import type { AgentSessionFailureWordsContext } from '../../../../shared/agent-session-failure-words'
import { structuredAgentSessionAttemptFailureParts } from '../../../../shared/structured-agent-session-send-disposition'
import { translate } from '@/i18n/i18n'
@@ -83,7 +88,8 @@ function deliveryNoticeText(
entry: StructuredAgentSessionOutboxEntry,
context: AgentSessionFailureWordsContext,
recorded: AgentJournalSubmission | undefined,
startFailures: readonly AgentSessionFailureFact[]
startFailures: readonly AgentSessionFailureFact[],
failedHere: ReadonlySet<string>
): string {
// A send attempted before a Stop and then interrupted may already be with the host.
const attemptedAcrossStop =
@@ -100,6 +106,14 @@ function deliveryNoticeText(
'Message was not sent.'
)
}
// An earlier attempt under the id the host forgot may already be in the chat.
if (structuredAgentSessionEntryIdExpired(entry)) {
return agentSessionWriteNoticeText(['outcomeUnknown'])
}
// Its cause may have cleared since it was saved; a Retry it still stops brings the cause back.
if (structuredAgentSessionEntryHeldForRetry(entry) && !failedHere.has(entry.clientMessageId)) {
return agentSessionWriteNoticeText(agentSessionWriteNotDoneParts('send'))
}
if (
entry.state === 'rejected' &&
agentSessionFailureStatedByStartRow(recorded?.rejection, startFailures)
@@ -115,35 +129,42 @@ function deliveryNoticeText(
)
}
/** Keyed by the message id the transcript renders each entry under. `blockedClientMessageId` is
* the entry a refusal stopped the queue on; `agentName` is the chat's agent, for the words. */
/** Keyed by the message id the transcript renders each entry under; `agentName` is the chat's
* agent, for the words. */
export function structuredAgentSessionDeliveryNotices(
outbox: readonly StructuredAgentSessionOutboxEntry[],
blockedClientMessageId: string | null,
agentName: string,
retry: (clientMessageId: string) => void,
/** The journal's rows, whose rejected ones carry more of a rejection than the message keeps. */
submissions: readonly AgentJournalSubmission[],
/** What the loaded start-failure rows state, from `structuredAgentSessionStartFailureFacts`. */
startFailures: readonly AgentSessionFailureFact[]
startFailures: readonly AgentSessionFailureFact[],
/** Ids whose send failed or was refused while this chat was open: only they word their cause. */
failedHere: ReadonlySet<string>
): ReadonlyMap<string, NativeChatDeliveryNotice> {
const admission = admitStructuredAgentSessionOutboxEntry(outbox, blockedClientMessageId)
const admission = admitStructuredAgentSessionOutboxEntry(outbox)
const held = admission.state === 'blocked' ? admission.entry.clientMessageId : null
const stalledFrom = admission.state === 'blocked' ? outbox.indexOf(admission.entry) : -1
const rejected = new Map(
submissions
.filter((submission) => submission.dispatchState === 'rejected')
.map((submission) => [submission.clientMessageId, submission])
)
const notices = new Map<string, NativeChatDeliveryNotice>()
for (const entry of outbox) {
if (entry.state === 'rejected' || entry.clientMessageId === held) {
for (const [index, entry] of outbox.entries()) {
if (
entry.state === 'rejected' ||
structuredAgentSessionEntryHeldForRetry(entry) ||
entry.clientMessageId === held
) {
// Its own Retry is the step, so the words leave out sending again.
const retryControl = held === null || entry.clientMessageId === held
const retryControl = stalledFrom === -1 || index <= stalledFrom
const text = deliveryNoticeText(
entry,
{ agentName, retryControl },
rejected.get(entry.clientMessageId),
startFailures
startFailures,
failedHere
)
notices.set(
agentJournalSubmissionKey(entry.clientMessageId),
@@ -3,10 +3,10 @@
import { describe, expect, it } from 'vitest'
import {
admitStructuredAgentSessionOutboxEntry,
createStructuredAgentSessionOutboxEntry,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
import { admitStructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox-admission'
import { requeueInterruptedStructuredAgentSessionDispatches } from './structured-agent-session-outbox-dispatch'
function dispatching(
@@ -33,8 +33,8 @@ describe('requeue after an owner change', () => {
[dispatching('stopped', { outlivedStop: true })],
1
)
expect(admitStructuredAgentSessionOutboxEntry([stopped], null).state).toBe('blocked')
expect(admitStructuredAgentSessionOutboxEntry([stopped]).state).toBe('blocked')
const [plain] = requeueInterruptedStructuredAgentSessionDispatches([dispatching('plain')], 1)
expect(admitStructuredAgentSessionOutboxEntry([plain], null).state).toBe('dispatch')
expect(admitStructuredAgentSessionOutboxEntry([plain]).state).toBe('dispatch')
})
})
@@ -80,14 +80,14 @@ export function requeueInterruptedStructuredAgentSessionDispatches(
export function dispatchStructuredAgentSessionOutboxEntry(args: {
next: StructuredAgentSessionOutboxEntry
persisted: readonly StructuredAgentSessionOutboxEntry[]
/** The outbox to stage `next` in: the latest, so a hold only memory keeps survives it. */
entries: readonly StructuredAgentSessionOutboxEntry[]
sessionId: string
target: RuntimeClientTarget
fence: number
dispatchGeneration: number
dispatchGenerationRef: MutableRef<number>
inFlightIdRef: MutableRef<string | null>
blockedIdRef: MutableRef<string | null>
setError: (error: string | null) => void
applyDisposition: (disposition: StructuredAgentSessionSendDisposition) => void
createOperationId: () => string
@@ -95,13 +95,22 @@ export function dispatchStructuredAgentSessionOutboxEntry(args: {
const start = async (): Promise<boolean> => {
args.inFlightIdRef.current = args.next.clientMessageId
const staged = updateStructuredAgentSessionOutboxEntry(
args.persisted,
args.entries,
args.next.clientMessageId,
(entry) => stageStructuredAgentSessionOutboxEntryForSend(entry, Date.now())
)
if (!commitStructuredAgentSessionOutbox(args.sessionId, staged, { onlyIfSaved: true })) {
args.inFlightIdRef.current = null
args.blockedIdRef.current = args.next.clientMessageId
// Held for Retry. Saved if storage takes this one write; otherwise only the open chat's
// outbox holds it.
commitStructuredAgentSessionOutbox(
args.sessionId,
updateStructuredAgentSessionOutboxEntry(
getStructuredAgentSessionOutbox(args.sessionId),
args.next.clientMessageId,
(entry) => ({ ...entry, lastFailure: { kind: 'failed' } })
)
)
args.setError('Message could not be saved to the outbox')
return false
}
@@ -118,7 +127,6 @@ export function dispatchStructuredAgentSessionOutboxEntry(args: {
disposeStructuredAgentSessionSendResult({
entries: getStructuredAgentSessionOutbox(args.sessionId),
entry: args.next,
blockedClientMessageId: args.blockedIdRef.current,
result,
createOperationId: args.createOperationId
})
@@ -138,11 +146,7 @@ export function dispatchStructuredAgentSessionOutboxEntry(args: {
if (args.dispatchGenerationRef.current !== args.dispatchGeneration) {
return false
}
const input = {
entries: getStructuredAgentSessionOutbox(args.sessionId),
entry: args.next,
blockedClientMessageId: args.blockedIdRef.current
}
const input = { entries: getStructuredAgentSessionOutbox(args.sessionId), entry: args.next }
const thrown = readAgentSessionErrorRefusal(caught)
const refusal = thrown ? agentSessionRefusalFailure(thrown) : undefined
args.applyDisposition(
@@ -2,7 +2,10 @@
// id when the recorded one can only ever replay a settled rejection.
import type { AgentJournalSubmission } from '../../../../shared/agent-session-journal-types'
import type { StructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
import {
structuredAgentSessionEntryIdExpired,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
import {
commitStructuredAgentSessionOutbox,
getStructuredAgentSessionOutbox
@@ -22,10 +25,16 @@ export function retryStructuredAgentSessionOutboxEntry(args: {
// The host settled this id as rejected, and reusing it only replays that forever, so rotate the
// id for a safe resend. Read from the message itself, which outlives a restart, or from a
// reconciliation that settled an earlier unknown before the outbox caught up. A refusal that
// settled the message already rotated it.
// settled the message already rotated it. An expired id is refused for good; its row told the
// user to check the chat first.
const recordedRejection =
current?.state === 'rejected' && current.lastFailure?.kind === 'rejected'
if (current && (recordedRejection || submission?.dispatchState === 'rejected')) {
if (
current &&
(recordedRejection ||
submission?.dispatchState === 'rejected' ||
structuredAgentSessionEntryIdExpired(current))
) {
const rotated = outbox.map((entry) =>
entry.clientMessageId === clientMessageId
? {
@@ -62,9 +71,10 @@ export function retryStructuredAgentSessionOutboxEntry(args: {
}
}
/** The user's own Retry is what a Stop left the entry waiting for. */
/** The user's own Retry is what a Stop, or a failure saved on the message, left it waiting for. */
function retriedByUser({
outlivedStop: _retried,
lastFailure: _sentAgain,
...entry
}: StructuredAgentSessionOutboxEntry): StructuredAgentSessionOutboxEntry {
return entry
@@ -186,21 +186,21 @@ describe('queued message cards', () => {
entry('a', { state: 'dispatching', lastAttemptAt: 2, sentDelivery: 'queue-if-active' }),
entry('b')
]
expect(ids(outboxOutsideQueuedCards(inFlight, [], true, null, QUEUEING))).toEqual([])
expect(ids(outboxOutsideQueuedCards(inFlight, [], false, null, QUEUEING))).toEqual(['a', 'b'])
// Refused and held for Retry: its text, and everything waiting behind it, stays in view.
expect(ids(outboxOutsideQueuedCards(inFlight, [], true, QUEUEING))).toEqual([])
expect(ids(outboxOutsideQueuedCards(inFlight, [], false, QUEUEING))).toEqual(['a', 'b'])
// Refused and held for Retry: its text stays in view, and what follows it is on its way.
const refused = [entry('a', { lastFailure: { kind: 'failed' } }), entry('b')]
expect(ids(outboxOutsideQueuedCards(refused, [], true, 'a', QUEUEING))).toEqual(['a', 'b'])
expect(ids(outboxOutsideQueuedCards(refused, [], true, QUEUEING))).toEqual(['a'])
// A rejected send holds nothing up: what follows it is still on its way to a card.
const rejected = [
entry('a', { state: 'rejected', lastFailure: { kind: 'rejected', reason: null } }),
entry('b')
]
expect(ids(outboxOutsideQueuedCards(rejected, [], true, null, QUEUEING))).toEqual(['a'])
expect(ids(outboxOutsideQueuedCards(rejected, [], true, QUEUEING))).toEqual(['a'])
const unconfirmed = [entry('a', { state: 'unconfirmed' }), entry('b')]
expect(ids(outboxOutsideQueuedCards(unconfirmed, [], true, null, QUEUEING))).toEqual(['a', 'b'])
expect(ids(outboxOutsideQueuedCards(unconfirmed, [], true, QUEUEING))).toEqual(['a', 'b'])
// Once the host visibly holds it, it is a card whatever this queue last heard.
expect(ids(outboxOutsideQueuedCards(unconfirmed, ['a'], true, null, QUEUEING))).toEqual(['b'])
expect(ids(outboxOutsideQueuedCards(unconfirmed, ['a'], true, QUEUEING))).toEqual(['b'])
})
it('hides a send only by what its request carries: a plain one is always a bubble', () => {
@@ -222,12 +222,12 @@ describe('queued message cards', () => {
entries.map((candidate) => candidate.clientMessageId)
// The capability is unknown: a new send goes out plain, so it stays in view.
const unknown = { capability: 'unknown', enabled: true } as const
expect(ids(outboxOutsideQueuedCards([entry('new')], [], true, null, unknown))).toEqual(['new'])
expect(ids(outboxOutsideQueuedCards([entry('new')], [], true, unknown))).toEqual(['new'])
// A send that went out plain stays a bubble; one that went out queued replays queued.
const sent = [
entry('plain', { state: 'dispatching', lastAttemptAt: 2, sentDelivery: null }),
entry('queued', { state: 'dispatching', lastAttemptAt: 2, sentDelivery: 'queue-if-active' })
]
expect(ids(outboxOutsideQueuedCards(sent, [], true, null, unknown))).toEqual(['plain'])
expect(ids(outboxOutsideQueuedCards(sent, [], true, unknown))).toEqual(['plain'])
})
})
@@ -9,10 +9,11 @@ import {
structuredAgentSessionEntryAsksToQueue,
type StructuredAgentSessionQueueDelivery
} from '../../../../shared/structured-agent-session-outbox-delivery'
import type { StructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
import {
admitStructuredAgentSessionOutboxEntry,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
structuredAgentSessionEntryHeldForRetry
} from '../../../../shared/structured-agent-session-outbox-admission'
/** Why a card is not on its way right now; decides the caption under the text. */
export type QueuedMessageCardHold =
@@ -103,23 +104,24 @@ export function newestSteerableQueuedMessageCard(
* read from what its request carries — otherwise it paints
* in the transcript until the queued answer retires it. A plain send stays a bubble. From
* the entry the drain is stopped on (read through the drain's own rule), nothing is on its
* way: those stay bubbles so their text is visible beside the Retry row.
* way, nor is one held for its Retry: those stay bubbles so their text is visible beside the
* Retry row.
*/
export function outboxOutsideQueuedCards(
outbox: readonly StructuredAgentSessionOutboxEntry[],
heldIds: readonly string[],
isWorking: boolean,
blockedClientMessageId: string | null,
host: StructuredAgentSessionQueueDelivery
): readonly StructuredAgentSessionOutboxEntry[] {
const held = new Set(heldIds)
const admission = admitStructuredAgentSessionOutboxEntry(outbox, blockedClientMessageId)
const admission = admitStructuredAgentSessionOutboxEntry(outbox)
const stalledFrom = admission.state === 'blocked' ? outbox.indexOf(admission.entry) : -1
const next = outbox.filter((entry, index) => {
const onItsWay =
isWorking &&
(stalledFrom === -1 || index < stalledFrom) &&
(entry.state === 'queued' || entry.state === 'dispatching') &&
!structuredAgentSessionEntryHeldForRetry(entry) &&
structuredAgentSessionEntryAsksToQueue(entry, host)
return !held.has(entry.clientMessageId) && !onItsWay
})
@@ -220,4 +220,30 @@ describe('native chat transcript virtualizer contract', () => {
expect(virtualizerMock.scrollToOffset).not.toHaveBeenCalled()
})
// Rows drawn after the window (a message shown as not sent) still fill the container.
it('follows the bottom of the container when no row is windowed', () => {
const container = document.createElement('div')
Object.defineProperty(container, 'scrollHeight', { configurable: true, value: 2000 })
virtualizerMock.scrollElement.current = container
const noSlots: NativeChatMessageSlot[] = []
const { result, rerender } = renderHook(
({ slots }) =>
useNativeChatTranscriptWindow({
scrollRef: { current: container },
slots,
isVisible: true,
revealIndex: -1
}),
{ initialProps: { slots: noSlots } }
)
result.current.scrollToEnd()
expect(virtualizerMock.scrollToEnd).not.toHaveBeenCalled()
expect(container.scrollTop).toBe(2000)
rerender({ slots: [slot('a')] })
result.current.scrollToEnd()
expect(virtualizerMock.scrollToEnd).toHaveBeenCalledOnce()
})
})
@@ -283,18 +283,27 @@ export function useNativeChatTranscriptWindow({
return
}
finishReaderTakeover()
if (virtualizer.scrollElement) {
// With no windowed row it would resolve the end from its own rows' height, 0, though rows
// drawn after the window (a message shown as not sent) still fill the container.
if (virtualizer.scrollElement && slots.length > 0) {
virtualizer.scrollToEnd({ behavior: 'auto' })
return
}
// No virtualizer yet (a container without layout): the document's own bottom
// is the same offset the virtualizer would resolve for the last row.
// No virtualizer yet (a container without layout), or no windowed row: the document's own
// bottom is the same offset the virtualizer would resolve for the last row.
const previous = container.scrollTop
container.scrollTop = container.scrollHeight
if (container.scrollTop !== previous) {
programmaticScrollMarks.mark(container.scrollTop)
}
}, [finishReaderTakeover, isVisible, programmaticScrollMarks, scrollRef, virtualizer])
}, [
finishReaderTakeover,
isVisible,
programmaticScrollMarks,
scrollRef,
slots.length,
virtualizer
])
const restoreScrollOffset = useCallback(
(offset: number) => {
@@ -33,7 +33,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: () => 'operation-1',
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: mocks.outboxSend,
retry: vi.fn()
@@ -38,7 +38,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: () => 'operation-1',
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -37,7 +37,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: vi.fn(),
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -23,6 +23,8 @@ import { enqueueStructuredAgentSessionLaunchPrompt } from './structured-agent-se
import { settleStructuredAgentLaunchPrompt } from '@/lib/structured-agent-session-launch-prompt'
import { structuredAgentSessionDeliveryNotices } from './structured-agent-session-delivery-notices'
import { agentJournalSubmissionKey } from '../../../../shared/agent-session-journal-item-key'
import type { StructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
import { structuredAgentSessionEntryHeldForRetry } from '../../../../shared/structured-agent-session-outbox-admission'
const LOCAL_TARGET = { kind: 'local' } as const
@@ -51,6 +53,13 @@ function sentTexts(): (string | undefined)[] {
return mocks.call.mock.calls.map((call) => requestText(call[2]))
}
/** The messages that wait for their own Retry. */
function heldIds(outbox: readonly StructuredAgentSessionOutboxEntry[]): string[] {
return outbox
.filter(structuredAgentSessionEntryHeldForRetry)
.map((entry) => entry.clientMessageId)
}
function submissionResult(
clientMessageId: string,
dispatchState: 'pending' | 'accepted',
@@ -155,17 +164,26 @@ describe('structured agent session outbox admission', () => {
expect(mocks.call).toHaveBeenCalledTimes(1)
})
it('keeps a refused head blocking the queue', async () => {
mocks.call.mockResolvedValue(refusedResult('agent_session_ownership_unknown'))
// A refused message lands only by its own Retry, so nothing it could be reordered around.
it('sends a message past a refused one, which waits for its own Retry', async () => {
mocks.call.mockImplementation((_target, _method, params) =>
Promise.resolve(
requestText(params) === 'first'
? refusedResult('agent_session_ownership_unknown')
: submissionResult(requestId(params), 'accepted', 20)
)
)
const { result } = renderOutbox()
act(() => expect(result.current.send('first')).toBe(true))
await waitFor(() => expect(result.current.blockedClientMessageId).not.toBeNull())
await waitFor(() => expect(heldIds(result.current.outbox)).toHaveLength(1))
const refusedId = heldIds(result.current.outbox)[0]
act(() => expect(result.current.send('second')).toBe(true))
await settleTimers(200)
await waitFor(() => expect(result.current.outbox).toHaveLength(1))
await settleTimers(50)
expect(sentTexts()).not.toContain('second')
expect(result.current.blockedClientMessageId).toBe(result.current.outbox[0]?.clientMessageId)
expect(sentTexts()).toEqual(['first', 'second'])
expect(heldIds(result.current.outbox)).toEqual([refusedId])
})
it.each(['accepted', 'pending'] as const)(
@@ -250,8 +268,8 @@ describe('structured agent session outbox admission', () => {
await waitFor(() => expect(result.current.outbox.at(-1)?.state).toBe('rejected'))
}
act(() => expect(result.current.send('held')).toBe(true))
await waitFor(() => expect(result.current.blockedClientMessageId).not.toBeNull())
const heldId = result.current.blockedClientMessageId
await waitFor(() => expect(heldIds(result.current.outbox)).toHaveLength(1))
const heldId = heldIds(result.current.outbox)[0]
const [x, a, c] = result.current.outbox
// x goes out and stays in flight, so a, retried after it, waits queued ahead of the held one.
@@ -262,21 +280,21 @@ describe('structured agent session outbox admission', () => {
// The drain would send a next, so the queue reads as moving and c offers its own Retry.
const notices = structuredAgentSessionDeliveryNotices(
result.current.outbox,
result.current.blockedClientMessageId,
'Claude',
() => {},
[],
[]
[],
result.current.failedHere
)
expect(notices.get(agentJournalSubmissionKey(c!.clientMessageId))?.onRetry).toBeDefined()
act(() => result.current.retry(c!.clientMessageId))
expect(result.current.blockedClientMessageId).toBe(heldId)
expect(heldIds(result.current.outbox)).toEqual([heldId])
await act(async () => inFlight.resolve(submissionResult(x!.clientMessageId, 'accepted', 20)))
await waitFor(() => expect(result.current.outbox).toHaveLength(1))
await settleTimers(50)
expect(sentTexts().filter((text) => text === 'held')).toHaveLength(1)
expect(result.current.blockedClientMessageId).toBe(heldId)
expect(heldIds(result.current.outbox)).toEqual([heldId])
})
it('keeps a refused message held when an unconfirmed message ahead of it is retried', async () => {
@@ -302,8 +320,8 @@ describe('structured agent session outbox admission', () => {
act(() => expect(result.current.send('first')).toBe(true))
await waitFor(() => expect(result.current.outbox[0]?.state).toBe('dispatching'))
act(() => expect(result.current.send('held')).toBe(true))
await waitFor(() => expect(result.current.blockedClientMessageId).not.toBeNull())
const heldId = result.current.blockedClientMessageId
await waitFor(() => expect(heldIds(result.current.outbox)).toHaveLength(1))
const heldId = heldIds(result.current.outbox)[0]
const firstId = String(result.current.outbox[0]?.clientMessageId)
const unknown: AgentJournalSubmission = {
...submissionResult(firstId, 'pending', 10).value.submission,
@@ -316,6 +334,6 @@ describe('structured agent session outbox admission', () => {
await waitFor(() => expect(result.current.outbox).toHaveLength(1))
await settleTimers(50)
expect(sentTexts()).toEqual(['first', 'held', 'first'])
expect(result.current.blockedClientMessageId).toBe(heldId)
expect(heldIds(result.current.outbox)).toEqual([heldId])
})
})
@@ -0,0 +1,63 @@
import { useCallback, useState } from 'react'
import type { StructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
const NONE: ReadonlySet<string> = new Set()
/**
* The messages whose send failed or was refused while this chat was open, in memory only. Only
* they word their saved cause: one read back from storage may have outlived it (Orca was updated
* since, say), and its Retry, refused again if the cause holds, brings the words back.
*/
export function useStructuredAgentSessionOutboxFailedHere(sessionId: string): {
failedHere: ReadonlySet<string>
/** Records the entries a send's answer gave a failure they did not carry before. */
recordFailures: (
before: readonly StructuredAgentSessionOutboxEntry[],
after: readonly StructuredAgentSessionOutboxEntry[]
) => void
forget: (clientMessageId: string) => void
} {
const [recorded, setRecorded] = useState({ sessionId, ids: NONE })
const recordFailures = useCallback(
(
before: readonly StructuredAgentSessionOutboxEntry[],
after: readonly StructuredAgentSessionOutboxEntry[]
): void => {
const previous = new Map(before.map((entry) => [entry.clientMessageId, entry.lastFailure]))
const failed = after.filter(
(entry) =>
entry.lastFailure !== undefined &&
previous.get(entry.clientMessageId) !== entry.lastFailure
)
if (failed.length === 0) {
return
}
setRecorded((current) => ({
sessionId,
ids: new Set([
...(current.sessionId === sessionId ? current.ids : NONE),
...failed.map((entry) => entry.clientMessageId)
])
}))
},
[sessionId]
)
const forget = useCallback(
(clientMessageId: string): void => {
setRecorded((current) => {
if (current.sessionId !== sessionId || !current.ids.has(clientMessageId)) {
return current
}
const ids = new Set(current.ids)
ids.delete(clientMessageId)
return { sessionId, ids }
})
},
[sessionId]
)
return {
failedHere: recorded.sessionId === sessionId ? recorded.ids : NONE,
recordFailures,
forget
}
}
@@ -106,12 +106,12 @@ describe('an outbox on a host that accepts a send before any agent has it', () =
expect(mocks.call).toHaveBeenCalledTimes(1)
})
it('keeps a blocked head blocked across a fence change; only Retry sends it', async () => {
it('keeps a failed send held across a fence change; only Retry sends it', async () => {
mocks.call.mockRejectedValueOnce(new Error('send failed')).mockResolvedValue({ ok: true })
const { result, rerender } = render()
act(() => expect(result.current.send('hello')).toBe(true))
await waitFor(() => expect(result.current.blockedClientMessageId).not.toBeNull())
await waitFor(() => expect(result.current.outbox[0]?.lastFailure).toBeDefined())
rerender({ fence: 2 })
await settle()
expect(mocks.call).toHaveBeenCalledTimes(1)
@@ -0,0 +1,232 @@
// @vitest-environment happy-dom
// A message the chat said was not sent goes out only on its Retry, on every host, however the
// chat's agent owner moves. A send still in flight when the owner moves, on a host not known to
// record sends before starting an agent, goes out again under its id.
import { act, cleanup, renderHook } from '@testing-library/react'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
const mocks = vi.hoisted(() => {
const regained: (() => void)[] = []
return { call: vi.fn(), probe: vi.fn(), regained }
})
vi.mock('@/runtime/structured-agent-session-client', () => ({
callStructuredAgentSession: mocks.call
}))
vi.mock('@/runtime/runtime-rpc-client', async (importOriginal) => ({
...(await importOriginal<Record<string, unknown>>()),
runtimeEnvironmentSupportsCapability: mocks.probe
}))
vi.mock('@/runtime/runtime-host-contact-regained', () => ({
subscribeRuntimeHostContactRegained: (_environmentId: string, listener: () => void) => {
mocks.regained.push(listener)
return () => {}
}
}))
import type { AgentSessionWireRefusalCode } from '../../../../shared/agent-session-wire'
import {
createStructuredAgentSessionOutboxEntry,
type StructuredAgentSessionOutboxEntry
} from '../../../../shared/structured-agent-session-outbox'
import { writeOutbox } from './structured-agent-session-outbox-storage'
import { useStructuredAgentSessionOutbox } from './use-structured-agent-session-outbox'
type SendRequest = { envelope?: { clientOperationId?: string; expectedRuntimeFence?: number } }
// Stable, as the view passes it: a new object each render would re-run the owner-change requeue.
const TARGET = { kind: 'environment', environmentId: 'env-1' } as const
function requestId(params: SendRequest | undefined): string {
return String(params?.envelope?.clientOperationId)
}
function sentIds(): string[] {
return mocks.call.mock.calls.map((call) => requestId(call[2]))
}
/** Each send as `id@fence`, so a resend on a new owner is told from the first attempt. */
function sentAttempts(): string[] {
return mocks.call.mock.calls.map((call) => {
const params: SendRequest | undefined = call[2]
return `${requestId(params)}@${String(params?.envelope?.expectedRuntimeFence)}`
})
}
const QUEUED = { ok: true, replayed: false, fence: 2, value: { queued: {} } }
function heldEntry(
clientMessageId: string,
code: AgentSessionWireRefusalCode
): StructuredAgentSessionOutboxEntry {
return {
...createStructuredAgentSessionOutboxEntry({
clientMessageId,
sessionId: 'session-1',
text: 'held',
attachments: [],
queuedAt: 1
}),
lastAttemptAt: 2,
lastFailure: { kind: 'refused', code }
}
}
async function settle(): Promise<void> {
await act(() => new Promise<void>((resolve) => setTimeout(resolve, 50)))
}
function mount(fence: number) {
return renderHook(
({ fence: current }: { fence: number }) =>
useStructuredAgentSessionOutbox({
sessionId: 'session-1',
target: TARGET,
fence: current,
submissions: []
}),
{ initialProps: { fence } }
)
}
function refusal(code: string) {
return { ok: false, refusal: { code, message: code } }
}
async function refusedOnce(code: string) {
mocks.call.mockResolvedValue(refusal(code))
const hook = mount(1)
await settle()
act(() => void hook.result.current.send('hello'))
await settle()
expect(sentIds()).toHaveLength(1)
expect(hook.result.current.outbox[0]?.lastFailure).toBeDefined()
return hook
}
describe('a refused message when the chat gets a new owner', () => {
beforeEach(() => {
localStorage.clear()
mocks.call.mockReset()
mocks.probe.mockReset()
mocks.regained.length = 0
})
afterEach(() => cleanup())
it('is not sent on a current host whose capability probe fails after the owner moved', async () => {
mocks.probe.mockResolvedValue(true)
const hook = await refusedOnce('agent_session_checkpoint_stale')
hook.rerender({ fence: 2 })
await settle()
mocks.probe.mockRejectedValueOnce(new Error('status timeout'))
act(() => mocks.regained.forEach((listener) => listener()))
await settle()
expect(sentIds()).toHaveLength(1)
expect(hook.result.current.outbox[0]?.lastFailure).toBeDefined()
})
it('is not sent while the capability probe has not answered', async () => {
mocks.probe.mockReturnValue(new Promise(() => {}))
const hook = await refusedOnce('agent_session_checkpoint_stale')
hook.rerender({ fence: 2 })
await settle()
expect(sentIds()).toHaveLength(1)
})
it.each([
'agent_session_owner_restart_failed',
'agent_session_checkpoint_stale',
'agent_session_conflict',
'agent_session_ownership_unknown',
'execution_owner_reconciling',
'agent_session_journal_unreadable',
'agent_session_operation_capacity'
] as const)(
'waits for its Retry on a host known to be older when the owner moves, after %s',
async (code) => {
mocks.probe.mockResolvedValue(false)
// Held under an id already sent once, so its Retry keeps the id.
writeOutbox('session-1', [heldEntry('op-held', code)])
mocks.call.mockResolvedValue(QUEUED)
const hook = mount(1)
await settle()
hook.rerender({ fence: 2 })
await settle()
expect(sentIds()).toEqual([])
expect(hook.result.current.outbox[0]?.lastFailure).toBeDefined()
act(() => hook.result.current.retry('op-held'))
await settle()
expect(sentIds()).toEqual(['op-held'])
expect(hook.result.current.outbox).toEqual([])
}
)
it('waits for its Retry on a host known to be older when the send never reached it', async () => {
mocks.probe.mockResolvedValue(false)
mocks.call.mockRejectedValueOnce(new Error('send failed')).mockResolvedValue(QUEUED)
const hook = mount(1)
await settle()
act(() => void hook.result.current.send('hello'))
await settle()
expect(hook.result.current.outbox[0]?.lastFailure).toEqual({ kind: 'failed' })
hook.rerender({ fence: 2 })
await settle()
expect(sentIds()).toHaveLength(1)
act(() => hook.result.current.retry(hook.result.current.outbox[0]!.clientMessageId))
await settle()
expect(sentIds()).toHaveLength(2)
expect(hook.result.current.outbox).toEqual([])
})
// The owner change's requeue lands in the same commit as the drain's next send.
it.each([
['a host known to be older', false],
['a host whose capability check has not answered', undefined]
])(
'sends again the send in flight when the owner moves on %s, with one queued behind it',
async (_label, supports) => {
mocks.probe.mockReturnValue(
supports === undefined ? new Promise(() => {}) : Promise.resolve(supports)
)
// The first owner never answers the first send.
mocks.call.mockReturnValueOnce(new Promise(() => {})).mockResolvedValue(QUEUED)
const hook = mount(1)
await settle()
act(() => void hook.result.current.send('first'))
await settle()
act(() => void hook.result.current.send('second'))
await settle()
const [first, second] = hook.result.current.outbox.map((entry) => entry.clientMessageId)
hook.rerender({ fence: 2 })
await settle()
await settle()
expect(sentAttempts()).toEqual([`${first}@1`, `${first}@2`, `${second}@2`])
expect(hook.result.current.outbox).toEqual([])
}
)
it('sends the send in flight again on a new owner, but not a message it refused', async () => {
mocks.probe.mockResolvedValue(false)
writeOutbox('session-1', [heldEntry('op-held', 'agent_session_checkpoint_stale')])
mocks.call.mockReturnValueOnce(new Promise(() => {})).mockResolvedValue(QUEUED)
const hook = mount(1)
await settle()
act(() => void hook.result.current.send('first'))
await settle()
act(() => void hook.result.current.send('second'))
await settle()
const [, first, second] = hook.result.current.outbox.map((entry) => entry.clientMessageId)
hook.rerender({ fence: 2 })
await settle()
await settle()
expect(sentAttempts()).toEqual([`${first}@1`, `${first}@2`, `${second}@2`])
expect(hook.result.current.outbox).toMatchObject([
{ clientMessageId: 'op-held', lastFailure: { code: 'agent_session_checkpoint_stale' } }
])
})
})
@@ -18,7 +18,6 @@ export function useStructuredAgentSessionOutboxOwnership(args: {
submissions: readonly AgentJournalSubmission[]
/** Ids of the host's published drafts; an entry with one of these ids is host-owned. */
queuedMessageIds: readonly string[] | undefined
blockedIdRef: { current: string | null }
/** The send in flight and its generation: a host that holds that send answered it. */
inFlightIdRef: { current: string | null }
dispatchGenerationRef: { current: number }
@@ -27,7 +26,7 @@ export function useStructuredAgentSessionOutboxOwnership(args: {
/** Stop's local step, before its RPC, so the drain has nothing left to send after it. */
withdrawUnsent: () => void
} {
const { blockedIdRef, queuedMessageIds, restoreWithdrawn, sessionId } = args
const { queuedMessageIds, restoreWithdrawn, sessionId } = args
const { dispatchGenerationRef, inFlightIdRef, submissions } = args
const withdrawUnsent = useCallback((): void => {
@@ -35,7 +34,6 @@ export function useStructuredAgentSessionOutboxOwnership(args: {
const next = withdrawUnsentStructuredAgentSessionOutboxEntries(
current,
submissions,
blockedIdRef.current,
inFlightIdRef.current
)
if (next.length === current.length && next.every((entry, index) => entry === current[index])) {
@@ -45,7 +43,7 @@ export function useStructuredAgentSessionOutboxOwnership(args: {
const kept = new Set(next.map((entry) => entry.clientMessageId))
restoreWithdrawn.byStop(current.filter((entry) => !kept.has(entry.clientMessageId)))
commitStructuredAgentSessionOutbox(sessionId, next)
}, [blockedIdRef, inFlightIdRef, restoreWithdrawn, sessionId, submissions])
}, [inFlightIdRef, restoreWithdrawn, sessionId, submissions])
// Drop host-owned entries without a restore: the published card is the text now.
const retire = useCallback(
@@ -111,7 +111,6 @@ describe('a send the host rejected because the agent never started', () => {
// Settled as not delivered: it waits for Retry and holds no later message up.
expect(result.current.error).toBeNull()
expect(result.current.outbox[0]?.state).toBe('rejected')
expect(result.current.blockedClientMessageId).toBeNull()
})
it('sends a new message past one the host could not start the agent for, without resending it', async () => {
@@ -184,7 +183,6 @@ describe('a send the host rejected because the agent never started', () => {
await waitFor(() => expect(result.current.outbox[0]?.state).toBe('rejected'))
expect(shownFailure(result.current.outbox[0])).toBe(reason)
expect(result.current.error).toBeNull()
expect(result.current.blockedClientMessageId).toBeNull()
// Retry is a new message with the same text: a fresh id, sent once.
act(() => result.current.retry(id))
@@ -401,7 +399,8 @@ function submission(
}
}
// An older host: it restarts the agent inside the send, so a new fence is its word to send again.
// An older host restarts the agent inside the send. A send it refused was shown as not sent, so it
// waits for the user's Retry even once the agent has a new owner, as on every host.
describe('a send refused while its agent restarted', () => {
beforeEach(() => {
vi.clearAllMocks()
@@ -414,7 +413,7 @@ describe('a send refused while its agent restarted', () => {
})
// The live order: the restart takes seconds, so the journal settles the resend before its reply.
it('leaves no error behind once a send refused before an agent restart is delivered', async () => {
it('waits for its Retry across an agent restart, and leaves no error once that Retry lands', async () => {
mocks.call
.mockResolvedValueOnce({
ok: false,
@@ -440,8 +439,13 @@ describe('a send refused while its agent restarted', () => {
expect(shownFailure(result.current.outbox[0])).toBe('Your message was not sent.')
)
// The pane learns the new owner and sends the same message again.
// The pane learns the new owner; the message still waits for its Retry.
rerender({ fence: 3, submissions: [] })
await act(() => new Promise<void>((resolve) => setTimeout(resolve, 50)))
expect(mocks.call).toHaveBeenCalledTimes(1)
expect(shownFailure(result.current.outbox[0])).toBe('Your message was not sent.')
act(() => result.current.retry(result.current.outbox[0]!.clientMessageId))
await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(2))
expect(result.current.outbox[0]).toMatchObject({ state: 'dispatching' })
expect(shownFailure(result.current.outbox[0])).toBeUndefined()
@@ -451,7 +455,7 @@ describe('a send refused while its agent restarted', () => {
rerender({ fence: 3, submissions: [submission(id, 'accepted')] })
await waitFor(() => expect(result.current.outbox).toHaveLength(0))
expect(result.current.error).toBeNull()
expect(result.current.blockedClientMessageId).toBeNull()
expect(mocks.call).toHaveBeenCalledTimes(2)
})
})
@@ -0,0 +1,567 @@
// @vitest-environment happy-dom
// A message the chat said was not sent, with a Retry beside it, waits for that Retry. The hold is
// read from the saved message, so quitting and reopening Orca (a new mount over the same storage)
// holds it exactly as the refusal did.
import { act, cleanup, renderHook, waitFor } from '@testing-library/react'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY } from '../../../../shared/protocol-version'
import type { AgentSessionWireRefusalCode } from '../../../../shared/agent-session-wire'
const mocks = vi.hoisted(() => ({ call: vi.fn() }))
vi.mock('@/runtime/structured-agent-session-client', () => ({
callStructuredAgentSession: mocks.call
}))
import { setLocalRuntimeCapabilitiesForTests } from '@/runtime/local-runtime-capabilities'
import { RuntimeRpcCallError } from '@/runtime/runtime-rpc-result'
import { agentJournalSubmissionKey } from '../../../../shared/agent-session-journal-item-key'
import { useStructuredAgentSessionOutbox } from './use-structured-agent-session-outbox'
import {
clearNativeChatDraftCacheForTests,
readNativeChatDraftCache
} from './native-chat-draft-cache'
import { createStructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
import { writeOutbox } from './structured-agent-session-outbox-storage'
import { structuredAgentSessionDeliveryNotices } from './structured-agent-session-delivery-notices'
const SESSION = 'session-1'
// Stable, as the view passes it: a new object each render would re-run the owner-change requeue.
const LOCAL_TARGET = { kind: 'local' } as const
// Read back from storage, the cause may have cleared since (the user updated Orca, say).
const NOT_SENT_WORDS = 'Your message was not sent.'
const NEWER_ORCA_WORDS =
'Chats were saved by a newer Orca. Your message was not sent. Update Orca to keep using them.'
type SendRequest = {
body?: { blocks?: { text?: string }[] }
envelope?: { clientOperationId?: string }
}
function requestId(params: SendRequest | undefined): string {
return String(params?.envelope?.clientOperationId)
}
function sentIds(): string[] {
return mocks.call.mock.calls.map((call) => requestId(call[2]))
}
function newerOrcaRefusal() {
return {
ok: false,
refusal: {
code: 'agent_session_journal_unreadable',
message: 'Chats were saved by a newer Orca. Update Orca to keep using them.',
details: { reason: 'journalWrittenByNewerOrca' }
}
}
}
function refusal(code: AgentSessionWireRefusalCode) {
return { ok: false, refusal: { code, message: code } }
}
function accepted(clientMessageId: string) {
return {
ok: true,
replayed: false,
fence: 1,
cursor: { epoch: 'epoch-1', sequence: 2 },
value: {
clientMessageId,
submission: {
clientMessageId,
fence: 1,
payloadFingerprint: 'fingerprint',
dispatchState: 'accepted',
providerItemId: `provider-${clientMessageId}`,
reason: null,
submittedAt: 1,
resolvedAt: 1
}
}
}
}
function hostAccepts(): void {
mocks.call.mockImplementation((_target, _method, params: SendRequest) =>
Promise.resolve(accepted(requestId(params)))
)
}
function mount(fence = 1) {
return renderHook(
({ fence: current }: { fence: number }) =>
useStructuredAgentSessionOutbox({
sessionId: SESSION,
target: LOCAL_TARGET,
fence: current,
submissions: []
}),
{ initialProps: { fence } }
)
}
type Outbox = ReturnType<typeof mount>['result']['current']
function notices(outbox: Outbox) {
return structuredAgentSessionDeliveryNotices(
outbox.outbox,
'Claude',
outbox.retry,
[],
[],
outbox.failedHere
)
}
function noticeFor(outbox: Outbox, clientMessageId: string) {
return notices(outbox).get(agentJournalSubmissionKey(clientMessageId))
}
/** Long enough for any effect a mount or a state change schedules to have sent. */
async function settle(): Promise<void> {
await act(() => new Promise<void>((resolve) => setTimeout(resolve, 50)))
}
describe('a message the host refused, across a relaunch', () => {
afterEach(() => {
cleanup()
setLocalRuntimeCapabilitiesForTests(null)
})
beforeEach(() => {
vi.clearAllMocks()
localStorage.clear()
setLocalRuntimeCapabilitiesForTests([AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY])
})
it.each(['returned', 'thrown'] as const)(
'is not sent on its own after a relaunch; its Retry sends it (refusal %s)',
async (refusalDelivery) => {
if (refusalDelivery === 'returned') {
mocks.call.mockResolvedValueOnce(newerOrcaRefusal())
} else {
// As `mapRuntimeError` sends a thrown refusal (pinned in `rpc/errors.test.ts`).
mocks.call.mockRejectedValueOnce(
new RuntimeRpcCallError({
id: 'req-1',
ok: false,
error: {
code: 'runtime_error',
message: 'agent_session_journal_unreadable',
data: {
refusal: {
code: 'agent_session_journal_unreadable',
details: { reason: 'journalWrittenByNewerOrca' }
}
}
}
})
)
}
const before = mount()
act(() => expect(before.result.current.send('hello')).toBe(true))
await waitFor(() => expect(before.result.current.outbox[0]?.lastFailure).toBeDefined())
const refusedId = sentIds()[0]!
expect(noticeFor(before.result.current, refusedId)?.text).toBe(NEWER_ORCA_WORDS)
before.unmount()
// The user updates Orca and opens the chat again; the host now takes sends.
hostAccepts()
const after = mount()
await settle()
expect(mocks.call).toHaveBeenCalledTimes(1)
const notice = noticeFor(after.result.current, refusedId)
expect(notice?.text).toBe(NOT_SENT_WORDS)
expect(notice?.onRetry).toBeDefined()
act(() => notice?.onRetry?.())
await waitFor(() => expect(after.result.current.outbox).toHaveLength(0))
// The refusal recorded nothing, so the Retry goes out under the refused id.
expect(sentIds()).toEqual([refusedId, refusedId])
}
)
it('sends a message typed again after a refusal once, not a second time after a relaunch', async () => {
// The chat's history is from a newer Orca: every send is refused until the user updates.
mocks.call.mockResolvedValue(newerOrcaRefusal())
const before = mount()
act(() => expect(before.result.current.send('hello')).toBe(true))
await waitFor(() => expect(before.result.current.outbox[0]?.lastFailure).toBeDefined())
// The user types the same message again rather than pressing Retry.
act(() => expect(before.result.current.send('hello')).toBe(true))
await settle()
const refusedBeforeUpdate = mocks.call.mock.calls.length
before.unmount()
hostAccepts()
const after = mount()
await settle()
// Nothing goes out on its own: each message still says it was not sent, with its own Retry.
expect(mocks.call).toHaveBeenCalledTimes(refusedBeforeUpdate)
expect(after.result.current.outbox).toHaveLength(2)
for (const entry of after.result.current.outbox) {
expect(noticeFor(after.result.current, entry.clientMessageId)?.onRetry).toBeDefined()
}
// One Retry delivers the message once; the other copy stays unsent.
act(() => after.result.current.retry(after.result.current.outbox[1]!.clientMessageId))
await waitFor(() => expect(after.result.current.outbox).toHaveLength(1))
await settle()
expect(mocks.call).toHaveBeenCalledTimes(refusedBeforeUpdate + 1)
})
it('keeps a send that never reached the host held across a relaunch', async () => {
mocks.call.mockRejectedValueOnce(new Error('send failed'))
const before = mount()
act(() => expect(before.result.current.send('hello')).toBe(true))
await waitFor(() =>
expect(before.result.current.outbox[0]?.lastFailure).toEqual({ kind: 'failed' })
)
const failedId = sentIds()[0]!
before.unmount()
hostAccepts()
const after = mount()
await settle()
expect(mocks.call).toHaveBeenCalledTimes(1)
expect(noticeFor(after.result.current, failedId)?.onRetry).toBeDefined()
})
it('holds a refused message an earlier Orca saved, in the shape that build writes', async () => {
// Exactly what a build from before this hold leaves behind after the refusal: no new field.
localStorage.setItem(
`orca:desktopStructuredAgentSessionOutbox:v1:${SESSION}`,
JSON.stringify([
{
clientMessageId: 'op-refused',
sessionId: SESSION,
body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'hello' }] },
previewUris: [],
state: 'queued',
queuedAt: 1,
lastAttemptAt: 2,
retryAfterUnknownSubmittedAt: null,
lastFailure: {
kind: 'refused',
code: 'agent_session_journal_unreadable',
details: { reason: 'journalWrittenByNewerOrca' }
}
}
])
)
hostAccepts()
const { result } = mount()
await settle()
expect(mocks.call).not.toHaveBeenCalled()
expect(noticeFor(result.current, 'op-refused')?.text).toBe(NOT_SENT_WORDS)
})
it("words the cause again once a Retry is refused for it; a relaunch's row only says not sent", async () => {
mocks.call.mockResolvedValue(newerOrcaRefusal())
const before = mount()
act(() => expect(before.result.current.send('hello')).toBe(true))
await waitFor(() => expect(before.result.current.outbox[0]?.lastFailure).toBeDefined())
const refusedId = sentIds()[0]!
before.unmount()
// The chat's history is still from a newer Orca.
const after = mount()
await settle()
expect(noticeFor(after.result.current, refusedId)).toMatchObject({ text: NOT_SENT_WORDS })
act(() => noticeFor(after.result.current, refusedId)?.onRetry?.())
await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(2))
await waitFor(() =>
expect(noticeFor(after.result.current, refusedId)?.text).toBe(NEWER_ORCA_WORDS)
)
expect(noticeFor(after.result.current, refusedId)?.onRetry).toBeDefined()
})
})
// An older host restarts the agent inside a send and refuses it unrecorded when that fails. The
// refused message was shown as not sent, so neither a relaunch nor a new owner sends it.
describe('a message an older host refused, across a relaunch', () => {
afterEach(() => {
cleanup()
setLocalRuntimeCapabilitiesForTests(null)
})
beforeEach(() => {
vi.clearAllMocks()
localStorage.clear()
setLocalRuntimeCapabilitiesForTests([])
})
it('is not sent on its own when the chat reopens on a moved fence', async () => {
mocks.call.mockResolvedValueOnce(refusal('agent_session_checkpoint_stale'))
const before = mount(1)
act(() => expect(before.result.current.send('hello')).toBe(true))
await waitFor(() => expect(before.result.current.outbox[0]?.lastFailure).toBeDefined())
before.unmount()
hostAccepts()
mount(3)
await settle()
expect(mocks.call).toHaveBeenCalledTimes(1)
})
it('is not sent when the fence moves while the chat is open; its Retry sends it once', async () => {
hostAccepts()
mocks.call.mockResolvedValueOnce(refusal('agent_session_checkpoint_stale'))
const { result, rerender } = mount(1)
act(() => expect(result.current.send('hello')).toBe(true))
await waitFor(() => expect(result.current.outbox[0]?.lastFailure).toBeDefined())
rerender({ fence: 3 })
await settle()
expect(mocks.call).toHaveBeenCalledTimes(1)
act(() => result.current.retry(result.current.outbox[0]!.clientMessageId))
await waitFor(() => expect(result.current.outbox).toHaveLength(0))
await settle()
expect(sentIds()).toHaveLength(2)
expect(sentIds()[1]).toBe(sentIds()[0])
})
})
// A message left on its way out is in doubt after a relaunch, whatever an earlier attempt failed
// with: the unconfirmed probe resends it under its id rather than holding it for a Retry.
describe('a message in doubt that carries an earlier failure', () => {
afterEach(() => {
cleanup()
setLocalRuntimeCapabilitiesForTests(null)
})
beforeEach(() => {
vi.clearAllMocks()
localStorage.clear()
setLocalRuntimeCapabilitiesForTests([AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY])
})
it('is probed and sent again, not held', async () => {
// What an earlier build left: on its way out, beside a failure it never cleared.
const inDoubt = {
...createStructuredAgentSessionOutboxEntry({
clientMessageId: 'op-in-doubt',
sessionId: SESSION,
text: 'hello',
attachments: [],
queuedAt: 1
}),
state: 'dispatching' as const,
lastAttemptAt: 2,
lastFailure: { kind: 'refused' as const, code: 'execution_owner_reconciling' as const }
}
writeOutbox(SESSION, [inDoubt])
hostAccepts()
const { result } = mount()
await waitFor(() => expect(result.current.outbox).toHaveLength(0), { timeout: 5000 })
expect(sentIds()).toEqual(['op-in-doubt'])
})
})
describe('a message whose send could not be saved before it went out', () => {
afterEach(() => {
cleanup()
vi.restoreAllMocks()
setLocalRuntimeCapabilitiesForTests(null)
})
beforeEach(() => {
vi.clearAllMocks()
localStorage.clear()
setLocalRuntimeCapabilitiesForTests([AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY])
})
it('waits for its Retry, and does not hold back the next message', async () => {
hostAccepts()
const save = localStorage.setItem.bind(localStorage)
// The send saves; the save that marks it on its way out fails.
const setItem = vi
.spyOn(localStorage, 'setItem')
.mockImplementationOnce(save)
.mockImplementationOnce(() => {
throw new Error('storage full')
})
const { result } = mount()
act(() => expect(result.current.send('first')).toBe(true))
await waitFor(() =>
expect(result.current.error).toBe('Message could not be saved to the outbox')
)
setItem.mockRestore()
const firstId = result.current.outbox[0]!.clientMessageId
act(() => expect(result.current.send('second')).toBe(true))
await waitFor(() => expect(result.current.outbox).toHaveLength(1))
await settle()
expect(mocks.call).toHaveBeenCalledOnce()
expect(sentIds()).not.toContain(firstId)
expect(noticeFor(result.current, firstId)?.onRetry).toBeDefined()
})
function failingWrites(failing: readonly number[]): void {
const save = localStorage.setItem.bind(localStorage)
let writes = 0
vi.spyOn(localStorage, 'setItem').mockImplementation((key: string, value: string) => {
writes += 1
if (failing.includes(writes)) {
throw new Error('storage full')
}
save(key, value)
})
}
it('stays held when the message behind it goes out, though no save recorded the hold', async () => {
hostAccepts()
// Both messages save; the save marking the first on its way out, and the one recording its
// hold, fail; storage then works again for the second.
failingWrites([3, 4])
const detached: { fence: number | null } = { fence: null }
const { result, rerender } = renderHook(
({ fence }: { fence: number | null }) =>
useStructuredAgentSessionOutbox({
sessionId: SESSION,
target: LOCAL_TARGET,
fence,
submissions: []
}),
{ initialProps: detached }
)
act(() => expect(result.current.send('first')).toBe(true))
act(() => expect(result.current.send('second')).toBe(true))
const [firstId, secondId] = result.current.outbox.map((entry) => entry.clientMessageId)
rerender({ fence: 1 })
await waitFor(() => expect(result.current.outbox).toHaveLength(1))
await settle()
expect(sentIds()).toEqual([secondId])
expect(result.current.outbox[0]).toMatchObject({
clientMessageId: firstId,
lastFailure: { kind: 'failed' }
})
expect(noticeFor(result.current, firstId!)?.onRetry).toBeDefined()
})
it('stays held across a relaunch once storage takes the hold', async () => {
hostAccepts()
// The send saves; the save marking it on its way out fails; the hold's own save goes through.
failingWrites([2])
const before = mount()
act(() => expect(before.result.current.send('first')).toBe(true))
await waitFor(() =>
expect(before.result.current.error).toBe('Message could not be saved to the outbox')
)
const firstId = before.result.current.outbox[0]!.clientMessageId
before.unmount()
vi.restoreAllMocks()
const after = mount()
await settle()
expect(mocks.call).not.toHaveBeenCalled()
expect(noticeFor(after.result.current, firstId)?.onRetry).toBeDefined()
})
})
// A host forgets an operation id a day after it was made and refuses it for good after that. An
// earlier attempt under that id may already be in the chat, so the message stays a held row that
// says so; only the user's Retry sends it, under a new id.
describe('a held message whose id expired', () => {
const DAY = 24 * 60 * 60 * 1000
const EXPIRED_WORDS = "Orca couldn't confirm what happened. Check the chat."
function expired() {
return {
ok: false,
refusal: {
code: 'agent_session_operation_expired',
message: 'Operation expired.',
details: { reason: 'operationExpired' }
}
}
}
afterEach(() => {
cleanup()
setLocalRuntimeCapabilitiesForTests(null)
})
beforeEach(() => {
vi.clearAllMocks()
localStorage.clear()
clearNativeChatDraftCacheForTests()
setLocalRuntimeCapabilitiesForTests([AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY])
})
it('stays a row that says it may be in the chat; its Retry sends it under a new id', async () => {
mocks.call.mockResolvedValueOnce(newerOrcaRefusal()).mockResolvedValueOnce(expired())
const before = mount()
act(() => expect(before.result.current.send('hello')).toBe(true))
await waitFor(() => expect(before.result.current.outbox[0]?.lastFailure).toBeDefined())
const keptId = sentIds()[0]!
act(() => before.result.current.retry(keptId))
await waitFor(() =>
expect(before.result.current.outbox[0]?.lastFailure).toMatchObject({
code: 'agent_session_operation_expired'
})
)
expect(before.result.current.error).toBeNull()
before.unmount()
hostAccepts()
const after = mount()
await settle()
expect(mocks.call).toHaveBeenCalledTimes(2)
const notice = noticeFor(after.result.current, keptId)
expect(notice?.text).toBe(EXPIRED_WORDS)
expect(notice?.onRetry).toBeDefined()
act(() => notice?.onRetry?.())
await waitFor(() => expect(after.result.current.outbox).toHaveLength(0))
expect(sentIds()).toHaveLength(3)
expect(sentIds()[2]).not.toBe(keptId)
})
it('stays in the chat, not the composer, when a relaunch resends it on its own', async () => {
// Quit mid-send two days ago: the first message was on its way, the second queued behind it.
const old = Date.now() - 2 * DAY
const inFlight = {
...createStructuredAgentSessionOutboxEntry({
clientMessageId: `${old}-${'1'.repeat(32)}`,
sessionId: SESSION,
text: 'first',
attachments: [],
queuedAt: old
}),
state: 'dispatching' as const,
lastAttemptAt: old
}
writeOutbox(SESSION, [inFlight])
mocks.call.mockResolvedValue(expired())
const { result } = renderHook(() =>
useStructuredAgentSessionOutbox({
sessionId: SESSION,
target: LOCAL_TARGET,
fence: 1,
submissions: [],
composerScopeKey: 'pane-1'
})
)
await waitFor(() => expect(mocks.call).toHaveBeenCalled(), { timeout: 5000 })
await waitFor(() =>
expect(result.current.outbox[0]?.lastFailure).toMatchObject({
code: 'agent_session_operation_expired'
})
)
await settle()
expect(readNativeChatDraftCache('pane-1')).toBe('')
expect(result.current.outbox).toMatchObject([
{ clientMessageId: inFlight.clientMessageId, state: 'queued' }
])
expect(noticeFor(result.current, inFlight.clientMessageId)?.text).toBe(EXPIRED_WORDS)
expect(sentIds()).toEqual([inFlight.clientMessageId])
})
})
@@ -9,8 +9,7 @@ import {
const UNCONFIRMED_PROBE_BASE_DELAY_MS = 1_000
/** No attempt ceiling: a transport outage outlives any fixed budget, and giving up
* restores the wedge this fixes. Growth caps the rate at one status query per 16s.
* A refusal that blocks the head still ends probing until a manual Retry (or, on an older
* host, a fence change), because the entry leaves `unconfirmed`. */
* A refusal still ends probing until a manual Retry, because the entry leaves `unconfirmed`. */
const UNCONFIRMED_PROBE_MAX_DELAY_MS = 16_000
/** Re-queues the entry holding the outbox in `unconfirmed`, with backoff, until the journal answers it. */
@@ -58,9 +57,14 @@ export function useStructuredAgentSessionOutboxUnconfirmedProbe(args: {
const timer = setTimeout(
() => {
probeAttemptsRef.current = { id: probeId, attempts: attempts + 1 }
const next = getStructuredAgentSessionOutbox(sessionId).map((entry) =>
entry.clientMessageId === probeId ? { ...entry, state: 'queued' as const } : entry
)
const next = getStructuredAgentSessionOutbox(sessionId).map((entry) => {
if (entry.clientMessageId !== probeId) {
return entry
}
// A saved failure would hold it for a Retry instead of resending it.
const { lastFailure: _probed, ...probed } = entry
return { ...probed, state: 'queued' as const }
})
commitStructuredAgentSessionOutbox(sessionId, next)
},
Math.min(UNCONFIRMED_PROBE_BASE_DELAY_MS * 2 ** attempts, UNCONFIRMED_PROBE_MAX_DELAY_MS)
@@ -120,7 +120,7 @@ describe('a Stop withdrawing what the host does not hold', () => {
]
expect(
withdrawUnsentStructuredAgentSessionOutboxEntries(entries, [pending('held')], null, null).map(
withdrawUnsentStructuredAgentSessionOutboxEntries(entries, [pending('held')], null).map(
(candidate) => candidate.clientMessageId
)
).toEqual(['held', 'refused'])
@@ -128,21 +128,19 @@ describe('a Stop withdrawing what the host does not hold', () => {
it('leaves every message that waits on its Retry, not only a refused one', () => {
const entries = [
entry('blocked', 'queued'),
{ ...entry('blocked', 'queued'), lastFailure: { kind: 'failed' as const } },
{ ...entry('retried-in-doubt', 'unconfirmed'), retryAfterUnknownSubmittedAt: 10 },
entry('probed-in-doubt', 'unconfirmed'),
entry('local', 'queued')
]
expect(
withdrawUnsentStructuredAgentSessionOutboxEntries(entries, [], 'blocked', null).map(
withdrawUnsentStructuredAgentSessionOutboxEntries(entries, [], null).map(
(candidate) => candidate.clientMessageId
)
).toEqual(['blocked', 'retried-in-doubt'])
expect(hasUnsentStructuredAgentSessionOutboxEntry(entries.slice(0, 2), [], 'blocked')).toBe(
false
)
expect(hasUnsentStructuredAgentSessionOutboxEntry(entries, [], 'blocked')).toBe(true)
expect(hasUnsentStructuredAgentSessionOutboxEntry(entries.slice(0, 2), [])).toBe(false)
expect(hasUnsentStructuredAgentSessionOutboxEntry(entries, [])).toBe(true)
})
it('keeps a message whose send failed for its Retry', async () => {
@@ -156,7 +154,7 @@ describe('a Stop withdrawing what the host does not hold', () => {
})
)
act(() => expect(result.current.send('first')).toBe(true))
await waitFor(() => expect(result.current.blockedClientMessageId).not.toBeNull())
await waitFor(() => expect(result.current.outbox[0]?.lastFailure).toEqual({ kind: 'failed' }))
act(() => result.current.withdrawUnsent())
@@ -132,7 +132,6 @@ describe('an outbox entry the host handed off as a queued draft', () => {
]
})
await waitFor(() => expect(view.result.current.outbox).toHaveLength(0))
expect(view.result.current.blockedClientMessageId).toBeNull()
expect(readNativeChatDraftCache('scope')).toBe('')
})
@@ -199,7 +198,6 @@ describe('an outbox entry the host handed off as a queued draft', () => {
await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(1))
await act(async () => new Promise((resolve) => setTimeout(resolve, 20)))
expect(view.result.current.outbox).toEqual([])
expect(view.result.current.blockedClientMessageId).toBeNull()
expect(view.result.current.error).toBeNull()
expect(readNativeChatDraftCache('scope')).toBe('')
act(() => {
@@ -315,7 +315,7 @@ async function attemptedQueueSend() {
{ initialProps: { capability: SUPPORTED } }
)
expect(view.result.current.send('follow-up')).toBe(true)
await waitFor(() => expect(view.result.current.blockedClientMessageId).not.toBeNull())
await waitFor(() => expect(view.result.current.outbox[0]?.lastFailure).toBeDefined())
expect(mocks.call.mock.calls[0]?.[2]?.delivery).toBe('queue-if-active')
mocks.call.mockImplementation(() => new Promise(() => {}))
return view
@@ -563,7 +563,6 @@ describe('useStructuredAgentSessionOutbox', () => {
expect(result.current.outbox).toHaveLength(1)
// Never sent, and never re-sent on its own: it waits for Retry and holds nothing up.
expect(result.current.outbox[0]?.state).toBe('rejected')
expect(result.current.blockedClientMessageId).toBeNull()
// Settled, not pending: the refused id never ran, so a Retry is a new operation.
const sentId: unknown = mocks.call.mock.calls[0]![2].envelope.clientOperationId
const retryId = result.current.outbox[0]!.clientMessageId
@@ -781,7 +780,6 @@ describe('useStructuredAgentSessionOutbox', () => {
// under the "delivery is unconfirmed" banner. The disposition tests pin its words.
await waitFor(() => expect(result.current.outbox[0]?.lastFailure?.kind).toBe('rejected'))
expect(result.current.outbox[0]?.state).toBe('rejected')
expect(result.current.blockedClientMessageId).toBeNull()
// Retry immediately, before the journal subscription can publish the rejected row.
act(() => result.current.retry(firstId))
@@ -8,7 +8,10 @@ import {
} from 'react'
import type { AgentJournalSubmission } from '../../../../shared/agent-session-journal-types'
import { createStructuredAgentSessionOperationId } from '../../../../shared/structured-agent-session-mutation'
import { admitStructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox'
import {
admitStructuredAgentSessionOutboxEntry,
structuredAgentSessionEntryHeldForRetry
} from '../../../../shared/structured-agent-session-outbox-admission'
import {
journalAnswersInFlightSend,
type StructuredAgentSessionSendDisposition
@@ -42,6 +45,7 @@ import {
type StructuredAgentSessionQueueDelivery
} from '../../../../shared/structured-agent-session-outbox-delivery'
import { retryStructuredAgentSessionOutboxEntry } from './structured-agent-session-outbox-retry'
import { useStructuredAgentSessionOutboxFailedHere } from './use-structured-agent-session-outbox-failed-here'
const NO_QUEUE_DELIVERY: StructuredAgentSessionQueueDelivery = {
capability: 'unsupported',
@@ -77,7 +81,7 @@ export function useStructuredAgentSessionOutbox(args: {
target
} = args
const { capability: queueCapability, enabled: queueEnabled } = queueDelivery
// What resends, unblocks and drops a send in flight besides a Retry or a new send; see the hook.
// What resends and drops a send in flight besides a Retry or a new send; see the hook.
const owner = useStructuredAgentSessionOutboxOwnerChange(target, fence)
const restoreWithdrawn = useStructuredAgentSessionWithdrawnRestore(sessionId, composerScopeKey)
// The outbox lives in the session's store, shared with every other writer; this view holds it
@@ -100,8 +104,9 @@ export function useStructuredAgentSessionOutbox(args: {
// "which entry" must never disagree: the journal can settle the tail while the head moves.
const inFlightIdRef = useRef<string | null>(null)
const dispatchGenerationRef = useRef(0)
const blockedIdRef = useRef<string | null>(null)
const [error, setError] = useState<string | null>(null)
const { failedHere, recordFailures, forget } =
useStructuredAgentSessionOutboxFailedHere(sessionId)
const [errorSession, setErrorSession] = useState(sessionId)
// Render-time reset (react.dev: adjusting state when a prop changes), so the
// old session's banner neither flashes for a frame nor resurrects on return.
@@ -113,7 +118,6 @@ export function useStructuredAgentSessionOutbox(args: {
useLayoutEffect(() => {
dispatchGenerationRef.current += 1
inFlightIdRef.current = null
blockedIdRef.current = null
}, [owner.ownerChange, owner.targetKey, sessionId])
useEffect(() => {
@@ -159,11 +163,12 @@ export function useStructuredAgentSessionOutbox(args: {
dispatchGenerationRef.current += 1
inFlightIdRef.current = null
}
if (blockedIdRef.current !== null && hostOwns.has(blockedIdRef.current)) {
blockedIdRef.current = null
setError(null)
} else if (
current.some((entry) => entry.state === 'unconfirmed' && hostOwns.has(entry.clientMessageId))
if (
current.some(
(entry) =>
(entry.state === 'unconfirmed' || structuredAgentSessionEntryHeldForRetry(entry)) &&
hostOwns.has(entry.clientMessageId)
)
) {
setError(null)
}
@@ -175,25 +180,28 @@ export function useStructuredAgentSessionOutbox(args: {
// Released here rather than in a `.finally`: the state write below is what re-runs the
// drain, so a later microtask would leave the queue with no trigger to move on.
inFlightIdRef.current = null
blockedIdRef.current = disposition.blockedClientMessageId
setError(disposition.error)
recordFailures(getStructuredAgentSessionOutbox(sessionId), disposition.entries)
commitStructuredAgentSessionOutbox(sessionId, disposition.entries)
},
[sessionId]
[recordFailures, sessionId]
)
const [drains, setDrains] = useState(0)
const drainAgain = useCallback(() => setDrains((count) => count + 1), [])
useEffect(() => {
const head = outbox[0]
// The shared outbox, not this render's copy: an effect earlier in this commit (an owner
// change's requeue, a reconcile) or a launch settlement may have written a newer one.
const current = getStructuredAgentSessionOutbox(sessionId)
const head = current[0]
if (!head || head.sessionId !== sessionId) {
return
}
// A launch settlement dispatches outside this hook's single-flight, so while its send is up
// nothing else may go out beside it and race it for the host's arrival order. Its writes land
// in the shared outbox, but it releases its marker after the last one, so drain again then.
const launching = outbox.find((entry) => entry.source === 'launch')
const launching = current.find((entry) => entry.source === 'launch')
const launchDispatch = launching
? getStructuredAgentLaunchPromptDispatch(
launching.sessionId,
@@ -205,33 +213,25 @@ export function useStructuredAgentSessionOutbox(args: {
void launchDispatch.then(drainAgain, drainAgain)
return
}
const admission = admitStructuredAgentSessionOutboxEntry(outbox, blockedIdRef.current)
const admission = admitStructuredAgentSessionOutboxEntry(current)
if (admission.state !== 'dispatch' || fence === null || inFlightIdRef.current !== null) {
return
}
const next = admission.entry
// A launch settlement may have already admitted this entry and cleared its in-flight marker
// before this render saw it; only dispatch while the shared outbox still holds it queued.
const current = getStructuredAgentSessionOutbox(sessionId)
const currentEntry = current.find((entry) => entry.clientMessageId === next.clientMessageId)
if (currentEntry?.state !== 'queued') {
return
}
// The request reads the capability; the entry keeps only what its first attempt sent.
const attempt = structuredAgentSessionEntryAttempt(currentEntry, {
const attempt = structuredAgentSessionEntryAttempt(next, {
capability: queueCapability,
enabled: queueEnabled
})
const dispatch = dispatchStructuredAgentSessionOutboxEntry({
next: attempt.wire,
persisted: current.map((entry) => (entry === currentEntry ? attempt.stored : entry)),
entries: current.map((entry) => (entry === next ? attempt.stored : entry)),
sessionId,
target,
fence,
dispatchGeneration: dispatchGenerationRef.current,
dispatchGenerationRef,
inFlightIdRef,
blockedIdRef,
setError,
applyDisposition,
createOperationId: structuredSessionOperationId
@@ -279,18 +279,14 @@ export function useStructuredAgentSessionOutbox(args: {
sessionId,
submissions,
queuedMessageIds,
blockedIdRef,
inFlightIdRef,
dispatchGenerationRef,
restoreWithdrawn
})
const retry = (clientMessageId: string): void => {
// Another message's Retry must not send the one the queue is held on.
if (blockedIdRef.current === clientMessageId) {
blockedIdRef.current = null
}
setError(null)
forget(clientMessageId)
retryStructuredAgentSessionOutboxEntry({
clientMessageId,
sessionId,
@@ -302,7 +298,7 @@ export function useStructuredAgentSessionOutbox(args: {
return {
outbox,
error,
blockedClientMessageId: blockedIdRef.current,
failedHere,
send,
retry,
withdrawUnsent
@@ -32,7 +32,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: mocks.operationId,
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -55,7 +55,6 @@ vi.mock('./use-structured-agent-session-read', () => ({
vi.mock('./use-structured-agent-session-outbox', () => ({
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -43,7 +43,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
mocks.outbox(args)
return {
outbox: [],
blockedClientMessageId: null,
error: null,
send: mocks.send,
retry: mocks.retry
@@ -34,7 +34,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: mocks.operationId,
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -20,7 +20,6 @@ const mocks = vi.hoisted(() => ({
let items: AgentJournalRenderItem[] = []
let submissions: AgentJournalSubmission[] = []
let outbox: StructuredAgentSessionOutboxEntry[] = []
let blockedClientMessageId: string | null = null
let fence = 3
vi.mock('@/runtime/structured-agent-session-client', () => ({
@@ -40,7 +39,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: () => `operation-${++mocks.operations}`,
useStructuredAgentSessionOutbox: () => ({
outbox,
blockedClientMessageId,
error: null,
send: vi.fn(),
retry: vi.fn(),
@@ -139,7 +137,6 @@ beforeEach(() => {
items = []
submissions = []
outbox = []
blockedClientMessageId = null
fence = 3
})
@@ -388,14 +385,12 @@ describe('Stop against a host that stops the conversation', () => {
})
it('is hidden with only a message that waits on its Retry', () => {
// A send that failed holds the queue until the user retries it.
outbox = [entry('queued')]
blockedClientMessageId = 'client-1'
// A send that failed waits, with its saved failure, until the user retries it.
outbox = [{ ...entry('queued'), lastFailure: { kind: 'failed' } }]
expect(render().result.current.canStop).toBe(false)
// A send the host restarted under is parked for the user, and the chat reads idle.
outbox = [{ ...entry('unconfirmed'), retryAfterUnknownSubmittedAt: -1 }]
blockedClientMessageId = null
submissions = [
submission({
dispatchState: 'unknown',
@@ -34,7 +34,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: () => 'operation-1',
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -54,7 +54,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
mocks.outboxArgs.push(args)
return {
outbox: outboxEntries,
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn(),
@@ -45,7 +45,6 @@ vi.mock('./use-structured-agent-session-outbox', () => ({
structuredSessionOperationId: mocks.operationId,
useStructuredAgentSessionOutbox: () => ({
outbox: [],
blockedClientMessageId: null,
error: null,
send: vi.fn(),
retry: vi.fn()
@@ -152,30 +152,16 @@ export function useStructuredAgentSession(args: {
transportState.turnId !== null ||
(stopsConversation &&
(transportState.isWorking ||
hasUnsentStructuredAgentSessionOutboxEntry(
outbox,
transportState.submissions,
outboxController.blockedClientMessageId
)))
hasUnsentStructuredAgentSessionOutboxEntry(outbox, transportState.submissions)))
// A queued send is a card, never a transcript bubble.
const isWorking = transportState.isWorking
const transcriptOutbox = useMemo(
() =>
outboxOutsideQueuedCards(
outbox,
queuedMessageIds,
isWorking,
outboxController.blockedClientMessageId,
{ capability: queueCapability, enabled: queueFollowUps }
),
[
isWorking,
outbox,
outboxController.blockedClientMessageId,
queueCapability,
queueFollowUps,
queuedMessageIds
]
outboxOutsideQueuedCards(outbox, queuedMessageIds, isWorking, {
capability: queueCapability,
enabled: queueFollowUps
}),
[isWorking, outbox, queueCapability, queueFollowUps, queuedMessageIds]
)
const messages = useStructuredAgentSessionMessages(
transportState.journalItems,
@@ -227,9 +213,9 @@ export function useStructuredAgentSession(args: {
loadOlder,
prompts,
outbox,
failedHere: outboxController.failedHere,
/** The journal's rows for sent messages, which carry a rejected message's whole fact. */
submissions: transportState.submissions,
blockedClientMessageId: outboxController.blockedClientMessageId,
// A message typed during a command queues behind it on the host.
send: (...input: Parameters<typeof outboxController.send>) =>
// Legacy: an older host refuses sends while a command runs; removable once those hosts age out.
@@ -19,11 +19,11 @@ export function useStructuredAgentSessionHostAcceptsSend(target: RuntimeClientTa
}
/**
* When an outbox treats its owner as changed: resending a send in flight under the same id,
* dropping that send's answer, and unblocking a refused head. An older host restarts the agent
* inside the send and refuses it, unrecorded, when that fails, so a new fence is its only word that
* another try may land. A host that accepts first records every send before it starts anything,
* so a moved fence means nothing there, and only a Retry or a new send goes out.
* When an outbox treats its owner as changed: resending a send in flight under the same id and
* dropping that send's answer. An older host restarts the agent inside the send, so a new fence is
* its only word that a send it never answered may land with the new owner. A send it refused was
* shown as not sent and waits for its Retry, as on every host. A host that accepts first records
* every send before it starts anything, so a moved fence means nothing there.
*/
export function useStructuredAgentSessionOutboxOwnerChange(
target: RuntimeClientTarget,
+11 -10
View File
@@ -3,10 +3,13 @@ import type { AgentSessionOwnerVerdict } from './agent-session-wire'
import { agentSessionOwnerVerdictAllowsFreshOperationId } from './agent-session-refusal-retry'
import { parseAgentSessionWriteFailure } from './agent-session-write-failure'
import {
admitStructuredAgentSessionOutboxEntry,
createStructuredAgentSessionOutboxEntry,
requeueStructuredAgentSessionSendRefusal
} from './structured-agent-session-outbox'
import {
admitStructuredAgentSessionOutboxEntry,
structuredAgentSessionEntryHeldForRetry
} from './structured-agent-session-outbox-admission'
import { disposeStructuredAgentSessionSendResult } from './structured-agent-session-send-disposition'
const entry = createStructuredAgentSessionOutboxEntry({
@@ -49,7 +52,7 @@ describe("the owner's verdict decides a retry's operation id only as a floor", (
expect(retried('exited')).toMatchObject({ clientMessageId: 'message-2', state: 'queued' })
})
it('keeps later sends behind the message after an exited refusal', () => {
it('holds the message for its Retry after an exited refusal, and sends later ones past it', () => {
const later = createStructuredAgentSessionOutboxEntry({
clientMessageId: 'message-3',
sessionId: 'session-1',
@@ -60,7 +63,6 @@ describe("the owner's verdict decides a retry's operation id only as a floor", (
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry, later],
entry,
blockedClientMessageId: null,
result: {
ok: false,
refusal: {
@@ -71,13 +73,12 @@ describe("the owner's verdict decides a retry's operation id only as a floor", (
},
createOperationId: () => 'message-2'
})
expect(disposition.blockedClientMessageId).toBe('message-2')
expect(
admitStructuredAgentSessionOutboxEntry(
disposition.entries,
disposition.blockedClientMessageId
)
).toMatchObject({ state: 'blocked', entry: { clientMessageId: 'message-2' } })
expect(disposition.entries[0]).toMatchObject({ clientMessageId: 'message-2', state: 'queued' })
expect(structuredAgentSessionEntryHeldForRetry(disposition.entries[0]!)).toBe(true)
expect(admitStructuredAgentSessionOutboxEntry(disposition.entries)).toMatchObject({
state: 'dispatch',
entry: { clientMessageId: 'message-3' }
})
})
it.each<AgentSessionOwnerVerdict | undefined>(['unverifiable', 'live', undefined])(
@@ -239,4 +239,23 @@ describe('a message the provider answered after the running turn', () => {
])
expect(membership.liveTurnKey).toBe('B')
})
// A message shown as not sent stays in the outbox, below every journal row.
it.each([
['states each row turn', THREAD],
['states no turn scope', null]
] as const)(
'puts a message shown as not sent in no turn, never the live one (host %s)',
(_host, scope) => {
const items = [user('u1', scope), turn('t1', 'u1', scope), user('u2', scope)]
const messages = [
...rows(items),
{ id: 'held', role: 'user' as const, unsent: true as const }
]
const membership = nativeChatTurnMembership(messages, { items, submissions: [] })
expect(membership.turnKeys).toEqual(['u1', 'u2', undefined])
// The send whose turn has not opened yet is live, not the message below it.
expect(membership.liveTurnKey).toBe('u2')
}
)
})
+32 -8
View File
@@ -125,29 +125,36 @@ export type NativeChatTurnMembership = {
drawOrder: readonly number[] | null
}
type NativeChatTurnMember = { id: string; role: NativeChatRole; unsent?: true }
/**
* Places each row in its turn. A user entry that anchors a turn, or is scoped to none, keys
* itself; one delivered into a running turn (a steer) takes that turn's key. A host that states no
* scope is read by journal order instead (`nativeChatJournalOrderTurnKeys`). `opensTurn` narrows
* which rows may key themselves, for rows that interleave a subagent's with the conversation's.
* which rows may key themselves, for rows that interleave a subagent's with the conversation's. A
* message shown as not sent is in no turn, so it is never the live one.
*/
export function nativeChatTurnMembership(
messages: readonly { id: string; role: NativeChatRole }[],
messages: readonly NativeChatTurnMember[],
journal?: NativeChatTurnJournal | null,
opensTurn: NativeChatOpensTurn = nativeChatUserRowOpensTurn
): NativeChatTurnMembership {
const opens: NativeChatOpensTurn = (message) => !isUnsent(message) && opensTurn(message)
if (!journal) {
const turnKeys = nativeChatRowTurnKeys(messages, null, opensTurn)
const turnKeys = withoutUnsent(messages, nativeChatRowTurnKeys(messages, null, opens))
return { turnKeys, liveTurnKey: newestUserTurnKey(messages, turnKeys), drawOrder: null }
}
const anchors = structuredAgentTurnAnchors(journal.items, journal.submissions)
const running = liveStructuredAgentSessionTurnScope(journal.items)
if (!hostStatesTurnScopes(journal.items)) {
const recordKeys = namedRecordKeys(journal.items, anchors)
const turnKeys = nativeChatRowTurnKeys(
const turnKeys = withoutUnsent(
messages,
nativeChatJournalOrderTurnKeys(journal.items, recordKeys),
opensTurn
nativeChatRowTurnKeys(
messages,
nativeChatJournalOrderTurnKeys(journal.items, recordKeys),
opens
)
)
const runningNamed = running.kind === 'turn' ? recordKeys.get(running.turnItemId) : null
return {
@@ -160,6 +167,9 @@ export function nativeChatTurnMembership(
const anchoring = new Set(anchors.values())
const scopes = new Map(journal.items.map((item) => [item.itemId, item.turnScope]))
const turnKeys = messages.map((message) => {
if (message.unsent === true) {
return undefined
}
const scope = scopes.get(message.id)
const turnKey =
scope?.kind === 'turn' && scope.turnItemId ? anchors.get(scope.turnItemId) : undefined
@@ -175,6 +185,18 @@ export function nativeChatTurnMembership(
}
}
function isUnsent(message: { id: string; role: NativeChatRole }): boolean {
return 'unsent' in message && message.unsent === true
}
/** An unsent row neither keys a turn nor inherits the one before it. */
function withoutUnsent(
messages: readonly NativeChatTurnMember[],
turnKeys: (string | undefined)[]
): (string | undefined)[] {
return turnKeys.map((turnKey, index) => (messages[index]?.unsent === true ? undefined : turnKey))
}
/** Each root record's anchor, or null for a record that names no opener (only older hosts write
* one): such a record owns no rows, which keep their position. */
function namedRecordKeys(
@@ -225,9 +247,11 @@ function commandTurnRunning(items: readonly AgentJournalRenderItem[]): boolean {
}
function newestUserTurnKey(
messages: readonly { role: NativeChatRole }[],
messages: readonly NativeChatTurnMember[],
turnKeys: readonly (string | undefined)[]
): string | undefined {
const index = messages.findLastIndex((message) => message.role === 'user')
const index = messages.findLastIndex(
(message) => message.role === 'user' && message.unsent !== true
)
return index === -1 ? undefined : turnKeys[index]
}
+3
View File
@@ -212,6 +212,9 @@ export type NativeChatMessage = AgentJournalProducerLinkage & {
sentAs?: AgentJournalMessageSendMode
/** Accepted but not yet handed to the agent: drawn after everything the agent has done. */
queued?: true
/** Shown as not sent, waiting for the user's Retry: in no turn, so a newer turn's bar and clock
* never land on it, and drawn after the conversation. */
unsent?: true
/** Set only by the structured projection, on rows the journal holds, and ranks
* them ahead of time. Terminal-backed messages never carry it, and worker reads strip it. */
journalPosition?: AgentJournalPosition
@@ -51,16 +51,13 @@ describe('a queued draft handed off under a fresh submission id', () => {
it('answers the send in flight and leaves nothing a Stop would withdraw', () => {
expect(journalAnswersInFlightSend([handOff('pending')], 'draft')).toBe(true)
expect(journalAnswersInFlightSend([handOff('pending')], null)).toBe(false)
expect(hasUnsentStructuredAgentSessionOutboxEntry([entry], [handOff('pending')], null)).toBe(
false
)
expect(hasUnsentStructuredAgentSessionOutboxEntry([entry], [handOff('pending')])).toBe(false)
})
it('settles a replayed send answered with the hand-off, with no notice', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
createOperationId: () => 'rotated',
result: {
ok: true,
@@ -70,6 +67,6 @@ describe('a queued draft handed off under a fresh submission id', () => {
value: { clientMessageId: 'hand-off', submission: handOff('rejected') }
}
})
expect(disposition).toEqual({ entries: [], error: null, blockedClientMessageId: null })
expect(disposition).toEqual({ entries: [], error: null })
})
})
@@ -5,6 +5,7 @@ import { isQueuedAgentJournalSubmission } from './agent-session-queued-submissio
import { collapseProviderRetryRuns } from './native-chat-provider-retry-runs'
import type { NativeChatMessage } from './native-chat-types'
import type { StructuredAgentSessionOutboxEntry } from './structured-agent-session-outbox'
import { structuredAgentSessionEntryHeldForRetry } from './structured-agent-session-outbox-admission'
import { reconcileStructuredAgentSessionOutboxWithQueue } from './structured-agent-session-draft-hand-off'
import { projectStructuredItemsToNativeChat } from './structured-agent-session-projection'
@@ -62,6 +63,9 @@ export function projectStructuredAgentSessionMessages(
source: 'transcript',
timestamp: entry.queuedAt,
blocks: entry.body.blocks,
...(entry.state === 'rejected' || structuredAgentSessionEntryHeldForRetry(entry)
? { unsent: true as const }
: {}),
// A send the journal recorded before refusing it keeps its place there.
...(recorded ? { journalPosition: agentJournalItemPosition(recorded) } : {})
}
@@ -0,0 +1,45 @@
// Which outbox entry goes out next, and which ones wait for the user.
import type { StructuredAgentSessionOutboxEntry } from './structured-agent-session-outbox'
/** A send the user was told did not go through, with a Retry: only that Retry sends it again. Read
* from the saved failure, so a relaunch holds it exactly as the refusal did. */
export function structuredAgentSessionEntryHeldForRetry(
entry: StructuredAgentSessionOutboxEntry
): boolean {
return entry.state === 'queued' && entry.lastFailure !== undefined
}
export type StructuredAgentSessionOutboxAdmission =
| { state: 'dispatch'; entry: StructuredAgentSessionOutboxEntry }
| { state: 'blocked'; entry: StructuredAgentSessionOutboxEntry }
| { state: 'idle'; entry: null }
/**
* What the queue does next. The drain and the Retry affordance both read it, so neither can
* disagree with the other about which entry is holding the queue.
*
* A `dispatching` entry is not a barrier: the host appended its journal row inside the
* per-session serialize chain before dispatching, so nothing behind it can overtake it, and
* waiting for its echo costs delivery of everything queued behind it. An `unconfirmed` entry
* is a barrier — sending past it would reorder around a message that may yet land. One the
* user was told did not go through is not: it lands only by its own Retry, so what the user
* sends after it goes out as they send it.
*/
export function admitStructuredAgentSessionOutboxEntry(
entries: readonly StructuredAgentSessionOutboxEntry[]
): StructuredAgentSessionOutboxAdmission {
for (const entry of entries) {
if (entry.state === 'rejected' || structuredAgentSessionEntryHeldForRetry(entry)) {
continue
}
if (entry.state === 'unconfirmed') {
return { state: 'blocked', entry }
}
// A queue send a Stop outlived goes again only on the user's Retry.
if (entry.state === 'queued') {
return { state: entry.outlivedStop === true ? 'blocked' : 'dispatch', entry }
}
}
return { state: 'idle', entry: null }
}
@@ -0,0 +1,103 @@
import { describe, expect, it } from 'vitest'
import type { AgentJournalSubmission } from './agent-session-journal-types'
import {
createStructuredAgentSessionOutboxEntry,
reconcileStructuredAgentSessionOutbox,
type StructuredAgentSessionOutboxEntry
} from './structured-agent-session-outbox'
import {
admitStructuredAgentSessionOutboxEntry,
structuredAgentSessionEntryHeldForRetry
} from './structured-agent-session-outbox-admission'
function entry(
clientMessageId: string,
patch: Partial<StructuredAgentSessionOutboxEntry> = {}
): StructuredAgentSessionOutboxEntry {
return {
...createStructuredAgentSessionOutboxEntry({
clientMessageId,
sessionId: 'session-1',
text: clientMessageId,
attachments: [],
queuedAt: 1
}),
...patch
}
}
const REFUSED = { kind: 'refused', code: 'agent_session_journal_unreadable' } as const
describe('a message held for its Retry', () => {
it('is one the user was told did not go through, and still queued', () => {
expect(structuredAgentSessionEntryHeldForRetry(entry('a', { lastFailure: REFUSED }))).toBe(true)
expect(
structuredAgentSessionEntryHeldForRetry(entry('a', { lastFailure: { kind: 'failed' } }))
).toBe(true)
expect(structuredAgentSessionEntryHeldForRetry(entry('a'))).toBe(false)
expect(
structuredAgentSessionEntryHeldForRetry(
entry('a', { state: 'rejected', lastFailure: { kind: 'rejected', reason: null } })
)
).toBe(false)
})
it('is passed over by the drain, which still stops on a message in doubt', () => {
const held = entry('held', { lastFailure: REFUSED })
expect(admitStructuredAgentSessionOutboxEntry([held])).toEqual({ state: 'idle', entry: null })
expect(admitStructuredAgentSessionOutboxEntry([held, entry('next')])).toMatchObject({
state: 'dispatch',
entry: { clientMessageId: 'next' }
})
expect(
admitStructuredAgentSessionOutboxEntry([
held,
entry('doubt', { state: 'unconfirmed' }),
entry('next')
])
).toMatchObject({ state: 'blocked', entry: { clientMessageId: 'doubt' } })
})
it('is released once the host shows it has the message after all', () => {
const submission: AgentJournalSubmission = {
clientMessageId: 'held',
fence: 1,
payloadFingerprint: 'fingerprint',
dispatchState: 'pending',
providerItemId: null,
reason: null,
submittedAt: 1,
resolvedAt: null
}
const [landed] = reconcileStructuredAgentSessionOutbox(
[entry('held', { lastFailure: REFUSED })],
[submission]
)
expect(landed?.state).toBe('dispatching')
expect(landed?.lastFailure).toBeUndefined()
})
// Any fresh word on where a message stands supersedes an earlier attempt's failure.
it('is in doubt, and holds the queue, once the host says it cannot tell whether it landed', () => {
const submission: AgentJournalSubmission = {
clientMessageId: 'held',
fence: 1,
payloadFingerprint: 'fingerprint',
dispatchState: 'unknown',
providerItemId: null,
reason: null,
submittedAt: 5,
resolvedAt: null
}
const reconciled = reconcileStructuredAgentSessionOutbox(
[entry('held', { lastAttemptAt: 2, lastFailure: REFUSED }), entry('next')],
[submission]
)
expect(reconciled[0]).toMatchObject({ state: 'unconfirmed' })
expect(reconciled[0]?.lastFailure).toBeUndefined()
expect(admitStructuredAgentSessionOutboxEntry(reconciled)).toMatchObject({
state: 'blocked',
entry: { clientMessageId: 'held' }
})
})
})
@@ -1,32 +1,28 @@
import type { AgentJournalSubmission } from './agent-session-journal-types'
import type { StructuredAgentSessionOutboxEntry } from './structured-agent-session-outbox'
import { structuredAgentSessionEntryHeldForRetry } from './structured-agent-session-outbox-admission'
import { handedOffQueuedMessageIds } from './structured-agent-session-draft-hand-off'
/** Whether only the user's Retry sends this entry again: a refused one, the one the drain stopped
* on, one a Stop outlived, or one in doubt the unconfirmed probe leaves alone. `NativeChatDeliveryRetry` offers it. */
function awaitsStructuredAgentSessionRetry(
entry: StructuredAgentSessionOutboxEntry,
blockedClientMessageId: string | null
): boolean {
/** Whether only the user's Retry sends this entry again: a rejected one, one whose send failed or
* was refused, one a Stop outlived, or one in doubt the unconfirmed probe leaves alone.
* `NativeChatDeliveryRetry` offers it. */
function awaitsStructuredAgentSessionRetry(entry: StructuredAgentSessionOutboxEntry): boolean {
return (
entry.state === 'rejected' ||
entry.clientMessageId === blockedClientMessageId ||
structuredAgentSessionEntryHeldForRetry(entry) ||
entry.outlivedStop === true ||
(entry.state === 'unconfirmed' && entry.retryAfterUnknownSubmittedAt !== null)
)
}
function unsentStructuredAgentSessionOutboxEntry(
submissions: readonly AgentJournalSubmission[],
blockedClientMessageId: string | null
submissions: readonly AgentJournalSubmission[]
): (entry: StructuredAgentSessionOutboxEntry) => boolean {
const held = handedOffQueuedMessageIds(submissions)
for (const submission of submissions) {
held.add(submission.clientMessageId)
}
return (entry) =>
!held.has(entry.clientMessageId) &&
!awaitsStructuredAgentSessionRetry(entry, blockedClientMessageId)
return (entry) => !held.has(entry.clientMessageId) && !awaitsStructuredAgentSessionRetry(entry)
}
/** A queue send that has gone out at least once and was not refused, in whatever state it now
@@ -63,10 +59,9 @@ function markedOutlivingStop(
export function withdrawUnsentStructuredAgentSessionOutboxEntries(
entries: readonly StructuredAgentSessionOutboxEntry[],
submissions: readonly AgentJournalSubmission[],
blockedClientMessageId: string | null,
inFlightClientMessageId: string | null
): StructuredAgentSessionOutboxEntry[] {
const unsent = unsentStructuredAgentSessionOutboxEntry(submissions, blockedClientMessageId)
const unsent = unsentStructuredAgentSessionOutboxEntry(submissions)
return entries
.filter(
(entry) =>
@@ -80,8 +75,7 @@ export function withdrawUnsentStructuredAgentSessionOutboxEntries(
/** Whether a Stop has something here to withdraw: a message that would still go out on its own. */
export function hasUnsentStructuredAgentSessionOutboxEntry(
entries: readonly StructuredAgentSessionOutboxEntry[],
submissions: readonly AgentJournalSubmission[],
blockedClientMessageId: string | null
submissions: readonly AgentJournalSubmission[]
): boolean {
return entries.some(unsentStructuredAgentSessionOutboxEntry(submissions, blockedClientMessageId))
return entries.some(unsentStructuredAgentSessionOutboxEntry(submissions))
}
@@ -5,7 +5,6 @@
import { describe, expect, it } from 'vitest'
import { structuredAgentSessionPayloadFingerprint } from './structured-agent-session-mutation'
import {
admitStructuredAgentSessionOutboxEntry,
createStructuredAgentSessionOutboxEntry,
parseStructuredAgentSessionOutboxEntry,
stageStructuredAgentSessionOutboxEntryForSend,
@@ -14,6 +13,7 @@ import {
type StructuredAgentSessionOutboxEntry,
type StructuredAgentSessionOutboxState
} from './structured-agent-session-outbox'
import { admitStructuredAgentSessionOutboxEntry } from './structured-agent-session-outbox-admission'
import { structuredAgentSessionEntryAttempt } from './structured-agent-session-outbox-delivery'
import { disposeStructuredAgentSessionSendResult } from './structured-agent-session-send-disposition'
import { withdrawUnsentStructuredAgentSessionOutboxEntries } from './structured-agent-session-outbox-stop-withdrawal'
@@ -104,7 +104,6 @@ describe('outbox queue delivery', () => {
at('sent-plain', 'dispatching', null)
],
[],
null,
null
)
// Its state is left to its answer: only the mark holds it back.
@@ -114,7 +113,7 @@ describe('outbox queue delivery', () => {
['probed', 'queued', true]
])
// The drain never admits the marked one it would otherwise send.
expect(admitStructuredAgentSessionOutboxEntry(next.slice(2), null)).toEqual({
expect(admitStructuredAgentSessionOutboxEntry(next.slice(2))).toEqual({
state: 'blocked',
entry: next[2]
})
@@ -138,16 +137,10 @@ describe('outbox queue delivery', () => {
disposeStructuredAgentSessionSendResult({
entries,
entry: attempt.wire,
blockedClientMessageId: null,
result: refusal,
createOperationId: () => operations.shift() ?? 'spent'
})
const stopped = withdrawUnsentStructuredAgentSessionOutboxEntries(
[staged],
[],
null,
'client-1'
)
const stopped = withdrawUnsentStructuredAgentSessionOutboxEntries([staged], [], 'client-1')
const withStop = answer(stopped)
const withoutStop = answer([staged])
expect(
@@ -156,6 +149,5 @@ describe('outbox queue delivery', () => {
expect(
withoutStop.entries.map((candidate) => [candidate.clientMessageId, candidate.state])
).toEqual([['rotated-2', 'rejected']])
expect(withStop.blockedClientMessageId).toBeNull()
})
})
+27 -41
View File
@@ -41,7 +41,8 @@ export type StructuredAgentSessionOutboxEntry = {
* what that request carries. */
sentDelivery?: 'queue-if-active' | null
/** Why the last attempt did not go through. Lives on the message so it goes when the message
* is sent again or delivered, instead of outliving it as a separate error. */
* is sent again or delivered, instead of outliving it as a separate error. On a `queued` entry
* it is also the hold (structured-agent-session-outbox-admission). */
lastFailure?: StructuredAgentSessionAttemptFailure
}
@@ -160,6 +161,19 @@ export function stageStructuredAgentSessionOutboxEntryForSend(
return { ...entry, state: 'dispatching', lastAttemptAt: now }
}
/** The host forgot this message's id, a day after it was made, and refuses it for good: only a new
* id sends it. The id was kept because an earlier attempt under it may already be in the chat; a
* first attempt's was replaced when it was refused. */
export function structuredAgentSessionEntryIdExpired(
entry: StructuredAgentSessionOutboxEntry
): boolean {
return (
entry.state === 'queued' &&
entry.lastFailure?.kind === 'refused' &&
entry.lastFailure.code === 'agent_session_operation_expired'
)
}
export function requeueStructuredAgentSessionSendRefusal(
entry: StructuredAgentSessionOutboxEntry,
refusal: AgentSessionWriteRefusal,
@@ -168,7 +182,7 @@ export function requeueStructuredAgentSessionSendRefusal(
): StructuredAgentSessionOutboxEntry {
const refusalSettled = agentSessionRefusalOperationState(refusal.code) === 'settled-rejected'
// An exited owner runs nothing under the old id, so a new one can't collide; the message still
// holds the head, since nothing recorded it.
// waits for its Retry, since nothing recorded it.
const ownerExited =
refusal.code === 'agent_session_ownership_unknown' &&
agentSessionOwnerVerdictAllowsFreshOperationId(refusal.details?.ownerVerdict)
@@ -181,8 +195,8 @@ export function requeueStructuredAgentSessionSendRefusal(
return { ...entry, state: 'queued' }
}
// Only here may the id rotate: an earlier attempt under this id, or one whose delivery was in
// doubt, may have landed, so those stay queued behind the block. Only a settled refusal proves
// the message never landed and so releases the queue.
// doubt, may have landed, so those keep it. Only a settled refusal proves the message never
// landed.
return {
...entry,
clientMessageId: createOperationId(),
@@ -209,7 +223,12 @@ export function reconcileStructuredAgentSessionOutbox(
return []
}
if (submission?.dispatchState === 'pending') {
return entry.state === 'dispatching' ? [entry] : [{ ...entry, state: 'dispatching' as const }]
if (entry.state === 'dispatching') {
return [entry]
}
// The host has it, so no failure of an earlier attempt describes it now.
const { lastFailure: _landed, ...landed } = entry
return [{ ...landed, state: 'dispatching' as const }]
}
// Accepted, then not delivered — the agent never started, or its start was refused. The text
// and why stay here for the user's Retry, and nothing queues behind it. `unconfirmed` is how a
@@ -231,47 +250,14 @@ export function reconcileStructuredAgentSessionOutbox(
entry.retryAfterUnknownSubmittedAt !== -1 &&
entry.retryAfterUnknownSubmittedAt !== submission.submittedAt
) {
return [{ ...entry, state: 'unconfirmed' as const }]
// In doubt now, not failed: the probe's resend decides it, as for any unconfirmed send.
const { lastFailure: _superseded, ...inDoubt } = entry
return [{ ...inDoubt, state: 'unconfirmed' as const }]
}
return [entry]
})
}
export type StructuredAgentSessionOutboxAdmission =
| { state: 'dispatch'; entry: StructuredAgentSessionOutboxEntry }
| { state: 'blocked'; entry: StructuredAgentSessionOutboxEntry }
| { state: 'idle'; entry: null }
/**
* What the queue does next. The drain and the Retry affordance both read it, so neither can
* disagree with the other about which entry is holding the queue.
*
* A `dispatching` entry is not a barrier: the host appended its journal row inside the
* per-session serialize chain before dispatching, so nothing behind it can overtake it, and
* waiting for its echo costs delivery of everything queued behind it. An `unconfirmed` entry,
* or one the user must act on, is a barrier — sending past either would reorder around a
* message that may yet land.
*/
export function admitStructuredAgentSessionOutboxEntry(
entries: readonly StructuredAgentSessionOutboxEntry[],
blockedClientMessageId: string | null
): StructuredAgentSessionOutboxAdmission {
for (const entry of entries) {
// It can no longer land, so nothing it could be reordered around; it waits for Retry.
if (entry.state === 'rejected') {
continue
}
if (entry.state === 'unconfirmed' || entry.clientMessageId === blockedClientMessageId) {
return { state: 'blocked', entry }
}
// A queue send a Stop outlived goes again only on the user's Retry.
if (entry.state === 'queued') {
return { state: entry.outlivedStop === true ? 'blocked' : 'dispatch', entry }
}
}
return { state: 'idle', entry: null }
}
export function parseStructuredAgentSessionOutboxEntry(
value: unknown,
sessionId: string
@@ -22,6 +22,7 @@ import {
reconcileStructuredAgentSessionOutbox,
type StructuredAgentSessionOutboxEntry
} from './structured-agent-session-outbox'
import { structuredAgentSessionEntryHeldForRetry } from './structured-agent-session-outbox-admission'
const entry: StructuredAgentSessionOutboxEntry = createStructuredAgentSessionOutboxEntry({
clientMessageId: 'client-1',
@@ -59,7 +60,6 @@ function notice(reason: string | null, rejection?: AgentSessionFailureFact): str
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: rejectedWith(reason, rejection ? { rejection } : {}),
createOperationId: () => 'unused'
})
@@ -76,7 +76,6 @@ describe('a queued draft answer', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: {
ok: true,
replayed: false,
@@ -97,7 +96,6 @@ describe('a queued draft answer', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: {
ok: true,
replayed: true,
@@ -158,15 +156,10 @@ describe('what a rejection shows the user', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: rejectedWith(DISPATCH_REJECTED_CANCELLED, { rejection, replayed }),
createOperationId: () => 'unused'
})
expect(disposition).toEqual({
entries: [],
error: null,
blockedClientMessageId: null
})
expect(disposition).toEqual({ entries: [], error: null })
}
})
@@ -180,7 +173,6 @@ describe('what a rejection shows the user', () => {
disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result,
createOperationId: () => 'unused'
}).error
@@ -263,7 +255,6 @@ describe('what a refusal shows the user', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: {
ok: false,
refusal: {
@@ -275,7 +266,7 @@ describe('what a refusal shows the user', () => {
})
expect(disposition.error).toBeNull()
expect(disposition.blockedClientMessageId).toBe(entry.clientMessageId)
expect(structuredAgentSessionEntryHeldForRetry(disposition.entries[0]!)).toBe(true)
expect(disposition.entries).toMatchObject([
{
clientMessageId: entry.clientMessageId,
@@ -294,7 +285,6 @@ describe('what a refusal shows the user', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: {
ok: false,
refusal: {
@@ -323,7 +313,6 @@ describe('what a refusal shows the user', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result: rejectedWith('An image on this message is empty, so the message was not sent.', {
rejection: {
kind: 'attachmentInvalid',
@@ -345,7 +334,6 @@ describe('what a refusal shows the user', () => {
const disposition = disposeStructuredAgentSessionSendFailure({
entries: [entry],
entry,
blockedClientMessageId: null,
cause: new Error('socket hang up: ECONNRESET 10.0.0.2:443'),
isDeliveryUnknown: () => false
})
@@ -367,7 +355,6 @@ describe('what a refusal shows the user', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [refused],
entry: refused,
blockedClientMessageId: null,
result,
createOperationId: () => 'unused'
})
@@ -381,25 +368,47 @@ describe('ambiguous operation refusals', () => {
it.each([
{ ...entry, state: 'unconfirmed' as const, lastAttemptAt: 10 },
{ ...entry, state: 'queued' as const, lastAttemptAt: 10, retryAfterUnknownSubmittedAt: 10 }
])('never rotates $state operation after its host tombstone expires', (ambiguous) => {
])(
'never rotates $state operation after its host tombstone expires; holds it for its Retry',
(ambiguous) => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [ambiguous],
entry: ambiguous,
result: {
ok: false,
refusal: {
code: 'agent_session_operation_expired',
message: 'Operation expired.'
}
},
createOperationId: () => 'fresh-id'
})
// An earlier attempt under the kept id may have landed, so no new id goes out on its own.
expect(disposition.entries).toMatchObject([
{
clientMessageId: entry.clientMessageId,
state: 'queued',
lastFailure: { kind: 'refused', code: 'agent_session_operation_expired' }
}
])
expect(structuredAgentSessionEntryHeldForRetry(disposition.entries[0]!)).toBe(true)
expect(disposition.error).toBeNull()
}
)
it('rotates a first attempt the host calls expired, which nothing can have delivered', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [ambiguous],
entry: ambiguous,
blockedClientMessageId: null,
entries: [entry],
entry,
result: {
ok: false,
refusal: {
code: 'agent_session_operation_expired',
message: 'Operation expired.'
}
refusal: { code: 'agent_session_operation_expired', message: 'Operation expired.' }
},
createOperationId: () => 'fresh-id'
})
expect(disposition.entries).toMatchObject([
{ clientMessageId: entry.clientMessageId, state: 'queued' }
])
expect(disposition.blockedClientMessageId).toBe(entry.clientMessageId)
expect(disposition.entries).toMatchObject([{ clientMessageId: 'fresh-id', state: 'rejected' }])
})
it('parks a recovered missing submission without polling forever', () => {
@@ -417,7 +426,6 @@ describe('ambiguous operation refusals', () => {
const disposition = disposeStructuredAgentSessionSendResult({
entries: [entry],
entry,
blockedClientMessageId: null,
result,
createOperationId: () => 'unused'
})
@@ -34,15 +34,11 @@ export type StructuredAgentSessionSendDisposition = {
entries: StructuredAgentSessionOutboxEntry[]
/** Only for an outcome with no entry left to carry it; a kept entry holds its own failure. */
error: string | null
/** The entry the queue is stuck on, or null when nothing blocks it. Always the
* next value, never "unchanged": the caller assigns it verbatim. */
blockedClientMessageId: string | null
}
type SendDispositionInput = {
entries: readonly StructuredAgentSessionOutboxEntry[]
entry: StructuredAgentSessionOutboxEntry
blockedClientMessageId: string | null
}
function replaceEntryState(
@@ -200,10 +196,8 @@ export function disposeStructuredAgentSessionSendRefusal(
createOperationId: () => string
}
): StructuredAgentSessionSendDisposition {
const refusedIndex = input.entries.findIndex(
(candidate) => candidate.clientMessageId === input.entry.clientMessageId
)
const entries = input.entries.map((candidate) =>
// The refusal saved on a message it keeps `queued` is what holds it for the user's Retry.
const entries: StructuredAgentSessionOutboxEntry[] = input.entries.map((candidate) =>
candidate.clientMessageId === input.entry.clientMessageId
? withLastFailure(
requeueStructuredAgentSessionSendRefusal(
@@ -216,18 +210,7 @@ export function disposeStructuredAgentSessionSendRefusal(
)
: candidate
)
const refused = entries[refusedIndex]
return {
entries,
error: null,
// Read back by index rather than from the input: a refusal can rotate the id, and the
// refused entry is not always the head now that an admitted one no longer holds the queue.
// A rejected one holds nothing: it can no longer land, and it keeps its own Retry.
blockedClientMessageId:
!refused || refused.state === 'rejected'
? input.blockedClientMessageId
: refused.clientMessageId
}
return { entries, error: null }
}
export function disposeStructuredAgentSessionSendResult(
@@ -249,8 +232,7 @@ export function disposeStructuredAgentSessionSendResult(
// and the draft card — not this queue — carries any later refusal.
return {
entries: dropEntry(input),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
const submission = result.value.submission
@@ -259,22 +241,19 @@ export function disposeStructuredAgentSessionSendResult(
if (submission.queuedMessageId === input.entry.clientMessageId) {
return {
entries: dropEntry(input),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
if (refusedRedelivery(input.entry, submission)) {
return {
entries: dropEntry(input),
error: 'Message delivery is unconfirmed and Orca will not send it again',
blockedClientMessageId: input.blockedClientMessageId
error: 'Message delivery is unconfirmed and Orca will not send it again'
}
}
if (submission.dispatchState === 'accepted') {
return {
entries: dropEntry(input),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
// A Stop's withdrawal failed nothing, first reply or replay: the entry leaves as the reconcile
@@ -285,8 +264,7 @@ export function disposeStructuredAgentSessionSendResult(
) {
return {
entries: dropEntry(input),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
if (submission.dispatchState === 'rejected') {
@@ -296,8 +274,7 @@ export function disposeStructuredAgentSessionSendResult(
'rejected',
structuredAgentSessionRejectedFailure(submission)
),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
if (submission.dispatchState === 'unknown' && submission.recovered) {
@@ -307,8 +284,7 @@ export function disposeStructuredAgentSessionSendResult(
? { ...candidate, state: 'unconfirmed', retryAfterUnknownSubmittedAt: -1 }
: candidate
),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
// `pending` is the host saying the message was written and is awaiting the
@@ -322,8 +298,7 @@ export function disposeStructuredAgentSessionSendResult(
input,
submission.dispatchState === 'unknown' ? 'unconfirmed' : 'dispatching'
),
error: null,
blockedClientMessageId: input.blockedClientMessageId
error: null
}
}
@@ -336,13 +311,11 @@ export function disposeStructuredAgentSessionSendFailure(
const failure = classifyStructuredAgentSessionSendFailure(input.cause, input.isDeliveryUnknown)
const deliveryUnknown = failure === 'delivery-unknown'
return {
// An unconfirmed entry's Retry row already says delivery is unconfirmed.
// An unconfirmed entry's Retry row already says delivery is unconfirmed; the saved failure
// holds the other for its Retry.
entries: deliveryUnknown
? replaceEntryState(input, 'unconfirmed')
: replaceEntryState(input, 'queued', { kind: 'failed' }),
error: null,
blockedClientMessageId: deliveryUnknown
? input.blockedClientMessageId
: input.entry.clientMessageId
error: null
}
}