diff --git a/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts b/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts index 810ffeca3be..91e1eba6151 100644 --- a/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts +++ b/src/main/runtime/orca-runtime-preserved-branch-cleanup.ts @@ -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({ 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 590f588fb0f..9cdd36b615a 100644 --- a/src/main/runtime/orca-runtime-subscribe-to-terminal-resize.ts +++ b/src/main/runtime/orca-runtime-subscribe-to-terminal-resize.ts @@ -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 } diff --git a/src/main/runtime/rpc/methods/orchestration.ts b/src/main/runtime/rpc/methods/orchestration.ts index ed89ae4519d..fbc8f263cd0 100644 --- a/src/main/runtime/rpc/methods/orchestration.ts +++ b/src/main/runtime/rpc/methods/orchestration.ts @@ -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) diff --git a/src/main/runtime/rpc/methods/orchestration/messaging/send-point-to-point.ts b/src/main/runtime/rpc/methods/orchestration/messaging/send-point-to-point.ts index 4ff72454397..52a33780be1 100644 --- a/src/main/runtime/rpc/methods/orchestration/messaging/send-point-to-point.ts +++ b/src/main/runtime/rpc/methods/orchestration/messaging/send-point-to-point.ts @@ -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 } diff --git a/src/main/runtime/rpc/methods/orchestration/worker/worker-release.ts b/src/main/runtime/rpc/methods/orchestration/worker/worker-release.ts index 4a81f05b4b9..46aa8a62175 100644 --- a/src/main/runtime/rpc/methods/orchestration/worker/worker-release.ts +++ b/src/main/runtime/rpc/methods/orchestration/worker/worker-release.ts @@ -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 -) diff --git a/src/main/runtime/rpc/methods/settled-worker-resume-fence-sweep.ts b/src/main/runtime/rpc/methods/settled-worker-resume-fence-sweep.ts index 33315a2f176..977007e7abc 100644 --- a/src/main/runtime/rpc/methods/settled-worker-resume-fence-sweep.ts +++ b/src/main/runtime/rpc/methods/settled-worker-resume-fence-sweep.ts @@ -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 + } + } +} diff --git a/src/main/runtime/runtime-legacy-worker-terminal-recovery-persistence.ts b/src/main/runtime/runtime-legacy-worker-terminal-recovery-persistence.ts index acfa675b5b1..836f9be1426 100644 --- a/src/main/runtime/runtime-legacy-worker-terminal-recovery-persistence.ts +++ b/src/main/runtime/runtime-legacy-worker-terminal-recovery-persistence.ts @@ -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() + 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, - changedHostIds: Set + changedHostIds: Set, + 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) diff --git a/src/main/runtime/runtime-legacy-worker-terminal-resume-fence.test.ts b/src/main/runtime/runtime-legacy-worker-terminal-resume-fence.test.ts index ec4eac5514d..ff078c52af2 100644 --- a/src/main/runtime/runtime-legacy-worker-terminal-resume-fence.test.ts +++ b/src/main/runtime/runtime-legacy-worker-terminal-resume-fence.test.ts @@ -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' + ) + }) +}) diff --git a/src/main/runtime/runtime-notifier-contract.ts b/src/main/runtime/runtime-notifier-contract.ts index a652f051935..aa2982082b4 100644 --- a/src/main/runtime/runtime-notifier-contract.ts +++ b/src/main/runtime/runtime-notifier-contract.ts @@ -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, diff --git a/src/main/window/runtime-window-lifecycle.ts b/src/main/window/runtime-window-lifecycle.ts index 78c5f2ec426..2a6ab95a07f 100644 --- a/src/main/window/runtime-window-lifecycle.ts +++ b/src/main/window/runtime-window-lifecycle.ts @@ -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, diff --git a/src/preload/api/agent-status-api.ts b/src/preload/api/agent-status-api.ts index 7aa6c21115d..89677022506 100644 --- a/src/preload/api/agent-status-api.ts +++ b/src/preload/api/agent-status-api.ts @@ -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 /** Drop a paneKey from the main-process hook cache and on-disk last-status file. Fire-and-forget. */ drop: (paneKey: string) => void diff --git a/src/preload/api/agent-status-bridge.ts b/src/preload/api/agent-status-bridge.ts index 3cc1654aaed..3c3415cd207 100644 --- a/src/preload/api/agent-status-bridge.ts +++ b/src/preload/api/agent-status-bridge.ts @@ -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 => 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. */ diff --git a/src/renderer/src/hooks/ipc-events/agent-status-listeners.ts b/src/renderer/src/hooks/ipc-events/agent-status-listeners.ts index 70b13e22a41..2426dcb932d 100644 --- a/src/renderer/src/hooks/ipc-events/agent-status-listeners.ts +++ b/src/renderer/src/hooks/ipc-events/agent-status-listeners.ts @@ -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) + } } diff --git a/src/renderer/src/hooks/useIpcEvents-lifecycle.test.ts b/src/renderer/src/hooks/useIpcEvents-lifecycle.test.ts index 6e0e7238c05..4f11e13a060 100644 --- a/src/renderer/src/hooks/useIpcEvents-lifecycle.test.ts +++ b/src/renderer/src/hooks/useIpcEvents-lifecycle.test.ts @@ -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', diff --git a/src/renderer/src/web/preload-api/web-agent-status-api.ts b/src/renderer/src/web/preload-api/web-agent-status-api.ts index d7c9740018c..1a07b6d6a6c 100644 --- a/src/renderer/src/web/preload-api/web-agent-status-api.ts +++ b/src/renderer/src/web/preload-api/web-agent-status-api.ts @@ -12,6 +12,7 @@ export function createWebAgentStatusApi(): Partial { onMigrationUnsupported: () => noopUnsubscribe, onMigrationUnsupportedClear: () => noopUnsubscribe, onLegacyWorkerTerminalRecovery: () => noopUnsubscribe, + onLegacyWorkerTerminalResumeFence: () => noopUnsubscribe, getMigrationUnsupportedSnapshot: () => Promise.resolve([]), drop: () => {}, dropPersisted: () => {},