From 2736a72edd192e998db1b08ac44b8a2793e7ea14 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Fri, 4 Sep 2026 04:21:39 -0400 Subject: [PATCH] fix(orchestration): bound the stopping-row exit credit and dedupe concurrent stops --- ...ca-runtime-subscribe-to-terminal-resize.ts | 5 +- .../worker/worker-stop-exit-race.test.ts | 36 +- .../orchestration/worker/worker-stop.ts | 336 ++++++++++-------- 3 files changed, 220 insertions(+), 157 deletions(-) diff --git a/src/main/runtime/orca-runtime-subscribe-to-terminal-resize.ts b/src/main/runtime/orca-runtime-subscribe-to-terminal-resize.ts index 9cdd36b615a..a169b1efd69 100644 --- a/src/main/runtime/orca-runtime-subscribe-to-terminal-resize.ts +++ b/src/main/runtime/orca-runtime-subscribe-to-terminal-resize.ts @@ -83,7 +83,10 @@ export class OrcaRuntimeWithSubscribeToTerminalResize extends OrcaRuntimeWithApp } // A process that dies while we are stopping it is that stop succeeding, not a failure: // settling it as `failed` here made the in-flight worker-stop report its own success as an error. - if (this._orchestrationDb.getWorkerDispatch?.(dispatch.id)?.state === 'stopping') { + // Only a stop begun in THIS runtime can claim the exit; a `stopping` row left durable by a + // killed process would otherwise absorb a much later crash as a clean stop. + const stopping = this._orchestrationDb.getWorkerDispatch?.(dispatch.id) + if (stopping?.state === 'stopping' && stopping.runtime_epoch === this.getRuntimeId()) { this._orchestrationDb.settleWorkerStop(dispatch.id) this.sweepSettledWorkerResumeFencesAfterExit() return diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-stop-exit-race.test.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-stop-exit-race.test.ts index 9135012590e..6f4dbcb3715 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-stop-exit-race.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-stop-exit-race.test.ts @@ -1,5 +1,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -import { OPERATOR_CLOSE_EXIT_CAUSE } from '../../../../../../shared/terminal-exit-cause' +import { + OPERATOR_CLOSE_EXIT_CAUSE, + type TerminalExitCause +} from '../../../../../../shared/terminal-exit-cause' import { createOrchestrationWorkerReleaseHarness } from './worker-release.test-support' const h = createOrchestrationWorkerReleaseHarness() @@ -8,17 +11,17 @@ afterEach(() => h.cleanup()) type StopReceipt = { state: string; alreadySettled: boolean; processAction: string } -function fireExit(handle: string): void { +function fireExit(handle: string, cause: TerminalExitCause = OPERATOR_CLOSE_EXIT_CAUSE): void { ;( h.runtime as unknown as { failActiveDispatchOnExit: ( handle: string, paneKey: string | null, exitCode: number, - cause: typeof OPERATOR_CLOSE_EXIT_CAUSE + cause: TerminalExitCause ) => void } - ).failActiveDispatchOnExit(handle, h.workerPaneKey, 0, OPERATOR_CLOSE_EXIT_CAUSE) + ).failActiveDispatchOnExit(handle, h.workerPaneKey, 0, cause) } describe('a worker whose process exits while its own stop is in flight', () => { @@ -60,4 +63,29 @@ describe('a worker whose process exits while its own stop is in flight', () => { fireExit('term_worker') expect(h.db.getWorkerDispatch(dispatchId)?.state).toBe('failed') }) + + it('certifies a later death instead of crediting a stopping row from a dead runtime', async () => { + const { dispatchId } = await h.startWorker() + // The stop RPC committed `stopping` in an earlier runtime and the app died before settling. + h.db.beginWorkerStop(dispatchId, 'runtime_from_a_previous_process') + expect(h.db.getWorkerDispatch(dispatchId)?.state).toBe('stopping') + + fireExit('term_worker', { kind: 'signaled', signal: 9 }) + + expect(h.db.getDispatchContextById(dispatchId)?.termination_reason).toBe('signaled') + expect(h.db.getWorkerDispatch(dispatchId)?.state).toBe('failed') + }) + + it('gives a second concurrent stop the first caller receipt, not dispatch_inactive', async () => { + const { dispatchId } = await h.startWorker() + + const [first, second] = await Promise.all([ + h.call('orchestration.workerStop', { dispatch: dispatchId }) as Promise, + h.call('orchestration.workerStop', { dispatch: dispatchId }) as Promise + ]) + + expect(first).toMatchObject({ state: 'stopped' }) + expect(second).toEqual(first) + expect(h.db.getWorkerDispatch(dispatchId)?.state).toBe('stopped') + }) }) diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-stop.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-stop.ts index 53320a4f36a..6438040e1bf 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-stop.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-stop.ts @@ -5,6 +5,7 @@ import { requiredString } from '../../../schemas' import { describeUnconfirmedAgentStop } from '../../../../../../shared/pty-liveness-verdict' import { ORCHESTRATION_WORKER_STOP_VERDICT_RUNTIME_CAPABILITY } from '../../../../../../shared/protocol-version' import type { RuntimeStatus } from '../../../../../../shared/runtime-types' +import type { OrcaRuntimeService } from '../../../../orca-runtime' import { inspectWorkerTerminal, resolvePinnedFederatedServer } from './worker-observation' const WorkerDispatchParams = z.object({ dispatch: requiredString('Missing --dispatch') }) @@ -13,182 +14,213 @@ export const ORCHESTRATION_WORKER_STOP_METHODS: RpcMethod[] = [ defineMethod({ name: 'orchestration.workerStop', params: WorkerDispatchParams, - handler: async (params, { runtime, orchestrationMutation }) => { - const db = runtime.getOrchestrationDb() - const federated = db.getFederatedDispatch(params.dispatch) - if (federated) { - if (!orchestrationMutation) { - throw new OrchestrationError( - 'invalid_argument', - 'Remote worker-stop requires a durable retry request.' - ) - } - const server = resolvePinnedFederatedServer(runtime, federated) - const begun = db.beginWorkerStop(params.dispatch, runtime.getRuntimeId()) - if (begun.disposition === 'already_settled') { - return settledReceipt(params.dispatch, begun.worker.state) - } - try { - const status = (await runtime.callOrchestrationWorkerServer( - server.environmentId, - 'status.get', - undefined, - 30_000, - undefined, - { expectedEnvironmentPairingRevision: server.pairingRevision } - )) as RuntimeStatus - if ( - !status.capabilities?.includes(ORCHESTRATION_WORKER_STOP_VERDICT_RUNTIME_CAPABILITY) - ) { + handler: (params, { runtime, orchestrationMutation }) => + dedupeWorkerStop(runtime, params.dispatch, async () => { + const db = runtime.getOrchestrationDb() + const federated = db.getFederatedDispatch(params.dispatch) + if (federated) { + if (!orchestrationMutation) { + throw new OrchestrationError( + 'invalid_argument', + 'Remote worker-stop requires a durable retry request.' + ) + } + const server = resolvePinnedFederatedServer(runtime, federated) + const begun = db.beginWorkerStop(params.dispatch, runtime.getRuntimeId()) + if (begun.disposition === 'already_settled') { + return settledReceipt(params.dispatch, begun.worker.state) + } + try { + const status = (await runtime.callOrchestrationWorkerServer( + server.environmentId, + 'status.get', + undefined, + 30_000, + undefined, + { expectedEnvironmentPairingRevision: server.pairingRevision } + )) as RuntimeStatus + if ( + !status.capabilities?.includes(ORCHESTRATION_WORKER_STOP_VERDICT_RUNTIME_CAPABILITY) + ) { + return unknownReceipt( + params.dispatch, + db.markWorkerStopUnknown( + params.dispatch, + `Connected server ${server.name} cannot prove the worker stop outcome.` + ), + 'none' + ) + } + const remote = (await runtime.callOrchestrationWorkerServer( + server.environmentId, + 'orchestration.federationStop', + { dispatchId: params.dispatch }, + 30_000, + { orchestrationRequestId: orchestrationMutation.requestId }, + { expectedEnvironmentPairingRevision: server.pairingRevision } + )) as RemoteStopReceipt + if (remote.state === 'stopped') { + const worker = db.reconcileFederatedWorkerStop(params.dispatch) + return { + dispatchId: params.dispatch, + state: worker.state, + alreadySettled: remote.alreadySettled, + processAction: remote.processAction, + close: remote.close + } + } + if (remote.state === 'succeeded' || remote.state === 'failed') { + db.resumeFederatedWorkerForTerminalRelay(params.dispatch) + await runtime + .syncOrchestrationFederatedDispatchAfterCurrent(params.dispatch) + .catch(() => undefined) + return { + dispatchId: params.dispatch, + state: db.getWorkerDispatch(params.dispatch)?.state ?? remote.state, + alreadySettled: true, + processAction: 'none' + } + } return unknownReceipt( params.dispatch, db.markWorkerStopUnknown( params.dispatch, - `Connected server ${server.name} cannot prove the worker stop outcome.` + remote.lastError ?? `The worker server returned ${remote.state}.` ), - 'none' + remote.processAction + ) + } catch (error) { + const reason = error instanceof Error ? error.message : String(error) + return unknownReceipt( + params.dispatch, + db.markWorkerStopUnknown(params.dispatch, reason), + 'unknown' ) } - const remote = (await runtime.callOrchestrationWorkerServer( - server.environmentId, - 'orchestration.federationStop', - { dispatchId: params.dispatch }, - 30_000, - { orchestrationRequestId: orchestrationMutation.requestId }, - { expectedEnvironmentPairingRevision: server.pairingRevision } - )) as RemoteStopReceipt - if (remote.state === 'stopped') { - const worker = db.reconcileFederatedWorkerStop(params.dispatch) - return { - dispatchId: params.dispatch, - state: worker.state, - alreadySettled: remote.alreadySettled, - processAction: remote.processAction, - close: remote.close - } + } + + const begun = db.beginWorkerStop(params.dispatch, runtime.getRuntimeId()) + if (begun.disposition === 'already_settled') { + return settledReceipt(params.dispatch, begun.worker.state) + } + if (begun.disposition === 'context_only') { + if (!begun.alreadySettled) { + runtime.notifyMessageArrived(`dispatch:${params.dispatch}`, 'status') } - if (remote.state === 'succeeded' || remote.state === 'failed') { - db.resumeFederatedWorkerForTerminalRelay(params.dispatch) - await runtime - .syncOrchestrationFederatedDispatchAfterCurrent(params.dispatch) - .catch(() => undefined) - return { - dispatchId: params.dispatch, - state: db.getWorkerDispatch(params.dispatch)?.state ?? remote.state, - alreadySettled: true, - processAction: 'none' - } + return { + dispatchId: params.dispatch, + state: begun.state, + alreadySettled: begun.alreadySettled, + processAction: 'none' as const, + warning: contextOnlyStopWarning(begun) } + } + const handle = begun.worker.agent_terminal_handle + if (!handle) { return unknownReceipt( params.dispatch, db.markWorkerStopUnknown( params.dispatch, - remote.lastError ?? `The worker server returned ${remote.state}.` + 'The Dispatch has no recorded agent terminal.' ), - remote.processAction - ) - } catch (error) { - const reason = error instanceof Error ? error.message : String(error) - return unknownReceipt( - params.dispatch, - db.markWorkerStopUnknown(params.dispatch, reason), 'unknown' ) } - } - - const begun = db.beginWorkerStop(params.dispatch, runtime.getRuntimeId()) - if (begun.disposition === 'already_settled') { - return settledReceipt(params.dispatch, begun.worker.state) - } - if (begun.disposition === 'context_only') { - if (!begun.alreadySettled) { - runtime.notifyMessageArrived(`dispatch:${params.dispatch}`, 'status') + const observation = await inspectWorkerTerminal(runtime, db, params.dispatch) + // Why `unverifiable` still proceeds: losing contact is a reason to report + // the outcome honestly, never a reason to stop trying to stop the worker. + if ( + !observation.exact || + (observation.status !== 'live' && observation.status !== 'unverifiable') + ) { + return unknownReceipt( + params.dispatch, + db.markWorkerStopUnknown( + params.dispatch, + `The recorded worker process is ${observation.status}; no terminal was closed.` + ), + 'none' + ) } + const resource = db.getWorkerTerminalResourceByOwner(params.dispatch) + if (!resource || resource.ownership_state !== 'owned') { + const ownership = resource?.ownership_state ?? 'unproven' + return unknownReceipt( + params.dispatch, + db.markWorkerStopUnknown( + params.dispatch, + `The worker terminal is ${ownership}; no terminal was closed.` + ), + 'none' + ) + } + const closed = await runtime + .closeTerminal(handle) + .then((close) => ({ close }) as const) + .catch( + (error: unknown) => + ({ error: error instanceof Error ? error.message : String(error) }) as const + ) + // The process exit can land mid-close and settle the stop from the exit path; that exit + // is this stop's proof of success, so do not re-settle it or report it as unknown. + if (db.getWorkerDispatch(params.dispatch)?.state !== 'stopped') { + if ('error' in closed) { + return unknownReceipt( + params.dispatch, + db.markWorkerStopUnknown(params.dispatch, closed.error), + 'unknown' + ) + } + if (!closed.close.ptyKilled) { + // The tab is retired, but the agent process was never confirmed stopped — + // settling here is the false success this receipt exists to prevent. + return unknownReceipt( + params.dispatch, + db.markWorkerStopUnknown(params.dispatch, describeUnconfirmedAgentStop(closed.close)), + 'closed_agent_terminal' + ) + } + db.settleWorkerStop(params.dispatch) + } + runtime.notifyMessageArrived(`dispatch:${params.dispatch}`, 'status') return { dispatchId: params.dispatch, - state: begun.state, - alreadySettled: begun.alreadySettled, - processAction: 'none' as const, - warning: contextOnlyStopWarning(begun) + state: db.getWorkerDispatch(params.dispatch)?.state ?? 'stopped', + alreadySettled: false, + processAction: 'closed_agent_terminal', + ...('close' in closed ? { close: closed.close } : {}) } - } - const handle = begun.worker.agent_terminal_handle - if (!handle) { - return unknownReceipt( - params.dispatch, - db.markWorkerStopUnknown(params.dispatch, 'The Dispatch has no recorded agent terminal.'), - 'unknown' - ) - } - const observation = await inspectWorkerTerminal(runtime, db, params.dispatch) - // Why `unverifiable` still proceeds: losing contact is a reason to report - // the outcome honestly, never a reason to stop trying to stop the worker. - if ( - !observation.exact || - (observation.status !== 'live' && observation.status !== 'unverifiable') - ) { - return unknownReceipt( - params.dispatch, - db.markWorkerStopUnknown( - params.dispatch, - `The recorded worker process is ${observation.status}; no terminal was closed.` - ), - 'none' - ) - } - const resource = db.getWorkerTerminalResourceByOwner(params.dispatch) - if (!resource || resource.ownership_state !== 'owned') { - const ownership = resource?.ownership_state ?? 'unproven' - return unknownReceipt( - params.dispatch, - db.markWorkerStopUnknown( - params.dispatch, - `The worker terminal is ${ownership}; no terminal was closed.` - ), - 'none' - ) - } - const closed = await runtime - .closeTerminal(handle) - .then((close) => ({ close }) as const) - .catch( - (error: unknown) => - ({ error: error instanceof Error ? error.message : String(error) }) as const - ) - // The process exit can land mid-close and settle the stop from the exit path; that exit - // is this stop's proof of success, so do not re-settle it or report it as unknown. - if (db.getWorkerDispatch(params.dispatch)?.state !== 'stopped') { - if ('error' in closed) { - return unknownReceipt( - params.dispatch, - db.markWorkerStopUnknown(params.dispatch, closed.error), - 'unknown' - ) - } - if (!closed.close.ptyKilled) { - // The tab is retired, but the agent process was never confirmed stopped — - // settling here is the false success this receipt exists to prevent. - return unknownReceipt( - params.dispatch, - db.markWorkerStopUnknown(params.dispatch, describeUnconfirmedAgentStop(closed.close)), - 'closed_agent_terminal' - ) - } - db.settleWorkerStop(params.dispatch) - } - runtime.notifyMessageArrived(`dispatch:${params.dispatch}`, 'status') - return { - dispatchId: params.dispatch, - state: db.getWorkerDispatch(params.dispatch)?.state ?? 'stopped', - alreadySettled: false, - processAction: 'closed_agent_terminal', - ...('close' in closed ? { close: closed.close } : {}) - } - } + }) }) ] +const activeStopByRuntime = new WeakMap>>() + +/** Two callers stopping one Dispatch: the second reached `beginWorkerStop` after the first moved + * the row to `stopping` and got `dispatch_inactive` instead of the first caller's receipt. */ +function dedupeWorkerStop( + runtime: OrcaRuntimeService, + dispatchId: string, + stop: () => Promise +): Promise { + let active = activeStopByRuntime.get(runtime) + if (!active) { + active = new Map() + activeStopByRuntime.set(runtime, active) + } + const inFlight = active.get(dispatchId) + if (inFlight) { + return inFlight + } + const started: Promise = stop().finally(() => { + if (active.get(dispatchId) === started) { + active.delete(dispatchId) + } + }) + active.set(dispatchId, started) + return started +} + type RemoteStopReceipt = { state: string alreadySettled: boolean