From 515ea7d88cfa1e0030021d91c57d33c40dd09acd Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Wed, 16 Sep 2026 20:44:41 -0700 Subject: [PATCH] fix(mobile): use launch receipts to authorize replay --- .../worktree-create-launch-operation.test.ts | 26 +++++++- mobile/src/tasks/worktree-create-retry.ts | 28 ++++----- .../agent-launch-mobile-replay.test.ts | 61 ++++++++++++++++--- 3 files changed, 89 insertions(+), 26 deletions(-) diff --git a/mobile/src/tasks/worktree-create-launch-operation.test.ts b/mobile/src/tasks/worktree-create-launch-operation.test.ts index 01788b0f491..353826bfc65 100644 --- a/mobile/src/tasks/worktree-create-launch-operation.test.ts +++ b/mobile/src/tasks/worktree-create-launch-operation.test.ts @@ -101,6 +101,7 @@ describe('agent.launch operation id', () => { attempts: Attempt[] replay?: boolean supported?: boolean + worktreeCreateIdempotency?: false mintLaunchOperationId?: () => string }): Promise { let minted = 0 @@ -108,7 +109,7 @@ describe('agent.launch operation id', () => { client: args.client, baseName: 'otter', buildParams: (name) => ({ repo: 'id:r', name }), - worktreeCreateIdempotency: IDEMPOTENT_CREATE_SUPPORT, + worktreeCreateIdempotency: args.worktreeCreateIdempotency ?? IDEMPOTENT_CREATE_SUPPORT, mintMutationId: () => 'key-launch', agentLaunch: { agent: 'claude', @@ -154,6 +155,29 @@ describe('agent.launch operation id', () => { expect(launchOperationIds(attempts)).toEqual(['op-1', 'op-1']) }) + it('uses launch replay support independently of worktree.create idempotency', async () => { + const attempts: Attempt[] = [] + const client = scriptedLaunchClient( + [{ throws: new LogicalClientCutoverError() }, { launched: 'wt-replay' }], + attempts + ) + + await expect( + launchRetry({ client, attempts, worktreeCreateIdempotency: false }) + ).resolves.toEqual({ worktreeId: 'wt-replay', name: 'otter' }) + expect(launchOperationIds(attempts)).toEqual(['op-1', 'op-1']) + }) + + it('bounds named timeout retries without minting another operation', async () => { + const attempts: Attempt[] = [] + const error = markRpcDeliveryUnknown(new Error('Request timed out')) + const client = scriptedLaunchClient([{ throws: error }], attempts) + + await expect(launchRetry({ client, attempts })).rejects.toBe(error) + expect(attempts).toHaveLength(3) + expect(launchOperationIds(attempts)).toEqual(['op-1', 'op-1', 'op-1']) + }) + // The invariant. `computeAgentLaunchFingerprint` folds `target` whole, so the candidate name is // inside the fingerprint: carrying one id onto the bumped name would meet its own row under a // different fingerprint and refuse `agent_session_operation_conflict`, failing the create. diff --git a/mobile/src/tasks/worktree-create-retry.ts b/mobile/src/tasks/worktree-create-retry.ts index 0aea698abfc..1c0fe73d1a0 100644 --- a/mobile/src/tasks/worktree-create-retry.ts +++ b/mobile/src/tasks/worktree-create-retry.ts @@ -209,6 +209,8 @@ async function sendWorktreeCreateResilient( params: WorkspaceCreateParams, worktreeCreateIdempotency: WorktreeCreateIdempotencySupport | false ): Promise<{ response: RpcResponse; replayed: boolean }> { + // Only the selected method's receipt can authorize replay after an ambiguous delivery. + const replaySupport = launchAgent ? launchOperationId : worktreeCreateIdempotency let migrationRetry = 0 let ambiguousRetry = 0 const firstSentAt = Date.now() @@ -229,7 +231,7 @@ async function sendWorktreeCreateResilient( // A refusal on the replacement connection says nothing about what the first call created. return { response, replayed: migrationRetry > 0 || ambiguousRetry > 0 } } catch (error) { - if (!worktreeCreateIdempotency) { + if (!replaySupport) { throw error } if (isLogicalClientCutoverError(error)) { @@ -244,29 +246,21 @@ async function sendWorktreeCreateResilient( if (!isRpcDeliveryUnknown(error) || ambiguousRetry >= WORKTREE_CREATE_AMBIGUOUS_MAX_RETRIES) { throw error } - // Why: every transport path that reports a *drop* leaves 'connected' before the - // rejection reaches us (rpc-client.ts:675/695/1213 set state first or reject via - // queueMicrotask; the relay's fail() publishes synchronously). So still being - // 'connected' here means the socket was healthy the whole time and only the - // response went missing — the request-timeout path, which surfaces after - // WORKTREE_CREATE_TIMEOUT_MS. That says nothing about when the host actually - // resolved, so the dedupe record may be long gone and a replay would build a - // second worktree instead of reconciling. Fail the create instead. - if (client.getState() === 'connected') { + // A legacy cache may expire before a request timeout; durable receipts refuse unsafe replay. + if (typeof replaySupport !== 'string' && client.getState() === 'connected') { throw error } - // Computed once: a later ambiguity reads a fresher lastInboundAt from the - // replacement session, which would push the deadline past the record it respects. - replayDeadlineAt ??= resolveReplayDeadline(client, firstSentAt, worktreeCreateIdempotency) + // Keep the legacy deadline fixed; the host itself refuses expired durable operation IDs. + replayDeadlineAt ??= + typeof replaySupport === 'string' + ? Infinity + : resolveReplayDeadline(client, firstSentAt, replaySupport) const remainingWindowMs = replayDeadlineAt - Date.now() if (remainingWindowMs <= 0) { throw error } ambiguousRetry += 1 - // Why: unlike a cutover, no replacement session exists yet — resending now - // would just hit the dead one, so wait for the transport to come back and - // surface the original ambiguity if it does not. Clamped to the window so the - // wait itself cannot carry the replay past the host's record. + // Disconnected transports must reconnect before resend; bound the wait even for durable IDs. if ( !(await waitForRpcClientReconnected( client, diff --git a/src/main/runtime/rpc/methods/agent-launch-mobile-replay.test.ts b/src/main/runtime/rpc/methods/agent-launch-mobile-replay.test.ts index 82b0684c5cc..45a63a5c19c 100644 --- a/src/main/runtime/rpc/methods/agent-launch-mobile-replay.test.ts +++ b/src/main/runtime/rpc/methods/agent-launch-mobile-replay.test.ts @@ -4,6 +4,7 @@ import { join } from 'node:path' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { createWorktreeWithNameRetry } from '../../../../../mobile/src/tasks/worktree-create-retry' import type { RpcClient } from '../../../../../mobile/src/transport/rpc-client' +import { markRpcDeliveryUnknown } from '../../../../../mobile/src/transport/rpc-delivery-ambiguity' import { LogicalClientCutoverError } from '../../../../../mobile/src/transport/stable-logical-rpc-client' import { AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS, @@ -40,13 +41,24 @@ afterEach(async () => { await rm(directory, { recursive: true, force: true }) }) -function mobileLaunch(args: { loseFirstReplyAfterMs?: number } = {}) { +function mobileLaunch( + args: { + loseFirstReplyAfterMs?: number + replay?: boolean + restartAfterReply?: boolean + replyLoss?: 'cutover' | 'timeout' + } = {} +) { const runtime = { ...runtimeStub(), getRuntimeId: () => 'runtime-1' } - const dispatcher = new RpcDispatcher({ - // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: this fixture implements the launch handler and dispatcher metadata dependencies. - runtime: runtime as unknown as OrcaRuntimeService, - methods: AGENT_LAUNCH_METHODS - }) + const runtimeAfterRestart = { ...runtimeStub(), getRuntimeId: () => 'runtime-2' } + const dispatchers = [runtime, runtimeAfterRestart].map( + (host) => + new RpcDispatcher({ + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: this fixture implements the launch handler and dispatcher metadata dependencies. + runtime: host as unknown as OrcaRuntimeService, + methods: AGENT_LAUNCH_METHODS + }) + ) const operationId = `${Date.now()}-000000000000000000000000000000aa` const attempts: unknown[] = [] // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the mobile retry loop reaches only these transport members; requests use the real host dispatcher. @@ -54,6 +66,7 @@ function mobileLaunch(args: { loseFirstReplyAfterMs?: number } = {}) { getState: () => 'connected', sendRequest: async (method: string, params: unknown) => { attempts.push(params) + const dispatcher = dispatchers[args.restartAfterReply && attempts.length > 1 ? 1 : 0]! const response = await dispatcher.dispatch({ id: `request-${attempts.length}`, authToken: 'token', @@ -63,6 +76,9 @@ function mobileLaunch(args: { loseFirstReplyAfterMs?: number } = {}) { if (attempts.length === 1 && args.loseFirstReplyAfterMs !== undefined) { const later = Date.now() + args.loseFirstReplyAfterMs vi.spyOn(Date, 'now').mockReturnValue(later) + if (args.replyLoss === 'timeout') { + throw markRpcDeliveryUnknown(new Error('Request timed out')) + } throw new LogicalClientCutoverError() } return response @@ -73,13 +89,32 @@ function mobileLaunch(args: { loseFirstReplyAfterMs?: number } = {}) { baseName: 'otter', buildParams: (name) => ({ repo: 'id:repo-1', name }), worktreeCreateIdempotency: { dedupeTtlMs: 60_000 }, - agentLaunch: { agent: 'claude', supported: { replay: true } }, + agentLaunch: { agent: 'claude', supported: { replay: args.replay !== false } }, mintLaunchOperationId: () => operationId }) - return { runtime, attempts, operationId, result } + return { runtime, runtimeAfterRestart, attempts, operationId, result } } describe('mobile launch retries through the host ledger', () => { + it('does not replay an unnamed launch after an older host loses its in-memory receipt', async () => { + const launch = mobileLaunch({ + replay: false, + restartAfterReply: true, + loseFirstReplyAfterMs: 1 + }) + + const outcome = await launch.result.catch((error: unknown) => error) + expect( + launch.runtime.createManagedWorktree.mock.calls.length + + launch.runtimeAfterRestart.createManagedWorktree.mock.calls.length + ).toBe(1) + expect(outcome).toBeInstanceOf(LogicalClientCutoverError) + expect(launch.runtime.createManagedWorktree).toHaveBeenCalledTimes(1) + expect(launch.runtimeAfterRestart.createManagedWorktree).not.toHaveBeenCalled() + expect(createStructuredSession).toHaveBeenCalledTimes(1) + expect(launch.attempts).toHaveLength(1) + }) + it.each([ 'agent_session_operation_capacity', 'agent_session_operation_invalid', @@ -116,4 +151,14 @@ describe('mobile launch retries through the host ledger', () => { expect(launch.attempts[1]).toEqual(launch.attempts[0]) expect(launch.runtime.dedupeWorktreeCreate).not.toHaveBeenCalled() }) + + it('recovers a named launch whose reply timed out on a connected transport', async () => { + const launch = mobileLaunch({ loseFirstReplyAfterMs: 10 * 60_000, replyLoss: 'timeout' }) + + await expect(launch.result).resolves.toEqual({ worktreeId: 'wt-new', name: 'otter' }) + expect(launch.runtime.createManagedWorktree).toHaveBeenCalledTimes(1) + expect(createStructuredSession).toHaveBeenCalledTimes(1) + expect(launch.attempts).toHaveLength(2) + expect(launch.attempts[1]).toEqual(launch.attempts[0]) + }) })