From c0bb243eb899b212a026c873fda2d97f100decde Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Wed, 9 Sep 2026 18:37:46 -0700 Subject: [PATCH] Fix structured chat pending-work lifecycle and mobile cancellation --- .../src/session/MobileNativeChatOverlay.tsx | 1 + .../src/session/MobileNativeChatView.test.ts | 13 +++++++ mobile/src/session/MobileNativeChatView.tsx | 8 ++--- .../mobile-native-chat-controller-contract.ts | 1 + .../use-mobile-native-chat-controller.test.ts | 35 +++++++++++++++++-- .../use-mobile-native-chat-controller.ts | 3 ++ .../use-mobile-structured-agent-session.ts | 2 +- .../journal-crash-boundary.test.ts | 25 +++++++++++++ .../journal-pending-submission-recovery.ts | 14 ++++++-- .../agent-session-journal/journal-reducer.ts | 3 ++ .../agent-session-journal/journal-store.ts | 7 ++-- .../structured-agent-session-host-handoff.ts | 7 ++++ ...red-agent-session-send-idempotency.test.ts | 34 ++++++++++++++++++ ...ructured-agent-session-status-feed.test.ts | 27 +++++++++++++- .../structured-agent-session-status-feed.ts | 3 +- ...red-agent-session-surface-lifetime.test.ts | 6 ++++ .../structured-agent-session-turns.ts | 12 +++++-- ...ured-agent-session-unexpected-exit.test.ts | 11 +++++- ...tructured-agent-session-unexpected-exit.ts | 13 +++++++ .../use-structured-agent-session.ts | 3 +- src/shared/agent-session-journal-types.ts | 1 + ...tructured-agent-session-projection.test.ts | 25 +++++++++---- .../structured-agent-session-projection.ts | 23 +++++++----- 23 files changed, 241 insertions(+), 36 deletions(-) diff --git a/mobile/src/session/MobileNativeChatOverlay.tsx b/mobile/src/session/MobileNativeChatOverlay.tsx index 357a089466e..dcb353e54f8 100644 --- a/mobile/src/session/MobileNativeChatOverlay.tsx +++ b/mobile/src/session/MobileNativeChatOverlay.tsx @@ -71,6 +71,7 @@ export function MobileNativeChatOverlay({ error={session.error} agent={controller.nativeChatAgent} agentWorking={controller.nativeChatAgentWorking} + canStop={controller.nativeChatCanStop} structuredActivityUi={controller.nativeChatStructured} streaming={streaming} onStop={controller.handleNativeChatStop} diff --git a/mobile/src/session/MobileNativeChatView.test.ts b/mobile/src/session/MobileNativeChatView.test.ts index d101c3f0ef6..63bfb715445 100644 --- a/mobile/src/session/MobileNativeChatView.test.ts +++ b/mobile/src/session/MobileNativeChatView.test.ts @@ -74,6 +74,7 @@ type Overrides = { pending?: Parameters[0]['pending'] structuredActivityUi?: boolean agentWorking?: boolean + canStop?: boolean sendSurfaceId?: string } @@ -118,6 +119,18 @@ describe('MobileNativeChatView', () => { } /** Ids of the rows the list is currently rendering. */ + it('keeps Stop hidden during a structured dispatch until a provider turn can be cancelled', async () => { + const props = { structuredActivityUi: true, agentWorking: true, canStop: false } + await render(props) + const stops = () => + renderer!.root.findAll((node) => node.props.accessibilityLabel === 'Stop the agent') + expect(stops()).toHaveLength(0) + await update({ ...props, canStop: true }) + expect(stops()).toHaveLength(1) + await update({ agentWorking: true }) + expect(stops()).toHaveLength(1) + }) + function listIds(): string[] { const list = renderer!.root.find((node) => node.type === 'FlatList') return (list.props.data as { id: string }[]).map((row) => row.id) diff --git a/mobile/src/session/MobileNativeChatView.tsx b/mobile/src/session/MobileNativeChatView.tsx index 70a67787de0..f58e2ad05b0 100644 --- a/mobile/src/session/MobileNativeChatView.tsx +++ b/mobile/src/session/MobileNativeChatView.tsx @@ -49,10 +49,11 @@ type Props = { /** Resolved agent for this chat; names the empty-state copy (desktop parity). */ agent?: string | null agentWorking?: boolean + canStop?: boolean /** Structured lane: per-turn "Working for N" status plus live tool progress, * replacing the bridge lane's static three-dot working row (desktop parity). */ structuredActivityUi?: boolean - /** Interrupt the agent mid-turn (shown as a Stop button on the working bar). */ + /** Interrupt a provider turn. */ onStop?: () => void /** Live partial assistant text to show as an in-progress bubble, already gated * by the overlay against the transcript catching up. */ @@ -129,6 +130,7 @@ export function MobileNativeChatView({ error, agent, agentWorking, + canStop = agentWorking, structuredActivityUi = false, onStop, streaming, @@ -396,8 +398,6 @@ export function MobileNativeChatView({ question={question} onAnswerQuestion={onAnswerQuestion} /> - {/* Chrome row above the composer: the working indicator and the global - tool-calls expand/collapse toggle on the left, Stop in the far corner. */} {agentWorking && !structuredActivityUi ? : null} @@ -414,7 +414,7 @@ export function MobileNativeChatView({ {toolsExpanded ? 'Collapse' : 'Tools'} - {agentWorking ? ( + {canStop ? ( [styles.stopButton, pressed && styles.pressed]} onPress={onStop} diff --git a/mobile/src/session/mobile-native-chat-controller-contract.ts b/mobile/src/session/mobile-native-chat-controller-contract.ts index 53187e0d6db..0c6ffd6a764 100644 --- a/mobile/src/session/mobile-native-chat-controller-contract.ts +++ b/mobile/src/session/mobile-native-chat-controller-contract.ts @@ -28,6 +28,7 @@ export type MobileNativeChatController = { /** Structured lane: drives the per-turn status row and live tool progress. */ nativeChatStructured: boolean nativeChatAgentWorking: boolean + nativeChatCanStop: boolean nativeChatStreamingText?: string /** Agent mid-turn, regardless of whether chat is the visible view. */ nativeChatStreamLive: boolean diff --git a/mobile/src/session/use-mobile-native-chat-controller.test.ts b/mobile/src/session/use-mobile-native-chat-controller.test.ts index 83f8075d914..c90033e2404 100644 --- a/mobile/src/session/use-mobile-native-chat-controller.test.ts +++ b/mobile/src/session/use-mobile-native-chat-controller.test.ts @@ -55,6 +55,7 @@ const structuredQuestion = { allowOther: true, optionTokens: ['choice-a', 'choice-b'] } +const structuredActivity = { isWorking: false, turnId: null as string | null } const structuredSessionState = { messages: [] as unknown[], status: 'ready', @@ -86,8 +87,7 @@ vi.mock('./use-mobile-native-chat-session', () => ({ vi.mock('./use-mobile-structured-agent-session', () => ({ useMobileStructuredAgentSession: () => ({ session: structuredSessionState, - isWorking: false, - turnId: null, + ...structuredActivity, sendWithOutcome: structuredSendWithOutcome, cancel: structuredCancel, permission: structuredPermission, @@ -341,6 +341,37 @@ describe('useMobileNativeChatController handleNativeChatSend', () => { expect(clientStub.sendRequest).not.toHaveBeenCalled() }) + it('separates structured working status from provider cancellation availability', async () => { + const props = { + tab: { + type: 'agent-session', + id: 'agent-tab-1', + title: 'Chat', + sessionId: 'session-structured', + agent: 'codex', + isActive: true + }, + activeHandle: null, + inputLeaseReady: false + } + structuredActivity.isWorking = true + try { + await act(async () => { + renderer?.update(createElement(Harness, props)) + }) + expect(controller?.nativeChatAgentWorking).toBe(true) + expect(controller?.nativeChatCanStop).toBe(false) + structuredActivity.turnId = 'provider-turn' + await act(async () => { + renderer?.update(createElement(Harness, props)) + }) + expect(controller?.nativeChatCanStop).toBe(true) + } finally { + structuredActivity.isWorking = false + structuredActivity.turnId = null + } + }) + it('exposes structured prompt cards and session options on structured tabs', async () => { await act(async () => { renderer?.update( diff --git a/mobile/src/session/use-mobile-native-chat-controller.ts b/mobile/src/session/use-mobile-native-chat-controller.ts index 729cec302c1..5195edf2124 100644 --- a/mobile/src/session/use-mobile-native-chat-controller.ts +++ b/mobile/src/session/use-mobile-native-chat-controller.ts @@ -299,6 +299,9 @@ export function useMobileNativeChatController(args: { /** Structured lane: drives the per-turn status row and live tool progress. */ nativeChatStructured: activeChatStructured, nativeChatAgentWorking, + nativeChatCanStop: activeChatStructured + ? structuredNativeChat.turnId !== null + : nativeChatAgentWorking, nativeChatStreamingText, nativeChatStreamLive, nativeChatStreamScopeKey: streamScopeKey, diff --git a/mobile/src/session/use-mobile-structured-agent-session.ts b/mobile/src/session/use-mobile-structured-agent-session.ts index 8d8db89b046..86d3f2c8e44 100644 --- a/mobile/src/session/use-mobile-structured-agent-session.ts +++ b/mobile/src/session/use-mobile-structured-agent-session.ts @@ -297,7 +297,7 @@ export function useMobileStructuredAgentSession(args: { // A dispatch the provider has not answered yet is already work — see the desktop hook. isWorking: activeStructuredAgentSessionTurnId(state.items) !== null || - hasUnansweredStructuredAgentSessionDispatch(state.submissions), + hasUnansweredStructuredAgentSessionDispatch(state.submissions, state.fence), turnId: activeStructuredAgentSessionTurnId(state.items), sendWithOutcome, cancel, diff --git a/src/main/native-chat/agent-session-journal/journal-crash-boundary.test.ts b/src/main/native-chat/agent-session-journal/journal-crash-boundary.test.ts index 6e33db68ee6..85cda398c1d 100644 --- a/src/main/native-chat/agent-session-journal/journal-crash-boundary.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-crash-boundary.test.ts @@ -15,6 +15,7 @@ import type { AgentJournalMessageItem, AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types' +import { hasUnansweredStructuredAgentSessionDispatch } from '../../../shared/structured-agent-session-projection' import { digestPayload } from './journal-payload-bounds' import { reconcileSubmissions, @@ -141,6 +142,30 @@ describe('crash between provider accept and journal commit', () => { expect(restarted.receiptFor('cm_1')?.providerItemId).toBe(agentJournalItemKey(outcome.identity)) }) + it('retires an ack timeout on restart without changing its delivery verdict', async () => { + const journal = await open() + await journal.appendSubmission({ + clientMessageId: 'cm_timeout', + payloadFingerprint: digestPayload('slow'), + body: userMessage('slow'), + fence: 1 + }) + await journal.resolveDispatch({ + clientMessageId: 'cm_timeout', + state: 'unknown', + reason: 'ack timeout', + fence: 1 + }) + expect(hasUnansweredStructuredAgentSessionDispatch(journal.submissions())).toBe(true) + const restarted = await open() + await restarted.markPendingSubmissionsUnknown(2) + expect(restarted.submissions()[0]?.dispatchState).toBe('unknown') + expect(hasUnansweredStructuredAgentSessionDispatch(restarted.submissions())).toBe(false) + const cursor = restarted.cursor() + await restarted.markPendingSubmissionsUnknown(2) + expect(restarted.cursor()).toEqual(cursor) + }) + it('reports a rejected submission as never delivered, and never re-sends it', async () => { const journal = await open() await journal.appendSubmission({ diff --git a/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts b/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts index 76bc00394f2..e09dd109d50 100644 --- a/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts +++ b/src/main/native-chat/agent-session-journal/journal-pending-submission-recovery.ts @@ -2,14 +2,22 @@ import type { AgentSessionJournal } from './journal-store' export async function markJournalPendingSubmissionsUnknown( journal: AgentSessionJournal, - fence: number + fence: number, + reason = 'host_restarted_before_acknowledgement' ): Promise { - const pending = journal.pendingSubmissions().map((entry) => entry.clientMessageId) + const pending = journal + .submissions() + .filter( + (entry) => + entry.dispatchState === 'pending' || + (entry.dispatchState === 'unknown' && entry.recovered !== true) + ) + .map((entry) => entry.clientMessageId) for (const clientMessageId of pending) { await journal.resolveDispatch({ clientMessageId, state: 'unknown', - reason: 'host_restarted_before_acknowledgement', + reason, fence, recovered: true }) diff --git a/src/main/native-chat/agent-session-journal/journal-reducer.ts b/src/main/native-chat/agent-session-journal/journal-reducer.ts index 97e10af950d..e01e7d6158f 100644 --- a/src/main/native-chat/agent-session-journal/journal-reducer.ts +++ b/src/main/native-chat/agent-session-journal/journal-reducer.ts @@ -256,12 +256,15 @@ function applyDispatch( if (submission.dispatchState === 'rejected' || submission.dispatchState === 'accepted') { return } + submission.fence = row.fence submission.dispatchState = row.state submission.providerItemId = row.providerItemId submission.reason = row.reason submission.resolvedAt = row.ts if (row.recovered) { submission.recovered = row.recovered + } else { + delete submission.recovered } if (row.state !== 'accepted' || !row.providerItemId) { return diff --git a/src/main/native-chat/agent-session-journal/journal-store.ts b/src/main/native-chat/agent-session-journal/journal-store.ts index 3c16800099b..812e2bca834 100644 --- a/src/main/native-chat/agent-session-journal/journal-store.ts +++ b/src/main/native-chat/agent-session-journal/journal-store.ts @@ -245,10 +245,9 @@ export class AgentSessionJournal { })) } - /** On restart every `pending` submission becomes `unknown` before the session - * accepts a writer. Orca never re-sends on the user's behalf. */ - async markPendingSubmissionsUnknown(fence: number): Promise { - return markJournalPendingSubmissionsUnknown(this, fence) + /** Retire unanswered sends after their execution owner ended, without assuming delivery. */ + async markPendingSubmissionsUnknown(fence: number, reason?: string): Promise { + return markJournalPendingSubmissionsUnknown(this, fence, reason) } /** The escape hatch for corruption, an unreconcilable prefix, a forked handle, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts index 3495409f127..113940ff0f4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-handoff.ts @@ -86,6 +86,13 @@ export function createStructuredAgentSessionHostHandoff( host.publishStatus?.(sessionId) try { await host.flush(sessionId) + const session = host.session(sessionId) + await session.journal.markPendingSubmissionsUnknown( + session.fence, + 'provider_exited_before_acknowledgement' + ) + host.subscribers.publish(sessionId, session.journal) + host.publishStatus?.(sessionId) host.eventSink(sessionId).unbind() return { state: 'stopped' } } catch (error) { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-send-idempotency.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-send-idempotency.test.ts index 581743633c0..be9386815dd 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-send-idempotency.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-send-idempotency.test.ts @@ -3,6 +3,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types' +import { hasUnansweredStructuredAgentSessionDispatch } from '../../../shared/structured-agent-session-projection' import { structuredAgentSessionPayloadFingerprint } from '../../../shared/structured-agent-session-mutation' import { createTrackedJournalOpener } from '../agent-session-journal/journal-store-test-open' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' @@ -34,6 +35,39 @@ afterEach(async () => { }) describe('structured send idempotency', () => { + it('publishes a recovered retry as working before waiting for its provider', async () => { + const body: AgentJournalMessageItem = { + kind: 'message', + role: 'user', + blocks: [{ type: 'text', text: 'retry' }] + } + const input = { clientMessageId: 'retry-id', payloadFingerprint: 'fingerprint', body } + await journal.appendSubmission({ ...input, fence: 1 }) + await journal.markPendingSubmissionsUnknown(2) + const originalItem = journal.snapshot().items[0] + const publish = vi.fn() + const dispatch = vi.fn(async () => { + expect(publish).toHaveBeenCalledOnce() + expect(hasUnansweredStructuredAgentSessionDispatch(journal.submissions(), 2)).toBe(true) + return { state: 'unknown' as const, reason: 'ack timeout' } + }) + await performSend( + { + sessionId: 'session-1', + journal, + fence: 2, + adapter: { dispatch } as unknown as StructuredAgentSessionAdapter, + persistOptions: async () => undefined, + resolvedBy: 'caller', + publish, + now: () => 1 + }, + { ...input, retryUnknown: true } + ) + expect(hasUnansweredStructuredAgentSessionDispatch(journal.submissions(), 2)).toBe(true) + expect(journal.snapshot().items).toEqual([originalItem]) + }) + it('does not redispatch one send id reused across caller ledgers', async () => { const body: AgentJournalMessageItem = { kind: 'message', diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.test.ts index 95d78ffc38d..7bd53742430 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.test.ts @@ -56,9 +56,11 @@ async function openJournal(sessionId = SESSION, now?: () => number) { function indexed(session: { journal: Awaited> hasProviderChild?: boolean + fence?: number }) { return { journal: session.journal, + fence: session.fence ?? 1, ...(session.hasProviderChild !== undefined ? { hasProviderChild: session.hasProviderChild } : {}), @@ -69,7 +71,7 @@ function indexed(session: { function feedFor( sessions: Map< string, - { journal: Awaited>; hasProviderChild?: boolean } + { journal: Awaited>; hasProviderChild?: boolean; fence?: number } >, record: Partial | null = null, onStatusChanged?: StructuredAgentSessionStatusFeedDeps['onStatusChanged'] @@ -153,6 +155,29 @@ describe('StructuredAgentSessionStatusFeed', () => { ]) }) + it('stops projecting an old-host unknown submission after the owner fence advances', async () => { + const journal = await openJournal() + const session = { journal, fence: 1 } + const { feed, events } = feedFor(new Map([[SESSION, session]])) + await journal.appendSubmission({ + clientMessageId: 'old-host', + payloadFingerprint: 'fp', + body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'slow' }] }, + fence: 1 + }) + await journal.resolveDispatch({ + clientMessageId: 'old-host', + state: 'unknown', + reason: 'ack timeout', + fence: 1 + }) + feed.publish(SESSION) + expect(events.at(-1)).toMatchObject({ session: { status: 'working' } }) + session.fence = 2 + feed.publish(SESSION) + expect(events.at(-1)).toMatchObject({ session: { status: 'idle' } }) + }) + it('publishes working from the pending submission, before the provider replays the turn', async () => { const journal = await openJournal() const { feed, events } = feedFor(new Map([[SESSION, { journal }]])) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.ts index 0a3e36ba453..3c4ce4b97a6 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-status-feed.ts @@ -31,6 +31,7 @@ type StatusFeedSession = { journal: AgentSessionJournal params: { location: { workspaceId: string }; provider: AgentSessionRecord['provider'] } hasProviderChild?: boolean + fence?: number } export type StructuredAgentSessionStatusFeedDeps = { @@ -166,7 +167,7 @@ export class StructuredAgentSessionStatusFeed { workspaceId: session.params.location.workspaceId, agent: session.params.provider, ...(session.hasProviderChild ? { hostExecutionOwned: true as const } : {}), - ...projectStructuredAgentSessionStatusSummary(items, submissions), + ...projectStructuredAgentSessionStatusSummary(items, submissions, session.fence), ...(record?.rewind?.phase === 'prepared' || record?.rewind?.phase === 'provider-succeeded' ? { rewindBlockedReason: 'outcome-unknown' as const } : {}), diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts index 230e39cdaef..f4e02f4491b 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-surface-lifetime.test.ts @@ -8,6 +8,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest' import type { AgentSessionOwnerProbe } from '../../../shared/agent-session-lease-adjudication' +import { hasUnansweredStructuredAgentSessionDispatch } from '../../../shared/structured-agent-session-projection' import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' import type { AgentSessionMutationEnvelope, @@ -334,6 +335,11 @@ describe('an unexpected provider exit', () => { acquisitionGeneration: 'generation-1' }) + const recoveredHistory = host.history({ sessionId: SESSION, direction: 'tail' }) + expect( + recoveredHistory.ok && + hasUnansweredStructuredAgentSessionDispatch(recoveredHistory.page.submissions) + ).toBe(false) expect(acquire).toHaveBeenCalledTimes(2) expect(dispatch).toHaveBeenCalledOnce() expect(store.getRecord(SESSION)?.lease).toMatchObject({ diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts index 3bd98a61735..ba83b7e54ac 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-turns.ts @@ -97,6 +97,15 @@ export async function performSend( if (!(input.retryUnknown && existing?.dispatchState === 'unknown')) { await ctx.journal.appendSubmission({ ...input, fence: ctx.fence }) ctx.publish() + } else { + // Retry resumes work without moving or duplicating the original message. + await ctx.journal.resolveDispatch({ + clientMessageId: input.clientMessageId, + state: 'unknown', + reason: 'dispatch_retry_in_progress', + fence: ctx.fence + }) + ctx.publish() } const outcome = await dispatchSafely(ctx, input.clientMessageId, input.body) @@ -124,8 +133,7 @@ export async function performSend( clientMessageId: input.clientMessageId, state: 'unknown', reason: 'dispatch_result_persistence_failed', - fence: ctx.fence, - recovered: true + fence: ctx.fence }) } catch { // Nothing further to record; the pending row is settled on the next attach. diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.test.ts index 084ffc7547d..2c2b9428a43 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.test.ts @@ -49,7 +49,11 @@ describe('provider-exit recovery tickets', () => { hasProviderChild: true, fence: 7, acquisitionGeneration: GENERATION, - journal: { snapshot: () => ({ items: [] }), appendLifecycleBatch } + journal: { + snapshot: () => ({ items: [] }), + appendLifecycleBatch, + markPendingSubmissionsUnknown: vi.fn(async () => []) + } } as unknown as StructuredAgentSessionHostSession const store = { getRecord: () => ({ @@ -89,6 +93,10 @@ describe('provider-exit recovery tickets', () => { expect(result).toMatchObject({ settlementRetryRequired: false, releasedFence: 8 }) expect(appendLifecycleBatch).toHaveBeenCalledOnce() + expect(session.journal.markPendingSubmissionsUnknown).toHaveBeenCalledWith( + 7, + 'provider_exited_before_acknowledgement' + ) expect(session.hasProviderChild).toBe(false) }) @@ -98,6 +106,7 @@ describe('provider-exit recovery tickets', () => { fence: 7, acquisitionGeneration: GENERATION, journal: { + markPendingSubmissionsUnknown: vi.fn(async () => []), snapshot: () => ({ items: [] }), appendLifecycleBatch: vi.fn(async () => { throw new Error('journal still unavailable') diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.ts index ddc4af9f2d3..43ba7dbbb2d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-unexpected-exit.ts @@ -81,6 +81,15 @@ export async function settleUnexpectedStructuredAgentSessionExit( settlementRetryRequired = true context.onBarrierError?.(unexpectedEvent.sessionId, error) } + try { + await session.journal.markPendingSubmissionsUnknown( + session.fence, + 'provider_exited_before_acknowledgement' + ) + } catch (error) { + settlementRetryRequired = true + context.onBarrierError?.(unexpectedEvent.sessionId, error) + } if (unexpectedEvent.settlementRetryRequired || settlementRetryRequired) { const retried = await retryUnexpectedExitSettlement({ context, @@ -170,6 +179,10 @@ export async function retryUnexpectedExitSettlement(input: { stableSettlementId: string }): Promise { try { + await input.session.journal.markPendingSubmissionsUnknown( + input.session.fence, + 'provider_exited_before_acknowledgement' + ) const mutations = unexpectedExitFallbackMutations( input.event, input.session, diff --git a/src/renderer/src/components/native-chat/use-structured-agent-session.ts b/src/renderer/src/components/native-chat/use-structured-agent-session.ts index 0e531924dfc..864b7e1f608 100644 --- a/src/renderer/src/components/native-chat/use-structured-agent-session.ts +++ b/src/renderer/src/components/native-chat/use-structured-agent-session.ts @@ -86,8 +86,7 @@ export function useStructuredAgentSession(args: { // A dispatch the provider has not answered is already work; Claude's running row trails the // send by seconds, and only a provider-minted turn is cancellable, so the two stay separate. const isWorking = - turnId !== null || - hasUnansweredStructuredAgentSessionDispatch(state.submissions) + turnId !== null || hasUnansweredStructuredAgentSessionDispatch(state.submissions, state.fence) const turnActivity = useMemo( () => selectStructuredAgentTurnActivity(state.items, turnId, state.activity), [state.activity, state.items, turnId] diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index f057af24ca5..c017a94b36e 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -193,6 +193,7 @@ export type AgentJournalDispatchState = (typeof AGENT_JOURNAL_DISPATCH_STATES)[n * the turn reads as delivery unconfirmed, never as sent and never as failed. */ export type AgentJournalSubmission = { clientMessageId: string + /** Execution fence of the latest dispatch attempt or recovery. */ fence: number payloadFingerprint: string dispatchState: AgentJournalDispatchState diff --git a/src/shared/structured-agent-session-projection.test.ts b/src/shared/structured-agent-session-projection.test.ts index 7c56d8806be..76cf6088cba 100644 --- a/src/shared/structured-agent-session-projection.test.ts +++ b/src/shared/structured-agent-session-projection.test.ts @@ -1,9 +1,6 @@ import { describe, expect, it } from 'vitest' import { AGENT_STATUS_MAX_FIELD_LENGTH } from './agent-status-field-normalization' -import type { - AgentJournalRenderItem, - AgentJournalSubmission -} from './agent-session-journal-types' +import type { AgentJournalRenderItem, AgentJournalSubmission } from './agent-session-journal-types' import { parsePaneKey } from './stable-pane-id' import { activeStructuredAgentSessionTurnId, @@ -180,6 +177,20 @@ describe('structured agent session status projection', () => { }) }) + it('does not resurrect old-host unknown work after its execution fence advances', () => { + const oldHostSubmission = { ...submission('m1', 'unknown'), fence: 2 } + expect(hasUnansweredStructuredAgentSessionDispatch([oldHostSubmission], 2)).toBe(true) + expect(hasUnansweredStructuredAgentSessionDispatch([oldHostSubmission], 3)).toBe(false) + }) + + it('recognizes recovery from an older host without the optional marker', () => { + expect( + hasUnansweredStructuredAgentSessionDispatch([ + { ...submission('m1', 'unknown'), reason: 'host_restarted_before_acknowledgement' } + ]) + ).toBe(false) + }) + it('stops reading a resolved dispatch as work, and lets a pending prompt outrank it', () => { const asked = item('asked', 1, { kind: 'message', @@ -205,9 +216,9 @@ describe('structured agent session status projection', () => { { ...submission('m1', 'unknown'), recovered: true } ]) ).toBe(false) - expect(projectStructuredAgentSessionStatus([asked, prompt], [submission('m1', 'pending')])).toBe( - 'attention' - ) + expect( + projectStructuredAgentSessionStatus([asked, prompt], [submission('m1', 'pending')]) + ).toBe('attention') expect(projectStructuredAgentSessionStatusSummary([], [])).toEqual({ status: null, latestPrompt: '' diff --git a/src/shared/structured-agent-session-projection.ts b/src/shared/structured-agent-session-projection.ts index 0b3afa22228..a48537d05b7 100644 --- a/src/shared/structured-agent-session-projection.ts +++ b/src/shared/structured-agent-session-projection.ts @@ -190,12 +190,17 @@ export function hasPersistedStructuredAgentSessionTurn( * that sent it, so there is nothing still running to report. */ export function hasUnansweredStructuredAgentSessionDispatch( - submissions: readonly AgentJournalSubmission[] + submissions: readonly AgentJournalSubmission[], + currentFence?: number | null ): boolean { return submissions.some( (submission) => - submission.dispatchState === 'pending' || - (submission.dispatchState === 'unknown' && submission.recovered !== true) + (currentFence == null || submission.fence >= currentFence) && + (submission.dispatchState === 'pending' || + (submission.dispatchState === 'unknown' && + submission.recovered !== true && + // Older hosts publish the recovery reason but omit the optional marker. + submission.reason !== 'host_restarted_before_acknowledgement')) ) } @@ -207,7 +212,8 @@ export function structuredAgentSessionTabId(sessionId: string): string { export function projectStructuredAgentSessionStatus( items: readonly AgentJournalRenderItem[], - submissions: readonly AgentJournalSubmission[] = [] + submissions: readonly AgentJournalSubmission[] = [], + currentFence?: number | null ): StructuredAgentSessionProjectedStatus { if ( items.some( @@ -219,7 +225,7 @@ export function projectStructuredAgentSessionStatus( return 'attention' } return activeStructuredAgentSessionTurnId(items) || - hasUnansweredStructuredAgentSessionDispatch(submissions) + hasUnansweredStructuredAgentSessionDispatch(submissions, currentFence) ? 'working' : 'idle' } @@ -299,17 +305,18 @@ export type StructuredAgentSessionStatusProjection = { * has to stay small even though the row only ever renders one line of it. */ export function projectStructuredAgentSessionStatusSummary( items: readonly AgentJournalRenderItem[], - submissions: readonly AgentJournalSubmission[] = [] + submissions: readonly AgentJournalSubmission[] = [], + currentFence?: number | null ): StructuredAgentSessionStatusProjection { // A first send has no journalled message until the provider replays it, so the pending // dispatch is also what makes a brand-new session listable at all. if ( !hasPersistedStructuredAgentSessionTurn(items) && - !hasUnansweredStructuredAgentSessionDispatch(submissions) + !hasUnansweredStructuredAgentSessionDispatch(submissions, currentFence) ) { return { status: null, latestPrompt: '' } } - const status = projectStructuredAgentSessionStatus(items, submissions) + const status = projectStructuredAgentSessionStatus(items, submissions, currentFence) const activeToolCall = status === 'working' ? activeStructuredAgentSessionToolCall(items) : null const toolName = activeToolCall ? normalizeOptionalField(activeToolCall.name, AGENT_STATUS_TOOL_NAME_MAX_LENGTH)