fix(orchestration): fence a settled worker's pane in the same session

setSleepingAgentAutomaticResumeBlocked had no production caller: the stamp only
ran at startup and the sweep only after release/retain/takeover, so a worker
settled by worker_done, stop, abandon or its own process exit left its pane
unfenced and reopening it in the same session respawned the agent. Every
settlement now sweeps, and the sweep pushes each fence change to the live
renderer over agentStatus:legacyWorkerTerminalResumeFence.
This commit is contained in:
Jinwoo-H
2026-09-04 03:35:25 -04:00
parent 5e1f38120c
commit d24ce644df
15 changed files with 188 additions and 34 deletions
@@ -131,7 +131,8 @@ export class OrcaRuntimeWithPreservedBranchCleanup extends OrcaRuntimeWithTermin
new RuntimeLegacyWorkerTerminalRecoveryPersistence(
() => this.store,
() => this.getOrchestrationDb(),
(worktreeId) => this.tryGetWorkspaceSessionHostIdForWorktree(worktreeId)
(worktreeId) => this.tryGetWorkspaceSessionHostIdForWorktree(worktreeId),
(paneKey, blocked) => this.notifier?.setLegacyWorkerTerminalResumeFence?.(paneKey, blocked)
)
protected readonly legacyWorkerRecovery = new RuntimeLegacyWorkerTerminalRecoveryController({
@@ -52,6 +52,16 @@ export class OrcaRuntimeWithSubscribeToTerminalResize extends OrcaRuntimeWithApp
// dispatch contexts immediately, rather than waiting for the coordinator's
// next poll cycle. This catches agent crashes and unexpected exits within
// milliseconds. The task is set back to 'pending' so it can be re-dispatched.
/** A worker settled by its own process exit makes its pane fenceable now, not at the next app
* start; a fence sweep must never fail the exit path behind it. */
private sweepSettledWorkerResumeFencesAfterExit(): void {
try {
this.prepareLegacyWorkerTerminalRecovery()
} catch (error) {
console.warn('[orchestration] settled worker resume fence sweep failed', error)
}
}
protected failActiveDispatchOnExit(
handle: string,
paneKey: string | null,
@@ -75,6 +85,7 @@ export class OrcaRuntimeWithSubscribeToTerminalResize extends OrcaRuntimeWithApp
// 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') {
this._orchestrationDb.settleWorkerStop(dispatch.id)
this.sweepSettledWorkerResumeFencesAfterExit()
return
}
@@ -83,6 +94,7 @@ export class OrcaRuntimeWithSubscribeToTerminalResize extends OrcaRuntimeWithApp
workerProcessExited: true,
terminationReason: cause.kind
})
this.sweepSettledWorkerResumeFencesAfterExit()
if (isDeliberateTerminalExit(cause)) {
return
}
@@ -1,4 +1,5 @@
import type { RpcMethod } from '../core'
import { sweepingSettledWorkerResumeFences } from './settled-worker-resume-fence-sweep'
import { ORCHESTRATION_RUN_METHODS } from './orchestration/runs/runs'
import { ORCHESTRATION_WORKER_METHODS } from './orchestration/worker/worker-methods'
import { ORCHESTRATION_FEDERATION_METHODS } from './orchestration/federation/federation-methods'
@@ -23,4 +24,4 @@ export const ORCHESTRATION_METHODS: RpcMethod[] = [
...ORCHESTRATION_ASK_METHODS,
...ORCHESTRATION_GATE_METHODS,
...ORCHESTRATION_RESET_METHODS
]
].map(sweepingSettledWorkerResumeFences)
@@ -6,6 +6,7 @@ import { isDispatchMutationMessageType, parseMessageTaskId } from '../schemas'
import type { SendParams } from '../schemas'
import { legacyWorkerDeliveryContract } from '../routing'
import { recordReceiptForPostCommitNudge } from './mutation-replay-nudge'
import { sweepSettledWorkerResumeFences } from '../../settled-worker-resume-fence-sweep'
import type { SendRecipientWarning } from './recipient-routing'
import type { z } from 'zod'
@@ -143,6 +144,11 @@ export function sendPointToPointMessage(args: {
? db.commitWorkerDoneMessageMutation(commitMessage)
: commitMessage()
committed.nudge()
if (messageType === 'worker_done') {
// Settlement is what makes the pane fenceable; without this the fence only appeared at the
// next app start and reopening the pane in the same session respawned the agent.
sweepSettledWorkerResumeFences(runtime)
}
return committed.receipt
}
@@ -13,29 +13,7 @@ import {
import { WorkerDispatchParams, WorkerRetainParams } from './worker-release-schemas'
import { sweepSettledWorkerResumeFences } from '../../settled-worker-resume-fence-sweep'
// Release and retain both drop the worker's row from the legacy recovery plan, and a fenced pane
// refuses a fresh spawn — so the sweep runs after every early return of both methods. workerList
// is a pure read and is deliberately absent.
const FENCE_SWEEPING_METHOD_NAMES = new Set([
'orchestration.workerRelease',
'orchestration.workerRetain'
])
function sweepingRetiredWorkerResumeFences(method: RpcMethod): RpcMethod {
if (!FENCE_SWEEPING_METHOD_NAMES.has(method.name)) {
return method
}
return {
...method,
handler: async (params, ctx) => {
const result = await method.handler(params, ctx)
sweepSettledWorkerResumeFences(ctx.runtime)
return result
}
}
}
const WORKER_RELEASE_METHODS: RpcMethod[] = [
export const ORCHESTRATION_WORKER_RELEASE_METHODS: RpcMethod[] = [
defineMethod({
name: 'orchestration.workerRelease',
params: WorkerDispatchParams,
@@ -170,7 +148,3 @@ const WORKER_RELEASE_METHODS: RpcMethod[] = [
}
})
]
export const ORCHESTRATION_WORKER_RELEASE_METHODS: RpcMethod[] = WORKER_RELEASE_METHODS.map(
sweepingRetiredWorkerResumeFences
)
@@ -1,4 +1,5 @@
import type { OrcaRuntimeService } from '../../orca-runtime'
import type { RpcMethod } from '../core'
/**
* One pass both stamps the automatic-resume fence on every settled worker pane and lifts it from
@@ -14,3 +15,28 @@ export function sweepSettledWorkerResumeFences(runtime: OrcaRuntimeService): voi
console.warn('[orchestration] settled worker resume fence sweep failed', error)
}
}
/** Settling a worker is what makes its pane fenceable, and release/retain/takeover are what make it
* unfenceable again — so every one of those has to sweep in the same call. Without the settlement
* half the fence only appeared at the next app start, and reopening the pane in the same session
* respawned the agent. */
const FENCE_SWEEPING_METHOD_NAMES = new Set([
'orchestration.workerRelease',
'orchestration.workerRetain',
'orchestration.workerStop',
'orchestration.workerAbandon'
])
export function sweepingSettledWorkerResumeFences(method: RpcMethod): RpcMethod {
if (!FENCE_SWEEPING_METHOD_NAMES.has(method.name)) {
return method
}
return {
...method,
handler: async (params, ctx) => {
const result = await method.handler(params, ctx)
sweepSettledWorkerResumeFences(ctx.runtime)
return result
}
}
}
@@ -18,7 +18,9 @@ export class RuntimeLegacyWorkerTerminalRecoveryPersistence {
constructor(
private readonly getStore: () => RuntimeStore | null,
private readonly getDb: () => OrchestrationDb,
private readonly getHostId: (worktreeId: string) => ExecutionHostId | null
private readonly getHostId: (worktreeId: string) => ExecutionHostId | null,
/** The store write only reaches the next app start; a live renderer holds its own copy. */
private readonly notifyFenceChanged?: (paneKey: string, blocked: boolean) => void
) {}
prepare(): LegacyWorkerTerminalRecoveryPlan {
@@ -41,6 +43,7 @@ export class RuntimeLegacyWorkerTerminalRecoveryPersistence {
{ current: WorkspaceSessionState; next: WorkspaceSessionState }
>()
const changedHostIds = new Set<ExecutionHostId>()
const fenceChanges: [string, boolean][] = []
for (const blocked of plan.blockedPanes) {
let hostIds: ExecutionHostId[]
try {
@@ -79,9 +82,10 @@ export class RuntimeLegacyWorkerTerminalRecoveryPersistence {
[blocked.paneKey]: { ...record, automaticResumeBlockedBy: 'legacy-orchestration-worker' }
}
changedHostIds.add(hostId)
fenceChanges.push([blocked.paneKey, true])
}
}
this.liftRetiredFences(store, plan, sessions, changedHostIds)
this.liftRetiredFences(store, plan, sessions, changedHostIds, fenceChanges)
const changed = [...sessions].filter(([hostId]) => changedHostIds.has(hostId))
if (changed.length === 0) {
return plan
@@ -92,6 +96,10 @@ export class RuntimeLegacyWorkerTerminalRecoveryPersistence {
}
} catch (error) {
console.warn('[orchestration] failed to stage legacy worker resume fence', error)
return plan
}
for (const [paneKey, blocked] of fenceChanges) {
this.notifyFenceChanged?.(paneKey, blocked)
}
return plan
}
@@ -103,7 +111,8 @@ export class RuntimeLegacyWorkerTerminalRecoveryPersistence {
store: RuntimeStore,
plan: LegacyWorkerTerminalRecoveryPlan,
sessions: Map<ExecutionHostId, { current: WorkspaceSessionState; next: WorkspaceSessionState }>,
changedHostIds: Set<ExecutionHostId>
changedHostIds: Set<ExecutionHostId>,
fenceChanges: [string, boolean][]
): void {
const blockedPaneKeys = new Set(plan.blockedPanes.map((blocked) => blocked.paneKey))
for (const hostId of store.getWorkspaceSessionHostIds?.() ?? [LOCAL_EXECUTION_HOST_ID]) {
@@ -130,6 +139,7 @@ export class RuntimeLegacyWorkerTerminalRecoveryPersistence {
for (const [paneKey, record] of retired) {
const { automaticResumeBlockedBy: _retired, ...unfenced } = record
next[paneKey] = unfenced
fenceChanges.push([paneKey, false])
}
state.next.sleepingAgentSessionsByPaneKey = next
changedHostIds.add(hostId)
@@ -3,6 +3,8 @@ import { getDefaultWorkspaceSession } from '../../shared/constants'
import { LOCAL_EXECUTION_HOST_ID } from '../../shared/execution-host'
import type { WorkspaceSessionState } from '../../shared/workspace-session-state-types'
import { OrchestrationDb } from './orchestration/db'
import { OrcaRuntimeService } from './orca-runtime'
import { ORCHESTRATION_METHODS } from './rpc/methods/orchestration'
import { RuntimeLegacyWorkerTerminalRecoveryPersistence } from './runtime-legacy-worker-terminal-recovery-persistence'
import type { RuntimeStore } from './runtime-store-contract'
@@ -34,7 +36,7 @@ describe('settled worker automatic-resume fence persistence', () => {
afterEach(() => db?.close())
function harness(): {
function harness(onFenceChanged?: (paneKey: string, blocked: boolean) => void): {
db: OrchestrationDb
taskId: string
dispatchId: string
@@ -77,7 +79,8 @@ describe('settled worker automatic-resume fence persistence', () => {
persistence: new RuntimeLegacyWorkerTerminalRecoveryPersistence(
() => store,
() => orchestrationDb,
() => LOCAL_EXECUTION_HOST_ID
() => LOCAL_EXECUTION_HOST_ID,
onFenceChanged
),
fence: () => session.sleepingAgentSessionsByPaneKey?.[PANE_KEY]?.automaticResumeBlockedBy
}
@@ -89,6 +92,16 @@ describe('settled worker automatic-resume fence persistence', () => {
).toBe('settled')
}
it('pushes the fence to the live renderer instead of waiting for the next app start', () => {
const fenceChanges: [string, boolean][] = []
const h = harness((paneKey, blocked) => fenceChanges.push([paneKey, blocked]))
settle(h.db, h.taskId, h.dispatchId)
h.persistence.prepare()
expect(fenceChanges).toEqual([[PANE_KEY, true]])
})
// The STA-4577 repro: worker_done, no release, restart, open the worktree — the pane still
// holds a resumable provider session and must not respawn `codex resume`.
it('fences a settled worker pane whose terminal was never released', () => {
@@ -164,3 +177,85 @@ describe('settled worker automatic-resume fence persistence', () => {
expect(plan.candidates).toEqual([expect.objectContaining({ dispatchId: h.dispatchId })])
})
})
// STA-4577's other half: settlement with no release and no restart. The stamp only ran at startup
// and after release/retain/takeover, so reopening the pane in the same session respawned the agent.
describe('worker_done without a release', () => {
let db: OrchestrationDb | undefined
afterEach(() => db?.close())
it('fences the pane in the same session', async () => {
const orchestrationDb = new OrchestrationDb(':memory:')
db = orchestrationDb
let session = sessionWithSleepingWorker()
const store = {
getWorkspaceSession: () => session,
setWorkspaceSession: (next: WorkspaceSessionState) => {
session = next
},
getWorkspaceSessionHostIds: () => [LOCAL_EXECUTION_HOST_ID],
flushOrThrow: vi.fn()
} as unknown as RuntimeStore
const runtime = new OrcaRuntimeService(store)
runtime.setOrchestrationDb(orchestrationDb)
vi.spyOn(runtime, 'getTerminalPaneKey').mockImplementation((handle) =>
handle === 'term_worker' ? PANE_KEY : 'tab_coord:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
)
vi.spyOn(runtime, 'getTerminalProcessIncarnation').mockReturnValue('runtime:pty:1')
vi.spyOn(runtime, 'notifyMessageArrived').mockImplementation(() => {})
const run = orchestrationDb.createRun({
objective: 'settle without release',
coordinatorHandle: 'term_coord',
coordinatorPaneKey: 'tab_coord:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
})
const task = orchestrationDb.createTask({ spec: 'settle without release', runId: run.id })
const started = orchestrationDb.createStartingWorkerDispatch({
creator: { kind: 'system' },
maxDepth: Number.MAX_SAFE_INTEGER,
taskId: task.id,
startOptions: {}
})
orchestrationDb.prepareStartingWorkerAuthority({
dispatchId: started.dispatch.id,
handle: 'term_worker',
paneKey: PANE_KEY,
processIncarnation: 'runtime:pty:1',
worktreeId: WORKTREE_ID,
setupState: 'not_applicable',
effects: [],
terminalOwnership: 'created'
})
orchestrationDb.markWorkerDispatchReady(started.dispatch.id)
const capability = orchestrationDb.mintDispatchCapability({
dispatchId: started.dispatch.id,
paneKey: PANE_KEY,
processIncarnation: 'runtime:pty:1'
})
expect(session.sleepingAgentSessionsByPaneKey?.[PANE_KEY]?.automaticResumeBlockedBy).toBe(
undefined
)
const send = ORCHESTRATION_METHODS.find((method) => method.name === 'orchestration.send')!
await send.handler(
send.params!.parse({
from: 'term_worker',
to: 'term_coord',
subject: 'Done',
type: 'worker_done',
payload: JSON.stringify({
taskId: task.id,
dispatchId: started.dispatch.id,
outcome: 'succeeded'
})
}),
{ runtime, orchestrationCapability: capability }
)
expect(orchestrationDb.getWorkerDispatch(started.dispatch.id)?.state).toBe('succeeded')
expect(session.sleepingAgentSessionsByPaneKey?.[PANE_KEY]?.automaticResumeBlockedBy).toBe(
'legacy-orchestration-worker'
)
})
})
@@ -79,6 +79,8 @@ export type RuntimeNotifier = {
resolution: 'adopted' | 'exited' | 'rolled_back',
ptyId?: string
): void
/** The fence lives in the workspace session, which a live renderer only re-reads at startup. */
setLegacyWorkerTerminalResumeFence?(paneKey: string, blocked: boolean): void
splitTerminal(
tabId: string,
paneRuntimeId: number,
@@ -149,6 +149,8 @@ export function registerRuntimeWindowLifecycle(
resolution,
...(ptyId ? { ptyId } : {})
}),
setLegacyWorkerTerminalResumeFence: (paneKey, blocked) =>
send('agentStatus:legacyWorkerTerminalResumeFence', { paneKey, blocked }),
splitTerminal: (tabId, paneRuntimeId, opts) => {
send('ui:splitTerminal', {
tabId,
+4
View File
@@ -28,6 +28,10 @@ export type AgentStatusApi = {
ptyId?: string
}) => void
) => () => void
/** Listen for the automatic-resume fence a settled worker's pane gains or loses mid-session. */
onLegacyWorkerTerminalResumeFence: (
callback: (data: { paneKey: string; blocked: boolean }) => void
) => () => void
getMigrationUnsupportedSnapshot: () => Promise<MigrationUnsupportedPtyEntry[]>
/** Drop a paneKey from the main-process hook cache and on-disk last-status file. Fire-and-forget. */
drop: (paneKey: string) => void
+10
View File
@@ -61,6 +61,16 @@ export const agentStatusApi = {
ipcRenderer.on('agentStatus:legacyWorkerTerminalRecovery', listener)
return () => ipcRenderer.removeListener('agentStatus:legacyWorkerTerminalRecovery', listener)
},
onLegacyWorkerTerminalResumeFence: (
callback: (data: { paneKey: string; blocked: boolean }) => void
): (() => void) => {
const listener = (
_event: Electron.IpcRendererEvent,
data: { paneKey: string; blocked: boolean }
) => callback(data)
ipcRenderer.on('agentStatus:legacyWorkerTerminalResumeFence', listener)
return () => ipcRenderer.removeListener('agentStatus:legacyWorkerTerminalResumeFence', listener)
},
getMigrationUnsupportedSnapshot: (): Promise<MigrationUnsupportedPtyEntry[]> =>
ipcRenderer.invoke('agentStatus:getMigrationUnsupportedSnapshot'),
/** Drop the cached hook status for a paneKey on both sides (memory + on-disk) so a relaunch can't resurrect a dismissed row. */
@@ -126,4 +126,12 @@ export function registerAgentStatusListeners(args: {
if (unsubscribeLegacyWorkerTerminalRecovery) {
unsubs.push(unsubscribeLegacyWorkerTerminalRecovery)
}
const unsubscribeResumeFence = window.api.agentStatus.onLegacyWorkerTerminalResumeFence?.(
({ paneKey, blocked }) => {
useAppStore.getState().setSleepingAgentAutomaticResumeBlocked(paneKey, blocked)
}
)
if (unsubscribeResumeFence) {
unsubs.push(unsubscribeResumeFence)
}
}
@@ -5,6 +5,7 @@ import { createHarnessStoreState } from './ipc-events-test-harness'
const EXPECTED_DIRECT_CALLBACK_METHODS = [
'agentStatus.onClear',
'agentStatus.onLegacyWorkerTerminalRecovery',
'agentStatus.onLegacyWorkerTerminalResumeFence',
'agentStatus.onMigrationUnsupported',
'agentStatus.onMigrationUnsupportedClear',
'agentStatus.onSet',
@@ -196,6 +197,7 @@ const EXPECTED_CALLBACK_REGISTRATION_SEQUENCE = [
'agentStatus.onMigrationUnsupported',
'agentStatus.onMigrationUnsupportedClear',
'agentStatus.onLegacyWorkerTerminalRecovery',
'agentStatus.onLegacyWorkerTerminalResumeFence',
'runtime.onTerminalFitOverrideChanged',
'runtime.onTerminalDriverChanged',
'runtime.onNativeChatLaunchDraftResolved',
@@ -12,6 +12,7 @@ export function createWebAgentStatusApi(): Partial<PreloadApi> {
onMigrationUnsupported: () => noopUnsubscribe,
onMigrationUnsupportedClear: () => noopUnsubscribe,
onLegacyWorkerTerminalRecovery: () => noopUnsubscribe,
onLegacyWorkerTerminalResumeFence: () => noopUnsubscribe,
getMigrationUnsupportedSnapshot: () => Promise.resolve([]),
drop: () => {},
dropPersisted: () => {},