diff --git a/src/main/runtime/orca-runtime.test.ts b/src/main/runtime/orca-runtime.test.ts index 6aae780db35..d72bb1fa390 100644 --- a/src/main/runtime/orca-runtime.test.ts +++ b/src/main/runtime/orca-runtime.test.ts @@ -22172,7 +22172,7 @@ describe('OrcaRuntimeService', () => { expect(harness.resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith( harness.workerPaneKey, 'rolled_back', - harness.ptyId + { ptyId: harness.ptyId } ) expect(harness.resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith( harness.workerPaneKey, @@ -22240,7 +22240,7 @@ describe('OrcaRuntimeService', () => { expect(harness.resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith( harness.workerPaneKey, 'rolled_back', - harness.ptyId + { ptyId: harness.ptyId } ) expect(harness.resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith( harness.workerPaneKey, @@ -22461,7 +22461,7 @@ describe('OrcaRuntimeService', () => { expect(resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith( workerPaneKey, 'rolled_back', - 'pty-missing-worker' + { ptyId: 'pty-missing-worker' } ) expect(resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith(workerPaneKey, 'exited') } finally { @@ -22570,7 +22570,7 @@ describe('OrcaRuntimeService', () => { expect(resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith( workerPaneKey, 'rolled_back', - 'pty-missing-retry' + { ptyId: 'pty-missing-retry' } ) expect(resolveLegacyWorkerTerminalRecovery).toHaveBeenCalledWith(workerPaneKey, 'exited') warn.mockRestore() diff --git a/src/main/runtime/orca-runtime.ts b/src/main/runtime/orca-runtime.ts index 50ae709f43d..46c250e5cac 100644 --- a/src/main/runtime/orca-runtime.ts +++ b/src/main/runtime/orca-runtime.ts @@ -55,7 +55,8 @@ import { type AgentStatusIpcPayload, type ParsedAgentStatusPayload, type AgentStatusOrchestrationContext, - type AgentStatusEntry + type AgentStatusEntry, + type LegacyWorkerTerminalRecoveryResolutionKind } from '../../shared/agent-status-types' import { terminalStatusPayloadMatchesHook } from '../../shared/agent-terminal-status-equivalence' import { indexAgentStatusRowsByPaneKey } from '../agent-hooks/agent-status-pane-index' @@ -598,7 +599,8 @@ import { getRepoIdFromWorktreeId, splitWorktreeId, splitWorktreeIdForFilesystem, - worktreeIdComparisonKey + worktreeIdComparisonKey, + worktreeIdsEqual } from '../../shared/worktree/id' import { getProjectIdForProviderIdentity } from '../../shared/project-host-setup-projection' import { @@ -2410,8 +2412,8 @@ type RuntimeNotifier = { | void resolveLegacyWorkerTerminalRecovery?( paneKey: string, - resolution: 'adopted' | 'exited' | 'rolled_back', - ptyId?: string + resolution: LegacyWorkerTerminalRecoveryResolutionKind, + identity?: { ptyId?: string; worktreeId?: string } ): void splitTerminal( tabId: string, @@ -4717,19 +4719,29 @@ export class OrcaRuntimeService { this.scheduleRestoredMessageRepoints() } - private getLegacyWorkerTerminalRecoveryPlan(): LegacyWorkerTerminalRecoveryPlan { + private getLegacyWorkerTerminalRecoveryPlan(): LegacyWorkerTerminalRecoveryPlan | null { try { return planLegacyWorkerTerminalRecovery( this.getOrchestrationDb().listLegacyWorkerTerminalRecoveryRows() ) } catch (error) { console.warn('[orchestration] failed to plan legacy worker terminal recovery', error) - return { blockedPanes: [], candidates: [], ambiguousDispatchIds: [] } + return null } } prepareLegacyWorkerTerminalRecovery(): LegacyWorkerTerminalRecoveryPlan { const plan = this.getLegacyWorkerTerminalRecoveryPlan() + if (!plan) { + return { blockedPanes: [], candidates: [], ambiguousDispatchIds: [] } + } + for (const blocked of plan.blockedPanes) { + if (blocked.settled) { + this.notifier?.resolveLegacyWorkerTerminalRecovery?.(blocked.paneKey, 'fenced', { + worktreeId: blocked.worktreeId + }) + } + } const store = this.store if ( !store?.getWorkspaceSession || @@ -4767,7 +4779,7 @@ export class OrcaRuntimeService { const record = state.next.sleepingAgentSessionsByPaneKey?.[blocked.paneKey] if ( !record || - !runtimeWorktreeIdsEqual(record.worktreeId, blocked.worktreeId) || + !worktreeIdsEqual(record.worktreeId, blocked.worktreeId) || record.automaticResumeBlockedBy === 'legacy-orchestration-worker' ) { continue @@ -4782,6 +4794,7 @@ export class OrcaRuntimeService { changedHostIds.add(hostId) } } + this.liftRetiredLegacyWorkerResumeFences(plan, sessions, changedHostIds) const changed = [...sessions].filter(([hostId]) => changedHostIds.has(hostId)) if (changed.length === 0) { return plan @@ -4796,6 +4809,49 @@ export class OrcaRuntimeService { return plan } + private liftRetiredLegacyWorkerResumeFences( + plan: LegacyWorkerTerminalRecoveryPlan, + sessions: Map, + changedHostIds: Set + ): void { + const store = this.store + if (!store?.getWorkspaceSession) { + return + } + const blockedPaneKeys = new Set(plan.blockedPanes.map((blocked) => blocked.paneKey)) + for (const hostId of store.getWorkspaceSessionHostIds?.() ?? [LOCAL_EXECUTION_HOST_ID]) { + const staged = sessions.get(hostId) + const session = staged?.next ?? store.getWorkspaceSession(hostId) + const retired = Object.entries(session?.sleepingAgentSessionsByPaneKey ?? {}).filter( + ([paneKey, record]) => + record.automaticResumeBlockedBy === 'legacy-orchestration-worker' && + !blockedPaneKeys.has(paneKey) + ) + if (retired.length === 0) { + continue + } + let state = staged + if (!state) { + const current = store.getWorkspaceSession(hostId) + if (!current) { + continue + } + state = { current, next: structuredClone(current) } + sessions.set(hostId, state) + } + const next = { ...state.next.sleepingAgentSessionsByPaneKey } + for (const [paneKey, record] of retired) { + const { automaticResumeBlockedBy: _retiredFence, ...unfenced } = record + next[paneKey] = unfenced + this.notifier?.resolveLegacyWorkerTerminalRecovery?.(paneKey, 'unfenced', { + worktreeId: record.worktreeId + }) + } + state.next.sleepingAgentSessionsByPaneKey = next + changedHostIds.add(hostId) + } + } + private async flushWorkspaceSessionOrThrowAsync(): Promise { const store = this.store if (store?.flushPendingOrThrowAsync) { @@ -5042,11 +5098,9 @@ export class OrcaRuntimeService { pty.tabId = null pty.paneKey = null } - this.notifier?.resolveLegacyWorkerTerminalRecovery?.( - candidate.paneKey, - 'rolled_back', - candidate.ptyId - ) + this.notifier?.resolveLegacyWorkerTerminalRecovery?.(candidate.paneKey, 'rolled_back', { + ptyId: candidate.ptyId + }) } private updateLegacyWorkerTerminalRecoveryRetry( diff --git a/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts b/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts index 904cc0d5b98..b056e75fe02 100644 --- a/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts +++ b/src/main/runtime/orchestration/db/worker-dispatch/worker-terminal-recovery.ts @@ -8,6 +8,7 @@ import { OrchestrationError } from '../../orchestration-error' import { DISPATCH_CIRCUIT_BREAK_FAILURES } from '../dispatch-context/dispatch-circuit-breaker' import type { OrchestrationDb } from '../orchestration-db' import { reconcileTaskAfterDispatchInterruption } from '../dispatch-context/task-dispatch-reconciliation' +import { WORKER_SETTLED_STATES } from '../../worker-terminal-ownership' export function listLegacyWorkerTerminalRecoveryRows( this: OrchestrationDb @@ -21,9 +22,16 @@ export function listLegacyWorkerTerminalRecoveryRows( FROM dispatch_contexts dc INNER JOIN worker_dispatches wd ON wd.dispatch_id = dc.id WHERE wd.state IN ('starting', 'ready', 'start_unknown', 'stopping', 'stop_unknown') + OR (wd.state IN (${WORKER_SETTLED_STATES.map(() => '?').join(', ')}) + AND EXISTS ( + SELECT 1 FROM worker_terminal_resources wtr + WHERE wtr.owner_dispatch_id = dc.id + AND wtr.ownership_state = 'owned' + AND wtr.release_state NOT IN ('released', 'retained') + )) ORDER BY dc.rowid` ) - .all() as LegacyWorkerTerminalRecoveryRow[] + .all(...WORKER_SETTLED_STATES) as LegacyWorkerTerminalRecoveryRow[] } export function reconcileMissingWorkerTerminal( diff --git a/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.test.ts b/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.test.ts index d425a52f2cb..4b032c480bb 100644 --- a/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.test.ts +++ b/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.test.ts @@ -30,7 +30,8 @@ describe('legacy worker terminal recovery planning', () => { { worktreeId: 'repo::/workspace', paneKey: `tab-worker:${LEAF_ID}`, - contractVersion: 0 + contractVersion: 0, + settled: false } ], candidates: [ @@ -52,7 +53,8 @@ describe('legacy worker terminal recovery planning', () => { { worktreeId: 'repo::/workspace', paneKey: `tab-worker:${LEAF_ID}`, - contractVersion: 0 + contractVersion: 0, + settled: false } ], candidates: [], @@ -60,6 +62,34 @@ describe('legacy worker terminal recovery planning', () => { }) }) + it('fences a settled worker pane without offering its terminal for adoption', () => { + const plan = planLegacyWorkerTerminalRecovery([ + recoveryRow({ worker_state: 'succeeded', dispatch_status: 'completed' }) + ]) + expect(plan).toEqual({ + blockedPanes: [ + { + worktreeId: 'repo::/workspace', + paneKey: `tab-worker:${LEAF_ID}`, + contractVersion: 0, + settled: true + } + ], + candidates: [], + ambiguousDispatchIds: [] + }) + }) + + it('does not let a settled row make a live worker identity ambiguous', () => { + const plan = planLegacyWorkerTerminalRecovery([ + recoveryRow({ dispatch_id: 'dispatch-settled', worker_state: 'succeeded' }), + recoveryRow({ dispatch_id: 'dispatch-live' }) + ]) + expect(plan.candidates).toEqual([expect.objectContaining({ dispatchId: 'dispatch-live' })]) + expect(plan.ambiguousDispatchIds).toEqual([]) + expect(plan.blockedPanes).toEqual([expect.objectContaining({ settled: false })]) + }) + it('fails closed when two Dispatches claim one terminal identity', () => { const plan = planLegacyWorkerTerminalRecovery([ recoveryRow(), diff --git a/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.ts b/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.ts index c994101aeda..eb4ce1306e0 100644 --- a/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.ts +++ b/src/main/runtime/orchestration/orchestration-legacy-worker-terminal-recovery.ts @@ -1,6 +1,7 @@ import { isPtyIncarnationId, type PtyIncarnationId } from '../../../shared/pty-incarnation' import { parsePaneKey } from '../../../shared/stable-pane-id' import type { LegacyWorkerTerminalRecoveryRow } from './types' +import { WORKER_SETTLED_STATES } from './worker-terminal-ownership' export type LegacyWorkerTerminalRecoveryCandidate = { dispatchId: string @@ -17,8 +18,15 @@ export type LegacyWorkerTerminalRecoveryCandidate = { incarnationId: PtyIncarnationId } +export type LegacyWorkerTerminalRecoveryBlockedPane = { + worktreeId: string + paneKey: string + contractVersion: number + settled: boolean +} + export type LegacyWorkerTerminalRecoveryPlan = { - blockedPanes: { worktreeId: string; paneKey: string; contractVersion: number }[] + blockedPanes: LegacyWorkerTerminalRecoveryBlockedPane[] candidates: LegacyWorkerTerminalRecoveryCandidate[] ambiguousDispatchIds: string[] } @@ -50,22 +58,26 @@ function countCandidateKeys( export function planLegacyWorkerTerminalRecovery( rows: readonly LegacyWorkerTerminalRecoveryRow[] ): LegacyWorkerTerminalRecoveryPlan { - const blockedPanes = new Map< - string, - { worktreeId: string; paneKey: string; contractVersion: number } - >() + const blockedPanes = new Map() const parsedCandidates: LegacyWorkerTerminalRecoveryCandidate[] = [] for (const row of rows) { const worktreeId = row.worktree_id?.trim() const paneKey = row.assignee_pane_key?.trim() const pane = paneKey ? parsePaneKey(paneKey) : null + const settled = WORKER_SETTLED_STATES.includes(row.worker_state) if (worktreeId && paneKey && pane) { - blockedPanes.set(`${worktreeId}\0${paneKey}`, { + const blockedKey = `${worktreeId}\0${paneKey}` + const alreadySettled = blockedPanes.get(blockedKey)?.settled + blockedPanes.set(blockedKey, { worktreeId, paneKey, - contractVersion: row.contract_version + contractVersion: row.contract_version, + settled: (alreadySettled ?? true) && settled }) } + if (settled) { + continue + } const terminalHandle = row.assignee_handle?.trim() const workerHandle = row.agent_terminal_handle?.trim() const processIncarnation = row.process_incarnation?.trim() diff --git a/src/main/runtime/orchestration/orchestration-settled-worker-resume-fence-db.test.ts b/src/main/runtime/orchestration/orchestration-settled-worker-resume-fence-db.test.ts new file mode 100644 index 00000000000..98baf9e5771 --- /dev/null +++ b/src/main/runtime/orchestration/orchestration-settled-worker-resume-fence-db.test.ts @@ -0,0 +1,91 @@ +import { afterEach, describe, expect, it } from 'vitest' +import { OrchestrationDb } from './db' +import type { WorkerTerminalResourceRow } from './worker-terminal-ownership' + +const PANE_KEY = 'tab_worker:33333333-3333-4333-8333-333333333333' + +describe('settled worker terminal resume fence rows', () => { + let db: OrchestrationDb | undefined + + afterEach(() => db?.close()) + + function createReadyWorker(): { db: OrchestrationDb; taskId: string; dispatchId: string } { + const d = new OrchestrationDb(':memory:') + db = d + const task = d.createTask({ spec: 'settled worker' }) + const started = d.createStartingWorkerDispatch({ + creator: { kind: 'system' }, + maxDepth: Number.MAX_SAFE_INTEGER, + taskId: task.id, + startOptions: {} + }) + d.prepareStartingWorkerAuthority({ + dispatchId: started.dispatch.id, + handle: 'term_worker', + paneKey: PANE_KEY, + processIncarnation: 'runtime:pty:1', + worktreeId: 'repo::worktree', + setupState: 'not_applicable', + effects: [], + terminalOwnership: 'created' + }) + d.markWorkerDispatchReady(started.dispatch.id) + return { db: d, taskId: task.id, dispatchId: started.dispatch.id } + } + + function requestRelease(d: OrchestrationDb, dispatchId: string): WorkerTerminalResourceRow { + const requested = d.requestWorkerTerminalRelease(dispatchId) + if (requested.disposition !== 'requested') { + throw new Error(`expected a release request, got ${requested.disposition}`) + } + return requested.resource + } + + function settle(d: OrchestrationDb, taskId: string, dispatchId: string): void { + expect( + d.settleWorkerReport({ + taskId, + dispatchId, + outcome: 'succeeded', + result: JSON.stringify({ provenance: 'worker_report', outcome: 'succeeded' }) + }).action + ).not.toBe('rejected') + } + + it('keeps a settled-but-unreleased worker terminal in recovery rows', () => { + const { db: d, taskId, dispatchId } = createReadyWorker() + settle(d, taskId, dispatchId) + expect(d.listLegacyWorkerTerminalRecoveryRows()).toEqual([ + expect.objectContaining({ + dispatch_id: dispatchId, + worker_state: 'succeeded', + assignee_pane_key: PANE_KEY + }) + ]) + }) + + it('keeps a settled worker whose release is unknown', () => { + const { db: d, taskId, dispatchId } = createReadyWorker() + settle(d, taskId, dispatchId) + const resource = requestRelease(d, dispatchId) + d.markWorkerTerminalReleaseUnknown(resource.id, 'terminal no longer resolves') + expect(d.listLegacyWorkerTerminalRecoveryRows()).toEqual([ + expect.objectContaining({ dispatch_id: dispatchId }) + ]) + }) + + it('drops a settled worker once its resource is released', () => { + const { db: d, taskId, dispatchId } = createReadyWorker() + settle(d, taskId, dispatchId) + const resource = requestRelease(d, dispatchId) + d.settleWorkerTerminalRelease(resource.id) + expect(d.listLegacyWorkerTerminalRecoveryRows()).toEqual([]) + }) + + it('drops a settled worker the user chose to retain', () => { + const { db: d, taskId, dispatchId } = createReadyWorker() + d.retainWorkerTerminalResource(dispatchId) + settle(d, taskId, dispatchId) + expect(d.listLegacyWorkerTerminalRecoveryRows()).toEqual([]) + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration-settlement-resume-fence.test.ts b/src/main/runtime/rpc/methods/orchestration-settlement-resume-fence.test.ts new file mode 100644 index 00000000000..a3a1ca28eb2 --- /dev/null +++ b/src/main/runtime/rpc/methods/orchestration-settlement-resume-fence.test.ts @@ -0,0 +1,144 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { ORCHESTRATION_METHODS } from './orchestration' +import type { RpcContext } from '../core' +import { OrchestrationDb } from '../../orchestration/db' +import { OrcaRuntimeService } from '../../orca-runtime' + +const COORDINATOR_PANE_KEY = 'tab_coord:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa' +const WORKER_PANE_KEY = 'tab_worker:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb' +const WORKTREE_ID = 'repo::worktree' + +describe('settled worker automatic-resume fence trigger', () => { + let db: OrchestrationDb | undefined + + afterEach(() => { + db?.close() + db = undefined + vi.restoreAllMocks() + }) + + function setup(): { + ctx: RpcContext + runId: string + taskId: string + dispatchId: string + fence: ReturnType + } { + const orchestrationDb = new OrchestrationDb(':memory:') + db = orchestrationDb + const runtime = new OrcaRuntimeService() + runtime.setOrchestrationDb(orchestrationDb) + vi.spyOn(runtime, 'getTerminalPaneKey').mockImplementation((handle) => + handle === 'term_coord' + ? COORDINATOR_PANE_KEY + : handle === 'term_worker' + ? WORKER_PANE_KEY + : null + ) + vi.spyOn(runtime, 'getTerminalProcessIncarnation').mockImplementation((handle) => + handle === 'term_worker' ? 'runtime_test:term_worker:1' : null + ) + vi.spyOn(runtime, 'notifyMessageArrived').mockImplementation(() => {}) + const fence = vi.fn() + runtime.setNotifier({ resolveLegacyWorkerTerminalRecovery: fence } as never) + const run = orchestrationDb.createRun({ + objective: 'Settlement fence run', + coordinatorHandle: 'term_coord', + coordinatorPaneKey: COORDINATOR_PANE_KEY + }) + const task = orchestrationDb.createTask({ spec: 'settle me', runId: run.id }) + const started = orchestrationDb.createStartingWorkerDispatch({ + creator: { kind: 'system' }, + maxDepth: Number.MAX_SAFE_INTEGER, + taskId: task.id, + startOptions: {} + }) + const orchestrationCapability = orchestrationDb.prepareStartingWorkerAuthority({ + dispatchId: started.dispatch.id, + handle: 'term_worker', + paneKey: WORKER_PANE_KEY, + processIncarnation: 'runtime_test:term_worker:1', + worktreeId: WORKTREE_ID, + setupState: 'not_applicable', + effects: [], + terminalOwnership: 'created' + }) + orchestrationDb.markWorkerDispatchReady(started.dispatch.id) + return { + ctx: { runtime, orchestrationCapability }, + runId: run.id, + taskId: task.id, + dispatchId: started.dispatch.id, + fence + } + } + + async function call( + name: string, + params: Record, + ctx: RpcContext + ): Promise { + const method = ORCHESTRATION_METHODS.find((entry) => entry.name === name) + if (!method) { + throw new Error(`Method not found: ${name}`) + } + return method.handler(method.params ? method.params.parse(params) : undefined, ctx) + } + + const workerDonePayload = (taskId: string, dispatchId: string) => + JSON.stringify({ taskId, dispatchId, outcome: 'succeeded' }) + + it('fences the pane when orchestration.send settles worker_done', async () => { + const fixture = setup() + const sent = (await call( + 'orchestration.send', + { + from: 'term_worker', + to: `run:${fixture.runId}`, + subject: 'Done', + type: 'worker_done', + payload: workerDonePayload(fixture.taskId, fixture.dispatchId), + run: fixture.runId + }, + fixture.ctx + )) as { lifecycle?: { action: string } } + expect(sent.lifecycle?.action).toBe('completed') + expect(fixture.fence).toHaveBeenCalledWith(WORKER_PANE_KEY, 'fenced', { + worktreeId: WORKTREE_ID + }) + }) + + it('fences when an unread check is first to settle worker_done', async () => { + const fixture = setup() + db?.insertMessage({ + from: 'term_worker', + to: 'term_watcher', + subject: 'Done', + type: 'worker_done', + payload: workerDonePayload(fixture.taskId, fixture.dispatchId), + senderPaneKey: WORKER_PANE_KEY, + runId: fixture.runId + }) + await call('orchestration.check', { terminal: 'term_watcher' }, fixture.ctx) + expect(fixture.fence).toHaveBeenCalledWith(WORKER_PANE_KEY, 'fenced', { + worktreeId: WORKTREE_ID + }) + }) + + it('does not fence a heartbeat', async () => { + const fixture = setup() + await call( + 'orchestration.send', + { + from: 'term_worker', + to: `run:${fixture.runId}`, + subject: 'Still working', + type: 'heartbeat', + payload: JSON.stringify({ dispatchId: fixture.dispatchId }), + run: fixture.runId + }, + fixture.ctx + ) + expect(fixture.fence).not.toHaveBeenCalled() + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration.ts b/src/main/runtime/rpc/methods/orchestration.ts index 876021e79e1..a8ec57bb5f3 100644 --- a/src/main/runtime/rpc/methods/orchestration.ts +++ b/src/main/runtime/rpc/methods/orchestration.ts @@ -15,7 +15,10 @@ import { MESSAGE_TYPES } from '../../orchestration/types' import { buildDispatchPreamble } from '../../orchestration/preamble' import { formatMessageBanner } from '../../orchestration/formatter' import { isGroupAddress, resolveGroupAddress } from '../../orchestration/groups' -import { reconcileLifecycleMessage } from '../../orchestration/lifecycle-reconciliation' +import { + reconcileLifecycleMessage, + type LifecycleReconciliationResult +} from '../../orchestration/lifecycle-reconciliation' import { waitForFederatedLifecycleSettlement } from '../../orchestration/federation-lifecycle-settlement' import { abbreviateOrchestrationTasks } from '../../../../shared/orchestration-task-summary' import { @@ -125,6 +128,15 @@ function isWorkerReportOutcome(value: unknown): value is 'succeeded' | 'failed' return value === 'succeeded' || value === 'failed' } +function fenceSettledWorkerPanes( + runtime: OrcaRuntimeService, + reconciled: readonly LifecycleReconciliationResult[] +): void { + if (reconciled.some((result) => result.action === 'completed' || result.action === 'failed')) { + runtime.prepareLegacyWorkerTerminalRecovery() + } +} + const SendParams = z .object({ to: OptionalString, @@ -801,6 +813,7 @@ export const ORCHESTRATION_METHODS: RpcMethod[] = [ runtime.notifyMessageArrived(rejection.to_handle, rejection.type) return withSendWarnings({ message: rejection, lifecycle: reconciled }) } + fenceSettledWorkerPanes(runtime, [reconciled]) runtime.notifyMessageArrived(msg.to_handle, msg.type) return withSendWarnings( msg.type === 'worker_done' ? { message: msg, lifecycle: reconciled } : { message: msg } @@ -1334,12 +1347,15 @@ export const ORCHESTRATION_METHODS: RpcMethod[] = [ let visibleMessages = messages if (consumeUnread && messages.length > 0) { // Why: unread check is an authoritative read path for worker_done/heartbeat, so reconcile lifecycle messages here too. + const reconciledMessages: LifecycleReconciliationResult[] = [] visibleMessages = messages.map((message) => { const reconciled = reconcileLifecycleMessage(db, message) + reconciledMessages.push(reconciled) return reconciled.action === 'rejected' ? (db.getMessageById(message.id) ?? message) : message }) + fenceSettledWorkerPanes(runtime, reconciledMessages) db.markAsRead(messages.map((m) => m.id)) } diff --git a/src/main/runtime/rpc/orchestration-legacy-lifecycle.ts b/src/main/runtime/rpc/orchestration-legacy-lifecycle.ts index ebf6cfc2869..ea000eafc67 100644 --- a/src/main/runtime/rpc/orchestration-legacy-lifecycle.ts +++ b/src/main/runtime/rpc/orchestration-legacy-lifecycle.ts @@ -108,6 +108,9 @@ export async function handleLegacyLifecycleSend(args: { } : { kind: 'message_only' } }) + if (committed.settlement?.action === 'settled') { + runtime.prepareLegacyWorkerTerminalRecovery() + } if (!committed.duplicate) { runtime.notifyMessageArrived(committed.message.to_handle, committed.message.type) } diff --git a/src/main/window/runtime-window-lifecycle.ts b/src/main/window/runtime-window-lifecycle.ts index 7d86c0bb9b4..655271f1077 100644 --- a/src/main/window/runtime-window-lifecycle.ts +++ b/src/main/window/runtime-window-lifecycle.ts @@ -143,11 +143,12 @@ export function registerRuntimeWindowLifecycle( reject(new Error('runtime_unavailable')) } }), - resolveLegacyWorkerTerminalRecovery: (paneKey, resolution, ptyId) => + resolveLegacyWorkerTerminalRecovery: (paneKey, resolution, identity) => send('agentStatus:legacyWorkerTerminalRecovery', { paneKey, resolution, - ...(ptyId ? { ptyId } : {}) + ...(identity?.ptyId ? { ptyId: identity.ptyId } : {}), + ...(identity?.worktreeId ? { worktreeId: identity.worktreeId } : {}) }), splitTerminal: (tabId, paneRuntimeId, opts) => { send('ui:splitTerminal', { diff --git a/src/preload/api/agent-status-api.ts b/src/preload/api/agent-status-api.ts index fd6da88e351..f6759099f1e 100644 --- a/src/preload/api/agent-status-api.ts +++ b/src/preload/api/agent-status-api.ts @@ -1,6 +1,7 @@ import type { AgentStatusClearIpcPayload, AgentStatusIpcPayload, + LegacyWorkerTerminalRecoveryEvent, MigrationUnsupportedPtyEntry } from '../../shared/agent-status-types' import type { AgentInterruptInferenceRequest } from '../../shared/agent-interrupt-intent' @@ -21,11 +22,7 @@ export type AgentStatusApi = { onMigrationUnsupported: (callback: (entry: MigrationUnsupportedPtyEntry) => void) => () => void onMigrationUnsupportedClear: (callback: (data: { ptyId: string }) => void) => () => void onLegacyWorkerTerminalRecovery: ( - callback: (data: { - paneKey: string - resolution: 'adopted' | 'exited' | 'rolled_back' - ptyId?: string - }) => void + callback: (data: LegacyWorkerTerminalRecoveryEvent) => void ) => () => void getMigrationUnsupportedSnapshot: () => Promise /** Drop a paneKey from the main-process hook cache and on-disk last-status file. Fire-and-forget. */ diff --git a/src/preload/index.ts b/src/preload/index.ts index 9799d4fe850..b88eb465a4e 100644 --- a/src/preload/index.ts +++ b/src/preload/index.ts @@ -253,6 +253,7 @@ import { import type { AgentStatusClearIpcPayload, AgentStatusIpcPayload, + LegacyWorkerTerminalRecoveryEvent, MigrationUnsupportedPtyEntry } from '../shared/agent-status-types' import type { AgentInterruptInferenceRequest } from '../shared/agent-interrupt-intent' @@ -5182,19 +5183,11 @@ const api = { return () => ipcRenderer.removeListener('agentStatus:migrationUnsupportedClear', listener) }, onLegacyWorkerTerminalRecovery: ( - callback: (data: { - paneKey: string - resolution: 'adopted' | 'exited' | 'rolled_back' - ptyId?: string - }) => void + callback: (data: LegacyWorkerTerminalRecoveryEvent) => void ): (() => void) => { const listener = ( _event: Electron.IpcRendererEvent, - data: { - paneKey: string - resolution: 'adopted' | 'exited' | 'rolled_back' - ptyId?: string - } + data: LegacyWorkerTerminalRecoveryEvent ) => callback(data) ipcRenderer.on('agentStatus:legacyWorkerTerminalRecovery', listener) return () => ipcRenderer.removeListener('agentStatus:legacyWorkerTerminalRecovery', listener) diff --git a/src/renderer/src/components/terminal-pane/pty-connection-main-side-effect-authority.test.ts b/src/renderer/src/components/terminal-pane/pty-connection-main-side-effect-authority.test.ts index 25b4027c545..147c79fa9b2 100644 --- a/src/renderer/src/components/terminal-pane/pty-connection-main-side-effect-authority.test.ts +++ b/src/renderer/src/components/terminal-pane/pty-connection-main-side-effect-authority.test.ts @@ -5,6 +5,7 @@ import { flushAsyncTicks } from './pty-connection-test-async' import { AGENT_TASK_COMPLETE_NOTIFICATION_MAX_WAIT_MS } from './pty-connection-test-constants' import { LEAF_1, + LEAF_2, createMockTransport, createPane, createManager, @@ -642,6 +643,67 @@ describe('connectPanePty', () => { expect(mockStoreState.setAgentStatus).not.toHaveBeenCalled() }) + it.each(['claude', 'codex'] as const)( + 'retires only the confirmed %s agent exit recovery record before PTY teardown', + async (agent) => { + enableMainAuthority() + const { connectPanePty } = await import('./pty-connection') + const handler = await import('./terminal-side-effect-facts-handler') + const transport = createMockTransport('pty-agent-exit') + transportFactoryQueue.push(transport) + const paneKey = 'tab-1:1' + const duplicatePaneKey = 'tab-1:2' + const siblingPaneKey = makePaneKey('tab-1', LEAF_2) + const record = { + paneKey, + tabId: 'tab-1', + worktreeId: 'wt-1', + agent, + providerSession: { key: 'session_id', id: 'session-1' }, + state: 'working' as const, + capturedAt: 1, + updatedAt: 1 + } + const duplicateRecord = { + ...record, + paneKey: duplicatePaneKey, + capturedAt: 2, + updatedAt: 2 + } + const siblingRecord = { + ...record, + paneKey: siblingPaneKey, + providerSession: { key: 'session_id', id: 'session-2' } + } + mockStoreState.sleepingAgentSessionsByPaneKey = { + [paneKey]: record, + [duplicatePaneKey]: duplicateRecord, + [siblingPaneKey]: siblingRecord + } + const deps = createDeps() + + connectPanePty(createPane(1) as never, createManager(1) as never, deps as never) + const onPtySpawn = createdTransportOptions[0]?.onPtySpawn as (ptyId: string) => void + onPtySpawn('pty-agent-exit') + handler._dispatchTerminalSideEffectBatchForTest({ + ptyId: 'pty-agent-exit', + seq: 1, + facts: [{ kind: 'agent-exited' }] + }) + + expect(mockStoreState.clearSleepingAgentSession).toHaveBeenCalledWith(paneKey) + expect(mockStoreState.clearSleepingAgentSession).toHaveBeenCalledWith(duplicatePaneKey) + expect(mockStoreState.sleepingAgentSessionsByPaneKey[paneKey]).toBeUndefined() + expect(mockStoreState.sleepingAgentSessionsByPaneKey[duplicatePaneKey]).toBeUndefined() + expect(mockStoreState.sleepingAgentSessionsByPaneKey[siblingPaneKey]).toBe(siblingRecord) + expect(deps.onAgentExitedRef.current).toHaveBeenCalledWith(LEAF_1) + const clearCallOrders = mockStoreState.clearSleepingAgentSession.mock.invocationCallOrder + expect(clearCallOrders.at(-1)).toBeLessThan( + deps.onAgentExitedRef.current.mock.invocationCallOrder[0] + ) + } + ) + it('honors the persisted kill switch for panes bound before settings hydrate', async () => { // Pre-hydration: settings not loaded but kill switch persisted off — pane registers byte parsers, not a fact consumer. mockStoreState.settings = null diff --git a/src/renderer/src/components/terminal-pane/pty-connection/agent-idle-working-handlers.ts b/src/renderer/src/components/terminal-pane/pty-connection/agent-idle-working-handlers.ts index 8516eab7f19..cde94f238e9 100644 --- a/src/renderer/src/components/terminal-pane/pty-connection/agent-idle-working-handlers.ts +++ b/src/renderer/src/components/terminal-pane/pty-connection/agent-idle-working-handlers.ts @@ -31,8 +31,13 @@ export function installAgentIdleWorkingHandlers(session: ConnectPanePtySession): } } session.onAgentExited = (): void => { - // Why: eligibility can disappear transiently during reconnect, but a - // confirmed shell-title transition is authoritative for native-chat exit. + // Why: a confirmed shell transition means the agent ended before its terminal. + // Retire exact resume authority so a later workspace open cannot resurrect it. + const state = useAppStore.getState() + const sleepingRecordEntry = session.getSleepingRecordForPane(state) + if (sleepingRecordEntry) { + session.clearSleepingRecordProviderDuplicates(state, sleepingRecordEntry) + } session.deps.onAgentExitedRef.current(session.pane.leafId) session.clearSuppressedTitleSideEffects() session.clearCommandInferredPaneAgent() 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..ea0c30a10e0 100644 --- a/src/renderer/src/hooks/ipc-events/agent-status-listeners.ts +++ b/src/renderer/src/hooks/ipc-events/agent-status-listeners.ts @@ -8,6 +8,7 @@ import { rollbackLegacyWorkerTerminalSurfaceInStore } from '../legacy-worker-terminal-recovery-event' import { useAppStore } from '../../store' +import { worktreeIdsEqual } from '../../../../shared/worktree/id' import { resolvePaneKey } from './agent-status-routing' import type { PendingAgentStatusEvent } from './agent-status-bridge-types' @@ -121,6 +122,12 @@ export function registerAgentStatusListeners(args: { rollbackLegacyWorkerTerminalSurfaceInStore(useAppStore.getState(), action.detail) } else if (action.kind === 'clear-sleeping') { useAppStore.getState().clearSleepingAgentSession(action.paneKey) + } else if (action.kind === 'set-automatic-resume-block') { + const store = useAppStore.getState() + const record = store.sleepingAgentSessionsByPaneKey[action.paneKey] + if (record && worktreeIdsEqual(record.worktreeId, action.worktreeId)) { + store.setSleepingAgentAutomaticResumeBlocked(action.paneKey, action.blocked) + } } }) if (unsubscribeLegacyWorkerTerminalRecovery) { diff --git a/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.test.ts b/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.test.ts index 6b6b3aa8862..00c6a3cf522 100644 --- a/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.test.ts +++ b/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.test.ts @@ -38,6 +38,39 @@ describe('legacy worker terminal recovery events', () => { }) }) + it('fences and unfences only the named workspace recovery record', () => { + expect( + resolveLegacyWorkerTerminalRecoveryAction({ + paneKey: `legacy-worker:${LEAF_ID}`, + resolution: 'fenced', + worktreeId: 'repo::/workspace' + }) + ).toEqual({ + kind: 'set-automatic-resume-block', + paneKey: `legacy-worker:${LEAF_ID}`, + worktreeId: 'repo::/workspace', + blocked: true + }) + expect( + resolveLegacyWorkerTerminalRecoveryAction({ + paneKey: `legacy-worker:${LEAF_ID}`, + resolution: 'unfenced', + worktreeId: 'repo::/workspace' + }) + ).toEqual({ + kind: 'set-automatic-resume-block', + paneKey: `legacy-worker:${LEAF_ID}`, + worktreeId: 'repo::/workspace', + blocked: false + }) + expect( + resolveLegacyWorkerTerminalRecoveryAction({ + paneKey: `legacy-worker:${LEAF_ID}`, + resolution: 'fenced' + }) + ).toEqual({ kind: 'ignore' }) + }) + it('removes an unmounted split surface only when its PTY identity still matches', () => { const setTabLayout = vi.fn() const clearTabPtyId = vi.fn() diff --git a/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.ts b/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.ts index 61940acedd7..d6597e70939 100644 --- a/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.ts +++ b/src/renderer/src/hooks/legacy-worker-terminal-recovery-event.ts @@ -2,24 +2,33 @@ import type { CloseTerminalPaneDetail } from '@/constants/terminal' import { detachTerminalLayoutLeaf } from '@/components/terminal-pane/terminal-layout-leaf-detach' import type { AppState } from '@/store' import { makePaneKey, parsePaneKey } from '../../../shared/stable-pane-id' - -type LegacyWorkerTerminalRecoveryEvent = { - paneKey: string - resolution: 'adopted' | 'exited' | 'rolled_back' - ptyId?: string -} +import type { LegacyWorkerTerminalRecoveryEvent } from '../../../shared/agent-status-types' export type LegacyWorkerTerminalRecoveryAction = | { kind: 'clear-sleeping'; paneKey: string } | { kind: 'rollback-surface'; detail: CloseTerminalPaneDetail } + | { kind: 'set-automatic-resume-block'; paneKey: string; worktreeId: string; blocked: boolean } | { kind: 'ignore' } export function resolveLegacyWorkerTerminalRecoveryAction( event: LegacyWorkerTerminalRecoveryEvent ): LegacyWorkerTerminalRecoveryAction { - if (event.resolution !== 'rolled_back') { + if (event.resolution === 'adopted' || event.resolution === 'exited') { return { kind: 'clear-sleeping', paneKey: event.paneKey } } + if (event.resolution === 'fenced' || event.resolution === 'unfenced') { + return event.worktreeId + ? { + kind: 'set-automatic-resume-block', + paneKey: event.paneKey, + worktreeId: event.worktreeId, + blocked: event.resolution === 'fenced' + } + : { kind: 'ignore' } + } + if (event.resolution !== 'rolled_back') { + return { kind: 'ignore' } + } const pane = parsePaneKey(event.paneKey) return pane && event.ptyId ? { diff --git a/src/renderer/src/store/slices/agent-status.ts b/src/renderer/src/store/slices/agent-status.ts index 5a6293ae5a4..049febf5b29 100644 --- a/src/renderer/src/store/slices/agent-status.ts +++ b/src/renderer/src/store/slices/agent-status.ts @@ -611,7 +611,7 @@ function sleepingRecordFromEntry(args: { return null } const tab = args.tab ?? findTabForAgentEntry(args.state, args.worktreeId, args.entry) - return { + const record: SleepingAgentSessionRecord = { paneKey: args.entry.paneKey, ...(tab ? { tabId: tab.id } : {}), worktreeId: args.worktreeId, @@ -632,6 +632,11 @@ function sleepingRecordFromEntry(args: { ...(args.entry.interrupted ? { interrupted: true } : {}), ...(args.origin ? { origin: args.origin } : {}) } + carryOverAutomaticResumeBlock( + record, + args.state.sleepingAgentSessionsByPaneKey[args.entry.paneKey] + ) + return record } type CollectSleepingAgentSessionRecordsOptions = { @@ -677,8 +682,7 @@ function manualSleepCaptureEntry(entry: AgentStatusEntry, capturedAt: number): A return { ...entry, updatedAt: capturedAt, interrupted: false } } -// Why: capture recreates a record the manual-sleep wipe would otherwise remove, so a deliberately -// blocked worker must not become auto-resumable at wake. +// Why: every re-derivation from a live entry must preserve a deliberate worker resume fence. function carryOverAutomaticResumeBlock( record: SleepingAgentSessionRecord, previous: SleepingAgentSessionRecord | undefined @@ -794,10 +798,6 @@ export function collectSleepingAgentSessionRecordsForWorktree( if (record) { if (isManualWorktreeSleep) { markManualSleepLazyRestore(record) - carryOverAutomaticResumeBlock( - record, - state.sleepingAgentSessionsByPaneKey[retained.entry.paneKey] - ) } records[record.paneKey] = record } @@ -831,7 +831,6 @@ export function collectSleepingAgentSessionRecordsForWorktree( if (record) { if (isManualWorktreeSleep) { markManualSleepLazyRestore(record) - carryOverAutomaticResumeBlock(record, state.sleepingAgentSessionsByPaneKey[paneKey]) } records[record.paneKey] = record } diff --git a/src/shared/agent-status-types.ts b/src/shared/agent-status-types.ts index 70509ea495a..f352683772f 100644 --- a/src/shared/agent-status-types.ts +++ b/src/shared/agent-status-types.ts @@ -421,3 +421,17 @@ export function parseAgentStatusPayload(json: string): ParsedAgentStatusPayload return null } } + +export type LegacyWorkerTerminalRecoveryResolutionKind = + | 'adopted' + | 'exited' + | 'rolled_back' + | 'fenced' + | 'unfenced' + +export type LegacyWorkerTerminalRecoveryEvent = { + paneKey: string + resolution: LegacyWorkerTerminalRecoveryResolutionKind + ptyId?: string + worktreeId?: string +} diff --git a/src/shared/worktree/id.ts b/src/shared/worktree/id.ts index 82220328eb3..f2bfe771998 100644 --- a/src/shared/worktree/id.ts +++ b/src/shared/worktree/id.ts @@ -43,6 +43,12 @@ export function worktreeIdComparisonKey(worktreeId: string): string | null { )}` } +/** Compare workspace identity while normalizing only runtime path spelling. */ +export function worktreeIdsEqual(left: string, right: string): boolean { + const leftKey = worktreeIdComparisonKey(left) + return leftKey === null ? left === right : leftKey === worktreeIdComparisonKey(right) +} + export function splitWorktreeId(worktreeId: string): ParsedWorktreeId | null { const separatorIdx = worktreeId.indexOf(WORKTREE_ID_SEPARATOR) if (separatorIdx === -1) { diff --git a/tests/e2e/completed-worker-retirement-resume.spec.ts b/tests/e2e/completed-worker-retirement-resume.spec.ts index 69f6e0af375..6935cbe1895 100644 --- a/tests/e2e/completed-worker-retirement-resume.spec.ts +++ b/tests/e2e/completed-worker-retirement-resume.spec.ts @@ -269,7 +269,7 @@ for (const closeMode of ['terminal-close-cli', 'worker-release'] as const) { const expectedRecovery = { origin: 'live', - state: 'working', + state: 'done', providerSessionId: PROVIDER_SESSION_ID } await expect