diff --git a/src/cli/handlers/orchestration-migration.test.ts b/src/cli/handlers/orchestration-migration.test.ts index 863720eaab1..84c1038e3d7 100644 --- a/src/cli/handlers/orchestration-migration.test.ts +++ b/src/cli/handlers/orchestration-migration.test.ts @@ -1,6 +1,9 @@ import { afterEach, describe, expect, it, vi } from 'vitest' import { ORCHESTRATION_HANDLERS } from './orchestration' +// worker_done carries a CLI-minted request id so its runtime_unavailable retries replay one mutation. +const WORKER_DONE_REQUEST = { orchestrationRequestId: expect.any(String) } + const originalPaneKey = process.env.ORCA_PANE_KEY afterEach(() => { @@ -34,20 +37,24 @@ describe('orchestration CLI migration recovery', () => { json: true } as never) - expect(call).toHaveBeenCalledWith('orchestration.send', { - from: 'term_worker', - to: undefined, - run: undefined, - subject: 'Done', - body: undefined, - type: 'worker_done', - priority: undefined, - threadId: undefined, - payload: undefined, - senderPaneKey: 'tab-worker:leaf-worker', - waitForLifecycleSettlement: true, - devMode: false - }) + expect(call).toHaveBeenCalledWith( + 'orchestration.send', + { + from: 'term_worker', + to: undefined, + run: undefined, + subject: 'Done', + body: undefined, + type: 'worker_done', + priority: undefined, + threadId: undefined, + payload: undefined, + senderPaneKey: 'tab-worker:leaf-worker', + waitForLifecycleSettlement: true, + devMode: false + }, + WORKER_DONE_REQUEST + ) expect(call).toHaveBeenCalledOnce() } ) diff --git a/src/cli/handlers/orchestration.test.ts b/src/cli/handlers/orchestration.test.ts index e285f431ad4..4ca66b603eb 100644 --- a/src/cli/handlers/orchestration.test.ts +++ b/src/cli/handlers/orchestration.test.ts @@ -1,6 +1,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const callMock = vi.fn() +// worker_done carries a CLI-minted request id so its runtime_unavailable retries replay one mutation. +const WORKER_DONE_REQUEST = { orchestrationRequestId: expect.any(String) } const getTerminalHandleMock = vi.hoisted(() => vi.fn()) const originalTerminalHandle = process.env.ORCA_TERMINAL_HANDLE const originalPaneKey = process.env.ORCA_PANE_KEY @@ -79,24 +81,28 @@ describe('orchestration send structured payload flags', () => { ]) ) - expect(callMock).toHaveBeenCalledWith('orchestration.send', { - from: 'term_worker', - to: 'term_coord', - subject: 'done', - body: undefined, - type: 'worker_done', - priority: undefined, - threadId: undefined, - payload: JSON.stringify({ - taskId: 'task_1', - dispatchId: 'ctx_1', - outcome: 'succeeded', - filesModified: ['src/a.ts', 'src/b.ts'], - reportPath: 'reports/done.md' - }), - waitForLifecycleSettlement: true, - devMode: false - }) + expect(callMock).toHaveBeenCalledWith( + 'orchestration.send', + { + from: 'term_worker', + to: 'term_coord', + subject: 'done', + body: undefined, + type: 'worker_done', + priority: undefined, + threadId: undefined, + payload: JSON.stringify({ + taskId: 'task_1', + dispatchId: 'ctx_1', + outcome: 'succeeded', + filesModified: ['src/a.ts', 'src/b.ts'], + reportPath: 'reports/done.md' + }), + waitForLifecycleSettlement: true, + devMode: false + }, + WORKER_DONE_REQUEST + ) }) it('forwards multiline message bodies without normalization', async () => { @@ -202,18 +208,22 @@ describe('orchestration send structured payload flags', () => { ]) ) - expect(callMock).toHaveBeenCalledWith('orchestration.send', { - from: 'term_worker', - to: 'term_coord', - subject: 'done', - body: undefined, - type: 'worker_done', - priority: undefined, - threadId: undefined, - payload: JSON.stringify({ outcome: 'succeeded' }), - waitForLifecycleSettlement: true, - devMode: false - }) + expect(callMock).toHaveBeenCalledWith( + 'orchestration.send', + { + from: 'term_worker', + to: 'term_coord', + subject: 'done', + body: undefined, + type: 'worker_done', + priority: undefined, + threadId: undefined, + payload: JSON.stringify({ outcome: 'succeeded' }), + waitForLifecycleSettlement: true, + devMode: false + }, + WORKER_DONE_REQUEST + ) }) it('sends lifecycle messages from ORCA_TERMINAL_HANDLE without a liveness probe', async () => { @@ -229,18 +239,22 @@ describe('orchestration send structured payload flags', () => { ) expect(callMock).toHaveBeenCalledTimes(1) - expect(callMock).toHaveBeenCalledWith('orchestration.send', { - from: 'term_worker_env', - to: 'term_coord', - subject: 'done', - body: undefined, - type: 'worker_done', - priority: undefined, - threadId: undefined, - payload: JSON.stringify({ outcome: 'succeeded' }), - waitForLifecycleSettlement: true, - devMode: false - }) + expect(callMock).toHaveBeenCalledWith( + 'orchestration.send', + { + from: 'term_worker_env', + to: 'term_coord', + subject: 'done', + body: undefined, + type: 'worker_done', + priority: undefined, + threadId: undefined, + payload: JSON.stringify({ outcome: 'succeeded' }), + waitForLifecycleSettlement: true, + devMode: false + }, + WORKER_DONE_REQUEST + ) }) it.each(['worker_done', 'heartbeat'] as const)( @@ -265,7 +279,8 @@ describe('orchestration send structured payload flags', () => { expect(callMock).toHaveBeenCalledTimes(1) expect(callMock).toHaveBeenCalledWith( 'orchestration.send', - expect.objectContaining({ from: 'term_worker_env' }) + expect.objectContaining({ from: 'term_worker_env' }), + ...(type === 'worker_done' ? [WORKER_DONE_REQUEST] : []) ) } ) @@ -285,7 +300,8 @@ describe('orchestration send structured payload flags', () => { expect(callMock).toHaveBeenCalledWith( 'orchestration.send', - expect.objectContaining({ senderPaneKey: 'tab_worker:leaf_worker' }) + expect.objectContaining({ senderPaneKey: 'tab_worker:leaf_worker' }), + WORKER_DONE_REQUEST ) }) diff --git a/src/cli/handlers/orchestration/message-send-handler.ts b/src/cli/handlers/orchestration/message-send-handler.ts index aa9bbaff47e..3b024de2cff 100644 --- a/src/cli/handlers/orchestration/message-send-handler.ts +++ b/src/cli/handlers/orchestration/message-send-handler.ts @@ -12,6 +12,8 @@ import { throwNoActiveSenderTerminal } from './terminal-identity' +const WORKER_DONE_UNAVAILABLE_RETRY_MS = 120_000 + type LifecycleSendResult = | { action: 'completed' | 'failed' @@ -110,7 +112,9 @@ export const ORCHESTRATION_SEND_HANDLER: Record = { flags, 'orchestration.send', sendParams, - dispatchCapability ? { orchestrationCapability: dispatchCapability } : undefined + dispatchCapability ? { orchestrationCapability: dispatchCapability } : undefined, + // Why: a worker reports once and ends its turn, so a brief app outage must delay worker_done, not drop it. + type === 'worker_done' ? WORKER_DONE_UNAVAILABLE_RETRY_MS : 0 ) await requireWorkerDoneSettlement(client, type, sendParams.payload, result.result) if ('lifecycle' in result.result && result.result.lifecycle?.action === 'rejected') { diff --git a/src/cli/handlers/orchestration/mutation-request.test.ts b/src/cli/handlers/orchestration/mutation-request.test.ts new file mode 100644 index 00000000000..24d98e75357 --- /dev/null +++ b/src/cli/handlers/orchestration/mutation-request.test.ts @@ -0,0 +1,240 @@ +import { createServer, type Server } from 'node:net' +import { mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { z } from 'zod' +import { ORCHESTRATION_CONTRACT_RUNTIME_CAPABILITY } from '../../../shared/protocol-version' +import { RuntimeClient, RuntimeClientError } from '../../runtime-client' +import { callOrchestrationMutation } from './mutation-request' + +const RETRY_MS = 120_000 +const WORKER_DONE = { + from: 'term_worker', + subject: 'done', + type: 'worker_done', + payload: '{"taskId":"task_1","dispatchId":"ctx_1","outcome":"succeeded"}' +} + +const RuntimeRequest = z.object({ + id: z.string(), + method: z.string(), + orchestrationRequestId: z.string().optional() +}) + +const servers = new Set() +const tempDirs = new Set() + +afterEach(async () => { + vi.useRealTimers() + await Promise.all([...servers].map((server) => new Promise((resolve) => server.close(resolve)))) + servers.clear() + for (const dir of tempDirs) { + rmSync(dir, { recursive: true, force: true }) + } + tempDirs.clear() +}) + +type CallOptions = { orchestrationRequestId?: string } + +function fakeClient(respond: (attempt: number, options?: CallOptions) => unknown): { + client: RuntimeClient + requestIds: (string | undefined)[] +} { + const requestIds: (string | undefined)[] = [] + const call = vi.fn(async (_method: string, _params: unknown, options?: CallOptions) => { + requestIds.push(options?.orchestrationRequestId) + return respond(requestIds.length, options) + }) + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: callOrchestrationMutation only uses client.call. + return { client: { call } as unknown as RuntimeClient, requestIds } +} + +function unavailable(options?: CallOptions): RuntimeClientError { + return new RuntimeClientError( + 'runtime_unavailable', + 'Could not connect to the running Orca app.', + { + orchestrationRequestId: options?.orchestrationRequestId, + originalCommand: ['orca', 'orchestration', 'send', '--type', 'worker_done'] + } + ) +} + +describe('callOrchestrationMutation runtime_unavailable retry', () => { + it.each([ + ['a fresh', new Map(), expect.stringMatching(/^[0-9a-f-]{36}$/)], + [ + 'an explicit --retry-request', + new Map([['retry-request', '11111111-2222-4333-8444-555555555555']]), + '11111111-2222-4333-8444-555555555555' + ] + ])('retries with %s request id until the runtime answers', async (_case, flags, requestId) => { + vi.useFakeTimers() + const { client, requestIds } = fakeClient((attempt, options) => { + if (attempt < 3) { + throw unavailable(options) + } + return { ok: true, result: 'sent' } + }) + const call = callOrchestrationMutation( + client, + flags, + 'orchestration.send', + WORKER_DONE, + undefined, + RETRY_MS + ) + await vi.advanceTimersByTimeAsync(3_000) + await expect(call).resolves.toEqual({ ok: true, result: 'sent' }) + expect(requestIds).toHaveLength(3) + expect(new Set(requestIds).size).toBe(1) + expect(requestIds[0]).toEqual(requestId) + }) + + it('does not retry an error other than runtime_unavailable', async () => { + const { client, requestIds } = fakeClient(() => { + throw new RuntimeClientError('runtime_timeout', 'Timed out.') + }) + await expect( + callOrchestrationMutation( + client, + new Map(), + 'orchestration.send', + WORKER_DONE, + undefined, + RETRY_MS + ) + ).rejects.toMatchObject({ code: 'runtime_timeout' }) + expect(requestIds).toHaveLength(1) + }) + + it('does not retry mutations that did not opt in', async () => { + const { client, requestIds } = fakeClient((_attempt, options) => { + throw unavailable(options) + }) + await expect( + callOrchestrationMutation(client, new Map(), 'orchestration.send', WORKER_DONE) + ).rejects.toMatchObject({ code: 'runtime_unavailable' }) + expect(requestIds).toEqual([undefined]) + }) + + it('keeps the recovery command when the last attempt fails before its request id is attached', async () => { + vi.useFakeTimers() + const { client, requestIds } = fakeClient((attempt, options) => { + throw attempt === 1 + ? unavailable(options) + : new RuntimeClientError('runtime_unavailable', 'No runtime metadata.') + }) + const settled = callOrchestrationMutation( + client, + new Map(), + 'orchestration.send', + WORKER_DONE, + undefined, + RETRY_MS + ).catch((error: unknown) => error) + await vi.advanceTimersByTimeAsync(RETRY_MS + 30_000) + expect(await settled).toMatchObject({ + data: { recovery: { orchestrationRequestId: requestIds[0] } } + }) + }) + + it('gives up after about two minutes and prints the recovery command with the same id', async () => { + vi.useFakeTimers() + const { client, requestIds } = fakeClient((_attempt, options) => { + throw unavailable(options) + }) + const call = callOrchestrationMutation( + client, + new Map(), + 'orchestration.send', + WORKER_DONE, + undefined, + RETRY_MS + ) + const settled = call.catch((error: unknown) => error) + await vi.advanceTimersByTimeAsync(RETRY_MS + 30_000) + const error = await settled + expect(requestIds.length).toBeGreaterThan(5) + expect(new Set(requestIds).size).toBe(1) + expect(error).toMatchObject({ + code: 'runtime_unavailable', + data: { + recovery: { + retryCommand: [ + 'orca', + 'orchestration', + 'send', + '--type', + 'worker_done', + '--retry-request', + requestIds[0] + ] + } + } + }) + }) + + // Why: the fake runtime listens on a Unix socket path, which Windows does not accept. + it.skipIf(process.platform === 'win32')( + 'reaches a runtime that comes back after the first attempt found it down', + async () => { + const userDataPath = mkdtempSync(join(tmpdir(), 'orca-worker-done-retry-')) + tempDirs.add(userDataPath) + const endpoint = join(userDataPath, 'runtime.sock') + const sendRequestIds: string[] = [] + const server = createServer((socket) => { + socket.setEncoding('utf8') + socket.once('data', (line: string) => { + const request = RuntimeRequest.parse(JSON.parse(line.trim())) + if (request.method === 'orchestration.send') { + sendRequestIds.push(String(request.orchestrationRequestId)) + if (sendRequestIds.length === 1) { + // A runtime that drops the connection mid-request leaves the outcome unknown. + socket.destroy() + return + } + } + socket.end( + `${JSON.stringify({ + id: request.id, + ok: true, + result: + request.method === 'status.get' + ? { capabilities: [ORCHESTRATION_CONTRACT_RUNTIME_CAPABILITY] } + : { message: { id: 'msg_1' } }, + _meta: { runtimeId: 'runtime-1' } + })}\n` + ) + }) + }) + servers.add(server) + const client = new RuntimeClient(userDataPath, 5_000, null, null, 'orca') + const call = callOrchestrationMutation( + client, + new Map(), + 'orchestration.send', + WORKER_DONE, + undefined, + RETRY_MS + ) + // The runtime is down for the first attempt: no metadata, so even the contract probe fails. + await new Promise((resolve) => server.listen(endpoint, resolve)) + writeFileSync( + join(userDataPath, 'orca-runtime.json'), + JSON.stringify({ + runtimeId: 'runtime-1', + pid: 1, + transports: [{ kind: 'unix', endpoint }], + authToken: 'token', + startedAt: 1 + }) + ) + await expect(call).resolves.toMatchObject({ result: { message: { id: 'msg_1' } } }) + expect(sendRequestIds).toHaveLength(2) + expect(sendRequestIds[1]).toBe(sendRequestIds[0]) + }, + 15_000 + ) +}) diff --git a/src/cli/handlers/orchestration/mutation-request.ts b/src/cli/handlers/orchestration/mutation-request.ts index c262a81b86c..2e9836c20ef 100644 --- a/src/cli/handlers/orchestration/mutation-request.ts +++ b/src/cli/handlers/orchestration/mutation-request.ts @@ -1,21 +1,56 @@ -import type { RuntimeClient } from '../../runtime-client' +import { randomUUID } from 'node:crypto' +import { RuntimeClientError, type RuntimeClient } from '../../runtime-client' import { readRetryRequestFlag } from '../../retry-request-flag' import { orchestrationMutationRecoveryError } from '../../orchestration-mutation-recovery' -export function callOrchestrationMutation( +const MAX_UNAVAILABLE_RETRY_DELAY_MS = 15_000 + +export async function callOrchestrationMutation( client: RuntimeClient, flags: Map, method: string, params: unknown, - options?: { timeoutMs?: number; orchestrationCapability?: string } + options?: { timeoutMs?: number; orchestrationCapability?: string }, + unavailableRetryMs = 0 ) { - const requestId = readRetryRequestFlag(flags) - const result = requestId - ? client.call(method, params, { ...options, orchestrationRequestId: requestId }) - : options - ? client.call(method, params, options) - : client.call(method, params) - return result.catch((error) => { - throw orchestrationMutationRecoveryError(error) - }) + // Why: every retry reuses one request id, so the host replays instead of applying the mutation twice. + const requestId = + readRetryRequestFlag(flags) ?? (unavailableRetryMs > 0 ? randomUUID() : undefined) + const deadline = Date.now() + unavailableRetryMs + let sentError: RuntimeClientError | undefined + for (let delayMs = 1_000; ; delayMs = Math.min(delayMs * 2, MAX_UNAVAILABLE_RETRY_DELAY_MS)) { + try { + return requestId + ? await client.call(method, params, { + ...options, + orchestrationRequestId: requestId + }) + : options + ? await client.call(method, params, options) + : await client.call(method, params) + } catch (error) { + const unavailable = + error instanceof RuntimeClientError && error.code === 'runtime_unavailable' + if (carriesRequestId(error)) { + sentError = error + } + if (!unavailable || Date.now() + delayMs > deadline) { + // Why: a later attempt can fail before its request id is attached, though an earlier one may have landed. + throw orchestrationMutationRecoveryError( + carriesRequestId(error) ? error : (sentError ?? error) + ) + } + await new Promise((resolve) => setTimeout(resolve, delayMs)) + } + } +} + +function carriesRequestId(error: unknown): error is RuntimeClientError { + const data: unknown = error instanceof RuntimeClientError ? error.data : undefined + return ( + typeof data === 'object' && + data !== null && + 'orchestrationRequestId' in data && + typeof data.orchestrationRequestId === 'string' + ) } diff --git a/src/cli/runtime/client.ts b/src/cli/runtime/client.ts index 68a099ff2f2..5efe54c7ba3 100644 --- a/src/cli/runtime/client.ts +++ b/src/cli/runtime/client.ts @@ -245,7 +245,13 @@ export class RuntimeClient { private async ensureOrchestrationContractCompatible(timeoutMs: number): Promise { if (!this.orchestrationContractCheck) { - this.orchestrationContractCheck = this.checkOrchestrationContractCompatibility(timeoutMs) + this.orchestrationContractCheck = this.checkOrchestrationContractCompatibility( + timeoutMs + ).catch((error: unknown) => { + // Why: a failed probe must not be cached, or a retry after a brief outage never reaches the app. + this.orchestrationContractCheck = null + throw error + }) } await this.orchestrationContractCheck }