fix(orchestration): retry worker_done while the Orca runtime is briefly unreachable (#23984)

* fix(orchestration): retry worker_done while the Orca runtime is briefly unreachable

A worker reports worker_done once and ends its turn, so a few-minute app
outage silently stranded finished work at dispatched. The CLI now retries
worker_done on runtime_unavailable for about two minutes with backoff,
reusing one request id so the host's mutation ledger replays rather than
double-applies it, then prints the existing recovery command.

The contract probe no longer caches a failed status.get, which otherwise
made every retry fail without reaching the app.

Part of STA-8833.

* refactor(cli): pass the worker_done retry window in the mutation options bag

* Revert "refactor(cli): pass the worker_done retry window in the mutation options bag"

The options bag is forwarded to client.call as-is; the retry window is not a client.call option, and folding it in needed a value scan to keep the no-options call shape.

* test(cli): fold the explicit retry-request case and drop a vacuous timing assert

* fix(cli): keep worker_done recovery when the last retry fails before sending

Also skip the Unix-socket retry test on Windows and remove its temp profile.
This commit is contained in:
Jinwoo Hong
2026-09-30 13:04:57 -04:00
committed by GitHub
parent e03870403e
commit 46d6b76ae8
6 changed files with 380 additions and 72 deletions
@@ -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()
}
)
+60 -44
View File
@@ -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
)
})
@@ -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<string, CommandHandler> = {
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') {
@@ -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<Server>()
const tempDirs = new Set<string>()
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<string, string | boolean>(), 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<void>((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
)
})
@@ -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<TResult>(
const MAX_UNAVAILABLE_RETRY_DELAY_MS = 15_000
export async function callOrchestrationMutation<TResult>(
client: RuntimeClient,
flags: Map<string, string | boolean>,
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<TResult>(method, params, { ...options, orchestrationRequestId: requestId })
: options
? client.call<TResult>(method, params, options)
: client.call<TResult>(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<TResult>(method, params, {
...options,
orchestrationRequestId: requestId
})
: options
? await client.call<TResult>(method, params, options)
: await client.call<TResult>(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'
)
}
+7 -1
View File
@@ -245,7 +245,13 @@ export class RuntimeClient {
private async ensureOrchestrationContractCompatible(timeoutMs: number): Promise<void> {
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
}