fix(orchestration): fence settled worker auto-resume

This commit is contained in:
Jinwoo-H
2026-08-30 23:29:42 -07:00
committed by Neil
parent c55121231a
commit d818fa0a0e
21 changed files with 546 additions and 62 deletions
+4 -4
View File
@@ -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()
+66 -12
View File
@@ -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<ExecutionHostId, { current: WorkspaceSessionState; next: WorkspaceSessionState }>,
changedHostIds: Set<ExecutionHostId>
): 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<void> {
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(
@@ -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(
@@ -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(),
@@ -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<string, LegacyWorkerTerminalRecoveryBlockedPane>()
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()
@@ -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([])
})
})
@@ -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<typeof vi.fn>
} {
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<string, unknown>,
ctx: RpcContext
): Promise<unknown> {
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()
})
})
+17 -1
View File
@@ -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))
}
@@ -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)
}
+3 -2
View File
@@ -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', {
+2 -5
View File
@@ -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<MigrationUnsupportedPtyEntry[]>
/** Drop a paneKey from the main-process hook cache and on-disk last-status file. Fire-and-forget. */
+3 -10
View File
@@ -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)
@@ -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
@@ -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()
@@ -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) {
@@ -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()
@@ -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
? {
@@ -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
}
+14
View File
@@ -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
}
+6
View File
@@ -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) {
@@ -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