Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions mobile/src/session/mobile-structured-agent-session-send.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<MobileNativeChatSendOutcome> {
const timeoutMs = timeoutForDeadline(input.deadline)
if (timeoutMs === null) {
Expand Down Expand Up @@ -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({
Expand Down
Original file line number Diff line number Diff line change
@@ -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<typeof useMobileStructuredAgentSession> | null = null
let listener: ((value: unknown) => void) | null = null
let storedOperations: Map<string, string>
const onSendError = vi.fn()
const sendRequest = vi.fn<RpcClient['sendRequest']>()
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<void> {
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<typeof ok> | Promise<ReturnType<typeof ok>>
): 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<unknown> = 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()
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
9 changes: 6 additions & 3 deletions mobile/src/session/use-mobile-structured-agent-session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 <TValue>(
Expand Down Expand Up @@ -219,11 +219,13 @@ export function useMobileStructuredAgentSession(args: {
text,
attachments: sendAttachments,
deadline,
onError: onSendError
onError: onSendError,
onAwaitingSettlement: awaitSend
})
},
[
agent,
awaitSend,
callerIdentity,
client,
conversationCommands,
Expand All @@ -233,7 +235,8 @@ export function useMobileStructuredAgentSession(args: {
optionSnapshot,
sessionId,
sessionKey,
setStructuredOption
setStructuredOption,
stateRef
]
)
const { groupedDraft, respondPermission, respondQuestion } = useMobileStructuredPromptResponses({
Expand Down
Original file line number Diff line number Diff line change
@@ -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<string>())
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]
)
}
12 changes: 12 additions & 0 deletions src/main/claude/claude-structured-session-close.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -94,6 +98,14 @@ async function finalizeClaudePublishedSession(
input: CloseClaudePublishedSessionInput,
session: ClaudeSession
): Promise<boolean> {
// 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.
Expand Down
Loading
Loading