diff --git a/mobile/src/session/mobile-structured-agent-session-send.ts b/mobile/src/session/mobile-structured-agent-session-send.ts index 82d237d4d4aa..feec2d5a1636 100644 --- a/mobile/src/session/mobile-structured-agent-session-send.ts +++ b/mobile/src/session/mobile-structured-agent-session-send.ts @@ -31,6 +31,8 @@ export async function sendMobileStructuredAgentSessionMessage(input: { attachments: readonly (StructuredAgentSessionAttachment & { contentFingerprint?: string })[] deadline?: number onError: (message: string) => void + /** The host answered `pending`; its settlement arrives on the stream, keyed by this id. */ + onAwaitingSettlement?: (clientMessageId: string) => void }): Promise { const timeoutMs = timeoutForDeadline(input.deadline) if (timeoutMs === null) { @@ -100,6 +102,9 @@ export async function sendMobileStructuredAgentSessionMessage(input: { timeoutMs }) const delivery = mobileStructuredSendDelivery(result, operation.retained) + if (result.status === 'accepted' && result.value.submission?.dispatchState === 'pending') { + input.onAwaitingSettlement?.(operation.operationId) + } if (delivery.operationIdSpent) { try { await clearMobileStructuredSendOperation({ diff --git a/mobile/src/session/use-mobile-structured-agent-session-late-rejection.test.tsx b/mobile/src/session/use-mobile-structured-agent-session-late-rejection.test.tsx new file mode 100644 index 000000000000..56f7508b301e --- /dev/null +++ b/mobile/src/session/use-mobile-structured-agent-session-late-rejection.test.tsx @@ -0,0 +1,243 @@ +// A send the host answered `pending` can still be rejected later, when the child holding it is +// stopped before it ever wrote it. Mobile has no outbox, so it must say so the way it says an +// immediate rejection, once, and a resend of the same text must be a new message. + +import { createElement } from 'react' +import { act, create, type ReactTestRenderer } from 'react-test-renderer' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { AgentJournalSubmission } from '../../../src/shared/agent-session-journal-types' +import type { AgentSessionSubscribeEvent } from '../../../src/shared/agent-session-wire' +import { DISPATCH_REJECTED_CANCELLED } from '../../../src/shared/structured-agent-session-dispatch-rejection' +import type { RpcClient } from '../transport/rpc-client' +import { resetMobileStructuredSendOperationJournalForTests } from './mobile-structured-send-operation-journal' +import { useMobileStructuredAgentSession } from './use-mobile-structured-agent-session' + +const asyncStorage = vi.hoisted(() => ({ + getItem: vi.fn(), + setItem: vi.fn(), + removeItem: vi.fn() +})) + +vi.mock('@react-native-async-storage/async-storage', () => ({ default: asyncStorage })) + +const REASON = 'Claude never finished starting, so Orca stopped it. Your message was not sent.' + +type SendParams = { envelope: { clientOperationId: string; payloadFingerprint: string } } + +function ok(result: unknown) { + return { id: 'rpc-1', ok: true as const, result, _meta: { runtimeId: 'runtime-1' } } +} + +function isSendParams(value: unknown): value is SendParams { + return typeof value === 'object' && value !== null && 'envelope' in value +} + +function submission( + params: SendParams, + dispatchState: AgentJournalSubmission['dispatchState'], + reason: string | null = null +): AgentJournalSubmission { + return { + clientMessageId: params.envelope.clientOperationId, + fence: 3, + payloadFingerprint: params.envelope.payloadFingerprint, + dispatchState, + providerItemId: null, + reason, + submittedAt: 10, + resolvedAt: dispatchState === 'pending' ? null : 11 + } +} + +function sendResult(params: SendParams, dispatchState: AgentJournalSubmission['dispatchState']) { + const row = submission(params, dispatchState, dispatchState === 'rejected' ? REASON : null) + return ok({ + ok: true, + replayed: false, + fence: 3, + cursor: { epoch: 'epoch-1', sequence: 1 }, + value: { clientMessageId: row.clientMessageId, submission: row } + }) +} + +function snapshotEvent(submissions: AgentJournalSubmission[] = []): AgentSessionSubscribeEvent { + return { + type: 'snapshot', + sessionId: 'session-1', + fence: 3, + page: { + sessionId: 'session-1', + epoch: 'epoch-1', + fence: 3, + direction: 'tail', + items: [], + removedItemIds: [], + submissions, + window: { oldest: null, newest: null, nextCursor: { epoch: 'epoch-1', sequence: 0 } }, + liveCursor: { epoch: 'epoch-1', sequence: 0 }, + hasOlder: false, + hasNewer: false + } + } +} + +describe('a mobile send the host rejects after answering pending', () => { + let renderer: ReactTestRenderer | null = null + let hook: ReturnType | null = null + let listener: ((value: unknown) => void) | null = null + let storedOperations: Map + const onSendError = vi.fn() + const sendRequest = vi.fn() + const client: RpcClient = { + sendRequest, + subscribe: (_method, _params, onData) => { + listener = onData + return () => {} + }, + updateTerminalSubscriptionViewport: () => {}, + getState: () => 'connected', + getReconnectAttempt: () => 0, + getLastConnectedAt: () => null, + onStateChange: () => () => {}, + notifyForeground: () => {}, + close: () => {} + } + + function Harness(): null { + hook = useMobileStructuredAgentSession({ + client, + sessionId: 'session-1', + sourceIdentity: 'host-a\0workspace-a', + enabled: true, + connected: true, + agent: 'claude', + onSendError + }) + return null + } + + async function mountSession(): Promise { + act(() => { + renderer = create(createElement(Harness)) + }) + await vi.waitFor(() => expect(listener).toEqual(expect.any(Function))) + act(() => listener?.(snapshotEvent())) + } + + function sends(): SendParams[] { + return sendRequest.mock.calls.flatMap(([method, params]) => + method === 'agentSession.send' && isSendParams(params) ? [params] : [] + ) + } + + function answerSends( + answer: (params: SendParams) => ReturnType | Promise> + ): void { + sendRequest.mockImplementation(async (method, params) => { + if (method !== 'agentSession.send' || !isSendParams(params)) { + return method === 'agentSession.options' ? ok({ models: [], current: {} }) : ok({}) + } + return answer(params) + }) + } + + beforeEach(() => { + vi.clearAllMocks() + resetMobileStructuredSendOperationJournalForTests() + storedOperations = new Map() + asyncStorage.getItem.mockImplementation( + async (key: string) => storedOperations.get(key) ?? null + ) + asyncStorage.setItem.mockImplementation(async (key: string, value: string) => { + storedOperations.set(key, value) + }) + asyncStorage.removeItem.mockImplementation(async (key: string) => { + storedOperations.delete(key) + }) + answerSends((params) => sendResult(params, 'pending')) + }) + + afterEach(() => { + act(() => renderer?.unmount()) + renderer = null + hook = null + listener = null + }) + + it('shows the reason once, and a resend of the same text is a new message', async () => { + await mountSession() + await act(async () => { + expect(await hook!.sendWithOutcome('held while starting')).toBe('accepted') + }) + const first = sends()[0]! + + act(() => listener?.(snapshotEvent([submission(first, 'rejected', REASON)]))) + act(() => listener?.(snapshotEvent([submission(first, 'rejected', REASON)]))) + + expect(onSendError).toHaveBeenCalledOnce() + expect(onSendError).toHaveBeenCalledWith(REASON) + await act(async () => { + expect(await hook!.sendWithOutcome('held while starting')).toBe('accepted') + }) + expect(sends()[1]!.envelope.clientOperationId).not.toBe(first.envelope.clientOperationId) + }) + + it('shows the reason when the stream rejects the send before its answer arrives', async () => { + await mountSession() + let answer: () => void = () => {} + answerSends( + (params) => + new Promise((resolve) => { + answer = () => resolve(sendResult(params, 'pending')) + }) + ) + let sent: Promise = Promise.resolve() + act(() => { + sent = hook!.sendWithOutcome('held while starting') + }) + await vi.waitFor(() => expect(sends()).toHaveLength(1)) + + act(() => listener?.(snapshotEvent([submission(sends()[0]!, 'rejected', REASON)]))) + expect(onSendError).not.toHaveBeenCalled() + await act(async () => { + answer() + await sent + }) + + expect(onSendError).toHaveBeenCalledOnce() + expect(onSendError).toHaveBeenCalledWith(REASON) + }) + + it('reports a rejection answered on the spot once, even when the stream repeats it', async () => { + answerSends((params) => sendResult(params, 'rejected')) + await mountSession() + await act(async () => { + expect(await hook!.sendWithOutcome('refused')).toBe('rejected') + }) + + act(() => listener?.(snapshotEvent([submission(sends()[0]!, 'rejected', REASON)]))) + + expect(onSendError).toHaveBeenCalledOnce() + expect(onSendError).toHaveBeenCalledWith(REASON) + }) + + it('stays quiet for a send the user cancelled and for one the provider took', async () => { + await mountSession() + await act(async () => { + await hook!.sendWithOutcome('never mind') + await hook!.sendWithOutcome('answered') + }) + const [cancelled, taken] = sends() + + act(() => + listener?.( + snapshotEvent([ + submission(cancelled!, 'rejected', DISPATCH_REJECTED_CANCELLED), + submission(taken!, 'accepted') + ]) + ) + ) + + expect(onSendError).not.toHaveBeenCalled() + }) +}) diff --git a/mobile/src/session/use-mobile-structured-agent-session-prompt-cancel.test.tsx b/mobile/src/session/use-mobile-structured-agent-session-prompt-cancel.test.tsx index 48b44e78ff5d..798c51c8c8b8 100644 --- a/mobile/src/session/use-mobile-structured-agent-session-prompt-cancel.test.tsx +++ b/mobile/src/session/use-mobile-structured-agent-session-prompt-cancel.test.tsx @@ -32,7 +32,7 @@ vi.mock('./use-mobile-structured-prompt-responses', () => ({ }) })) vi.mock('./use-mobile-structured-send-operation-reconciliation', () => ({ - useMobileStructuredSendOperationReconciliation: vi.fn() + useMobileStructuredSendOperationReconciliation: () => vi.fn() })) import { useMobileStructuredAgentSession } from './use-mobile-structured-agent-session' diff --git a/mobile/src/session/use-mobile-structured-agent-session.ts b/mobile/src/session/use-mobile-structured-agent-session.ts index 2a24d38a38e7..4991c25c162f 100644 --- a/mobile/src/session/use-mobile-structured-agent-session.ts +++ b/mobile/src/session/use-mobile-structured-agent-session.ts @@ -99,7 +99,7 @@ export function useMobileStructuredAgentSession(args: { useEffect(() => () => operationIdsRef.current.clear(), []) const stateArgs = { client, sessionId, sessionKey, enabled, connected } const { state, stateRef, loadingOlder, loadEarlier } = useMobileStructuredAgentState(stateArgs) - useMobileStructuredSendOperationReconciliation(state.submissions) + const awaitSend = useMobileStructuredSendOperationReconciliation(state.submissions, onSendError) const mutate = useCallback( async ( @@ -219,11 +219,13 @@ export function useMobileStructuredAgentSession(args: { text, attachments: sendAttachments, deadline, - onError: onSendError + onError: onSendError, + onAwaitingSettlement: awaitSend }) }, [ agent, + awaitSend, callerIdentity, client, conversationCommands, @@ -233,7 +235,8 @@ export function useMobileStructuredAgentSession(args: { optionSnapshot, sessionId, sessionKey, - setStructuredOption + setStructuredOption, + stateRef ] ) const { groupedDraft, respondPermission, respondQuestion } = useMobileStructuredPromptResponses({ diff --git a/mobile/src/session/use-mobile-structured-send-operation-reconciliation.ts b/mobile/src/session/use-mobile-structured-send-operation-reconciliation.ts index 1b18613bc7c9..700615b681ae 100644 --- a/mobile/src/session/use-mobile-structured-send-operation-reconciliation.ts +++ b/mobile/src/session/use-mobile-structured-send-operation-reconciliation.ts @@ -1,11 +1,46 @@ -import { useEffect } from 'react' +import { useCallback, useEffect, useRef } from 'react' import type { AgentJournalSubmission } from '../../../src/shared/agent-session-journal-types' +import { settleAwaitedStructuredAgentSessionSends } from '../../../src/shared/structured-agent-session-send-disposition' import { clearMobileStructuredSettledSendOperations } from './mobile-structured-send-operation-journal' +/** Settles sends from the stream: releases their ids, and reports a send the host answered + * `pending` and later rejected, with the notice an immediate rejection gets. */ export function useMobileStructuredSendOperationReconciliation( - submissions: readonly AgentJournalSubmission[] -): void { + submissions: readonly AgentJournalSubmission[], + onSendError: (message: string) => void +): (clientMessageId: string) => void { + const awaitedRef = useRef(new Set()) + const submissionsRef = useRef(submissions) + const onSendErrorRef = useRef(onSendError) useEffect(() => { + onSendErrorRef.current = onSendError + }, [onSendError]) + + const reportSettled = useCallback((current: readonly AgentJournalSubmission[]) => { + const { settled, notices } = settleAwaitedStructuredAgentSessionSends( + awaitedRef.current, + current + ) + for (const clientMessageId of settled) { + awaitedRef.current.delete(clientMessageId) + } + for (const notice of notices) { + onSendErrorRef.current(notice) + } + }, []) + + useEffect(() => { + submissionsRef.current = submissions + reportSettled(submissions) void clearMobileStructuredSettledSendOperations({ submissions }).catch(() => undefined) - }, [submissions]) + }, [reportSettled, submissions]) + + // Registers a send to await; the stream can settle it before its RPC answer arrives. + return useCallback( + (clientMessageId: string) => { + awaitedRef.current.add(clientMessageId) + reportSettled(submissionsRef.current) + }, + [reportSettled] + ) } diff --git a/src/main/claude/claude-structured-session-close.ts b/src/main/claude/claude-structured-session-close.ts index 2d1b1f255d85..c7bed7d04b5a 100644 --- a/src/main/claude/claude-structured-session-close.ts +++ b/src/main/claude/claude-structured-session-close.ts @@ -18,6 +18,10 @@ import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session import { closeProcessRegistry } from '../../shared/child-process/close-process-registry' import { retireClaudeDispatchWaiters } from './claude-structured-dispatch' import { settledClaudeTurnEndLeaf } from './claude-structured-resume-point' +import { + CLAUDE_STARTUP_ABANDONED_REJECTION, + failClaudeStartupGate +} from './claude-structured-session-startup-gate' /** The root's own exit was seen first-hand; only its descendants went unverified. */ export function claudeRootExitObserved( @@ -94,6 +98,14 @@ async function finalizeClaudePublishedSession( input: CloseClaudePublishedSessionInput, session: ClaudeSession ): Promise { + // Orca is stopping a start that never landed; the CLI did not fail, so the reason says so. + if (session.startup.state === 'pending') { + failClaudeStartupGate( + session, + new Error('claude session closed before startup completed'), + CLAUDE_STARTUP_ABANDONED_REJECTION + ) + } retireClaudeDispatchWaiters(session) // Settle every in-flight permission callback so closing leaves no dangling promise; `null` // writes no response, and the SDK ignores any post-cleanup answer regardless. diff --git a/src/main/claude/claude-structured-session-startup-gate.ts b/src/main/claude/claude-structured-session-startup-gate.ts index 020d625b5257..b5d5074e566f 100644 --- a/src/main/claude/claude-structured-session-startup-gate.ts +++ b/src/main/claude/claude-structured-session-startup-gate.ts @@ -144,12 +144,20 @@ export function rejectClaudeStartupWrites(session: ClaudeSession, reason: string return held.length > 0 } +/** Why a held prompt was rejected when Orca itself stopped a child still starting. */ +export const CLAUDE_STARTUP_ABANDONED_REJECTION = + 'Claude never finished starting, so Orca stopped it. Your message was not sent.' + /** Startup cannot land any more; nothing held was written, so all of it is rejected. */ -export function failClaudeStartupGate(session: ClaudeSession, error: Error): void { +export function failClaudeStartupGate( + session: ClaudeSession, + error: Error, + rejection = providerStartupFailureRejection(error) +): void { const gate = session.startup if (gate.state === 'pending') { gate.state = 'failed' gate.failure = error } - rejectClaudeStartupWrites(session, providerStartupFailureRejection(error)) + rejectClaudeStartupWrites(session, rejection) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts index aee8a8825766..6abfd06dc58d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-lifetime.ts @@ -190,14 +190,12 @@ export function createStructuredAgentSessionHolds( }, evict: close, hasProviderChild: (sessionId) => hasProviderChild(context, sessionId), - // A send pending while the child is still starting is held for that start; evicting would - // refuse it. Any other pending send may wait on an echo that never comes, so eviction retires it. + // A pending send is not owed work: one held for a start that never lands, or waiting on an echo + // that never comes, would keep the child forever. Eviction settles it instead. hasOwedWork: (sessionId) => { const session = context.sessions.get(sessionId) return session - ? activeStructuredAgentSessionTurnId(session.journal.snapshot().items) !== null || - (session.providerChildPhase === 'starting' && - session.journal.pendingSubmissions().length > 0) + ? activeStructuredAgentSessionTurnId(session.journal.snapshot().items) !== null : false }, onError: (error) => context.deps.onEventSinkError?.(error), diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-release-clock.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-release-clock.ts index 1cce910c09d6..de84d772efed 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-release-clock.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-release-clock.ts @@ -8,9 +8,10 @@ // already asked for must finish: stopping the child mid-answer strands the open turn marker. // // So the clock arms when the last holder leaves, every journal write while it is armed starts it -// again, and a tick that finds work still owed — a turn running, or a message sent but not yet -// taken by the provider — re-arms instead of evicting. The child goes only after a full window -// with no holder and no owed work. Quit still stops every child at once. +// again, and a tick that finds a turn running re-arms instead of evicting. A sent message the +// provider has not taken is not owed work: a start that never lands would otherwise keep its +// child forever, and eviction settles that message instead. The child goes only after a full +// window with no holder and no turn running. Quit still stops every child at once. export const STRUCTURED_AGENT_SESSION_RELEASE_GRACE_MS = 30 * 60_000 diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-starting-release.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-starting-release.test.ts index df6e6cec4bf1..faa5f81f0207 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-starting-release.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-starting-release.test.ts @@ -1,14 +1,16 @@ // A Claude chat is published before its CLI answers initialize, and a message sent in that window -// is held until it does. Switching away from the chat starts the release clock; the clock must -// treat that held message as work still owed, exactly as it treats a running turn, or it evicts -// the session and refuses a message the user already sent. +// is held until it does. Switching away from the chat starts the release clock. A held message is +// not work the clock waits on: a start that never lands would otherwise keep its child for the +// host's lifetime. The full idle window still passes first, and the start landing restarts it. import { mkdtemp, rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' +import type { AgentSessionSubscribeEvent } from '../../../shared/agent-session-wire' import { ClaudeStructuredSessionAdapter } from '../../claude/claude-structured-session-adapter' +import { CLAUDE_STARTUP_ABANDONED_REJECTION } from '../../claude/claude-structured-session-startup-gate' import { fakeClaude, PROVIDER_SESSION_ID @@ -136,28 +138,37 @@ function dispatchState(clientMessageId: string): string | undefined { .submissions.find((entry) => entry.clientMessageId === clientMessageId)?.dispatchState } -/** Long enough for several grace windows to elapse, so "not evicted" means the clock declined. */ -function waitOutSeveralGraceWindows(): Promise { - return new Promise((resolve) => setTimeout(resolve, GRACE_MS * 20)) +function publishedSubmissions(events: AgentSessionSubscribeEvent[]) { + return events.flatMap((event) => (event.type === 'batch' ? event.batch.submissions : [])) } describe('a chat left while its Claude CLI is still starting', () => { - it('keeps the session for a message it is holding, and delivers it once startup lands', async () => { + it('is released after the idle window when its start never lands, and rejects the message it held', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) await attachStarting() + const events: AgentSessionSubscribeEvent[] = [] + const unsubscribe = host.subscribe({ + id: 'pane', + sessionId: SESSION, + emit: (event) => events.push(event) + }) const held = await send('sent while starting') host.release(SESSION, SURFACE) - await waitOutSeveralGraceWindows() - + await vi.advanceTimersByTimeAsync(GRACE_MS - 1) expect(host.hasSession(SESSION)).toBe(true) expect(claude.connections[0].closeCount).toBe(0) - expect(dispatchState(held)).toBe('pending') - landInit() - await adapter.drainStartup(SESSION) + await vi.advanceTimersByTimeAsync(1) + await vi.waitFor(() => expect(host.hasSession(SESSION)).toBe(false)) + unsubscribe() - expect(claude.connections[0].sent).toEqual([expect.objectContaining({ type: 'user' })]) - await vi.waitFor(() => expect(dispatchState(held)).toBe('accepted')) + expect(claude.connections[0].closeCount).toBe(1) + expect(claude.connections[0].sent).toEqual([]) + // Never written, so it is refused with why rather than left in doubt; the user can resend it. + expect( + publishedSubmissions(events).findLast((entry) => entry.clientMessageId === held) + ).toMatchObject({ dispatchState: 'rejected', reason: CLAUDE_STARTUP_ABANDONED_REJECTION }) }) it('gives the message it wrote at startup a full grace to open its turn', async () => { @@ -170,10 +181,10 @@ describe('a chat left while its Claude CLI is still starting', () => { connection.sent.push(message) } host.release(SESSION, SURFACE) - await vi.advanceTimersByTimeAsync(GRACE_MS * 3 - 1) + await vi.advanceTimersByTimeAsync(GRACE_MS - 1) expect(host.hasSession(SESSION)).toBe(true) - // Startup lands just before the clock's next tick. + // A slow start lands just before the window would close. landInit() await adapter.drainStartup(SESSION) await Promise.all(lifecycle) diff --git a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox-late-rejection.test.tsx b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox-late-rejection.test.tsx new file mode 100644 index 000000000000..5a3bdfd34ef6 --- /dev/null +++ b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox-late-rejection.test.tsx @@ -0,0 +1,218 @@ +// @vitest-environment happy-dom + +// A send the host answered `pending` can still be rejected later, when the child holding it is +// stopped before it ever wrote it. The chat must treat that exactly like a send rejected on the +// spot: the reason on screen, the queue stopped on it, and Retry sending it under a new id. + +import { act, renderHook, waitFor } from '@testing-library/react' +import { beforeEach, describe, expect, it, vi } from 'vitest' +import type { AgentJournalSubmission } from '../../../../shared/agent-session-journal-types' +import { DISPATCH_REJECTED_CANCELLED } from '../../../../shared/structured-agent-session-dispatch-rejection' + +const mocks = vi.hoisted(() => ({ + call: vi.fn<(target: unknown, method: string, params: SendParams) => Promise>() +})) + +vi.mock('@/runtime/structured-agent-session-client', () => ({ + callStructuredAgentSession: mocks.call +})) + +import { useStructuredAgentSessionOutbox } from './use-structured-agent-session-outbox' + +const LOCAL_TARGET = { kind: 'local' } as const +const REASON = 'Claude never finished starting, so Orca stopped it. Your message was not sent.' + +type SendParams = { envelope: { clientOperationId: string }; retryUnknown?: true } + +function submission( + clientMessageId: string, + dispatchState: AgentJournalSubmission['dispatchState'], + reason: string | null = null +): AgentJournalSubmission { + return { + clientMessageId, + fence: 1, + payloadFingerprint: 'fingerprint', + dispatchState, + providerItemId: null, + reason, + submittedAt: 10, + resolvedAt: dispatchState === 'pending' ? null : 11 + } +} + +function sendResult( + clientMessageId: string, + dispatchState: AgentJournalSubmission['dispatchState'], + reason: string | null = null +) { + return { + ok: true, + replayed: false, + fence: 1, + cursor: { epoch: 'epoch-1', sequence: 10 }, + value: { clientMessageId, submission: submission(clientMessageId, dispatchState, reason) } + } +} + +function answerEach(dispatchState: AgentJournalSubmission['dispatchState'], reason?: string) { + mocks.call.mockImplementation(async (_target, _method, params) => + sendResult(params.envelope.clientOperationId, dispatchState, reason ?? null) + ) +} + +function sentIds(): string[] { + return mocks.call.mock.calls.map(([, , params]) => params.envelope.clientOperationId) +} + +function renderOutbox(submissions: readonly AgentJournalSubmission[] = []) { + return renderHook( + (props: { submissions: readonly AgentJournalSubmission[] }) => + useStructuredAgentSessionOutbox({ + sessionId: 'session-1', + target: LOCAL_TARGET, + fence: 1, + submissions: props.submissions + }), + { initialProps: { submissions } } + ) +} + +/** Lets the drain and the probe timer run, so "nothing was sent" means nothing would be. */ +async function settleEffects(): Promise { + await act(async () => { + await new Promise((resolve) => setTimeout(resolve, 50)) + }) +} + +describe('a send the host rejects after answering pending', () => { + beforeEach(() => { + vi.clearAllMocks() + localStorage.clear() + }) + + it('shows the reason and blocks on it while the chat is open, and Retry sends it fresh', async () => { + answerEach('pending') + const { result, rerender } = renderOutbox() + act(() => expect(result.current.send('held while starting')).toBe(true)) + await waitFor(() => expect(result.current.outbox[0]?.state).toBe('dispatching')) + const id = result.current.outbox[0]!.clientMessageId + + rerender({ submissions: [submission(id, 'rejected', REASON)] }) + + await waitFor(() => expect(result.current.error).toBe(REASON)) + expect(result.current.outbox[0]).toMatchObject({ clientMessageId: id, state: 'queued' }) + expect(result.current.blockedClientMessageId).toBe(id) + await settleEffects() + expect(mocks.call).toHaveBeenCalledOnce() + + answerEach('accepted') + act(() => result.current.retry(id)) + await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(2)) + expect(sentIds()[1]).not.toBe(id) + expect(mocks.call.mock.calls[1]![2].retryUnknown).toBeUndefined() + expect(result.current.error).toBeNull() + }) + + it('shows the reason when the chat reopens with the rejection already recorded', async () => { + answerEach('pending') + const first = renderOutbox() + act(() => expect(first.result.current.send('held while starting')).toBe(true)) + await waitFor(() => expect(first.result.current.outbox[0]?.state).toBe('dispatching')) + const id = first.result.current.outbox[0]!.clientMessageId + first.unmount() + + const reopened = renderOutbox([submission(id, 'rejected', REASON)]) + + await waitFor(() => expect(reopened.result.current.error).toBe(REASON)) + expect(reopened.result.current.outbox[0]).toMatchObject({ + clientMessageId: id, + state: 'queued' + }) + expect(reopened.result.current.blockedClientMessageId).toBe(id) + await settleEffects() + expect(mocks.call).toHaveBeenCalledOnce() + }) + + it('takes the journal answer over a send result that arrives after it', async () => { + let answer: (value: unknown) => void = () => {} + mocks.call.mockImplementationOnce( + () => + new Promise((resolve) => { + answer = resolve + }) + ) + const { result, rerender } = renderOutbox() + act(() => expect(result.current.send('held while starting')).toBe(true)) + await waitFor(() => expect(mocks.call).toHaveBeenCalledOnce()) + const id = sentIds()[0]! + + rerender({ submissions: [submission(id, 'rejected', REASON)] }) + await waitFor(() => expect(result.current.error).toBe(REASON)) + await act(async () => answer(sendResult(id, 'pending'))) + await settleEffects() + + expect(result.current.error).toBe(REASON) + expect(result.current.outbox[0]).toMatchObject({ clientMessageId: id, state: 'queued' }) + expect(result.current.blockedClientMessageId).toBe(id) + }) + + it('handles a rejection answered on the spot once, even when the journal repeats it', async () => { + answerEach('rejected', REASON) + const { result, rerender } = renderOutbox() + act(() => expect(result.current.send('refused')).toBe(true)) + await waitFor(() => expect(result.current.error).toBe(REASON)) + const id = result.current.outbox[0]!.clientMessageId + + rerender({ submissions: [submission(id, 'rejected', REASON)] }) + await settleEffects() + expect(result.current.blockedClientMessageId).toBe(id) + expect(mocks.call).toHaveBeenCalledOnce() + + answerEach('accepted') + act(() => result.current.retry(id)) + await waitFor(() => expect(result.current.outbox).toHaveLength(0)) + expect(mocks.call).toHaveBeenCalledTimes(2) + expect(sentIds()[1]).not.toBe(id) + }) + + it('keeps the reason when a send still in flight answers after another entry is rejected', async () => { + answerEach('pending') + const { result, rerender } = renderOutbox() + act(() => expect(result.current.send('held while starting')).toBe(true)) + await waitFor(() => expect(result.current.outbox[0]?.state).toBe('dispatching')) + const heldId = result.current.outbox[0]!.clientMessageId + let answer: (value: unknown) => void = () => {} + mocks.call.mockImplementationOnce( + () => + new Promise((resolve) => { + answer = resolve + }) + ) + act(() => expect(result.current.send('sent behind it')).toBe(true)) + await waitFor(() => expect(mocks.call).toHaveBeenCalledTimes(2)) + const inFlightId = sentIds()[1]! + + rerender({ submissions: [submission(heldId, 'rejected', REASON)] }) + await waitFor(() => expect(result.current.error).toBe(REASON)) + await act(async () => answer(sendResult(inFlightId, 'pending'))) + await settleEffects() + + expect(result.current.error).toBe(REASON) + expect(result.current.blockedClientMessageId).toBe(heldId) + }) + + it('drops a send the user cancelled, with nothing on screen', async () => { + answerEach('pending') + const { result, rerender } = renderOutbox() + act(() => expect(result.current.send('never mind')).toBe(true)) + await waitFor(() => expect(result.current.outbox[0]?.state).toBe('dispatching')) + const id = result.current.outbox[0]!.clientMessageId + + rerender({ submissions: [submission(id, 'rejected', DISPATCH_REJECTED_CANCELLED)] }) + + await waitFor(() => expect(result.current.outbox).toHaveLength(0)) + expect(result.current.error).toBeNull() + expect(result.current.blockedClientMessageId).toBeNull() + }) +}) diff --git a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts index e38709b9fed7..67d6bf7b7ae7 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts +++ b/src/renderer/src/components/native-chat/use-structured-agent-session-outbox.ts @@ -4,10 +4,12 @@ import { createStructuredAgentSessionOperationId } from '../../../../shared/stru import { admitStructuredAgentSessionOutboxEntry, createStructuredAgentSessionOutboxEntry, - reconcileStructuredAgentSessionOutbox, type StructuredAgentSessionOutboxEntry } from '../../../../shared/structured-agent-session-outbox' -import type { StructuredAgentSessionSendDisposition } from '../../../../shared/structured-agent-session-send-disposition' +import { + foldStructuredAgentSessionJournal, + type StructuredAgentSessionSendDisposition +} from '../../../../shared/structured-agent-session-send-disposition' import type { RuntimeClientTarget } from '@/runtime/runtime-rpc-client' import { readOutbox, writeOutbox } from './structured-agent-session-outbox-storage' import { @@ -93,18 +95,15 @@ export function useStructuredAgentSessionOutbox(args: { useEffect(() => { const current = outboxRef.current - const hostOwns = new Set( - submissions - .filter( - (submission) => - submission.dispatchState === 'pending' || submission.dispatchState === 'accepted' - ) - .map((submission) => submission.clientMessageId) - ) - const next = reconcileStructuredAgentSessionOutbox(current, submissions) - const admittedInFlight = inFlightIdRef.current !== null && hostOwns.has(inFlightIdRef.current) + const fold = foldStructuredAgentSessionJournal({ + entries: current, + submissions, + blockedClientMessageId: blockedIdRef.current, + inFlightClientMessageId: inFlightIdRef.current + }) + const next = fold.entries if ( - admittedInFlight || + fold.answeredInFlight || next.some((entry, index) => entry !== current[index]) || next.length !== current.length ) { @@ -113,21 +112,22 @@ export function useStructuredAgentSessionOutbox(args: { writeOutbox(sessionId, next) } // Keyed on the entry actually in flight, which is no longer always the head: the journal - // owning it outranks a send promise that has not settled, so release single-flight and make + // answering it outranks a send promise that has not settled, so release single-flight and make // that promise a no-op. Keying on the head would discard the tail's unsettled send instead, // and with it a refusal only that send can report. - if (admittedInFlight) { + if (fold.answeredInFlight) { 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)) - ) { + blockedIdRef.current = fold.blockedClientMessageId + if (fold.clearsError) { setError(null) } + if (fold.lateRejection) { + blockedIdRef.current = fold.lateRejection.blockedClientMessageId + retryWithFreshClientMessageIdRef.current = fold.lateRejection.retryWithFreshClientMessageId + setError(fold.lateRejection.error) + } }, [sessionId, submissions]) // The one place that owns the refs, the React state and the storage write. @@ -136,9 +136,16 @@ 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 + // A clean answer that leaves another entry blocked says nothing about it, so its reason stays. + const keepsBlockedReason = + disposition.error === null && + disposition.blockedClientMessageId !== null && + disposition.blockedClientMessageId === blockedIdRef.current blockedIdRef.current = disposition.blockedClientMessageId retryWithFreshClientMessageIdRef.current = disposition.retryWithFreshClientMessageId - setError(disposition.error) + if (!keepsBlockedReason) { + setError(disposition.error) + } outboxRef.current = disposition.entries setOutbox(disposition.entries) writeOutbox(sessionId, disposition.entries) diff --git a/src/shared/structured-agent-session-send-disposition.ts b/src/shared/structured-agent-session-send-disposition.ts index 16765a024714..b638b76ffa18 100644 --- a/src/shared/structured-agent-session-send-disposition.ts +++ b/src/shared/structured-agent-session-send-disposition.ts @@ -7,13 +7,16 @@ // the refs, the React state and the storage write, and nothing else decides an // entry's state. +import type { AgentJournalSubmission } from './agent-session-journal-types' import type { AgentSessionMutationResult, AgentSessionSendResult } from './agent-session-wire' import { + DISPATCH_REJECTED_CANCELLED, dispatchRejectionReasonIsInternal, dispatchRejectionWasTransportWriteFailure } from './structured-agent-session-dispatch-rejection' import { classifyStructuredAgentSessionSendFailure, + reconcileStructuredAgentSessionOutbox, requeueStructuredAgentSessionSendRefusal, type StructuredAgentSessionOutboxEntry } from './structured-agent-session-outbox' @@ -103,6 +106,144 @@ export function structuredAgentSessionRejectionNotice(reason: string | null): st : reason } +/** The message provably did not happen: it parks with its reason, and Retry sends it under a new id. */ +function rejectedSubmissionDisposition( + entries: StructuredAgentSessionOutboxEntry[], + clientMessageId: string, + reason: string | null +): StructuredAgentSessionSendDisposition { + return { + entries, + error: structuredAgentSessionRejectionNotice(reason), + blockedClientMessageId: clientMessageId, + retryWithFreshClientMessageId: clientMessageId + } +} + +/** A rejection the user is told about; a cancel is their own withdrawal and drops silently. */ +function reportsRejection(submission: AgentJournalSubmission): boolean { + return ( + submission.dispatchState === 'rejected' && submission.reason !== DISPATCH_REJECTED_CANCELLED + ) +} + +/** + * The journal rejected a send after its dispatch answered `pending`, or while the pane was away. + * That is the same fact as a rejected send result, so it gets the same disposition. Every such + * entry is requeued so none is skipped as in flight, and the queue stops on the first. Null while + * another answer already holds the queue, or when nothing was rejected. + */ +export function disposeStructuredAgentSessionLateRejection(input: { + entries: readonly StructuredAgentSessionOutboxEntry[] + submissions: readonly AgentJournalSubmission[] + blockedClientMessageId: string | null +}): StructuredAgentSessionSendDisposition | null { + if (input.blockedClientMessageId !== null) { + return null + } + const rejected = new Map( + input.submissions + .filter(reportsRejection) + .map((submission) => [submission.clientMessageId, submission.reason]) + ) + const head = input.entries.find((entry) => rejected.has(entry.clientMessageId)) + if (!head) { + return null + } + const entries = input.entries.map((entry) => + rejected.has(entry.clientMessageId) && entry.state !== 'queued' + ? { ...entry, state: 'queued' as const } + : entry + ) + return rejectedSubmissionDisposition( + entries, + head.clientMessageId, + rejected.get(head.clientMessageId) ?? null + ) +} + +/** What one journal update does to the outbox; the hook applies it and owns the refs and storage. */ +export type StructuredAgentSessionJournalFold = { + entries: StructuredAgentSessionOutboxEntry[] + /** The journal answered the send in flight, so its promise must not apply. */ + answeredInFlight: boolean + /** The host now owns the blocked or an unconfirmed entry, so the banner no longer applies. */ + clearsError: boolean + /** Always the next value, never "unchanged". */ + blockedClientMessageId: string | null + /** Applied last: it names the new blocked entry, its error and the id Retry rotates. */ + lateRejection: StructuredAgentSessionSendDisposition | null +} + +export function foldStructuredAgentSessionJournal(input: { + entries: readonly StructuredAgentSessionOutboxEntry[] + submissions: readonly AgentJournalSubmission[] + blockedClientMessageId: string | null + inFlightClientMessageId: string | null +}): StructuredAgentSessionJournalFold { + const hostOwns = new Set( + input.submissions + .filter( + (submission) => + submission.dispatchState === 'pending' || submission.dispatchState === 'accepted' + ) + .map((submission) => submission.clientMessageId) + ) + const unblocks = + input.blockedClientMessageId !== null && hostOwns.has(input.blockedClientMessageId) + const blockedClientMessageId = unblocks ? null : input.blockedClientMessageId + const reconciled = reconcileStructuredAgentSessionOutbox(input.entries, input.submissions) + const lateRejection = disposeStructuredAgentSessionLateRejection({ + entries: reconciled, + submissions: input.submissions, + blockedClientMessageId + }) + const entries = lateRejection?.entries ?? reconciled + const inFlight = input.inFlightClientMessageId + return { + entries, + // A late rejection that requeued the entry in flight answered it too. + answeredInFlight: + inFlight !== null && + (hostOwns.has(inFlight) || + entries.some( + (entry, index) => entry.clientMessageId === inFlight && entry !== reconciled[index] + )), + clearsError: + unblocks || + input.entries.some( + (entry) => entry.state === 'unconfirmed' && hostOwns.has(entry.clientMessageId) + ), + blockedClientMessageId, + lateRejection + } +} + +/** + * For a client with no outbox: which awaited sends the journal has settled, and the notice for each + * one it rejected, so a late rejection reads exactly like an immediate one. + */ +export function settleAwaitedStructuredAgentSessionSends( + awaited: ReadonlySet, + submissions: readonly AgentJournalSubmission[] +): { settled: string[]; notices: string[] } { + const settled: string[] = [] + const notices: string[] = [] + for (const submission of submissions) { + if ( + !awaited.has(submission.clientMessageId) || + (submission.dispatchState !== 'accepted' && submission.dispatchState !== 'rejected') + ) { + continue + } + settled.push(submission.clientMessageId) + if (reportsRejection(submission)) { + notices.push(structuredAgentSessionRejectionNotice(submission.reason)) + } + } + return { settled, notices } +} + export function disposeStructuredAgentSessionSendResult( input: SendDispositionInput & { result: AgentSessionMutationResult @@ -151,12 +292,11 @@ export function disposeStructuredAgentSessionSendResult( } } if (submission.dispatchState === 'rejected') { - return { - entries: replaceEntryState(input, 'queued'), - error: structuredAgentSessionRejectionNotice(submission.reason), - blockedClientMessageId: input.entry.clientMessageId, - retryWithFreshClientMessageId: input.entry.clientMessageId - } + return rejectedSubmissionDisposition( + replaceEntryState(input, 'queued'), + input.entry.clientMessageId, + submission.reason + ) } if (submission.dispatchState === 'unknown' && submission.recovered) { return {