mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(orchestration): bound the stopping-row exit credit and dedupe concurrent stops
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<StopReceipt>,
|
||||
h.call('orchestration.workerStop', { dispatch: dispatchId }) as Promise<StopReceipt>
|
||||
])
|
||||
|
||||
expect(first).toMatchObject({ state: 'stopped' })
|
||||
expect(second).toEqual(first)
|
||||
expect(h.db.getWorkerDispatch(dispatchId)?.state).toBe('stopped')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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<OrcaRuntimeService, Map<string, Promise<unknown>>>()
|
||||
|
||||
/** 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<unknown>
|
||||
): Promise<unknown> {
|
||||
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<unknown> = stop().finally(() => {
|
||||
if (active.get(dispatchId) === started) {
|
||||
active.delete(dispatchId)
|
||||
}
|
||||
})
|
||||
active.set(dispatchId, started)
|
||||
return started
|
||||
}
|
||||
|
||||
type RemoteStopReceipt = {
|
||||
state: string
|
||||
alreadySettled: boolean
|
||||
|
||||
Reference in New Issue
Block a user