mirror of
https://github.com/stablyai/orca.git
synced 2026-10-08 16:02:37 +00:00
fix(mobile): use launch receipts to authorize replay
This commit is contained in:
@@ -101,6 +101,7 @@ describe('agent.launch operation id', () => {
|
||||
attempts: Attempt[]
|
||||
replay?: boolean
|
||||
supported?: boolean
|
||||
worktreeCreateIdempotency?: false
|
||||
mintLaunchOperationId?: () => string
|
||||
}): Promise<WorktreeCreateResult> {
|
||||
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.
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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])
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user