fix(orchestration): reland process-incarnation reap for stale worker terminal handles (#23583)

* fix(orchestration): reland process-incarnation reap for stale worker terminal handles

Relands the orchestration part of #18790, which #22601 reverted in full
because that squash commit also bundled an unannounced agent. Only the
worker terminal-reap fix returns; no agent catalog changes.

When a worker's saved terminal handle stops resolving while its process is
still running, worker-show/read, worker-stop and worker-release now re-mint
a live handle from the recorded process incarnation (exact pty id +
incarnation id, same host scope) and act on it, instead of reporting the
terminal missing and leaving the process running on the host.

Fixes STA-8493

* test(orchestration): spy on the public incarnation methods instead of casting the runtime
This commit is contained in:
Brennan Benson
2026-09-28 02:24:04 -07:00
committed by GitHub
parent 2f8f4f576d
commit 5813cb640e
13 changed files with 1344 additions and 60 deletions
@@ -22,6 +22,7 @@ import { getAppEnvironment } from '../../shared/app-environment'
import type { FleetAgentStatusEvidence } from '../../shared/orchestration-fleet-agent-status-evidence'
import { readOrchestrationFleetAgentStatusSnapshot } from './orchestration-fleet-agent-status-snapshot'
import { resolveStructuredWorkerAuthority } from './structured-worker-authority'
import { matchesProcessIncarnation } from './orchestration/worker-terminal-process-liveness'
export class OrcaRuntimeWithGetOrchestrationDispatchAuthority extends OrcaRuntimeWithVerifyOrchestrationCompatibilityCaller {
/** Every pane key this PTY could be addressed by, including restored receipts. */
@@ -270,4 +271,38 @@ export class OrcaRuntimeWithGetOrchestrationDispatchAuthority extends OrcaRuntim
: undefined
})
}
/**
* Recover a live terminal handle for a worker whose durable handle stopped resolving
* (renderer graph epoch bump / handle invalidation) while its PTY is still tracked. Fences on
* the recorded process incarnation EXACTLY — never a bare ptyId, worktree, or pane — so a
* reused ptyId belonging to a different process can never be closed, and fails closed on an
* unknown host scope (this is also consumed by workerShow, which does no lease re-check).
* Returns a freshly minted live handle, or null when no live PTY carries that exact incarnation.
*/
resolveTerminalHandleByProcessIncarnation(
processIncarnation: string,
serializedHostScope: string | null
): string | null {
if (!processIncarnation || !serializedHostScope) {
return null
}
// Scan by the incarnation itself (startsWith + exact equality, mirroring
// classifyWorkerTerminalProcessIncarnation) rather than splitting on a colon, so relay/SSH
// ptyIds and colon-bearing incarnationIds still match. A pty with no incarnationId can never
// match, so the legacy `${runtimeId}:${ptyId}:${ptyGeneration}` fence stays fail-closed.
for (const [ptyId, pty] of this.ptysById) {
if (!matchesProcessIncarnation(ptyId, pty.incarnationId, processIncarnation)) {
continue
}
const hostScope = this.getOrchestrationCompatibilityHostScope(pty)
if (!hostScope || JSON.stringify(hostScope) !== serializedHostScope) {
// Keep scanning: a colon-ambiguous decoy pty in a different host scope that this
// incarnation string happens to prefix-match must not suppress the genuine same-scope
// pty later in ptysById. The scope check still fences the real match below.
continue
}
return this.issuePtyHandle(pty)
}
return null
}
}
@@ -1,11 +1,13 @@
import { describe, expect, it, vi } from 'vitest'
import { OrcaRuntimeService } from './orca-runtime'
import { makePaneKey } from '../../shared/stable-pane-id'
const PTY_ID = 'ssh:target@@relay-pty'
const WORKTREE_ID = 'repo::/worktree'
const TAB_ID = 'tab-terminal'
const LEAF_ID = '11111111-1111-4111-8111-111111111111'
/** A runtime whose pty controller captures every write for assertion. */
function makeRuntime(): { runtime: OrcaRuntimeService; writes: string[] } {
const writes: string[] = []
const runtime = new OrcaRuntimeService(null)
@@ -20,6 +22,7 @@ function makeRuntime(): { runtime: OrcaRuntimeService; writes: string[] } {
return { runtime, writes }
}
/** Attach a window and publish a one-tab, one-leaf terminal graph. */
function syncGraph(runtime: OrcaRuntimeService): void {
runtime.attachWindow(1)
runtime.syncWindowGraph(1, {
@@ -44,6 +47,7 @@ function syncGraph(runtime: OrcaRuntimeService): void {
})
}
/** (Re)register the fixture pty leaf under the given incarnation id. */
function register(runtime: OrcaRuntimeService, incarnationId: string): void {
runtime.registerPty(PTY_ID, WORKTREE_ID, 'target', {
tabId: TAB_ID,
@@ -66,6 +70,8 @@ describe('runtime terminal handle incarnation fencing', () => {
inspectProcess
})
expect(runtime.markRendererReloading(1)).not.toBeNull()
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
expect((runtime as unknown as { handles: Map<string, unknown> }).handles.has(handle)).toBe(
false
)
@@ -89,6 +95,8 @@ describe('runtime terminal handle incarnation fencing', () => {
})
})
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
it('treats a null-to-known incarnation as the same un-fenced PTY', async () => {
const { runtime } = makeRuntime()
const handle = runtime.preAllocateHandleForPty(PTY_ID)
@@ -132,6 +140,8 @@ describe('runtime terminal handle incarnation fencing', () => {
// Rotate the record directly so reconcile is the only fence exercised.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: test reaches the runtime's protected pty record map to bypass the registerPty fence.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const internals = runtime as unknown as {
ptysById: Map<string, { incarnationId: string | null }>
}
@@ -219,6 +229,8 @@ describe('runtime terminal handle incarnation fencing', () => {
runtime.registerPreAllocatedHandleForPty(PTY_ID, replacementHandle)
syncGraph(runtime)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const internals = runtime as unknown as {
handles: Map<string, unknown>
handleByLeafKey: Map<string, string>
@@ -241,3 +253,130 @@ describe('runtime terminal handle incarnation fencing', () => {
})
})
})
describe('resolveTerminalHandleByProcessIncarnation direct fencing', () => {
const LOCAL_SCOPE = JSON.stringify({ kind: 'local', hostId: 'local' })
// PTY_ID ('ssh:target@@relay-pty') derives connectionId 'target' via parseAppSshPtyId, so its
// host scope is ssh:target rather than local.
const SSH_TARGET_SCOPE = JSON.stringify({ kind: 'ssh', targetId: 'target' })
/** Register a pty leaf directly, optionally under a connection and incarnation. */
function seedPty(
runtime: OrcaRuntimeService,
ptyId: string,
incarnationId: string | null,
connectionId: string | null = null
): void {
runtime.registerPty(ptyId, WORKTREE_ID, connectionId, {
tabId: TAB_ID,
leafId: LEAF_ID,
...(incarnationId ? { incarnationId } : {})
})
}
function resolve(
runtime: OrcaRuntimeService,
processIncarnation: string,
serializedHostScope: string | null
): string | null {
return runtime.resolveTerminalHandleByProcessIncarnation(
processIncarnation,
serializedHostScope
)
}
it('mints a live handle for a pty whose exact incarnation and host scope match', () => {
const { runtime } = makeRuntime()
seedPty(runtime, PTY_ID, 'incarnation-1')
const handle = resolve(runtime, `${PTY_ID}:incarnation-1`, SSH_TARGET_SCOPE)
expect(handle).toMatch(/^term_/)
// Re-minting the same live pty is idempotent.
expect(resolve(runtime, `${PTY_ID}:incarnation-1`, SSH_TARGET_SCOPE)).toBe(handle)
})
it('returns null when the recorded incarnationId no longer matches the live pty', () => {
const { runtime } = makeRuntime()
seedPty(runtime, PTY_ID, 'incarnation-1')
expect(resolve(runtime, `${PTY_ID}:incarnation-2`, LOCAL_SCOPE)).toBeNull()
})
it('returns null when the live pty has no incarnationId (legacy runtimeId:ptyId:gen never reaps)', () => {
const { runtime } = makeRuntime()
seedPty(runtime, PTY_ID, null)
// A modern-shaped probe cannot match a pty that carries no incarnationId...
expect(resolve(runtime, `${PTY_ID}:incarnation-1`, LOCAL_SCOPE)).toBeNull()
// ...and the legacy `${runtimeId}:${ptyId}:${ptyGeneration}` shape stays fail-closed rather
// than reap on a ptyGeneration guess.
expect(resolve(runtime, `runtime_test:${PTY_ID}:0`, LOCAL_SCOPE)).toBeNull()
})
it('returns null when the host scope does not match', () => {
const { runtime } = makeRuntime()
seedPty(runtime, PTY_ID, 'incarnation-1')
expect(
resolve(runtime, `${PTY_ID}:incarnation-1`, JSON.stringify({ kind: 'ssh', targetId: 'nope' }))
).toBeNull()
})
it('returns null when no host scope is supplied (fails closed)', () => {
const { runtime } = makeRuntime()
seedPty(runtime, PTY_ID, 'incarnation-1')
expect(resolve(runtime, `${PTY_ID}:incarnation-1`, null)).toBeNull()
})
it('resolves a colon-bearing SSH/relay incarnationId that lastIndexOf would mis-split', () => {
const { runtime } = makeRuntime()
// PTY_ID already carries ':' and '@@'; the incarnationId itself also carries colons.
seedPty(runtime, PTY_ID, 'relay:conn-3:incarnation-9', 'target')
const handle = resolve(
runtime,
`${PTY_ID}:relay:conn-3:incarnation-9`,
JSON.stringify({ kind: 'ssh', targetId: 'target' })
)
expect(handle).toMatch(/^term_/)
})
it('resolves a Windows repo::C:\\path@@1 ptyId', () => {
const { runtime } = makeRuntime()
const windowsPtyId = 'repo::C:\\path@@1'
seedPty(runtime, windowsPtyId, 'incarnation-win')
const handle = resolve(runtime, `${windowsPtyId}:incarnation-win`, LOCAL_SCOPE)
expect(handle).toMatch(/^term_/)
})
it('keeps scanning past a colon-ambiguous decoy in another host scope to the genuine match', () => {
const { runtime } = makeRuntime()
// One process-incarnation string, two colon-ambiguous pty ids that both satisfy the exact
// matcher: 'relay' + 'conn:incarnation-1' AND 'relay:conn' + 'incarnation-1' both stringify to
// 'relay:conn:incarnation-1'. The decoy is registered FIRST and lives in a different host scope
// (ssh:decoy); the genuine pty is local and registered later.
const processIncarnation = 'relay:conn:incarnation-1'
const DECOY_TAB_ID = 'tab-decoy'
const DECOY_LEAF_ID = '22222222-2222-4222-8222-222222222222'
const GENUINE_TAB_ID = 'tab-genuine'
const GENUINE_LEAF_ID = '33333333-3333-4333-8333-333333333333'
runtime.registerPty('relay', WORKTREE_ID, 'decoy', {
tabId: DECOY_TAB_ID,
leafId: DECOY_LEAF_ID,
incarnationId: 'conn:incarnation-1'
})
runtime.registerPty('relay:conn', WORKTREE_ID, null, {
tabId: GENUINE_TAB_ID,
leafId: GENUINE_LEAF_ID,
incarnationId: 'incarnation-1'
})
// Before the fix, the earlier decoy's host-scope mismatch returned null and suppressed the
// genuine same-scope pty entirely. `continue` lets the scan reach it.
const handle = resolve(runtime, processIncarnation, LOCAL_SCOPE)
expect(handle).toMatch(/^term_/)
// getTerminalProcessIncarnation cannot separate the two (both stringify to
// 'relay:conn:incarnation-1'), so pin the remint to the GENUINE pane. The paneKey is the only
// terminal-specific observable that distinguishes the genuine leaf from the same-incarnation
// decoy; this fails if the resolver ever mints the decoy's handle instead of the genuine pty's.
const genuinePaneKey = makePaneKey(GENUINE_TAB_ID, GENUINE_LEAF_ID)
const decoyPaneKey = makePaneKey(DECOY_TAB_ID, DECOY_LEAF_ID)
expect(runtime.getTerminalPaneKey(handle!)).toBe(genuinePaneKey)
expect(runtime.getTerminalPaneKey(handle!)).not.toBe(decoyPaneKey)
})
})
@@ -0,0 +1,70 @@
import { describe, expect, it } from 'vitest'
import { matchesProcessIncarnation } from './worker-terminal-process-liveness'
describe('matchesProcessIncarnation', () => {
it.each([
{
name: 'an exact ptyId:incarnation match',
ptyId: 'pty-1',
incarnationId: 'incarnation-1',
processIncarnation: 'pty-1:incarnation-1',
expected: true
},
{
name: 'a Windows repo::C:\\path@@1 ptyId whose incarnation matches',
ptyId: 'repo::C:\\path@@1',
incarnationId: 'incarnation-win',
processIncarnation: 'repo::C:\\path@@1:incarnation-win',
expected: true
},
{
name: 'a colon-bearing relay incarnationId that a lastIndexOf split would mangle',
ptyId: 'relay-pty',
incarnationId: 'relay:conn-3:incarnation-9',
processIncarnation: 'relay-pty:relay:conn-3:incarnation-9',
expected: true
},
{
name: 'a whitespace-dirty incarnationId (lost contact is never a match)',
ptyId: 'pty-1',
incarnationId: ' incarnation-1',
processIncarnation: 'pty-1: incarnation-1',
expected: false
},
{
name: 'an empty-string incarnationId',
ptyId: 'pty-1',
incarnationId: '',
processIncarnation: 'pty-1:',
expected: false
},
{
name: 'an absent (null) incarnationId',
ptyId: 'pty-1',
incarnationId: null,
processIncarnation: 'pty-1:incarnation-1',
expected: false
},
{
name: 'an absent (undefined) incarnationId',
ptyId: 'pty-1',
incarnationId: undefined,
processIncarnation: 'pty-1:incarnation-1',
expected: false
},
{
name: 'a prefix-decoy ptyId (@@1) against a longer live pty (@@10)',
ptyId: 'repo::C:\\path@@1',
incarnationId: 'inc-1',
processIncarnation: 'repo::C:\\path@@10:inc-1',
expected: false
}
])('returns $expected for $name', ({ ptyId, incarnationId, processIncarnation, expected }) => {
expect(matchesProcessIncarnation(ptyId, incarnationId, processIncarnation)).toBe(expected)
})
it('rejects a partial prefix that is not colon-delimited', () => {
// 'pty-10' is not the pty 'pty-1'; the `${ptyId}:` fence keeps them distinct.
expect(matchesProcessIncarnation('pty-1', 'inc', 'pty-10:inc')).toBe(false)
})
})
@@ -5,6 +5,24 @@ import type { PtyProcessInfo } from '../../providers/pty-process-info'
export type { WorkerTerminalHostScope } from '../../../shared/worker-terminal-host-scope'
export { parseWorkerTerminalHostScope } from '../../../shared/worker-terminal-host-scope'
/**
* Does `processIncarnation` name exactly this pty's live incarnation? Requires exact
* `${ptyId}:${incarnationId}` equality, so it is immune to colons on either side (relay/SSH
* ptyIds, colon-bearing relay incarnationIds). A pty with no (or a whitespace-dirty)
* incarnationId can never match — the exact-incarnation fence stays intact.
*/
export function matchesProcessIncarnation(
ptyId: string,
incarnationId: string | null | undefined,
processIncarnation: string
): boolean {
if (!incarnationId || incarnationId !== incarnationId.trim()) {
return false
}
return `${ptyId}:${incarnationId}` === processIncarnation
}
/** Classify a recorded incarnation against live sessions: live on exact match, unverifiable when a candidate pty has a dirty or absent incarnationId (lost contact is never a death certificate), else exited. */
export function classifyWorkerTerminalProcessIncarnation(
processIncarnation: string,
sessions: readonly PtyProcessInfo[]
@@ -13,13 +31,9 @@ export function classifyWorkerTerminalProcessIncarnation(
processIncarnation.startsWith(`${session.id}:`)
)
if (
possibleMatches.some((session) => {
const incarnationId = session.incarnationId
if (!incarnationId || incarnationId !== incarnationId.trim()) {
return false
}
return `${session.id}:${incarnationId}` === processIncarnation
})
possibleMatches.some((session) =>
matchesProcessIncarnation(session.id, session.incarnationId, processIncarnation)
)
) {
return 'live'
}
@@ -0,0 +1,534 @@
import { afterEach, describe, expect, it, vi, type MockInstance } from 'vitest'
import { ORCHESTRATION_METHODS } from './orchestration'
import { eraseRpcMethods, type RpcContext } from '../core'
import { OrchestrationDb } from '../../orchestration/db'
import { OrcaRuntimeService } from '../../orca-runtime'
import { completeWorkerTerminalRelease } from './orchestration/worker/worker-release-completion'
describe('orchestration worker release incarnation fallback', () => {
let db: OrchestrationDb
let dbOpen = false
let runtime: OrcaRuntimeService
let ctx: RpcContext
let activeRunId: string
let inspectProcessLiveness: MockInstance<
OrcaRuntimeService['inspectTerminalProcessIncarnationLiveness']
>
const coordinatorPaneKey = 'tab_coord:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
const workerPaneKey = 'tab_worker:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb'
/** Fresh in-memory db and a fully-stubbed runtime for one worker-release scenario. */
function setup(): void {
db = new OrchestrationDb(':memory:')
dbOpen = true
runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(db)
inspectProcessLiveness = vi
.spyOn(runtime, 'inspectTerminalProcessIncarnationLiveness')
.mockResolvedValue('live')
vi.spyOn(runtime, 'getTerminalPaneKey').mockImplementation((handle) =>
handle === 'term_coord'
? coordinatorPaneKey
: handle === 'term_worker' || handle === 'term_reminted'
? workerPaneKey
: null
)
vi.spyOn(runtime, 'getTerminalProcessIncarnation').mockImplementation((handle) =>
handle === 'term_worker' || handle === 'term_reminted' ? 'runtime_test:term_worker:1' : null
)
vi.spyOn(runtime, 'getOrchestrationDispatchAuthority').mockImplementation((handle) =>
handle === 'term_worker' || handle === 'term_reminted'
? // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This test fixture is deliberately shaped to exercise the private/runtime boundary.
({
terminalHandle: handle,
paneKey: workerPaneKey,
processIncarnation: 'runtime_test:term_worker:1',
hostScope: { kind: 'local', hostId: 'local' }
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
} as never)
: null
)
vi.spyOn(runtime, 'validateOrchestrationAgentLauncher').mockImplementation(() => {})
vi.spyOn(runtime, 'showTerminal').mockImplementation(
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
async (handle) => ({ handle, worktreeId: 'repo::worktree', status: 'running' }) as never
)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
vi.spyOn(runtime, 'showManagedTerminalWorkspace').mockResolvedValue({
id: 'repo::worktree'
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
} as never)
vi.spyOn(runtime, 'createTerminal').mockResolvedValue({
handle: 'term_worker',
worktreeId: 'repo::worktree',
title: 'worker'
})
vi.spyOn(runtime, 'waitForTerminal').mockResolvedValue({
handle: 'term_worker',
condition: 'tui-idle',
satisfied: true,
status: 'running',
exitCode: null
})
vi.spyOn(runtime, 'getTerminalOrchestrationCliCommand').mockReturnValue('orca')
vi.spyOn(runtime, 'sendTerminalAgentPrompt').mockResolvedValue({
handle: 'term_worker',
accepted: true,
bytesWritten: 1
})
vi.spyOn(runtime, 'isTerminalRunningAgent').mockResolvedValue(true)
vi.spyOn(runtime, 'getExactWorkerProviderSession').mockReturnValue(null)
vi.spyOn(runtime, 'readTerminal').mockResolvedValue({
handle: 'term_worker',
status: 'running',
tail: ['worker output line 1', 'worker output line 2'],
truncated: false,
nextCursor: '2'
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
})
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
vi.spyOn(runtime, 'closeTerminal').mockResolvedValue({
handle: 'term_worker',
tabId: 'tab-worker',
ptyKilled: true
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
} as never)
vi.spyOn(runtime, 'notifyMessageArrived').mockImplementation(() => {})
activeRunId = db.createRun({
objective: 'Release test Run',
coordinatorHandle: 'term_coord',
coordinatorPaneKey
}).id
ctx = { runtime }
}
afterEach(() => {
if (dbOpen) {
dbOpen = false
db.close()
}
vi.restoreAllMocks()
})
/** Look up a registered orchestration RPC method by name. */
function findMethod(name: string) {
const method = eraseRpcMethods(ORCHESTRATION_METHODS).find((m) => m.name === name)
if (!method) {
throw new Error(`Method not found: ${name}`)
}
return method
}
/** Parse a method's params and invoke its handler against the shared ctx. */
async function call(name: string, params: Record<string, unknown>) {
const method = findMethod(name)
const parsed = method.params ? method.params.parse(params) : undefined
return method.handler(parsed, ctx)
}
/** Start a ready worker on a fresh task off the coordinator terminal. */
async function startWorker(options: { terminal?: string } = {}): Promise<{
taskId: string
dispatchId: string
}> {
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const task = db.createTask({ spec: 'release fixture task', runId: activeRunId })
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const result = (await call('orchestration.workerStart', {
task: task.id,
from: 'term_coord',
...(options.terminal ? { terminal: options.terminal } : { agent: 'codex' })
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
})) as { dispatchId: string; state: string }
expect(result.state).toBe('ready')
return { taskId: task.id, dispatchId: result.dispatchId }
}
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
/** Record the worker's report as settled — the precondition for release. */
function settle(taskId: string, dispatchId: string, outcome: 'succeeded' | 'failed'): void {
const settlement = db.settleWorkerReport({
taskId,
dispatchId,
outcome,
result: `worker ${outcome}`
})
expect(settlement.action).toBe('settled')
}
/** Start a worker and settle its report, the state a release acts on. */
async function startSettledWorker(
outcome: 'succeeded' | 'failed' = 'succeeded',
options: { terminal?: string } = {}
): Promise<{ taskId: string; dispatchId: string }> {
const worker = await startWorker(options)
settle(worker.taskId, worker.dispatchId, outcome)
return worker
}
it('closes a live worker terminal whose durable handle no longer resolves but whose process incarnation still matches', async () => {
setup()
const { dispatchId } = await startSettledWorker()
// The durable handle stops resolving (renderer graph epoch bump / handle invalidation)...
vi.mocked(runtime.showTerminal).mockImplementation(async (handle) =>
handle === 'term_worker'
? Promise.reject(new Error('terminal_handle_stale'))
: // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This test fixture is deliberately shaped to exercise the private/runtime boundary.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
({ handle, worktreeId: 'repo::worktree', status: 'running' } as never)
)
// ...but the recorded process incarnation still names a live PTY, re-minted to a fresh handle.
const resolveByIncarnation = vi.fn().mockReturnValue('term_reminted')
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
processAction: string
}
expect(receipt).toMatchObject({ state: 'released', processAction: 'closed_agent_terminal' })
expect(resolveByIncarnation).toHaveBeenCalledWith(
'runtime_test:term_worker:1',
JSON.stringify({ kind: 'local', hostId: 'local' })
)
// The close targeted exactly the re-minted live handle for that PTY, never the stale one.
expect(runtime.closeTerminal).toHaveBeenCalledTimes(1)
expect(runtime.closeTerminal).toHaveBeenCalledWith('term_reminted')
expect(db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({
ownership_state: 'released',
release_state: 'released'
})
})
it('stays release_unknown and closes nothing when the recorded incarnation no longer matches a live pty', async () => {
setup()
const { dispatchId } = await startSettledWorker()
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
// A reused ptyId now belongs to a different process: the incarnation mismatch refuses a close.
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
}
expect(receipt.state).toBe('release_unknown')
expect(resolveByIncarnation).toHaveBeenCalledWith(
'runtime_test:term_worker:1',
JSON.stringify({ kind: 'local', hostId: 'local' })
)
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)?.release_state).toBe('unknown')
})
it('does not plain-settle an exited missing worker when the archive was never committed', async () => {
setup()
const { dispatchId } = await startSettledWorker()
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
inspectProcessLiveness.mockResolvedValue('exited')
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
processAction: string
}
expect(receipt).toMatchObject({ state: 'release_unknown', processAction: 'none' })
expect(inspectProcessLiveness).toHaveBeenCalled()
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({
ownership_state: 'owned',
release_state: 'unknown'
})
})
it('settles released without a close when exited missing and an archive is already committed', async () => {
setup()
const { dispatchId } = await startSettledWorker()
const requested = db.requestWorkerTerminalRelease(dispatchId)
if (requested.disposition !== 'requested') {
throw new Error(`expected requested, got ${requested.disposition}`)
}
db.commitWorkerTerminalArchiveForRelease({
dispatchId,
resourceId: requested.resource.id,
kind: 'terminal_tail',
content: JSON.stringify({ lines: [] }),
archiveSource: 'terminal',
archiveStatus: 'empty'
})
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
inspectProcessLiveness.mockResolvedValue('exited')
const receipt = await completeWorkerTerminalRelease({
runtime,
db,
dispatchId,
resource: db.getWorkerTerminalResource(requested.resource.id)!,
mode: 'interactive'
})
expect(receipt).toMatchObject({ state: 'released', processAction: 'none' })
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({
ownership_state: 'released',
release_state: 'released'
})
})
it('reaches settleDead before the lease check when a gone worker is exited with no live authority', async () => {
// Without a committed archive settleDead retains; lease must not run first and force
// retained/identity_unproven. Disposition is release_unknown (interactive) after the retain.
setup()
const { dispatchId } = await startSettledWorker()
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockReturnValue(null)
inspectProcessLiveness.mockResolvedValue('exited')
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
processAction: string
}
expect(receipt).toMatchObject({ state: 'release_unknown', processAction: 'none' })
expect(receipt.state).not.toBe('retained')
expect(inspectProcessLiveness).toHaveBeenCalled()
expect(runtime.closeTerminal).not.toHaveBeenCalled()
})
it('concedes release_unknown before the lease check when a gone worker has no live authority and liveness is unproven', async () => {
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
setup()
const { dispatchId } = await startSettledWorker()
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockReturnValue(null)
// Liveness is unresolvable/not-exited: the process may have been re-homed, so concede rather
// than retain or guess at a live process.
inspectProcessLiveness.mockResolvedValue('live')
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
}
expect(receipt.state).toBe('release_unknown')
expect(receipt.state).not.toBe('retained')
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)?.release_state).toBe('unknown')
})
it('workerStop closes a live worker via the reminted handle when the durable handle is stale', async () => {
setup()
const { dispatchId } = await startWorker()
// The durable handle stops resolving, but the recorded incarnation still names a live PTY.
vi.mocked(runtime.showTerminal).mockImplementation(async (handle) =>
handle === 'term_worker'
? Promise.reject(new Error('terminal_handle_stale'))
: // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This test fixture is deliberately shaped to exercise the private/runtime boundary.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
({ handle, worktreeId: 'repo::worktree', status: 'running' } as never)
)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const resolveByIncarnation = vi.fn().mockReturnValue('term_reminted')
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerStop', { dispatch: dispatchId })) as {
processAction: string
}
expect(receipt.processAction).toBe('closed_agent_terminal')
// The kill targeted exactly the reminted live handle, never the stale durable one — closing
// the stale handle would throw terminal_handle_stale and leak the PTY.
expect(runtime.closeTerminal).toHaveBeenCalledTimes(1)
expect(runtime.closeTerminal).toHaveBeenCalledWith('term_reminted')
})
it('workerRead reads a live worker via the reminted handle when the durable handle is stale', async () => {
setup()
const { dispatchId } = await startWorker()
vi.mocked(runtime.showTerminal).mockImplementation(
async (handle) =>
handle === 'term_worker'
? Promise.reject(new Error('terminal_handle_stale'))
: // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This test fixture is deliberately shaped to exercise the private/runtime boundary.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
({ handle, worktreeId: 'repo::worktree', status: 'running' } as never)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
)
const resolveByIncarnation = vi.fn().mockReturnValue('term_reminted')
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const output = (await call('orchestration.workerRead', { dispatch: dispatchId })) as {
terminal?: { tail: string[] }
}
// Both the exact-session probe and the terminal read addressed the reminted handle.
expect(runtime.getExactWorkerProviderSession).toHaveBeenCalledWith(
'term_reminted',
expect.any(Number)
)
expect(runtime.readTerminal).toHaveBeenCalledWith('term_reminted', expect.anything())
expect(output.terminal?.tail).toEqual(['worker output line 1', 'worker output line 2'])
})
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
it('recovery-mode: defers when exited missing has no archive rather than plain-settling', async () => {
// Proof of death still runs settleDead first; when it retains (no archive), recovery must
// stay release_pending — not plain-settle, not unknown.
setup()
const { dispatchId } = await startSettledWorker()
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
inspectProcessLiveness.mockResolvedValue('exited')
const requested = db.requestWorkerTerminalRelease(dispatchId)
if (requested.disposition !== 'requested') {
throw new Error(`expected a requested release, got ${requested.disposition}`)
}
const receipt = await completeWorkerTerminalRelease({
runtime,
db,
dispatchId,
resource: requested.resource,
mode: 'recovery'
})
expect(receipt).toMatchObject({ state: 'release_pending', processAction: 'none' })
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({
ownership_state: 'owned',
release_state: 'requested'
})
})
it('recovery-mode: settles released before the defer when exited missing already has an archive', async () => {
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
setup()
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const { dispatchId } = await startSettledWorker()
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
inspectProcessLiveness.mockResolvedValue('exited')
const requested = db.requestWorkerTerminalRelease(dispatchId)
if (requested.disposition !== 'requested') {
throw new Error(`expected a requested release, got ${requested.disposition}`)
}
db.commitWorkerTerminalArchiveForRelease({
dispatchId,
resourceId: requested.resource.id,
kind: 'terminal_tail',
content: JSON.stringify({ lines: [] }),
archiveSource: 'terminal',
archiveStatus: 'empty'
})
const receipt = await completeWorkerTerminalRelease({
runtime,
db,
dispatchId,
resource: db.getWorkerTerminalResource(requested.resource.id)!,
mode: 'recovery'
})
expect(receipt).toMatchObject({ state: 'released', processAction: 'none' })
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({
ownership_state: 'released',
release_state: 'released'
})
})
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
it('recovery-mode: defers release_pending when liveness is unverifiable rather than provably exited', async () => {
setup()
const { dispatchId } = await startSettledWorker()
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
vi.mocked(runtime.showTerminal).mockRejectedValue(new Error('terminal_handle_stale'))
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const resolveByIncarnation = vi.fn().mockReturnValue(null)
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
resolveByIncarnation
)
// Not a death certificate: inventory may still be incomplete, so recovery must defer.
inspectProcessLiveness.mockResolvedValue('unverifiable')
const requested = db.requestWorkerTerminalRelease(dispatchId)
if (requested.disposition !== 'requested') {
throw new Error(`expected a requested release, got ${requested.disposition}`)
}
const receipt = await completeWorkerTerminalRelease({
runtime,
db,
dispatchId,
resource: requested.resource,
mode: 'recovery'
})
expect(receipt.state).toBe('release_pending')
expect(runtime.closeTerminal).not.toHaveBeenCalled()
expect(db.getWorkerTerminalResourceByOwner(dispatchId)?.release_state).not.toBe('released')
})
})
@@ -0,0 +1,358 @@
// PRB-0219 / upstream #18737 — FUNCTIONAL (integration) reproduction of the steady-state
// worker-release reap LEAK.
//
// This is the "functional/integration" tier of the PRB-0219 test strategy. It wires the REAL
// orchestration RPC surface (orchestration.workerStart / workerRelease / workerList), the REAL
// OrchestrationDb, and the REAL release completion + observation modules against a fake runtime
// that faithfully models the two identity planes involved in the bug:
//
// * VOLATILE, epoch-fenced handle table — `handleTable` keyed by terminal handle, each entry
// stamped with the `rendererGraphEpoch` at which it was issued. `showTerminal` resolves a
// handle ONLY while its stamped epoch matches the current epoch (mirrors getLiveLeafForHandle
// throwing 'terminal_handle_stale' on a rendererGraphEpoch bump). A relay reconnect / renderer
// remount on headless serve bumps the epoch.
// * DURABLE, incarnation-addressed process table — `ptysById` keyed by the ptyId embedded in the
// worker's persisted process_incarnation (`${ptyId}:${incarnationId}`). The pty stays LIVE
// across an epoch bump; nothing about the graph epoch kills the process.
//
// The bug (mechanism): when the epoch bumps, the durable db handle stops resolving through
// showTerminal WHILE THE PTY IS STILL ALIVE. inspectWorkerTerminal swallows the throw and reports
// `missing`; completeWorkerTerminalRelease then commits `release_unknown` and returns WITHOUT ever
// calling runtime.closeTerminal — so the process/PTY leaks. On mtl-02 those orphans accumulate in
// the orca-serve@factory cgroup until TasksMax=4096 is hit and Bun/omp abort() on EAGAIN.
//
// The observable leak signal asserted here: state 'release_unknown' + processAction 'none' +
// closeTerminal NEVER called + the pty STILL alive in ptysById + worker-list terminalState stuck
// on 'release_unknown' (never 'released').
import { afterEach, describe, expect, it, vi } from 'vitest'
import { ORCHESTRATION_METHODS } from './orchestration'
import { eraseRpcMethods, type RpcContext } from '../core'
import { OrchestrationDb } from '../../orchestration/db'
import { OrcaRuntimeService } from '../../orca-runtime'
describe('PRB-0219 worker-release reap FIX (functional verification)', () => {
let db: OrchestrationDb
let dbOpen = false
let runtime: OrcaRuntimeService
let ctx: RpcContext
let activeRunId: string
const coordinatorPaneKey = 'tab_coord:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
const workerPaneKey = 'tab_worker:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb'
// Dispatch/automation ptyId shape: `${repoId}::${worktreePath}@@${suffix}` (executionHostId:null).
const PTY_ID = 'repo-7f3a::/data/wt/factory-task-1@@a1b2c3d4'
const INCARNATION_ID = 1
const PROCESS_INCARNATION = `${PTY_ID}:${INCARNATION_ID}`
// ---- fake runtime process/handle planes (module-level so spies can mutate them) ----
let ptysById: Map<string, { incarnationId: number; alive: boolean }>
let handleTable: Map<string, { ptyId: string; epoch: number }>
let rendererGraphEpoch: number
let closedPtyIds: string[]
/** Resolve a handle to its live pty only while its stamped graph epoch is current. */
function resolveHandleToLivePty(
handle: string
): { ptyId: string; pty: { incarnationId: number; alive: boolean } } | null {
const entry = handleTable.get(handle)
if (!entry) {
return null
}
if (entry.epoch !== rendererGraphEpoch) {
// rendererGraphEpoch fence: the durable handle no longer resolves to a live leaf.
return null
}
const pty = ptysById.get(entry.ptyId)
if (!pty || !pty.alive) {
return null
}
return { ptyId: entry.ptyId, pty }
}
/** Wire the real RPC surface and db against the two-plane fake runtime. */
function setup(): void {
ptysById = new Map([
['coord-pty', { incarnationId: 1, alive: true }],
[PTY_ID, { incarnationId: INCARNATION_ID, alive: true }]
])
handleTable = new Map([
['term_coord', { ptyId: 'coord-pty', epoch: 0 }],
['term_worker', { ptyId: PTY_ID, epoch: 0 }]
])
rendererGraphEpoch = 0
closedPtyIds = []
db = new OrchestrationDb(':memory:')
dbOpen = true
runtime = new OrcaRuntimeService()
runtime.setOrchestrationDb(db)
// Incarnation-addressed liveness probe: reads the DURABLE plane, so it stays 'live' across the
// epoch bump (the process really is still running). Matches the real asymmetry.
vi.spyOn(runtime, 'inspectTerminalProcessIncarnationLiveness').mockImplementation(
async (incarnation: string) => {
const idx = incarnation.lastIndexOf(':')
const ptyId = incarnation.slice(0, idx)
const pty = ptysById.get(ptyId)
return pty?.alive ? 'live' : 'exited'
}
)
// The incarnation-addressed re-resolution primitive the fix adds. Present here as a spy so the
// same harness proves BOTH tiers: pre-fix completion never calls it (leak); post-fix completion
// calls it to remint a live handle (reap). Fence: EXACT incarnationId match only.
vi.spyOn(runtime, 'resolveTerminalHandleByProcessIncarnation').mockImplementation(
(incarnation: string): string | null => {
const idx = incarnation.lastIndexOf(':')
const ptyId = incarnation.slice(0, idx)
const inc = Number(incarnation.slice(idx + 1))
const pty = ptysById.get(ptyId)
if (!pty || !pty.alive) {
return null
}
if (pty.incarnationId !== inc) {
// Fence: a reused ptyId with a different incarnation must NOT resolve.
return null
}
// Remint a live handle at the current graph epoch.
handleTable.set('term_reminted', { ptyId, epoch: rendererGraphEpoch })
return 'term_reminted'
}
)
// Identity plane (durable) — answers for the original AND any reminted handle. Independent of
// the graph epoch, exactly like the real getTerminal* accessors that read dispatch authority.
const knownWorkerHandle = (handle: string): boolean =>
handle === 'term_worker' || handle === 'term_reminted'
vi.spyOn(runtime, 'getTerminalPaneKey').mockImplementation((handle) =>
handle === 'term_coord'
? coordinatorPaneKey
: knownWorkerHandle(handle)
? workerPaneKey
: null
)
vi.spyOn(runtime, 'getTerminalProcessIncarnation').mockImplementation((handle) =>
knownWorkerHandle(handle) ? PROCESS_INCARNATION : null
)
vi.spyOn(runtime, 'getOrchestrationDispatchAuthority').mockImplementation((handle) =>
knownWorkerHandle(handle)
? // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: This test fixture is deliberately shaped to exercise the private/runtime boundary.
({
terminalHandle: handle,
paneKey: workerPaneKey,
processIncarnation: PROCESS_INCARNATION,
hostScope: { kind: 'local', hostId: 'local' }
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
} as never)
: null
)
vi.spyOn(runtime, 'validateOrchestrationAgentLauncher').mockImplementation(() => {})
// Volatile handle resolution — epoch-fenced. Throws 'terminal_handle_stale' once the epoch
// moves past the epoch at which the handle was issued.
vi.spyOn(runtime, 'showTerminal').mockImplementation(async (handle) => {
if (!resolveHandleToLivePty(handle)) {
throw new Error('terminal_handle_stale')
}
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
return { handle, worktreeId: 'repo::worktree', status: 'running' } as never
})
// The reap: closing a handle kills exactly the pty it resolves to.
vi.spyOn(runtime, 'closeTerminal').mockImplementation(async (handle) => {
const live = resolveHandleToLivePty(handle)
if (!live) {
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
return { handle, tabId: null, ptyKilled: false } as never
}
live.pty.alive = false
ptysById.delete(live.ptyId)
closedPtyIds.push(live.ptyId)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
return { handle, tabId: `tab:${live.ptyId}`, ptyKilled: true } as never
})
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// Remaining runtime surface required to start + settle a worker (mirrors the unit harness).
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
vi.spyOn(runtime, 'showManagedTerminalWorkspace').mockResolvedValue({
id: 'repo::worktree'
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
} as never)
vi.spyOn(runtime, 'createTerminal').mockResolvedValue({
handle: 'term_worker',
worktreeId: 'repo::worktree',
title: 'worker'
})
vi.spyOn(runtime, 'waitForTerminal').mockResolvedValue({
handle: 'term_worker',
condition: 'tui-idle',
satisfied: true,
status: 'running',
exitCode: null
})
vi.spyOn(runtime, 'getTerminalOrchestrationCliCommand').mockReturnValue('orca')
vi.spyOn(runtime, 'sendTerminalAgentPrompt').mockResolvedValue({
handle: 'term_worker',
accepted: true,
bytesWritten: 1
})
vi.spyOn(runtime, 'isTerminalRunningAgent').mockResolvedValue(true)
vi.spyOn(runtime, 'getExactWorkerProviderSession').mockReturnValue(null)
vi.spyOn(runtime, 'readTerminal').mockResolvedValue({
handle: 'term_worker',
status: 'running',
tail: ['worker output line 1', 'worker output line 2'],
truncated: false,
nextCursor: '2'
})
vi.spyOn(runtime, 'notifyMessageArrived').mockImplementation(() => {})
activeRunId = db.createRun({
objective: 'PRB-0219 reap leak fixture',
coordinatorHandle: 'term_coord',
coordinatorPaneKey
}).id
ctx = { runtime }
}
afterEach(() => {
if (dbOpen) {
dbOpen = false
db.close()
}
vi.restoreAllMocks()
})
/** Look up a registered orchestration RPC method by name. */
function findMethod(name: string) {
const method = eraseRpcMethods(ORCHESTRATION_METHODS).find((m) => m.name === name)
if (!method) {
throw new Error(`Method not found: ${name}`)
}
return method
}
/** Parse a method's params and invoke its handler against the shared ctx. */
async function call(name: string, params: Record<string, unknown>) {
const method = findMethod(name)
const parsed = method.params ? method.params.parse(params) : undefined
return method.handler(parsed, ctx)
}
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
/** Start a worker and settle its report, the state a release acts on. */
async function startSettledWorker(): Promise<{ taskId: string; dispatchId: string }> {
const task = db.createTask({ spec: 'reap-leak fixture task', runId: activeRunId })
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const result = (await call('orchestration.workerStart', {
task: task.id,
from: 'term_coord',
agent: 'codex'
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
})) as { dispatchId: string; state: string }
expect(result.state).toBe('ready')
const settlement = db.settleWorkerReport({
taskId: task.id,
dispatchId: result.dispatchId,
outcome: 'succeeded',
result: 'worker succeeded'
})
expect(settlement.action).toBe('settled')
return { taskId: task.id, dispatchId: result.dispatchId }
}
/** The terminalState workerList projects for a dispatch, or null. */
async function workerTerminalState(dispatchId: string): Promise<string | null> {
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const listed = (await call('orchestration.workerList', { run: activeRunId })) as {
workers: { dispatchId: string; terminalState: string | null }[]
}
return listed.workers.find((w) => w.dispatchId === dispatchId)?.terminalState ?? null
}
it('REAP (fixed): a rendererGraphEpoch bump strands the durable handle, but the incarnation fallback remints a live handle and reaps the process', async () => {
setup()
const { dispatchId } = await startSettledWorker()
const resource = db.getWorkerTerminalResourceByOwner(dispatchId)
expect(resource?.process_incarnation).toBe(PROCESS_INCARNATION)
// Relay reconnect / renderer remount bumps the graph epoch: the durable handle goes stale
// while the PTY stays alive.
rendererGraphEpoch = 1
await expect(runtime.showTerminal('term_worker')).rejects.toThrow('terminal_handle_stale')
expect(ptysById.get(PTY_ID)?.alive).toBe(true)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
processAction: string
}
// The fix: inspectWorkerTerminal re-resolved the live PTY by process incarnation, reminted a
// handle, and completion closed THAT handle — the process is actually reaped.
expect(receipt.state).toBe('released')
expect(receipt.processAction).toBe('closed_agent_terminal')
expect(runtime.closeTerminal).toHaveBeenCalledWith('term_reminted')
expect(closedPtyIds).toEqual([PTY_ID])
expect(ptysById.has(PTY_ID)).toBe(false)
expect(await workerTerminalState(dispatchId)).toBe('released')
})
it('CONTROL: with the graph epoch intact the same release reaps exactly that PTY', async () => {
setup()
const { dispatchId } = await startSettledWorker()
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
processAction: string
}
expect(receipt.state).toBe('released')
expect(receipt.processAction).toBe('closed_agent_terminal')
expect(closedPtyIds).toEqual([PTY_ID])
expect(ptysById.has(PTY_ID)).toBe(false)
expect(await workerTerminalState(dispatchId)).toBe('released')
})
it('FENCE (fixed): a reused ptyId carrying a different incarnation must NOT remint or close — stays release_unknown', async () => {
setup()
const { dispatchId } = await startSettledWorker()
// The graph epoch bumps AND the ptyId has been reused by a different process (incarnation 2),
// while the worker's recorded incarnation is still 1. The exact-incarnation fence must refuse.
rendererGraphEpoch = 1
const reused = ptysById.get(PTY_ID)
if (reused) {
reused.incarnationId = 2
}
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: Test fixture crosses a private/runtime boundary with a verified shape.
const receipt = (await call('orchestration.workerRelease', { dispatch: dispatchId })) as {
state: string
processAction: string
}
expect(receipt.state).toBe('release_unknown')
expect(receipt.processAction).toBe('none')
expect(runtime.closeTerminal).not.toHaveBeenCalled()
// The other lane's live process is left untouched (never reaped by an over-broad match).
expect(ptysById.get(PTY_ID)?.alive).toBe(true)
expect(await workerTerminalState(dispatchId)).toBe('release_unknown')
})
})
@@ -4,15 +4,15 @@ import type { WorkerTerminalResourceRow } from '../../../../orchestration/worker
import { stopStructuredWorker } from '../../orchestration-structured-worker-lifecycle'
import type { StructuredWorkerIdentity } from '../../../../structured-worker-identity'
import { archiveSummary } from './worker-terminal-resource-presentation'
import type { WorkerReleaseReceipt } from './worker-release-completion'
import { releaseUnknownRecovery, type WorkerReleaseReceipt } from './worker-release-completion'
/**
* The close half of a release for a worker that IS a structured session.
*
* Separate from the PTY close for the same reason the delivery lane is: there is no terminal to
* close and no exit to observe, so the host's own settlement is the only proof available. Only a
* proven close may settle; an unproven one reports `release_unknown` and stays retryable under the
* same request id.
* proven close may settle; an unproven one reports `release_unknown` and stays retryable, but only
* under a fresh request id — replaying the prior one just returns this same receipt.
*/
export async function stopStructuredWorkerForRelease(args: {
structured: StructuredWorkerIdentity
@@ -36,7 +36,7 @@ export async function stopStructuredWorkerForRelease(args: {
processAction: stop.closeAttempted ? 'closed_agent_terminal' : 'none',
archive: { source: args.archiveSource, status: args.archiveStatus },
lastError: unknown.release_error ?? stop.reason,
recovery: `Inspect with: orca orchestration worker-show --dispatch ${dispatchId} --json — then repeat worker-release with the same --retry-request.`
recovery: releaseUnknownRecovery(dispatchId)
}
}
const settled = db.settleWorkerTerminalRelease(resource.id)
@@ -143,6 +143,9 @@ export const ORCHESTRATION_WORKER_CONTROL_METHODS = [
`Worker Dispatch ${params.dispatch} no longer resolves to its exact process.`
)
}
// Read via the handle inspectWorkerTerminal proved live: the durable one, or a handle
// re-minted from the recorded incarnation after the durable handle went stale.
const liveHandle = observation.terminalHandle ?? terminalHandle
const structured = await readStructuredWorkerOutput({
db,
dispatchId: params.dispatch,
@@ -163,7 +166,7 @@ export const ORCHESTRATION_WORKER_CONTROL_METHODS = [
const output = await readExactWorkerOutput({
runtime,
dispatchId: params.dispatch,
terminalHandle,
terminalHandle: liveHandle,
workerState: worker?.state ?? 'unsupervised',
terminalStatus:
observation.status === 'exited'
@@ -16,6 +16,7 @@ import type {
WorkerDispatchRow
} from '../../../../orchestration/types'
/** Observe a worker terminal, re-minting a live handle from the recorded process incarnation when the durable handle went stale, so a still-running worker is never reported missing and leaked. */
export async function inspectWorkerTerminal(
runtime: OrcaRuntimeService,
db: OrchestrationDb,
@@ -31,12 +32,15 @@ export async function inspectWorkerTerminal(
/** Structured workers only: whether mail still reaches it — at rest included, since the mail
* starts it. Absent when ownership cannot be read. `status` stays the process verdict. */
addressable?: boolean
/** The handle that actually resolved: the durable one, or a live handle re-minted from the
* recorded process incarnation after the durable handle went stale. Null when none resolved. */
terminalHandle: string | null
}> {
const worker = db.getWorkerDispatch(dispatchId)
const terminalHandle =
worker?.agent_terminal_handle ?? db.getDispatchContextById(dispatchId)?.assignee_handle
if (!terminalHandle) {
return { terminal: null, exact: false, status: 'unattached' }
return { terminal: null, exact: false, status: 'unattached', terminalHandle: null }
}
const structured = resolveStructuredWorkerForDispatch(db, dispatchId)
if (structured) {
@@ -63,20 +67,44 @@ export async function inspectWorkerTerminal(
exact,
status: exact ? observation.status : 'identity_changed',
...(exact && observation.reason ? { reason: observation.reason } : {}),
...(exact && addressable !== null ? { addressable } : {})
...(exact && addressable !== null ? { addressable } : {}),
terminalHandle: null
}
}
const terminal = await runtime.showTerminal(terminalHandle).catch(() => null)
let effectiveHandle = terminalHandle
let terminal = await runtime.showTerminal(effectiveHandle).catch(() => null)
if (!terminal) {
return { terminal: null, exact: false, status: 'missing' }
// Why: the durable handle resolves nowhere after a renderer graph epoch bump or handle
// invalidation, yet the recorded process incarnation may still name a live PTY. Re-mint a
// live handle (incarnation-fenced) so worker-show and release act on the still-running
// process instead of reporting it missing — which would leak the agent process tree.
// (workerList is a pure DB projection and never calls inspectWorkerTerminal.)
const resource = db.getWorkerTerminalResourceByOwner(dispatchId)
const reminted = resource?.process_incarnation
? runtime.resolveTerminalHandleByProcessIncarnation(
resource.process_incarnation,
resource.host_scope
)
: null
if (reminted) {
const remintedTerminal = await runtime.showTerminal(reminted).catch(() => null)
if (remintedTerminal) {
effectiveHandle = reminted
terminal = remintedTerminal
}
}
}
if (!terminal) {
// The re-mint above failed, so no live handle resolved; report the durable handle unresolved.
return { terminal: null, exact: false, status: 'missing', terminalHandle: null }
}
const exact = db.isDispatchProcessCurrent({
dispatchId,
paneKey: runtime.getTerminalPaneKey(terminalHandle),
processIncarnation: runtime.getTerminalProcessIncarnation(terminalHandle)
paneKey: runtime.getTerminalPaneKey(effectiveHandle),
processIncarnation: runtime.getTerminalProcessIncarnation(effectiveHandle)
})
if (!exact) {
return { terminal, exact, status: 'identity_changed' }
return { terminal, exact, status: 'identity_changed', terminalHandle: effectiveHandle }
}
// Why: the aggregate inventory only iterates registered providers, so a dropped
// relay clears `connected` for every remote PTY at once. Lost contact is not a
@@ -86,38 +114,48 @@ export async function inspectWorkerTerminal(
// Exact-gated by the early return above: a replaced process's prompt would attribute another
// lane's blocker to this worker.
const agentWait = terminal.agentWait
const verdict = runtime.getTerminalLivenessVerdict?.(terminalHandle) ?? null
const verdict = runtime.getTerminalLivenessVerdict?.(effectiveHandle) ?? null
if (verdict?.status === 'unverifiable') {
return { terminal, exact, status: 'unverifiable', reason: verdict.reason, agentWait }
return {
terminal,
exact,
status: 'unverifiable',
reason: verdict.reason,
agentWait,
terminalHandle: effectiveHandle
}
}
if (verdict?.status === 'live') {
return { terminal, exact, status: 'live', agentWait }
return { terminal, exact, status: 'live', agentWait, terminalHandle: effectiveHandle }
}
if (!verdict) {
const dispatch = db.getDispatchContextById?.(dispatchId)
const persistedHostScope = parseWorkerTerminalHostScope(dispatch?.host_scope ?? null)
const currentHostScope = runtime.getOrchestrationDispatchAuthority?.(terminalHandle)?.hostScope
const currentHostScope = runtime.getOrchestrationDispatchAuthority?.(effectiveHandle)?.hostScope
if (persistedHostScope?.kind === 'ssh' || currentHostScope?.kind === 'ssh') {
return {
terminal,
exact,
status: 'unverifiable',
reason: 'missing_liveness_verdict',
agentWait
agentWait,
terminalHandle: effectiveHandle
}
}
return {
terminal,
exact,
status: terminal.connected === false ? 'exited' : 'live',
agentWait
agentWait,
terminalHandle: effectiveHandle
}
}
return {
terminal,
exact,
status: 'exited',
agentWait
agentWait,
terminalHandle: effectiveHandle
}
}
@@ -120,6 +120,9 @@ async function completeWorkerTerminalReleaseOnce(
}
}
const observation = await inspectWorkerTerminal(runtime, db, dispatchId)
// The live handle to act on: the durable one, or a handle re-minted from the recorded process
// incarnation when the durable handle went stale (inspectWorkerTerminal proved it live).
const terminalHandle = observation.terminalHandle ?? resource.terminal_handle
if (observation.status === 'identity_changed') {
const retained = db.revertWorkerTerminalReleaseToRetained(resource.id, 'identity_unproven')
return {
@@ -131,33 +134,39 @@ async function completeWorkerTerminalReleaseOnce(
}
}
if (observation.status === 'missing' || observation.status === 'unattached') {
if (args.mode === 'recovery') {
// A close can succeed before the process crashes, leaving `releasing` durable state while
// terminal inventory no longer resolves the handle. Only a positive host liveness verdict
// may settle that exact incarnation; contact loss remains pending/unverifiable.
if (resource.process_incarnation) {
const processLiveness = await runtime.inspectTerminalProcessIncarnationLiveness(
resource.process_incarnation,
resource.host_scope
)
if (processLiveness === 'exited') {
const reconciled = db.settleDeadWorkerTerminalRelease({
requestingDispatchId: dispatchId,
resourceId: resource.id,
processIncarnation: resource.process_incarnation
})
if (reconciled.disposition === 'released') {
runtime.notifyMessageArrived(`dispatch:${dispatchId}`, 'status')
return {
dispatchId,
state: 'released',
processAction: 'closed_exited_terminal',
archive: archiveSummary(reconciled.resource)
}
// Re-resolution by process incarnation (inspectWorkerTerminal) already failed, so no live PTY
// carries this worker's exact incarnation. If that incarnation is provably gone, settle
// released BEFORE the recovery defer: proof of death outranks deferral, so a provably-exited
// worker never languishes in release_pending across recovery passes.
if (resource.process_incarnation) {
const processLiveness = await runtime.inspectTerminalProcessIncarnationLiveness(
resource.process_incarnation,
resource.host_scope
)
if (processLiveness === 'exited') {
// Prefer incarnation-fenced settle (dispatch relation + process_incarnation CAS).
const reconciled = db.settleDeadWorkerTerminalRelease({
requestingDispatchId: dispatchId,
resourceId: resource.id,
processIncarnation: resource.process_incarnation
})
if (reconciled.disposition === 'released') {
runtime.notifyMessageArrived(`dispatch:${dispatchId}`, 'status')
return {
dispatchId,
state: 'released',
processAction: 'none',
archive: archiveSummary(reconciled.resource)
}
}
// settleDead retains when the archive is still mandatory and missing (e.g. requested but
// never committed). Do NOT plain-settle: that would discard output and break recovery's
// "archive is mandatory" invariant. Fall through to recovery pending / unknown instead.
}
// Inventory may still be incomplete during startup/reconnect discovery; defer.
}
if (args.mode === 'recovery') {
// No death certificate yet: inventory may still be incomplete during startup/reconnect
// discovery, so defer instead of guessing.
return {
dispatchId,
state: 'release_pending',
@@ -183,7 +192,7 @@ async function completeWorkerTerminalReleaseOnce(
}
}
if (!workerTerminalLeaseIsCurrent(runtime, db, dispatchId, resource)) {
if (!workerTerminalLeaseIsCurrent(runtime, db, dispatchId, resource, terminalHandle)) {
const retained = db.revertWorkerTerminalReleaseToRetained(resource.id, 'identity_unproven')
return {
dispatchId,
@@ -202,7 +211,7 @@ async function completeWorkerTerminalReleaseOnce(
const captured = await captureWorkerOutputArchive({
runtime,
dispatchId,
terminalHandle: resource.terminal_handle,
terminalHandle,
attachedAtMs: orchestrationTimestampToMs(worker.created_at),
structuredWorker: structured
})
@@ -230,7 +239,7 @@ async function completeWorkerTerminalReleaseOnce(
archive: archiveSummary(releasing)
}
}
if (!workerTerminalLeaseIsCurrent(runtime, db, dispatchId, releasing)) {
if (!workerTerminalLeaseIsCurrent(runtime, db, dispatchId, releasing, terminalHandle)) {
const retained = db.revertWorkerTerminalReleaseToRetained(resource.id, 'identity_unproven')
return {
dispatchId,
@@ -253,7 +262,7 @@ async function completeWorkerTerminalReleaseOnce(
archiveStatus
})
}
const close = await runtime.closeTerminal(resource.terminal_handle)
const close = await runtime.closeTerminal(terminalHandle)
if (!close.ptyKilled) {
const reason = describeUnconfirmedAgentStop(close)
const unknown = db.markWorkerTerminalReleaseUnknown(resource.id, reason)
@@ -0,0 +1,79 @@
import { afterEach, describe, expect, it, vi } from 'vitest'
import { inspectWorkerTerminal } from './worker-observation'
import { createOrchestrationWorkerReleaseHarness } from './worker-release.test-support'
import { makePaneKey } from '../../../../../../shared/stable-pane-id'
const PTY_ID = 'runtime_test:term_worker'
const LEAF_ID = 'bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb'
describe('worker release through runtime incarnation recovery', () => {
const harness = createOrchestrationWorkerReleaseHarness()
afterEach(() => harness.cleanup())
it.each(['release', 'stop', 'read', 'replacement', 'disconnected'] as const)(
'%s uses runtime identity and close paths',
async (scenario) => {
harness.setup()
const { dispatchId } =
scenario === 'stop' ? await harness.startWorker() : await harness.startSettledWorker()
const runtime = harness.runtime
vi.mocked(runtime.showTerminal).mockRestore()
vi.mocked(runtime.getTerminalPaneKey).mockRestore()
vi.mocked(runtime.getTerminalProcessIncarnation).mockRestore()
vi.mocked(runtime.getOrchestrationDispatchAuthority).mockRestore()
vi.mocked(runtime.readTerminal).mockRestore()
vi.mocked(runtime.closeTerminal).mockRestore()
const kill = vi.fn(() => true)
runtime.setPtyController({ write: () => true, kill, getForegroundProcess: async () => null })
runtime.registerPty(PTY_ID, 'repo::worktree', null, {
tabId: 'tab_worker',
leafId: LEAF_ID,
incarnationId: scenario === 'replacement' ? '2' : '1'
})
await expect(runtime.showTerminal('term_worker')).rejects.toThrow()
if (scenario === 'disconnected') {
vi.spyOn(runtime, 'getTerminalLivenessVerdict').mockReturnValue({
status: 'unverifiable',
reason: 'transport unavailable'
})
const observed = await inspectWorkerTerminal(runtime, harness.db, dispatchId)
expect(observed).toMatchObject({ exact: true, status: 'unverifiable' })
expect(kill).not.toHaveBeenCalled()
return
}
if (scenario === 'stop') {
const receipt = await harness.call('orchestration.workerStop', { dispatch: dispatchId })
expect(receipt).toMatchObject({ state: 'stopped', processAction: 'closed_agent_terminal' })
expect(kill).toHaveBeenCalledExactlyOnceWith(PTY_ID)
return
}
if (scenario === 'read') {
const readTerminal = vi.spyOn(runtime, 'readTerminal')
const output = await harness.call('orchestration.workerRead', { dispatch: dispatchId })
// Pins the read to the reminted terminal: the handle it reached must resolve to the
// registered pane and incarnation, which the stale durable handle never does.
const [[readHandle]] = readTerminal.mock.calls
expect(runtime.getTerminalPaneKey(readHandle)).toBe(makePaneKey('tab_worker', LEAF_ID))
expect(runtime.getTerminalProcessIncarnation(readHandle)).toBe(`${PTY_ID}:1`)
expect(output).toMatchObject({
source: 'terminal',
fallbackReason: 'session_not_reported'
})
expect(kill).not.toHaveBeenCalled()
return
}
const receipt = await harness.call('orchestration.workerRelease', { dispatch: dispatchId })
if (scenario === 'replacement') {
expect(receipt).toMatchObject({ state: 'release_unknown', processAction: 'none' })
expect(kill).not.toHaveBeenCalled()
return
}
expect(receipt).toMatchObject({ state: 'released', processAction: 'closed_agent_terminal' })
expect(kill).toHaveBeenCalledExactlyOnceWith(PTY_ID)
expect(harness.db.getWorkerTerminalResourceByOwner(dispatchId)).toMatchObject({
ownership_state: 'released',
release_state: 'released'
})
}
)
})
@@ -141,6 +141,7 @@ export const ORCHESTRATION_WORKER_STOP_METHODS = [
})
}
const observation = await inspectWorkerTerminal(runtime, db, params.dispatch)
const liveHandle = observation.terminalHandle ?? handle
// The host exit can settle this stop while terminal inspection is awaiting inventory.
if (db.getWorkerDispatch(params.dispatch)?.state === 'stopped') {
runtime.notifyMessageArrived(`dispatch:${params.dispatch}`, 'status')
@@ -201,7 +202,7 @@ export const ORCHESTRATION_WORKER_STOP_METHODS = [
}
}
const closed = await runtime
.closeTerminal(handle)
.closeTerminal(liveHandle)
.then((close) => ({ close }) as const)
.catch(
(error: unknown) =>
@@ -9,23 +9,27 @@ export function workerTerminalLeaseIsCurrent(
runtime: OrcaRuntimeService,
db: OrchestrationDb,
dispatchId: string,
resource: WorkerTerminalResourceRow
resource: WorkerTerminalResourceRow,
// Live (possibly reminted) handle for runtime probes. Durable identity still compares
// worker.agent_terminal_handle to resource.terminal_handle; after beginGraphReload the durable
// string is absent from the handle table, so authority/pane/incarnation must use the live one.
liveTerminalHandle: string
): boolean {
const worker = db.getWorkerDispatch(dispatchId)
if (isStructuredWorkerHandle(resource.terminal_handle)) {
return structuredWorkerTerminalLeaseIsCurrent(db, dispatchId, worker, resource)
}
const authority = runtime.getOrchestrationDispatchAuthority(resource.terminal_handle)
const authority = runtime.getOrchestrationDispatchAuthority(liveTerminalHandle)
// Exited PTYs retain identity and host evidence but no longer mint launch authority.
return Boolean(
worker?.agent_terminal_handle === resource.terminal_handle &&
(authority
? resource.host_scope === JSON.stringify(authority.hostScope)
: runtime.getTerminalLivenessVerdict(resource.terminal_handle)?.status === 'exited') &&
: runtime.getTerminalLivenessVerdict(liveTerminalHandle)?.status === 'exited') &&
db.isDispatchProcessCurrent({
dispatchId,
paneKey: runtime.getTerminalPaneKey(resource.terminal_handle),
processIncarnation: runtime.getTerminalProcessIncarnation(resource.terminal_handle)
paneKey: runtime.getTerminalPaneKey(liveTerminalHandle),
processIncarnation: runtime.getTerminalProcessIncarnation(liveTerminalHandle)
}) &&
!db.workerTerminalResourceHasIdentityConflict(resource.id)
)