mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(orchestration): stop three PTY-only probes answering for structured sessions
Three defects, one shape: a probe that enumerates PTYs or resolves a pane was standing in for a question that is not about panes at all. `worktree rm` destroyed a live structured worker. `killAllProcessesForWorktree` sweeps the renderer graph, the provider session list and the local pty-registry, and a structured session is registered on none of them — so all three counted zero, nothing errored, and removal deleted the checkout out from under a running provider child, which kept running with its `cwd` gone while the dispatch still reported the worker live and exact. A fourth sweep now asks what the other three cannot: membership by `location.workspaceId`, which covers a plain chat session as well as a dispatched worker, and liveness by the same `live`/`unverifiable`/`exited` observation the rest of the structured surface uses. It REFUSES a destructive removal rather than auto-closing, on the same bargain and the same `--force` escape hatch as the unstopped-PTY gate — this is the verb that deletes a user's work, and a running agent is exactly what they would want to be told about. Force closes the sessions properly instead of orphaning a child. Best-effort reconciliation callers are excluded: they repair state, delete nothing, and must never be failed closed. Twelve coordinator verbs failed for a structured worker running as itself. `isLiveTerminalHandle` validated `ORCA_TERMINAL_HANDLE` with `terminal.show`, a PTY verb whose leaf lookup misses for a session that never had a pane; the pane remint that would have recovered it needs `ORCA_PANE_KEY`, which a structured child deliberately does not carry, so every one of them died on `no_active_sender_terminal` — including the ones the worker's own dispatch preamble tells it to run. The identity question gets its own probe, `terminal.resolveIdentity`: a handle and a boolean and nothing writable. `terminal.show` still refuses a structured handle, because synthesising ptyId/leafId/paneRuntimeId would hand every public terminal verb something that looks writable and is not. The PTY half is byte-for-byte today's check, `getLiveLeafForHandle` included, so its `rendererGraphEpoch` re-check still runs — that check is the whole reason the sender is validated at all, and a cheaper probe would have quietly started passing stale post-reload handles. A host that predates the method answers `method_not_found` and the client falls back to `terminal.show`, which is correct for that host: one without the identity probe has no structured workers to miss. `dispatch --inject` reported `no_agent_detected` for a structured worker, because `isTerminalRunningAgent` reaches `getLiveLeaf`, throws, and the catch returns false. A structured session IS the agent; there is no foreground process to recognise, so it answers before the PTY probes rather than through them. Also: a Run whose coordinator is structured now gets its `run:` mail. Both lanes declined and neither logged — the PTY lane because the owner is structured, the structured lane because the mailbox was not `dispatch:` — so each half believed the other owned it. The PTY lane's reasoning (a coordinator blocks in `check --wait`, where a waiter preempts pointer delivery) does not transfer: a structured coordinator is a chat session whose turn ends. Its `run:` deliveries take the `hasOutstandingRunDelivery` gate the PTY lane applies for exactly that mailbox, and only for that mailbox. The test that would have caught the twelve drives the CLI with `ORCA_TERMINAL_HANDLE=structworker_…` and no `--from`. Every existing orchestration CLI test passes `--from` explicitly, so the resolver a real worker goes through was never exercised — which is why the suite stayed green while the preamble failed on its first line. Two files crossed their line ceiling and are split rather than waived: `worktree-teardown.ts` sheds its two PTY-surface sweeps and the deadline arithmetic they share, and `orchestration.test.ts` — which sat exactly on 800 — sheds the two caller-identity suites this change rewrote.
This commit is contained in:
@@ -0,0 +1,377 @@
|
||||
/**
|
||||
* How the orchestration CLI decides WHO is speaking.
|
||||
*
|
||||
* Split out of `orchestration.test.ts`, which sat exactly on the test-file line ceiling: these two
|
||||
* suites are one subject — the coordinator and task-creator identity a command carries — and both
|
||||
* exercise the env-handle validation and pane-remint chain rather than flag-to-param mapping.
|
||||
*/
|
||||
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
const callMock = vi.fn()
|
||||
const getTerminalHandleMock = vi.hoisted(() => vi.fn())
|
||||
const originalTerminalHandle = process.env.ORCA_TERMINAL_HANDLE
|
||||
const originalPaneKey = process.env.ORCA_PANE_KEY
|
||||
// Why: isolate the handler's flag-to-param mapping; printResult only writes output.
|
||||
vi.mock('../format', () => ({ printResult: vi.fn() }))
|
||||
vi.mock('../selectors', () => ({ getTerminalHandle: getTerminalHandleMock }))
|
||||
|
||||
import { ORCHESTRATION_HANDLERS } from './orchestration'
|
||||
import { RuntimeClientError } from '../runtime-client'
|
||||
|
||||
function staleHandleError(): RuntimeClientError {
|
||||
return new RuntimeClientError('terminal_handle_stale', 'terminal_handle_stale')
|
||||
}
|
||||
|
||||
// Queues the stale-handle remint chain shared by coordinator commands:
|
||||
// `terminal.resolveIdentity` answers not-live → resolvePane returns liveHandle → downstream RPC.
|
||||
function stubStaleHandleRemint(liveHandle: string, downstream: unknown): void {
|
||||
callMock
|
||||
.mockResolvedValueOnce(notLiveIdentity())
|
||||
.mockResolvedValueOnce({ result: { terminal: { handle: liveHandle } } })
|
||||
.mockResolvedValueOnce(downstream)
|
||||
}
|
||||
|
||||
// Queues a not-live identity followed by a resolvePane remint that fails with `error`.
|
||||
function stubStaleHandleRemintFailure(error: RuntimeClientError): void {
|
||||
callMock.mockResolvedValueOnce(notLiveIdentity()).mockRejectedValueOnce(error)
|
||||
}
|
||||
|
||||
/** What the runtime answers for a handle whose leaf check reports `terminal_handle_stale`. */
|
||||
function notLiveIdentity(): { result: { identity: { live: false } } } {
|
||||
return { result: { identity: { live: false } } }
|
||||
}
|
||||
|
||||
function liveIdentity(handle: string): { result: { identity: { handle: string; live: true } } } {
|
||||
return { result: { identity: { handle, live: true } } }
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
getTerminalHandleMock.mockReset()
|
||||
if (originalTerminalHandle === undefined) {
|
||||
delete process.env.ORCA_TERMINAL_HANDLE
|
||||
} else {
|
||||
process.env.ORCA_TERMINAL_HANDLE = originalTerminalHandle
|
||||
}
|
||||
if (originalPaneKey === undefined) {
|
||||
delete process.env.ORCA_PANE_KEY
|
||||
} else {
|
||||
process.env.ORCA_PANE_KEY = originalPaneKey
|
||||
}
|
||||
})
|
||||
|
||||
describe('orchestration dispatch coordinator handle', () => {
|
||||
beforeEach(() => {
|
||||
callMock.mockReset()
|
||||
getTerminalHandleMock.mockReset()
|
||||
delete process.env.ORCA_TERMINAL_HANDLE
|
||||
delete process.env.ORCA_PANE_KEY
|
||||
})
|
||||
|
||||
const invokeDispatch = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration dispatch']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
const invokeDispatchShow = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration dispatch-show']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
const invokeRun = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration coordinator-start']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
it('remints a stale coordinator env handle from the caller pane key', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
stubStaleHandleRemint('term_live_coord', {
|
||||
result: { dispatch: { id: 'ctx_1', task_id: 'task_1', status: 'dispatched' } }
|
||||
})
|
||||
getTerminalHandleMock.mockRejectedValue(new Error('active terminal fallback is unsafe'))
|
||||
|
||||
await invokeDispatch(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['to', 'term_worker'],
|
||||
['inject', true]
|
||||
])
|
||||
)
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_coord:leaf_coord'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenNthCalledWith(3, 'orchestration.dispatch', {
|
||||
task: 'task_1',
|
||||
to: 'term_worker',
|
||||
from: 'term_live_coord',
|
||||
inject: true,
|
||||
dryRun: undefined,
|
||||
returnPreamble: undefined,
|
||||
devMode: false
|
||||
})
|
||||
})
|
||||
|
||||
it('rejects stale coordinator env handles when the caller pane cannot be proven', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
callMock.mockRejectedValueOnce(staleHandleError())
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeDispatch(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['to', 'term_worker']
|
||||
])
|
||||
)
|
||||
).rejects.toMatchObject({
|
||||
code: 'no_active_sender_terminal'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('propagates unexpected caller pane remint failures for coordinator commands', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
stubStaleHandleRemintFailure(
|
||||
new RuntimeClientError('runtime_unavailable', 'runtime_unavailable')
|
||||
)
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeDispatch(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['to', 'term_worker']
|
||||
])
|
||||
)
|
||||
).rejects.toMatchObject({
|
||||
code: 'runtime_unavailable'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_coord:leaf_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenCalledTimes(2)
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('uses a live coordinator handle for dispatch-show preamble previews', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
stubStaleHandleRemint('term_live_coord', {
|
||||
result: { dispatch: null, preamble: 'preamble' }
|
||||
})
|
||||
getTerminalHandleMock.mockRejectedValue(new Error('active terminal fallback is unsafe'))
|
||||
|
||||
await invokeDispatchShow(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['preamble', true]
|
||||
])
|
||||
)
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_coord:leaf_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(3, 'orchestration.dispatchShow', {
|
||||
task: 'task_1',
|
||||
preamble: true,
|
||||
from: 'term_live_coord',
|
||||
devMode: false
|
||||
})
|
||||
})
|
||||
|
||||
it('retires the legacy coordinator command without runtime effects', async () => {
|
||||
await expect(
|
||||
invokeRun(new Map<string, string | boolean>([['spec', 'run the plan']]))
|
||||
).rejects.toMatchObject({
|
||||
code: 'orchestration_migration_required',
|
||||
data: {
|
||||
reason: 'command_retired',
|
||||
effectsApplied: false,
|
||||
nextCommandArgs: ['skills', 'get', 'orchestration', '--full']
|
||||
}
|
||||
})
|
||||
expect(callMock).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('orchestration task-create caller handle', () => {
|
||||
beforeEach(() => {
|
||||
callMock.mockReset()
|
||||
getTerminalHandleMock.mockReset()
|
||||
delete process.env.ORCA_TERMINAL_HANDLE
|
||||
delete process.env.ORCA_PANE_KEY
|
||||
})
|
||||
|
||||
const invokeTaskCreate = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration task-create']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
it('records a live env terminal handle as task creator', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_creator'
|
||||
callMock
|
||||
.mockResolvedValueOnce(liveIdentity('term_creator'))
|
||||
.mockResolvedValueOnce({ result: { task: { id: 'task_1', status: 'ready' } } })
|
||||
|
||||
await invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_creator'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'orchestration.taskCreate', {
|
||||
spec: 'do work',
|
||||
taskTitle: undefined,
|
||||
displayName: undefined,
|
||||
deps: undefined,
|
||||
parent: undefined,
|
||||
run: undefined,
|
||||
callerTerminalHandle: 'term_creator'
|
||||
})
|
||||
})
|
||||
|
||||
it('fails closed when a stale task creator handle cannot be reminted', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
callMock.mockRejectedValueOnce(staleHandleError())
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({ code: 'no_active_sender_terminal' })
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('propagates runtime unavailability while proving the bound coordinator', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_creator'
|
||||
callMock.mockRejectedValueOnce(
|
||||
new RuntimeClientError('runtime_unavailable', 'runtime_unavailable')
|
||||
)
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({ code: 'runtime_unavailable' })
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_creator'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('propagates runtime unavailability while reminting the bound coordinator', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_creator:leaf_creator'
|
||||
stubStaleHandleRemintFailure(
|
||||
new RuntimeClientError('runtime_unavailable', 'runtime_unavailable')
|
||||
)
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({ code: 'runtime_unavailable' })
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_creator:leaf_creator'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('propagates unexpected caller pane remint failures for task creation', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_creator:leaf_creator'
|
||||
stubStaleHandleRemintFailure(new RuntimeClientError('permission_denied', 'denied'))
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({
|
||||
code: 'permission_denied'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_creator:leaf_creator'
|
||||
})
|
||||
expect(callMock).toHaveBeenCalledTimes(2)
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('propagates unexpected env handle validation failures', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_creator'
|
||||
callMock.mockRejectedValueOnce(new RuntimeClientError('permission_denied', 'denied'))
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({
|
||||
code: 'permission_denied'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('remints a stale task creator env handle from the caller pane key', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_creator:leaf_creator'
|
||||
stubStaleHandleRemint('term_live', {
|
||||
result: { task: { id: 'task_1', status: 'ready' } }
|
||||
})
|
||||
getTerminalHandleMock.mockRejectedValue(new Error('active terminal fallback is unsafe'))
|
||||
|
||||
await invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.resolveIdentity', {
|
||||
terminal: 'term_stale'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_creator:leaf_creator'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenNthCalledWith(3, 'orchestration.taskCreate', {
|
||||
spec: 'do work',
|
||||
taskTitle: undefined,
|
||||
displayName: undefined,
|
||||
deps: undefined,
|
||||
parent: undefined,
|
||||
run: undefined,
|
||||
callerTerminalHandle: 'term_live'
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -72,7 +72,7 @@ describe('orchestration gate commands carry caller identity', () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_coord'
|
||||
queueFixtures(
|
||||
callMock,
|
||||
okFixture('req_show', { terminal: { handle: 'term_coord' } }),
|
||||
okFixture('req_identity', { identity: { handle: 'term_coord', live: true } }),
|
||||
okFixture('req_gate', { gate: { id: 'gate_1', task_id: 'task_1', status: 'pending' } })
|
||||
)
|
||||
|
||||
@@ -91,8 +91,8 @@ describe('orchestration gate commands carry caller identity', () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
callMock.mockImplementation(async (method: string) => {
|
||||
if (method === 'terminal.show') {
|
||||
throw new RuntimeClientError('terminal_handle_stale', 'stale')
|
||||
if (method === 'terminal.resolveIdentity') {
|
||||
return okFixture('req_identity', { identity: { handle: 'term_stale', live: false } })
|
||||
}
|
||||
if (method === 'terminal.resolvePane') {
|
||||
return okFixture('req_pane', { terminal: { handle: 'term_live' } })
|
||||
@@ -146,7 +146,7 @@ describe('orchestration gate commands carry caller identity', () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_coord'
|
||||
queueFixtures(
|
||||
callMock,
|
||||
okFixture('req_show', { terminal: { handle: 'term_coord' } }),
|
||||
okFixture('req_identity', { identity: { handle: 'term_coord', live: true } }),
|
||||
okFixture('req_list', { gates: [], count: 0 })
|
||||
)
|
||||
|
||||
@@ -199,7 +199,9 @@ describe('orchestration gate commands carry caller identity', () => {
|
||||
it('reports idempotent recovery when a mutation connection drops', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_coord'
|
||||
callMock
|
||||
.mockResolvedValueOnce(okFixture('req_show', { terminal: { handle: 'term_coord' } }))
|
||||
.mockResolvedValueOnce(
|
||||
okFixture('req_identity', { identity: { handle: 'term_coord', live: true } })
|
||||
)
|
||||
.mockRejectedValueOnce(
|
||||
new RuntimeClientError(
|
||||
'runtime_unavailable',
|
||||
|
||||
@@ -27,7 +27,7 @@ describe('orchestration task-create CLI mapping', () => {
|
||||
|
||||
it('passes PowerShell-stripped deps through to the runtime', async () => {
|
||||
callMock
|
||||
.mockResolvedValueOnce({ result: { terminal: { handle: 'term_creator' } } })
|
||||
.mockResolvedValueOnce({ result: { identity: { handle: 'term_creator', live: true } } })
|
||||
.mockResolvedValueOnce({ result: { task: { id: 'task_2', status: 'pending' } } })
|
||||
|
||||
await ORCHESTRATION_HANDLERS['orchestration task-create']({
|
||||
|
||||
@@ -16,24 +16,6 @@ import { ORCHESTRATION_HANDLERS } from './orchestration'
|
||||
import { RuntimeClientError } from '../runtime-client'
|
||||
import { printResult } from '../format'
|
||||
|
||||
function staleHandleError(): RuntimeClientError {
|
||||
return new RuntimeClientError('terminal_handle_stale', 'terminal_handle_stale')
|
||||
}
|
||||
|
||||
// Queues the stale-handle remint chain shared by coordinator commands:
|
||||
// stale terminal.show → resolvePane returns liveHandle → downstream RPC result.
|
||||
function stubStaleHandleRemint(liveHandle: string, downstream: unknown): void {
|
||||
callMock
|
||||
.mockRejectedValueOnce(staleHandleError())
|
||||
.mockResolvedValueOnce({ result: { terminal: { handle: liveHandle } } })
|
||||
.mockResolvedValueOnce(downstream)
|
||||
}
|
||||
|
||||
// Queues a stale terminal.show followed by a resolvePane remint that fails with `error`.
|
||||
function stubStaleHandleRemintFailure(error: RuntimeClientError): void {
|
||||
callMock.mockRejectedValueOnce(staleHandleError()).mockRejectedValueOnce(error)
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
getTerminalHandleMock.mockReset()
|
||||
if (originalTerminalHandle === undefined) {
|
||||
@@ -329,309 +311,6 @@ describe('orchestration send structured payload flags', () => {
|
||||
)
|
||||
})
|
||||
|
||||
describe('orchestration dispatch coordinator handle', () => {
|
||||
beforeEach(() => {
|
||||
callMock.mockReset()
|
||||
getTerminalHandleMock.mockReset()
|
||||
delete process.env.ORCA_TERMINAL_HANDLE
|
||||
delete process.env.ORCA_PANE_KEY
|
||||
})
|
||||
|
||||
const invokeDispatch = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration dispatch']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
const invokeDispatchShow = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration dispatch-show']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
const invokeRun = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration coordinator-start']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
it('remints a stale coordinator env handle from the caller pane key', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
stubStaleHandleRemint('term_live_coord', {
|
||||
result: { dispatch: { id: 'ctx_1', task_id: 'task_1', status: 'dispatched' } }
|
||||
})
|
||||
getTerminalHandleMock.mockRejectedValue(new Error('active terminal fallback is unsafe'))
|
||||
|
||||
await invokeDispatch(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['to', 'term_worker'],
|
||||
['inject', true]
|
||||
])
|
||||
)
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', {
|
||||
terminal: 'term_stale_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_coord:leaf_coord'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenNthCalledWith(3, 'orchestration.dispatch', {
|
||||
task: 'task_1',
|
||||
to: 'term_worker',
|
||||
from: 'term_live_coord',
|
||||
inject: true,
|
||||
dryRun: undefined,
|
||||
returnPreamble: undefined,
|
||||
devMode: false
|
||||
})
|
||||
})
|
||||
|
||||
it('rejects stale coordinator env handles when the caller pane cannot be proven', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
callMock.mockRejectedValueOnce(staleHandleError())
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeDispatch(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['to', 'term_worker']
|
||||
])
|
||||
)
|
||||
).rejects.toMatchObject({
|
||||
code: 'no_active_sender_terminal'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('propagates unexpected caller pane remint failures for coordinator commands', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
stubStaleHandleRemintFailure(
|
||||
new RuntimeClientError('runtime_unavailable', 'runtime_unavailable')
|
||||
)
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeDispatch(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['to', 'term_worker']
|
||||
])
|
||||
)
|
||||
).rejects.toMatchObject({
|
||||
code: 'runtime_unavailable'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', {
|
||||
terminal: 'term_stale_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_coord:leaf_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenCalledTimes(2)
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('uses a live coordinator handle for dispatch-show preamble previews', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale_coord'
|
||||
process.env.ORCA_PANE_KEY = 'tab_coord:leaf_coord'
|
||||
stubStaleHandleRemint('term_live_coord', {
|
||||
result: { dispatch: null, preamble: 'preamble' }
|
||||
})
|
||||
getTerminalHandleMock.mockRejectedValue(new Error('active terminal fallback is unsafe'))
|
||||
|
||||
await invokeDispatchShow(
|
||||
new Map<string, string | boolean>([
|
||||
['task', 'task_1'],
|
||||
['preamble', true]
|
||||
])
|
||||
)
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', {
|
||||
terminal: 'term_stale_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_coord:leaf_coord'
|
||||
})
|
||||
expect(callMock).toHaveBeenNthCalledWith(3, 'orchestration.dispatchShow', {
|
||||
task: 'task_1',
|
||||
preamble: true,
|
||||
from: 'term_live_coord',
|
||||
devMode: false
|
||||
})
|
||||
})
|
||||
|
||||
it('retires the legacy coordinator command without runtime effects', async () => {
|
||||
await expect(
|
||||
invokeRun(new Map<string, string | boolean>([['spec', 'run the plan']]))
|
||||
).rejects.toMatchObject({
|
||||
code: 'orchestration_migration_required',
|
||||
data: {
|
||||
reason: 'command_retired',
|
||||
effectsApplied: false,
|
||||
nextCommandArgs: ['skills', 'get', 'orchestration', '--full']
|
||||
}
|
||||
})
|
||||
expect(callMock).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('orchestration task-create caller handle', () => {
|
||||
beforeEach(() => {
|
||||
callMock.mockReset()
|
||||
getTerminalHandleMock.mockReset()
|
||||
delete process.env.ORCA_TERMINAL_HANDLE
|
||||
delete process.env.ORCA_PANE_KEY
|
||||
})
|
||||
|
||||
const invokeTaskCreate = (flags: Map<string, string | boolean>) =>
|
||||
ORCHESTRATION_HANDLERS['orchestration task-create']({
|
||||
flags,
|
||||
client: { call: callMock },
|
||||
cwd: '/tmp/repo',
|
||||
json: true
|
||||
} as never)
|
||||
|
||||
it('records a live env terminal handle as task creator', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_creator'
|
||||
callMock
|
||||
.mockResolvedValueOnce({ result: { terminal: { handle: 'term_creator' } } })
|
||||
.mockResolvedValueOnce({ result: { task: { id: 'task_1', status: 'ready' } } })
|
||||
|
||||
await invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', { terminal: 'term_creator' })
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'orchestration.taskCreate', {
|
||||
spec: 'do work',
|
||||
taskTitle: undefined,
|
||||
displayName: undefined,
|
||||
deps: undefined,
|
||||
parent: undefined,
|
||||
run: undefined,
|
||||
callerTerminalHandle: 'term_creator'
|
||||
})
|
||||
})
|
||||
|
||||
it('fails closed when a stale task creator handle cannot be reminted', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
callMock.mockRejectedValueOnce(staleHandleError())
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({ code: 'no_active_sender_terminal' })
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', { terminal: 'term_stale' })
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('propagates runtime unavailability while proving the bound coordinator', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_creator'
|
||||
callMock.mockRejectedValueOnce(
|
||||
new RuntimeClientError('runtime_unavailable', 'runtime_unavailable')
|
||||
)
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({ code: 'runtime_unavailable' })
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', { terminal: 'term_creator' })
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('propagates runtime unavailability while reminting the bound coordinator', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_creator:leaf_creator'
|
||||
stubStaleHandleRemintFailure(
|
||||
new RuntimeClientError('runtime_unavailable', 'runtime_unavailable')
|
||||
)
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({ code: 'runtime_unavailable' })
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', { terminal: 'term_stale' })
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_creator:leaf_creator'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('propagates unexpected caller pane remint failures for task creation', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_creator:leaf_creator'
|
||||
stubStaleHandleRemintFailure(new RuntimeClientError('permission_denied', 'denied'))
|
||||
getTerminalHandleMock.mockResolvedValue('term_wrong_active')
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({
|
||||
code: 'permission_denied'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', { terminal: 'term_stale' })
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_creator:leaf_creator'
|
||||
})
|
||||
expect(callMock).toHaveBeenCalledTimes(2)
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('propagates unexpected env handle validation failures', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_creator'
|
||||
callMock.mockRejectedValueOnce(new RuntimeClientError('permission_denied', 'denied'))
|
||||
|
||||
await expect(
|
||||
invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
).rejects.toMatchObject({
|
||||
code: 'permission_denied'
|
||||
})
|
||||
|
||||
expect(callMock).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('remints a stale task creator env handle from the caller pane key', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
process.env.ORCA_PANE_KEY = 'tab_creator:leaf_creator'
|
||||
stubStaleHandleRemint('term_live', {
|
||||
result: { task: { id: 'task_1', status: 'ready' } }
|
||||
})
|
||||
getTerminalHandleMock.mockRejectedValue(new Error('active terminal fallback is unsafe'))
|
||||
|
||||
await invokeTaskCreate(new Map<string, string | boolean>([['spec', 'do work']]))
|
||||
|
||||
expect(callMock).toHaveBeenNthCalledWith(1, 'terminal.show', { terminal: 'term_stale' })
|
||||
expect(callMock).toHaveBeenNthCalledWith(2, 'terminal.resolvePane', {
|
||||
paneKey: 'tab_creator:leaf_creator'
|
||||
})
|
||||
expect(getTerminalHandleMock).not.toHaveBeenCalled()
|
||||
expect(callMock).toHaveBeenNthCalledWith(3, 'orchestration.taskCreate', {
|
||||
spec: 'do work',
|
||||
taskTitle: undefined,
|
||||
displayName: undefined,
|
||||
deps: undefined,
|
||||
parent: undefined,
|
||||
run: undefined,
|
||||
callerTerminalHandle: 'term_live'
|
||||
})
|
||||
})
|
||||
})
|
||||
describe('orchestration timeout flag validation', () => {
|
||||
const invalidTimeoutValues: [string, string | boolean][] = [
|
||||
['missing', true],
|
||||
|
||||
@@ -35,7 +35,38 @@ export async function resolveOrchestrationTerminalHandle(
|
||||
return await getTerminalHandle(flags, cwd, client)
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether the handle this process was born with still names a live identity.
|
||||
*
|
||||
* `terminal.resolveIdentity`, never `terminal.show`: `show` is a PTY verb, so it missed for a
|
||||
* structured worker and reported `terminal_handle_stale` for a handle that was perfectly live —
|
||||
* which then failed every coordinator verb, because the pane remint below needs an `ORCA_PANE_KEY`
|
||||
* a structured child deliberately does not carry.
|
||||
*/
|
||||
async function isLiveTerminalHandle(handle: string, client: RuntimeClient): Promise<boolean> {
|
||||
try {
|
||||
const response = await client.call<{ identity?: { live?: boolean } }>(
|
||||
'terminal.resolveIdentity',
|
||||
{ terminal: handle }
|
||||
)
|
||||
const live = response.result?.identity?.live
|
||||
// An unrecognised shape is an older host answering something else, not a dead handle.
|
||||
return typeof live === 'boolean' ? live : await showResolvesTerminalHandle(handle, client)
|
||||
} catch (err) {
|
||||
if (isStaleTerminalIdentityError(err)) {
|
||||
return false
|
||||
}
|
||||
if (getClientErrorCode(err) === 'method_not_found') {
|
||||
// Clients and remote hosts update independently, so a host that predates the identity probe
|
||||
// is the normal mixed-version state. Fall back to what it does have — which is correct for
|
||||
// that host, because a host without the probe also has no structured workers to miss.
|
||||
return await showResolvesTerminalHandle(handle, client)
|
||||
}
|
||||
throw err
|
||||
}
|
||||
}
|
||||
|
||||
async function showResolvesTerminalHandle(handle: string, client: RuntimeClient): Promise<boolean> {
|
||||
try {
|
||||
await client.call('terminal.show', { terminal: handle })
|
||||
return true
|
||||
|
||||
@@ -0,0 +1,151 @@
|
||||
/**
|
||||
* The env-handle path, with NO `--from`.
|
||||
*
|
||||
* Every other orchestration CLI test passes `--from term_coord` explicitly, so the resolver a real
|
||||
* worker actually goes through — `ORCA_TERMINAL_HANDLE` plus `validateEnvHandle` — was never
|
||||
* exercised. That is why twelve coordinator verbs could fail for a structured worker while the
|
||||
* whole suite stayed green, and why the worker's own preamble (which tells it to run these with no
|
||||
* `--from`) failed on its first line.
|
||||
*/
|
||||
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
|
||||
const {
|
||||
callMock,
|
||||
runtimeClientConstructorMock,
|
||||
serveOrcaAppMock,
|
||||
getDefaultUserDataPathMock,
|
||||
addEnvironmentFromPairingCodeMock,
|
||||
listEnvironmentsMock,
|
||||
spawnMock
|
||||
} = vi.hoisted(() => ({
|
||||
callMock: vi.fn(),
|
||||
runtimeClientConstructorMock: vi.fn(),
|
||||
serveOrcaAppMock: vi.fn(),
|
||||
getDefaultUserDataPathMock: vi.fn(() => '/tmp/orca-user-data'),
|
||||
addEnvironmentFromPairingCodeMock: vi.fn(),
|
||||
listEnvironmentsMock: vi.fn(),
|
||||
spawnMock: vi.fn()
|
||||
}))
|
||||
|
||||
vi.mock('./runtime-client', async () => {
|
||||
const { createRuntimeClientModuleMock } = await import('./index-test-harness.js')
|
||||
return createRuntimeClientModuleMock({
|
||||
callMock,
|
||||
runtimeClientConstructorMock,
|
||||
serveOrcaAppMock,
|
||||
getDefaultUserDataPathMock
|
||||
})
|
||||
})
|
||||
|
||||
vi.mock('./runtime/environments', () => ({
|
||||
addEnvironmentFromPairingCode: addEnvironmentFromPairingCodeMock,
|
||||
listEnvironments: listEnvironmentsMock,
|
||||
removeEnvironment: vi.fn(),
|
||||
resolveEnvironment: vi.fn()
|
||||
}))
|
||||
|
||||
vi.mock('child_process', async () => {
|
||||
const { createChildProcessModuleMock } = await import('./index-test-harness.js')
|
||||
return createChildProcessModuleMock(spawnMock)
|
||||
})
|
||||
|
||||
import { main } from './index'
|
||||
import { useWorktreeAwarenessEnvironment } from './index-test-harness'
|
||||
|
||||
const STRUCTURED_HANDLE = 'structworker_a1b2c3d4-e5f6-4a7b-8c9d-0e1f2a3b4c5d'
|
||||
|
||||
/**
|
||||
* Every command below is one of the twelve that resolve their sender through
|
||||
* `resolveCoordinatorTerminalHandle`. All twelve share ONE resolver and one liveness probe, so the
|
||||
* set is sufficient if it covers the distinct shapes that reach it: a read, a mutation, a
|
||||
* run-scoped verb, a dispatch verb and a worker-start. A per-verb sweep would pin the argv parsing
|
||||
* of twelve handlers and still tell us nothing more about the seam that actually broke.
|
||||
*/
|
||||
const SENDER_VERBS: { argv: string[]; method: string }[] = [
|
||||
{ argv: ['orchestration', 'run-current'], method: 'orchestration.runCurrent' },
|
||||
{ argv: ['orchestration', 'run-create', '--objective', 'x'], method: 'orchestration.runCreate' },
|
||||
{ argv: ['orchestration', 'task-list'], method: 'orchestration.taskList' },
|
||||
{ argv: ['orchestration', 'gate-list'], method: 'orchestration.gateList' },
|
||||
{ argv: ['orchestration', 'dispatch-show', '--task', 't1'], method: 'orchestration.dispatchShow' }
|
||||
]
|
||||
|
||||
describe('a structured worker running orchestration commands as itself', () => {
|
||||
useWorktreeAwarenessEnvironment({
|
||||
callMock,
|
||||
serveOrcaAppMock,
|
||||
getDefaultUserDataPathMock,
|
||||
addEnvironmentFromPairingCodeMock,
|
||||
listEnvironmentsMock,
|
||||
spawnMock
|
||||
})
|
||||
|
||||
function answerCalls(): void {
|
||||
callMock.mockImplementation(async (method: string) => {
|
||||
if (method === 'terminal.resolveIdentity') {
|
||||
return {
|
||||
id: 'req',
|
||||
ok: true,
|
||||
result: { identity: { handle: STRUCTURED_HANDLE, live: true } },
|
||||
_meta: { runtimeId: 'runtime-1' }
|
||||
}
|
||||
}
|
||||
return { id: 'req', ok: true, result: {}, _meta: { runtimeId: 'runtime-1' } }
|
||||
})
|
||||
}
|
||||
|
||||
it.each(SENDER_VERBS)(
|
||||
'resolves its own identity for $method with no --from',
|
||||
async ({ argv, method }) => {
|
||||
// The defect this pins: the sender resolver validated the env handle with `terminal.show`, a
|
||||
// PTY verb that misses for a structured worker and answers `terminal_handle_stale`. The pane
|
||||
// remint that would have recovered it needs `ORCA_PANE_KEY`, which a structured child
|
||||
// deliberately does not carry, so the command died on `no_active_sender_terminal`.
|
||||
process.env.ORCA_TERMINAL_HANDLE = STRUCTURED_HANDLE
|
||||
answerCalls()
|
||||
vi.spyOn(console, 'log').mockImplementation(() => {})
|
||||
await expect(main(argv)).resolves.not.toThrow()
|
||||
const called = callMock.mock.calls.map((call) => call[0] as string)
|
||||
expect(called).toContain(method)
|
||||
// Never through `terminal.show`: teaching that verb structured handles would hand every
|
||||
// public terminal verb something that looks writable and is not.
|
||||
expect(called).not.toContain('terminal.show')
|
||||
}
|
||||
)
|
||||
|
||||
it('sends the structured handle as the sender, not a guessed sibling', async () => {
|
||||
process.env.ORCA_TERMINAL_HANDLE = STRUCTURED_HANDLE
|
||||
answerCalls()
|
||||
vi.spyOn(console, 'log').mockImplementation(() => {})
|
||||
await main(['orchestration', 'run-create', '--objective', 'x'])
|
||||
const create = callMock.mock.calls.find((call) => call[0] === 'orchestration.runCreate')
|
||||
expect((create?.[1] as { from?: string } | undefined)?.from).toBe(STRUCTURED_HANDLE)
|
||||
})
|
||||
|
||||
it('still refuses a handle the runtime reports dead, with no pane key to remint from', async () => {
|
||||
// The invariant the fix must not break: a stale `ORCA_TERMINAL_HANDLE` in a long-lived shell
|
||||
// must keep failing rather than being baked into a coordinator preamble.
|
||||
process.env.ORCA_TERMINAL_HANDLE = 'term_stale'
|
||||
callMock.mockImplementation(async (method: string) => {
|
||||
if (method === 'terminal.resolveIdentity') {
|
||||
return {
|
||||
id: 'req',
|
||||
ok: true,
|
||||
result: { identity: { handle: 'term_stale', live: false } },
|
||||
_meta: { runtimeId: 'runtime-1' }
|
||||
}
|
||||
}
|
||||
return { id: 'req', ok: true, result: {}, _meta: { runtimeId: 'runtime-1' } }
|
||||
})
|
||||
const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => {})
|
||||
const priorExitCode = process.exitCode
|
||||
await main(['orchestration', 'run-create', '--objective', 'x'])
|
||||
expect(process.exitCode).toBe(1)
|
||||
expect(errorSpy.mock.calls.flat().join(' ')).toMatch(
|
||||
/no_active_sender_terminal|sender terminal/i
|
||||
)
|
||||
expect(callMock.mock.calls.map((call) => call[0])).not.toContain('orchestration.runCreate')
|
||||
process.exitCode = priorExitCode
|
||||
errorSpy.mockRestore()
|
||||
})
|
||||
})
|
||||
@@ -8,6 +8,10 @@ import { agentSessionPtyWriteGate } from './agent-session-pty-write-gate'
|
||||
import { resolveStructuredWorkerAuthority } from './structured-worker-authority'
|
||||
import { isSettledNativeOwner } from './orchestration/structured-session-pointer-delivery'
|
||||
import type { StructuredPointerTarget } from './orchestration/structured-mailbox-pointer-delivery'
|
||||
import {
|
||||
resolveTerminalIdentityFromProbes,
|
||||
type RuntimeTerminalIdentity
|
||||
} from './terminal-identity-probe'
|
||||
|
||||
export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneMobileSessionTabGroupLayout {
|
||||
protected getPtyRecordForPaneKey(paneKey: string): RuntimePtyWorktreeRecord | null {
|
||||
@@ -144,6 +148,23 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The identity seam: whether this handle still names a live agent identity, in either lane.
|
||||
*
|
||||
* Read-only by construction — a handle and a boolean — so it can serve the CLI's sender
|
||||
* validation without `terminal.show`'s writable-looking pane payload.
|
||||
*/
|
||||
resolveTerminalIdentity(handle: string): RuntimeTerminalIdentity {
|
||||
return resolveTerminalIdentityFromProbes(handle, {
|
||||
isLiveStructuredWorker: () =>
|
||||
Boolean(resolveStructuredWorkerAuthority(handle, this._orchestrationDb)),
|
||||
hasLivePty: () => Boolean(this.getLivePtyForHandle(handle)),
|
||||
assertLiveLeaf: () => {
|
||||
this.getLiveLeafForHandle(handle)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
deliverPendingMessagesForHandle(handle: string, reservedTypes?: ReadonlySet<string>): void {
|
||||
this.orchestrationMailboxNotifications.deliverForHandle(handle, reservedTypes)
|
||||
}
|
||||
@@ -161,12 +182,16 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
/**
|
||||
* The session a mailbox must be nudged through, or null when a live PTY can take the bytes.
|
||||
*
|
||||
* Two mailbox shapes reach a structured worker: its `dispatch:` address, and its own bearer
|
||||
* handle, which is how agents mail each other outside a dispatch. `run:` mail is coordinator
|
||||
* mail and stays out of scope: a coordinator blocks in `check --wait`, where a waiter preempts
|
||||
* pointer delivery anyway.
|
||||
* All THREE address forms a structured session can own resolve here — its `dispatch:` address,
|
||||
* its `run:` mailbox when it coordinates, and its own bearer handle for peer mail outside a
|
||||
* dispatch. `run:` was the one that fell in a hole: the PTY lane declines because the owner is
|
||||
* structured, and this lane used to decline anything that was not `dispatch:`, so each half
|
||||
* believed the other owned it and a structured coordinator was never nudged.
|
||||
*/
|
||||
protected resolveStructuredMailboxTarget(mailboxHandle: string): StructuredPointerTarget | null {
|
||||
if (mailboxHandle.startsWith('run:')) {
|
||||
return this.resolveStructuredCoordinatorMailboxTarget(mailboxHandle.slice('run:'.length))
|
||||
}
|
||||
if (!mailboxHandle.startsWith('dispatch:')) {
|
||||
return this.resolveStructuredWorkerDirectMailboxTarget(mailboxHandle)
|
||||
}
|
||||
@@ -182,6 +207,25 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
return this.resolveAdoptedStructuredMailboxTarget(assignee, dispatchId)
|
||||
}
|
||||
|
||||
/**
|
||||
* A Run's own mailbox, when the coordinator holding it is a structured session.
|
||||
*
|
||||
* A structured coordinator does NOT block in `check --wait` the way a PTY one does — it is a
|
||||
* chat session, and its turn ends — so the waiter that used to preempt pointer delivery is not
|
||||
* there to cover for the missing nudge. Session-scoped: a coordinator's run mailbox has no
|
||||
* dispatch, and needs none, since the ledger bucket is all a dispatch id ever supplied.
|
||||
*/
|
||||
protected resolveStructuredCoordinatorMailboxTarget(
|
||||
runId: string
|
||||
): StructuredPointerTarget | null {
|
||||
const coordinator = this._orchestrationDb?.getRun?.(runId)?.coordinator_handle
|
||||
if (!coordinator) {
|
||||
return null
|
||||
}
|
||||
const identity = resolveStructuredWorkerAuthority(coordinator, this._orchestrationDb)?.identity
|
||||
return identity ? { sessionId: identity.sessionId, dispatchId: null } : null
|
||||
}
|
||||
|
||||
/**
|
||||
* Direct peer mail, addressed to the worker's own handle rather than to a dispatch.
|
||||
*
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
import { OrchestrationStructuredMailboxPointerDelivery } from './orchestration/structured-mailbox-pointer-delivery'
|
||||
import { createStructuredMailboxPointerHost } from './orchestration/structured-mailbox-pointer-host'
|
||||
import { isStructuredWorkerHandle } from './structured-worker-identity'
|
||||
import { resolveStructuredWorkerAuthority } from './structured-worker-authority'
|
||||
import { OrcaRuntimeWithRuntimeId } from './orca-runtime-runtime-id'
|
||||
import { RuntimeTerminalAgentPresence } from './runtime-terminal-agent-presence'
|
||||
import type { RuntimeNotifier } from './runtime-notifier-contract'
|
||||
@@ -40,6 +41,8 @@ export class OrcaRuntimeWithStopRequestedPtyIds extends OrcaRuntimeWithRuntimeId
|
||||
protected readonly ptyExitListenersByPtyId = new Map<string, Set<() => void>>()
|
||||
|
||||
protected readonly terminalAgentPresence = new RuntimeTerminalAgentPresence({
|
||||
isLiveStructuredAgent: (handle) =>
|
||||
Boolean(resolveStructuredWorkerAuthority(handle, this._orchestrationDb)),
|
||||
getLivePty: (handle) => this.getLivePtyForHandle(handle)?.pty ?? null,
|
||||
getLiveLeaf: (handle) => this.getLiveLeafForHandle(handle).leaf,
|
||||
getPrimaryLeaf: (ptyId) => this.getLeavesForPty(ptyId)[0] ?? null,
|
||||
|
||||
@@ -7,8 +7,9 @@
|
||||
* everything below it is different — the nudge is a session turn, the idle edge is the journal,
|
||||
* and only an `accepted` dispatch may consume mail.
|
||||
*
|
||||
* Coordinators are deliberately out of scope: `run:` mail routes through a coordinator handle, and
|
||||
* a coordinator blocks in `check --wait`, where a waiter preempts pointer delivery anyway.
|
||||
* Coordinators are in scope here, unlike the PTY lane's reasoning: a PTY coordinator blocks in
|
||||
* `check --wait`, where a waiter preempts pointer delivery, but a structured coordinator is a chat
|
||||
* session whose turn ends — so nothing else would ever wake it for its own `run:` mail.
|
||||
*/
|
||||
|
||||
import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types'
|
||||
@@ -159,11 +160,18 @@ export class OrchestrationStructuredMailboxPointerDelivery<
|
||||
if (!db || this.inFlight.has(mailboxHandle)) {
|
||||
return
|
||||
}
|
||||
// No `hasOutstandingRunDelivery` gate, unlike the PTY lane: there it guards a COORDINATOR's
|
||||
// own `run:` mailbox against re-notifying a batch already handed to that coordinator. This
|
||||
// lane never resolves a `run:` mailbox, and a delivery row exists only for a `run:` address,
|
||||
// so the run's outstanding delivery belongs to the coordinator that is replying — gating on
|
||||
// it suppresses exactly the nudges a coordinator sends its workers.
|
||||
// The `hasOutstandingRunDelivery` gate applies to a `run:` mailbox and ONLY to one, exactly as
|
||||
// in the PTY lane: it guards a coordinator's own mailbox against re-notifying a batch already
|
||||
// handed to that coordinator. A delivery row exists only for a `run:` address, so for a
|
||||
// `dispatch:` or bare-handle mailbox the run's outstanding delivery belongs to the coordinator
|
||||
// that is replying — gating there would suppress exactly the nudges a coordinator sends its
|
||||
// workers.
|
||||
if (
|
||||
mailboxHandle.startsWith('run:') &&
|
||||
db.hasOutstandingRunDelivery?.(mailboxHandle.slice('run:'.length))
|
||||
) {
|
||||
return
|
||||
}
|
||||
const unread = selectOrchestrationPointerBatch({
|
||||
db,
|
||||
mailboxHandle,
|
||||
|
||||
@@ -55,15 +55,16 @@ function registerWorker(): string {
|
||||
return handle
|
||||
}
|
||||
|
||||
function probe(activeDispatch: { id: string } | undefined) {
|
||||
function probe(activeDispatch: { id: string } | undefined, run?: { coordinator_handle: string }) {
|
||||
const findActiveDispatchForAssignee = vi.fn(() => activeDispatch)
|
||||
const getRun = vi.fn(() => run)
|
||||
const instance = Object.assign(Object.create(MailboxTargetProbe.prototype), {
|
||||
_orchestrationDb: { findActiveDispatchForAssignee },
|
||||
_orchestrationDb: { findActiveDispatchForAssignee, getRun },
|
||||
getLiveLeafForHandle: () => {
|
||||
throw new Error('no leaf backs a native-born structured worker')
|
||||
}
|
||||
}) as MailboxTargetProbe
|
||||
return { instance, findActiveDispatchForAssignee }
|
||||
return { instance, findActiveDispatchForAssignee, getRun }
|
||||
}
|
||||
|
||||
describe('the mailbox target for direct peer mail to a structured worker', () => {
|
||||
@@ -109,11 +110,43 @@ describe('the mailbox target for direct peer mail to a structured worker', () =>
|
||||
expect(probe({ id: 'd1' }).instance.probeResolveTarget(handle)).toBeNull()
|
||||
})
|
||||
|
||||
it('claims neither a PTY handle nor a coordinator mailbox', () => {
|
||||
it('claims a PTY handle for neither lane', () => {
|
||||
installRecord({ runtimeKind: 'native', claimStatus: 'live' })
|
||||
const { instance, findActiveDispatchForAssignee } = probe({ id: 'd1' })
|
||||
expect(instance.probeResolveTarget('term_abc')).toBeNull()
|
||||
expect(instance.probeResolveTarget('run:run_1')).toBeNull()
|
||||
expect(findActiveDispatchForAssignee).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('the mailbox target for a Run whose coordinator is structured', () => {
|
||||
beforeEach(() => {
|
||||
structuredWorkerIdentities.clear()
|
||||
hostRef.current = null
|
||||
})
|
||||
|
||||
it('owns the run mailbox, which neither lane used to claim', () => {
|
||||
// The defect this pins: the PTY lane declines because the owner is structured, and this lane
|
||||
// used to decline anything that was not `dispatch:`. Each half believed the other owned it, so
|
||||
// a structured coordinator was never nudged for its own Run mail and nothing logged. A PTY
|
||||
// coordinator is covered by blocking in `check --wait`; a chat session's turn just ends.
|
||||
const handle = registerWorker()
|
||||
installRecord({ runtimeKind: 'native', claimStatus: 'live' })
|
||||
expect(
|
||||
probe(undefined, { coordinator_handle: handle }).instance.probeResolveTarget('run:run_1')
|
||||
).toEqual({ sessionId: SESSION_ID, dispatchId: null })
|
||||
})
|
||||
|
||||
it('leaves the run mailbox of a PTY coordinator to the PTY lane', () => {
|
||||
installRecord({ runtimeKind: 'native', claimStatus: 'live' })
|
||||
expect(
|
||||
probe(undefined, { coordinator_handle: 'term_coord' }).instance.probeResolveTarget(
|
||||
'run:run_1'
|
||||
)
|
||||
).toBeNull()
|
||||
})
|
||||
|
||||
it('claims nothing for a run that does not resolve', () => {
|
||||
installRecord({ runtimeKind: 'native', claimStatus: 'live' })
|
||||
expect(probe(undefined).instance.probeResolveTarget('run:run_1')).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -43,8 +43,10 @@ describe('orchestration CLI/runtime boundary', () => {
|
||||
isRemote: false,
|
||||
/** Preserves terminal-handle validation while routing other calls through runtime RPC. */
|
||||
async call<T>(method: string, params?: unknown): Promise<{ result: T }> {
|
||||
if (method === 'terminal.show') {
|
||||
return { result: { terminal: { handle: objectParams(params).terminal } } as T }
|
||||
if (method === 'terminal.resolveIdentity') {
|
||||
return {
|
||||
result: { identity: { handle: objectParams(params).terminal, live: true } } as T
|
||||
}
|
||||
}
|
||||
return { result: (await callRpc(method, objectParams(params))) as T }
|
||||
}
|
||||
|
||||
@@ -9,7 +9,6 @@
|
||||
|
||||
import type { AgentType, NativeChatMessage } from '../../../../shared/native-chat-types'
|
||||
import type { OrchestrationWorkerReadTranscriptResult } from '../../../../shared/orchestration-worker-output'
|
||||
import { getStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import type { OrchestrationDb } from '../../orchestration/db'
|
||||
import { OrchestrationError } from '../../orchestration/orchestration-error'
|
||||
@@ -35,10 +34,10 @@ import {
|
||||
structuredWorkerTerminalState,
|
||||
type StructuredWorkerObservation
|
||||
} from '../../structured-worker-authority'
|
||||
import { retireSettledStructuredWorkerTab } from '../../structured-agent-session-tab-retirement'
|
||||
import type { StructuredWorkerIdentity } from '../../structured-worker-identity'
|
||||
import type { WorkerTerminalReleaseState } from '../../orchestration/worker-terminal-ownership'
|
||||
import { releaseStructuredWorkerSession } from './orchestration-structured-worker-session'
|
||||
import { closeStructuredAgentSessionChild } from '../../structured-agent-session-close'
|
||||
|
||||
export { observeStructuredWorker, type StructuredWorkerObservation }
|
||||
|
||||
@@ -75,43 +74,12 @@ export async function stopStructuredWorker(
|
||||
'forgetStructuredSessionMail' | 'retireStructuredAgentSessionTabFromSnapshot'
|
||||
>
|
||||
): Promise<StructuredWorkerStopOutcome> {
|
||||
const host = getStructuredAgentSessionHost()
|
||||
if (!host) {
|
||||
// Nothing was reached, so nothing was acted on; the receipt must not claim a close.
|
||||
return {
|
||||
stopped: false,
|
||||
closeAttempted: false,
|
||||
reason: 'The structured agent-session host is not installed; no session was closed.'
|
||||
}
|
||||
}
|
||||
// Set only once the close is actually issued: `setSessionTabVisibility` throwing first leaves a
|
||||
// running child, and a receipt that still said `closed_agent_terminal` for it would be the
|
||||
// close-that-never-happened this flag exists to rule out.
|
||||
let closeAttempted = false
|
||||
try {
|
||||
await host.setSessionTabVisibility?.(identity.sessionId, false)
|
||||
closeAttempted = true
|
||||
await host.close(identity.sessionId)
|
||||
} catch (error) {
|
||||
return {
|
||||
stopped: false,
|
||||
closeAttempted,
|
||||
reason: error instanceof Error ? error.message : String(error)
|
||||
}
|
||||
}
|
||||
releaseStructuredWorkerSession(dispatchId, runtime)
|
||||
const after = observeStructuredWorker(identity)
|
||||
if (after.status === 'live') {
|
||||
return {
|
||||
stopped: false,
|
||||
closeAttempted: true,
|
||||
reason: 'The structured session is still attached after close.'
|
||||
}
|
||||
}
|
||||
// Only past the proof, and structurally unable to throw: the worker's chat tab is retired from
|
||||
// the live snapshot, which `setSessionTabVisibility(false)` above does not do.
|
||||
retireSettledStructuredWorkerTab(identity.sessionId, runtime)
|
||||
return { stopped: true, closeAttempted: true }
|
||||
return closeStructuredAgentSessionChild(identity.sessionId, {
|
||||
...(runtime ? { runtime } : {}),
|
||||
// Between the close and the proof, never after: an unsettled close returns early, and a
|
||||
// surviving hold keeps the provider child un-evictable for the life of the app.
|
||||
afterClose: () => releaseStructuredWorkerSession(dispatchId, runtime)
|
||||
})
|
||||
}
|
||||
|
||||
/** The structured half of `worker-read`, or null when a PTY worker owns the dispatch. */
|
||||
|
||||
@@ -13,6 +13,7 @@ const METHOD_CASES: readonly (readonly [string, unknown, boolean])[] = [
|
||||
['terminal.resolvePane', { paneKey: 'pane' }, false],
|
||||
['terminal.recoverPane', { paneKey: 'pane', worktreeId: 'worktree' }, false],
|
||||
['terminal.show', { terminal: 'term' }, false],
|
||||
['terminal.resolveIdentity', { terminal: 'term' }, false],
|
||||
['terminal.read', { terminal: 'term' }, false],
|
||||
['terminal.inspectProcess', { terminal: 'term' }, false],
|
||||
['terminal.isRunningAgent', { terminal: 'term' }, false],
|
||||
@@ -65,11 +66,11 @@ async function invoke(name: string, params: unknown, runtime: Partial<OrcaRuntim
|
||||
|
||||
describe('terminal RPC manifest characterization', () => {
|
||||
it('preserves all method names, order, streaming flags, and parseable minimum inputs', () => {
|
||||
expect(TERMINAL_METHODS).toHaveLength(34)
|
||||
expect(TERMINAL_METHODS).toHaveLength(35)
|
||||
expect(TERMINAL_METHODS.map((method) => [method.name, 'stream' in method])).toEqual(
|
||||
METHOD_CASES.map(([name, _params, stream]) => [name, stream])
|
||||
)
|
||||
expect(new Set(TERMINAL_METHODS.map((method) => method.name)).size).toBe(34)
|
||||
expect(new Set(TERMINAL_METHODS.map((method) => method.name)).size).toBe(35)
|
||||
for (const [name, params] of METHOD_CASES) {
|
||||
expect(() => schemaFor(name).parse(params), name).not.toThrow()
|
||||
}
|
||||
|
||||
@@ -56,6 +56,15 @@ export const TERMINAL_QUERY_METHODS: RpcAnyMethod[] = [
|
||||
terminal: await runtime.showTerminal(params.terminal)
|
||||
})
|
||||
}),
|
||||
defineMethod({
|
||||
// Read-only identity probe. Deliberately NOT `terminal.show`: this one resolves a structured
|
||||
// worker too, and must therefore never hand back anything that looks writable.
|
||||
name: 'terminal.resolveIdentity',
|
||||
params: TerminalHandle,
|
||||
handler: async (params, { runtime }) => ({
|
||||
identity: runtime.resolveTerminalIdentity(params.terminal)
|
||||
})
|
||||
}),
|
||||
defineMethod({
|
||||
name: 'terminal.read',
|
||||
params: TerminalRead,
|
||||
|
||||
@@ -20,6 +20,8 @@ const WRAPPER_RETRY_INTERVAL_MS = 150
|
||||
const WRAPPER_RETRY_TIMEOUT_MS = 6_500
|
||||
|
||||
type RuntimeTerminalAgentPresenceDependencies = {
|
||||
/** A structured agent session of this runtime; it has no pane, so no PTY probe can see it. */
|
||||
isLiveStructuredAgent?(handle: string): boolean
|
||||
getLivePty(handle: string): RuntimePtyWorktreeRecord | null
|
||||
getLiveLeaf(handle: string): RuntimeLeafRecord
|
||||
getPrimaryLeaf(ptyId: string): RuntimeLeafRecord | null
|
||||
@@ -41,6 +43,13 @@ export class RuntimeTerminalAgentPresence {
|
||||
handle: string,
|
||||
options: RuntimeTerminalAgentPresenceOptions = {}
|
||||
): Promise<boolean> {
|
||||
// Before every PTY probe below, because none of them can answer for a session that has no
|
||||
// pane: `getLiveLeaf` threw, the catch turned that into `false`, and a coordinator running
|
||||
// `dispatch --inject` concluded its structured worker was a bare shell — `no_agent_detected`.
|
||||
// A structured session IS the agent; there is no foreground process to recognise.
|
||||
if (this.deps.isLiveStructuredAgent?.(handle)) {
|
||||
return true
|
||||
}
|
||||
try {
|
||||
const pty = this.deps.getLivePty(handle)
|
||||
if (pty) {
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
/**
|
||||
* Closing a structured agent session's provider child, and proving it went.
|
||||
*
|
||||
* Extracted from `stopStructuredWorker` so that orchestration settlement and worktree teardown
|
||||
* close a session the SAME way rather than one of them inventing a shorter version. Everything
|
||||
* dispatch-shaped — dropping the hold, the redrive subscription and the parked mail — stays with
|
||||
* the caller that has a dispatch; this is only the child.
|
||||
*
|
||||
* `host.close` returns void and keeps a failed close indexed for retry, so the only settlement
|
||||
* evidence is the observation AFTER it: a session the host no longer holds and whose lease is no
|
||||
* longer live is proven gone. Anything else is retained rather than settled.
|
||||
*/
|
||||
|
||||
import { getStructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { OrcaRuntimeService } from './orca-runtime'
|
||||
import { retireSettledStructuredWorkerTab } from './structured-agent-session-tab-retirement'
|
||||
import { observeStructuredWorker } from './structured-worker-authority'
|
||||
|
||||
export type StructuredAgentSessionCloseOutcome = {
|
||||
stopped: boolean
|
||||
/** Whether a close was actually issued; a receipt must not claim one that never happened. */
|
||||
closeAttempted: boolean
|
||||
reason?: string
|
||||
}
|
||||
|
||||
export type StructuredAgentSessionCloseOptions = {
|
||||
runtime?: Pick<
|
||||
OrcaRuntimeService,
|
||||
'forgetStructuredSessionMail' | 'retireStructuredAgentSessionTabFromSnapshot'
|
||||
>
|
||||
/**
|
||||
* Runs after the close is issued and BEFORE the proof is read.
|
||||
*
|
||||
* Not after: an unsettled close returns early, so a dispatch that released its hold there would
|
||||
* keep the child un-evictable for the life of the app. Every settlement has to reach it.
|
||||
*/
|
||||
afterClose?: () => void
|
||||
}
|
||||
|
||||
export async function closeStructuredAgentSessionChild(
|
||||
sessionId: string,
|
||||
options: StructuredAgentSessionCloseOptions = {}
|
||||
): Promise<StructuredAgentSessionCloseOutcome> {
|
||||
const host = getStructuredAgentSessionHost()
|
||||
if (!host) {
|
||||
// Nothing was reached, so nothing was acted on; the receipt must not claim a close.
|
||||
return {
|
||||
stopped: false,
|
||||
closeAttempted: false,
|
||||
reason: 'The structured agent-session host is not installed; no session was closed.'
|
||||
}
|
||||
}
|
||||
// Set only once the close is actually issued: `setSessionTabVisibility` throwing first leaves a
|
||||
// running child, and a receipt that still said `closed_agent_terminal` for it would be the
|
||||
// close-that-never-happened this flag exists to rule out.
|
||||
let closeAttempted = false
|
||||
try {
|
||||
await host.setSessionTabVisibility?.(sessionId, false)
|
||||
closeAttempted = true
|
||||
await host.close(sessionId)
|
||||
} catch (error) {
|
||||
return {
|
||||
stopped: false,
|
||||
closeAttempted,
|
||||
reason: error instanceof Error ? error.message : String(error)
|
||||
}
|
||||
}
|
||||
options.afterClose?.()
|
||||
if (observeStructuredWorker({ sessionId }).status === 'live') {
|
||||
return {
|
||||
stopped: false,
|
||||
closeAttempted: true,
|
||||
reason: 'The structured session is still attached after close.'
|
||||
}
|
||||
}
|
||||
// Only past the proof, and structurally unable to throw: the session's chat tab is retired from
|
||||
// the live snapshot, which `setSessionTabVisibility(false)` above does not do.
|
||||
retireSettledStructuredWorkerTab(sessionId, options.runtime)
|
||||
return { stopped: true, closeAttempted: true }
|
||||
}
|
||||
@@ -0,0 +1,143 @@
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentSessionRecord } from '../../shared/agent-session-record'
|
||||
|
||||
const hostRef: { current: unknown } = { current: null }
|
||||
|
||||
vi.mock('../native-chat/agent-session-wire/structured-agent-session-registry', () => ({
|
||||
getStructuredAgentSessionHost: () => hostRef.current
|
||||
}))
|
||||
|
||||
const { killAllProcessesForWorktree } = await import('./worktree-teardown')
|
||||
const { listLiveStructuredSessionsForWorktree } =
|
||||
await import('./structured-session-worktree-teardown')
|
||||
|
||||
const WORKTREE = 'repo_1::/tmp/wt-a'
|
||||
const OTHER_WORKTREE = 'repo_1::/tmp/wt-b'
|
||||
|
||||
function record(sessionId: string, workspaceId: string): AgentSessionRecord {
|
||||
return {
|
||||
sessionId,
|
||||
provider: 'claude',
|
||||
location: { executionHostId: 'local', wslDistro: null, workspaceId, workspaceKind: 'folder' },
|
||||
lease: {
|
||||
sessionId,
|
||||
runtimeKind: 'native',
|
||||
claimStatus: 'live',
|
||||
handoffStage: null,
|
||||
runtimeFence: 1,
|
||||
deathEvidence: null
|
||||
}
|
||||
} as unknown as AgentSessionRecord
|
||||
}
|
||||
|
||||
function installHost(options: {
|
||||
records: AgentSessionRecord[]
|
||||
/** Sessions the host still holds; a close removes one unless it is listed as stuck. */
|
||||
stuck?: Set<string>
|
||||
}): { closed: string[] } {
|
||||
const held = new Set(options.records.map((entry) => entry.sessionId))
|
||||
const closed: string[] = []
|
||||
hostRef.current = {
|
||||
deps: { store: { listRecords: () => options.records, getRecord: () => null } },
|
||||
hasSession: (sessionId: string) => held.has(sessionId),
|
||||
setSessionTabVisibility: async () => {},
|
||||
close: async (sessionId: string) => {
|
||||
closed.push(sessionId)
|
||||
if (!options.stuck?.has(sessionId)) {
|
||||
held.delete(sessionId)
|
||||
}
|
||||
}
|
||||
}
|
||||
// `observeStructuredWorker` reads the record through the same host, so keep them consistent.
|
||||
;(
|
||||
hostRef.current as { deps: { store: { getRecord: (id: string) => unknown } } }
|
||||
).deps.store.getRecord = (sessionId: string) =>
|
||||
options.records.find((entry) => entry.sessionId === sessionId) ?? null
|
||||
return { closed }
|
||||
}
|
||||
|
||||
const localProvider = {
|
||||
listProcesses: async () => [],
|
||||
shutdown: async () => {}
|
||||
} as never
|
||||
|
||||
function destructiveDeps(extra: { allowUnverifiedStop?: boolean } = {}) {
|
||||
return {
|
||||
localProvider,
|
||||
requirePhysicalStop: true,
|
||||
includeProviderInventory: false as const,
|
||||
includeLocalRegistry: false as const,
|
||||
...extra
|
||||
}
|
||||
}
|
||||
|
||||
describe('worktree teardown and structured agent sessions', () => {
|
||||
beforeEach(() => {
|
||||
hostRef.current = null
|
||||
})
|
||||
|
||||
it('finds sessions by workspace, and ignores a sibling worktree', () => {
|
||||
installHost({ records: [record('s1', WORKTREE), record('s2', OTHER_WORKTREE)] })
|
||||
expect(listLiveStructuredSessionsForWorktree(WORKTREE)).toEqual([
|
||||
{ sessionId: 's1', agent: 'claude' }
|
||||
])
|
||||
})
|
||||
|
||||
it('refuses a destructive removal rather than deleting the checkout under a live child', async () => {
|
||||
// The defect this pins: all three PTY sweeps enumerate leaves, provider sessions and the local
|
||||
// registry, and a structured session is on NONE of them. Every sweep answered zero, nothing
|
||||
// errored, and removal proceeded — leaving the provider child running with its `cwd` deleted
|
||||
// and the dispatch still reporting the worker live and exact.
|
||||
installHost({ records: [record('s1', WORKTREE)] })
|
||||
await expect(killAllProcessesForWorktree(WORKTREE, destructiveDeps())).rejects.toThrow(
|
||||
/1 running agent session/
|
||||
)
|
||||
})
|
||||
|
||||
it('names the force escape hatch in the refusal, like the unstopped-PTY gate', async () => {
|
||||
installHost({ records: [record('s1', WORKTREE)] })
|
||||
await expect(killAllProcessesForWorktree(WORKTREE, destructiveDeps())).rejects.toThrow(/force/i)
|
||||
})
|
||||
|
||||
it('closes them under force instead of orphaning the child', async () => {
|
||||
const host = installHost({ records: [record('s1', WORKTREE), record('s2', WORKTREE)] })
|
||||
const result = await killAllProcessesForWorktree(
|
||||
WORKTREE,
|
||||
destructiveDeps({ allowUnverifiedStop: true })
|
||||
)
|
||||
expect(host.closed).toEqual(['s1', 's2'])
|
||||
expect(result.structuredStopped).toBe(2)
|
||||
})
|
||||
|
||||
it('still removes under force when a close does not settle, and says so', async () => {
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
installHost({ records: [record('s1', WORKTREE)], stuck: new Set(['s1']) })
|
||||
const result = await killAllProcessesForWorktree(
|
||||
WORKTREE,
|
||||
destructiveDeps({ allowUnverifiedStop: true })
|
||||
)
|
||||
expect(result.structuredStopped).toBeUndefined()
|
||||
expect(warn).toHaveBeenCalledWith(expect.stringContaining('still attached'))
|
||||
warn.mockRestore()
|
||||
})
|
||||
|
||||
it('leaves the best-effort reconciliation paths alone', async () => {
|
||||
// Those callers repair state and delete nothing, so a refusal there would wedge a repair.
|
||||
installHost({ records: [record('s1', WORKTREE)] })
|
||||
await expect(
|
||||
killAllProcessesForWorktree(WORKTREE, {
|
||||
localProvider,
|
||||
includeProviderInventory: false,
|
||||
includeLocalRegistry: false
|
||||
})
|
||||
).resolves.toMatchObject({ runtimeStopped: 0 })
|
||||
})
|
||||
|
||||
it('does not block removal when no structured host is installed', async () => {
|
||||
// Not being able to look is not evidence a child is there, and reading the persisted store
|
||||
// directly would force-install the host as a side effect of a teardown.
|
||||
await expect(killAllProcessesForWorktree(WORKTREE, destructiveDeps())).resolves.toMatchObject({
|
||||
runtimeStopped: 0
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,97 @@
|
||||
/**
|
||||
* The structured half of worktree teardown.
|
||||
*
|
||||
* `killAllProcessesForWorktree` sweeps three PTY surfaces — the renderer graph, the provider's
|
||||
* session list, and the local pty-registry — and a structured agent session appears on NONE of
|
||||
* them. It has no PTY, no leaf, and no provider session row. So every sweep counted zero, no error
|
||||
* was raised, and removal deleted the checkout out from under a running provider child: the child
|
||||
* kept running with its `cwd` gone, the durable record and chat tab survived to republish at the
|
||||
* next launch pointing at a deleted worktree, and `worker-show` still reported the worker live.
|
||||
*
|
||||
* Membership is `location.workspaceId`, which every structured session carries — so this covers a
|
||||
* plain chat session in the worktree as well as a dispatched worker. Liveness is
|
||||
* `observeStructuredWorker`, the same `live` / `unverifiable` / `exited` vocabulary the rest of the
|
||||
* structured surface uses; only a PROVEN live child is worth refusing a removal over.
|
||||
*/
|
||||
|
||||
import { getStructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import { observeStructuredWorker } from './structured-worker-authority'
|
||||
import { closeStructuredAgentSessionChild } from './structured-agent-session-close'
|
||||
import type { OrcaRuntimeService } from './orca-runtime'
|
||||
|
||||
export type LiveStructuredSessionInWorkspace = {
|
||||
sessionId: string
|
||||
agent: 'claude' | 'codex'
|
||||
}
|
||||
|
||||
export type StructuredWorktreeSweepRuntime = Pick<
|
||||
OrcaRuntimeService,
|
||||
'forgetStructuredSessionMail' | 'retireStructuredAgentSessionTabFromSnapshot'
|
||||
>
|
||||
|
||||
/**
|
||||
* Structured sessions with a proven-live child in this worktree.
|
||||
*
|
||||
* An uninstalled host answers empty rather than throwing: no host in this generation means no
|
||||
* provider child was started by this process, and the three PTY sweeps fall through the same way
|
||||
* when their surface is unavailable. It is deliberately NOT read through the persisted store
|
||||
* directly — that would force-install the host, which is itself a side effect on a teardown path.
|
||||
*/
|
||||
export function listLiveStructuredSessionsForWorktree(
|
||||
worktreeId: string
|
||||
): LiveStructuredSessionInWorkspace[] {
|
||||
const host = getStructuredAgentSessionHost()
|
||||
if (!host) {
|
||||
return []
|
||||
}
|
||||
let records: ReturnType<typeof host.deps.store.listRecords>
|
||||
try {
|
||||
records = host.deps.store.listRecords()
|
||||
} catch {
|
||||
return []
|
||||
}
|
||||
return records
|
||||
.filter(
|
||||
(record) =>
|
||||
record.location.workspaceId === worktreeId &&
|
||||
observeStructuredWorker({ sessionId: record.sessionId }).status === 'live'
|
||||
)
|
||||
.map((record) => ({ sessionId: record.sessionId, agent: record.provider }))
|
||||
}
|
||||
|
||||
export function describeLiveStructuredSessions(
|
||||
worktreeId: string,
|
||||
sessions: readonly LiveStructuredSessionInWorkspace[]
|
||||
): string {
|
||||
const noun = sessions.length === 1 ? 'agent session' : 'agent sessions'
|
||||
return `${sessions.length} running ${noun} in ${worktreeId} (${sessions
|
||||
.map((session) => `${session.agent}:${session.sessionId}`)
|
||||
.join(', ')})`
|
||||
}
|
||||
|
||||
/**
|
||||
* Closes every live structured session in the worktree, and reports what stayed.
|
||||
*
|
||||
* Force is the documented escape hatch, so it closes rather than orphaning: a child left running
|
||||
* against a deleted `cwd` is the exact outcome this whole sweep exists to prevent.
|
||||
*/
|
||||
export async function closeStructuredSessionsForWorktree(
|
||||
worktreeId: string,
|
||||
runtime?: StructuredWorktreeSweepRuntime
|
||||
): Promise<{ closed: number; unstopped: LiveStructuredSessionInWorkspace[] }> {
|
||||
const sessions = listLiveStructuredSessionsForWorktree(worktreeId)
|
||||
const unstopped: LiveStructuredSessionInWorkspace[] = []
|
||||
let closed = 0
|
||||
for (const session of sessions) {
|
||||
const outcome = await closeStructuredAgentSessionChild(
|
||||
session.sessionId,
|
||||
runtime ? { runtime } : {}
|
||||
)
|
||||
if (outcome.stopped) {
|
||||
closed += 1
|
||||
} else {
|
||||
unstopped.push(session)
|
||||
}
|
||||
}
|
||||
return { closed, unstopped }
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
resolveTerminalIdentityFromProbes,
|
||||
TERMINAL_HANDLE_STALE_ERROR
|
||||
} from './terminal-identity-probe'
|
||||
|
||||
function probes(overrides: { structured?: boolean; livePty?: boolean; leafError?: Error | null }) {
|
||||
const assertLiveLeaf = vi.fn(() => {
|
||||
if (overrides.leafError) {
|
||||
throw overrides.leafError
|
||||
}
|
||||
})
|
||||
return {
|
||||
calls: { assertLiveLeaf },
|
||||
probes: {
|
||||
isLiveStructuredWorker: () => overrides.structured ?? false,
|
||||
hasLivePty: () => overrides.livePty ?? false,
|
||||
assertLiveLeaf
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
describe('the terminal identity probe', () => {
|
||||
it('answers live for a structured worker without touching the PTY graph', () => {
|
||||
// The defect this pins: the sender validator asked `terminal.show`, whose leaf lookup misses
|
||||
// for a session that never had a pane, and reported a live worker's own handle as stale.
|
||||
const { calls, probes: p } = probes({ structured: true })
|
||||
expect(resolveTerminalIdentityFromProbes('structworker_1', p)).toEqual({
|
||||
handle: 'structworker_1',
|
||||
live: true
|
||||
})
|
||||
expect(calls.assertLiveLeaf).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('answers live for a PTY handle the runtime still holds', () => {
|
||||
const { probes: p } = probes({ livePty: true })
|
||||
expect(resolveTerminalIdentityFromProbes('term_1', p).live).toBe(true)
|
||||
})
|
||||
|
||||
it('runs the full leaf check for a handle with no live PTY', () => {
|
||||
// `getLiveLeafForHandle` is the one that re-checks `rendererGraphEpoch`, and that check is the
|
||||
// entire reason the sender is validated: a long-lived shell keeps a stale
|
||||
// `ORCA_TERMINAL_HANDLE` across a window reload. A cheaper probe would start passing it.
|
||||
const { calls, probes: p } = probes({})
|
||||
expect(resolveTerminalIdentityFromProbes('term_1', p).live).toBe(true)
|
||||
expect(calls.assertLiveLeaf).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('answers not-live for a stale handle', () => {
|
||||
const { probes: p } = probes({ leafError: new Error(TERMINAL_HANDLE_STALE_ERROR) })
|
||||
expect(resolveTerminalIdentityFromProbes('term_1', p)).toEqual({
|
||||
handle: 'term_1',
|
||||
live: false
|
||||
})
|
||||
})
|
||||
|
||||
it('propagates "could not look" rather than reporting it as a dead handle', () => {
|
||||
// A graph that is not ready yet is not evidence the handle died, and `terminal.show` lets that
|
||||
// error through today. Answering `live: false` here would make a command refuse its own sender
|
||||
// during startup instead of failing loudly.
|
||||
const { probes: p } = probes({ leafError: new Error('graph_not_ready') })
|
||||
expect(() => resolveTerminalIdentityFromProbes('term_1', p)).toThrow('graph_not_ready')
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,55 @@
|
||||
/**
|
||||
* "Is this handle a live orchestration identity?" — answered for BOTH lanes.
|
||||
*
|
||||
* The CLI asked that question by calling `terminal.show`, which is a PTY verb: it resolves a pane,
|
||||
* a ptyId and a preview. A structured worker has none of those, so `showTerminal` missed, threw
|
||||
* `terminal_handle_stale`, and the caller concluded the handle the child was BORN with was dead —
|
||||
* failing twelve coordinator verbs for a worker whose own preamble tells it to run them.
|
||||
*
|
||||
* So the identity question gets its own probe, returning a handle and a boolean and nothing
|
||||
* writable. `terminal.show` deliberately still refuses a structured handle: synthesising
|
||||
* `ptyId`/`leafId`/`paneRuntimeId` would hand every public terminal verb something that looks
|
||||
* writable and is not.
|
||||
*
|
||||
* The PTY half is EXACTLY today's `terminal.show` liveness test, `getLiveLeafForHandle` included,
|
||||
* so its `rendererGraphEpoch` re-check still runs. That check is the whole point of validating at
|
||||
* all — a long-lived shell keeps a stale `ORCA_TERMINAL_HANDLE` across a window reload — and a
|
||||
* cheaper probe that skipped it (`getPaneKeyForTerminalHandle`, say) would quietly start passing
|
||||
* handles that fail today.
|
||||
*/
|
||||
|
||||
export type RuntimeTerminalIdentity = {
|
||||
handle: string
|
||||
live: boolean
|
||||
}
|
||||
|
||||
/** The one error code that means "not live" rather than "could not look". */
|
||||
export const TERMINAL_HANDLE_STALE_ERROR = 'terminal_handle_stale'
|
||||
|
||||
export type TerminalIdentityProbes = {
|
||||
/** A structured worker of THIS runtime, proven through its durable record. */
|
||||
isLiveStructuredWorker: () => boolean
|
||||
hasLivePty: () => boolean
|
||||
/** Today's leaf check; throws `terminal_handle_stale` for a stale or reloaded handle. */
|
||||
assertLiveLeaf: () => void
|
||||
}
|
||||
|
||||
export function resolveTerminalIdentityFromProbes(
|
||||
handle: string,
|
||||
probes: TerminalIdentityProbes
|
||||
): RuntimeTerminalIdentity {
|
||||
if (probes.isLiveStructuredWorker() || probes.hasLivePty()) {
|
||||
return { handle, live: true }
|
||||
}
|
||||
try {
|
||||
probes.assertLiveLeaf()
|
||||
return { handle, live: true }
|
||||
} catch (error) {
|
||||
if (error instanceof Error && error.message === TERMINAL_HANDLE_STALE_ERROR) {
|
||||
return { handle, live: false }
|
||||
}
|
||||
// Anything else — a graph that is not ready yet — is "could not look", and must propagate
|
||||
// exactly as it does through `terminal.show` today rather than being read as a dead handle.
|
||||
throw error
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,140 @@
|
||||
/**
|
||||
* The two PTY-surface sweeps `killAllProcessesForWorktree` fans out to.
|
||||
*
|
||||
* Split from the teardown entry point so that file stays under the line ceiling once the
|
||||
* structured-session sweep joined it. Each function owns one registration surface: the installed
|
||||
* provider's session list, and the local pty-registry.
|
||||
*/
|
||||
|
||||
import type { IPtyProvider } from '../providers/types'
|
||||
import { listRegisteredPtys } from '../memory/pty-registry'
|
||||
import { isPathInsideOrEqual } from '../../shared/cross-platform-path'
|
||||
import { splitWorktreeId, splitWorktreeIdForFilesystem } from '../../shared/worktree/id'
|
||||
import { mapWithConcurrency } from '../../shared/map-with-concurrency'
|
||||
import { teardownRpcDeadline } from './worktree-teardown-deadline'
|
||||
|
||||
// Why: normal inventories still coalesce into one process scan, while a stale
|
||||
// or pathological inventory cannot fan out unbounded provider/RPC shutdowns.
|
||||
const WORKTREE_TEARDOWN_CONCURRENCY = 32
|
||||
|
||||
export type WorktreeTeardownStopPty = (
|
||||
ptyId: string,
|
||||
stop: () => Promise<boolean>
|
||||
) => Promise<{ stopped: boolean; owner: boolean }>
|
||||
|
||||
export async function sweepProviderByPrefix(
|
||||
worktreeId: string,
|
||||
provider: IPtyProvider,
|
||||
deadline: number,
|
||||
stopPty: (
|
||||
ptyId: string,
|
||||
stop: () => Promise<boolean>
|
||||
) => Promise<{ stopped: boolean; owner: boolean }>,
|
||||
onPtyStopped?: (ptyId: string) => void,
|
||||
failClosed = false
|
||||
): Promise<number> {
|
||||
const prefix = `${worktreeId}@@`
|
||||
// Why (#10252): the cwd fallback only proves ownership when the filesystem path
|
||||
// is the *whole* worktree path. A folder-workspace instance strips its
|
||||
// `::workspace:<uuid>` suffix to a checkout dir shared with sibling instances,
|
||||
// so leave the fallback unset whenever stripping shortened the path — else
|
||||
// deleting one instance would sweep the others.
|
||||
const fullWorktreePath = splitWorktreeId(worktreeId)?.worktreePath
|
||||
const cwdFallbackPath =
|
||||
splitWorktreeIdForFilesystem(worktreeId)?.worktreePath === fullWorktreePath
|
||||
? fullWorktreePath
|
||||
: undefined
|
||||
const rpcDeadline = teardownRpcDeadline(deadline)
|
||||
const sessions = failClosed
|
||||
? await provider.listProcesses({ deadlineMs: rpcDeadline })
|
||||
: await provider.listProcesses({ deadlineMs: rpcDeadline }).catch(() => [])
|
||||
const ownedSessions = sessions.filter((session) => {
|
||||
// Why: older daemon/relay process rows may omit cwd; their established ID
|
||||
// and authoritative worktree ownership must remain usable during teardown.
|
||||
const cwdOwned =
|
||||
cwdFallbackPath !== undefined &&
|
||||
session.worktreeId === undefined &&
|
||||
typeof session.cwd === 'string' &&
|
||||
session.cwd.length > 0 &&
|
||||
isPathInsideOrEqual(cwdFallbackPath, session.cwd)
|
||||
return session.id.startsWith(prefix) || session.worktreeId === worktreeId || cwdOwned
|
||||
})
|
||||
// Why: agent shutdown snapshots coalesce only when requests begin together;
|
||||
// bounded concurrency avoids serial process scans without unbounded fanout.
|
||||
const stopped = await mapWithConcurrency(
|
||||
ownedSessions,
|
||||
WORKTREE_TEARDOWN_CONCURRENCY,
|
||||
async (session) => {
|
||||
if (Date.now() >= deadline) {
|
||||
return 0
|
||||
}
|
||||
const stopResult = await stopPty(session.id, async () => {
|
||||
if (Date.now() >= deadline) {
|
||||
return false
|
||||
}
|
||||
try {
|
||||
await provider.shutdown(session.id, { immediate: true, deadlineMs: rpcDeadline })
|
||||
return Date.now() < deadline
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
})
|
||||
if (stopResult.owner && Date.now() < deadline) {
|
||||
clearStoppedPtyState(session.id, onPtyStopped)
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
)
|
||||
return stopped.reduce<number>((count, value) => count + value, 0)
|
||||
}
|
||||
|
||||
export async function sweepRegistryForWorktree(
|
||||
worktreeId: string,
|
||||
localProvider: IPtyProvider,
|
||||
deadline: number,
|
||||
stopPty: (
|
||||
ptyId: string,
|
||||
stop: () => Promise<boolean>
|
||||
) => Promise<{ stopped: boolean; owner: boolean }>,
|
||||
onPtyStopped?: (ptyId: string) => void
|
||||
): Promise<number> {
|
||||
const rpcDeadline = teardownRpcDeadline(deadline)
|
||||
const entries = listRegisteredPtys().filter((r) => r.worktreeId === worktreeId)
|
||||
const stopped = await mapWithConcurrency(
|
||||
entries,
|
||||
WORKTREE_TEARDOWN_CONCURRENCY,
|
||||
async (entry) => {
|
||||
if (Date.now() >= deadline) {
|
||||
return 0
|
||||
}
|
||||
const stopResult = await stopPty(entry.ptyId, async () => {
|
||||
if (Date.now() >= deadline) {
|
||||
return false
|
||||
}
|
||||
try {
|
||||
await localProvider.shutdown(entry.ptyId, { immediate: true, deadlineMs: rpcDeadline })
|
||||
return Date.now() < deadline
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
})
|
||||
if (stopResult.owner && Date.now() < deadline) {
|
||||
clearStoppedPtyState(entry.ptyId, onPtyStopped)
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
)
|
||||
return stopped.reduce<number>((count, value) => count + value, 0)
|
||||
}
|
||||
|
||||
export function clearStoppedPtyState(ptyId: string, onPtyStopped?: (ptyId: string) => void): void {
|
||||
try {
|
||||
// Why: daemon shutdown does not always fan a local pty:exit event back
|
||||
// through pty.ts, but removed worktrees must immediately drop memory rows.
|
||||
onPtyStopped?.(ptyId)
|
||||
} catch {
|
||||
/* cleanup is best-effort and must not block git-level removal */
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
/**
|
||||
* The one deadline arithmetic worktree teardown shares.
|
||||
*
|
||||
* Its own module because both the teardown entry point and the PTY-surface sweeps need it, and a
|
||||
* sweep importing the entry point back would be a cycle.
|
||||
*/
|
||||
|
||||
// Why: keep each bounded stop RPC settling before the sweep deadline itself, so
|
||||
// a wedged provider surfaces as a stop failure rather than as the outer timeout.
|
||||
// (The recheck this margin once also reserved time for now runs on its own
|
||||
// budget — see verifyUnstoppedPtys — because sharing this one wedged #11960.)
|
||||
export const WORKTREE_TEARDOWN_RPC_MARGIN_MS = 500
|
||||
|
||||
// Absolute deadline (epoch ms) threaded into provider RPCs on the destructive
|
||||
// path; each RPC leaf converts it to the remaining time when it actually issues,
|
||||
// so sequential RPCs share one budget without any relative-timeout bookkeeping.
|
||||
export function teardownRpcDeadline(sweepDeadline: number): number {
|
||||
return sweepDeadline - WORKTREE_TEARDOWN_RPC_MARGIN_MS
|
||||
}
|
||||
@@ -1,15 +1,21 @@
|
||||
import type { IPtyProvider } from '../providers/types'
|
||||
import type { OrcaRuntimeService } from './orca-runtime'
|
||||
import { listRegisteredPtys } from '../memory/pty-registry'
|
||||
import { isPathInsideOrEqual } from '../../shared/cross-platform-path'
|
||||
import { splitWorktreeId, splitWorktreeIdForFilesystem } from '../../shared/worktree/id'
|
||||
import { mapWithConcurrency } from '../../shared/map-with-concurrency'
|
||||
import {
|
||||
isUnstoppedPtyRemovalError,
|
||||
WORKTREE_TEARDOWN_FORCE_HINT,
|
||||
WORKTREE_TEARDOWN_TIMEOUT_PREFIX
|
||||
} from '../../shared/worktree/removal'
|
||||
import { settleBeforeDeadline } from './settle-before-deadline'
|
||||
import {
|
||||
clearStoppedPtyState,
|
||||
sweepProviderByPrefix,
|
||||
sweepRegistryForWorktree
|
||||
} from './worktree-pty-surface-sweeps'
|
||||
import {
|
||||
closeStructuredSessionsForWorktree,
|
||||
describeLiveStructuredSessions,
|
||||
listLiveStructuredSessionsForWorktree
|
||||
} from './structured-session-worktree-teardown'
|
||||
import { createWorktreeSweepTracker, settleSweepsForForcedRemoval } from './forced-sweep-settlement'
|
||||
import {
|
||||
describeError,
|
||||
@@ -18,10 +24,6 @@ import {
|
||||
resolveUnstoppedPtyVerdict
|
||||
} from './unstopped-pty-verification'
|
||||
|
||||
// Why: normal inventories still coalesce into one process scan, while a stale
|
||||
// or pathological inventory cannot fan out unbounded provider/RPC shutdowns.
|
||||
const WORKTREE_TEARDOWN_CONCURRENCY = 32
|
||||
|
||||
export type WorktreeTeardownDeps = {
|
||||
runtime?: OrcaRuntimeService
|
||||
/** Authoritative id for callers whose selector no longer resolves (orphaned workspace). */
|
||||
@@ -44,22 +46,13 @@ export type WorktreeTeardownResult = {
|
||||
runtimeStopped: number
|
||||
providerStopped: number
|
||||
registryStopped: number
|
||||
/** Structured agent sessions closed by the force path; absent when none were found. */
|
||||
structuredStopped?: number
|
||||
}
|
||||
|
||||
export const WORKTREE_PROCESS_SWEEP_TIMEOUT_MS = 10_000
|
||||
|
||||
// Why: keep each bounded stop RPC settling before the sweep deadline itself, so
|
||||
// a wedged provider surfaces as a stop failure rather than as the outer timeout.
|
||||
// (The recheck this margin once also reserved time for now runs on its own
|
||||
// budget — see verifyUnstoppedPtys — because sharing this one wedged #11960.)
|
||||
export const WORKTREE_TEARDOWN_RPC_MARGIN_MS = 500
|
||||
|
||||
// Absolute deadline (epoch ms) threaded into provider RPCs on the destructive
|
||||
// path; each RPC leaf converts it to the remaining time when it actually issues,
|
||||
// so sequential RPCs share one budget without any relative-timeout bookkeeping.
|
||||
export function teardownRpcDeadline(sweepDeadline: number): number {
|
||||
return sweepDeadline - WORKTREE_TEARDOWN_RPC_MARGIN_MS
|
||||
}
|
||||
export { WORKTREE_TEARDOWN_RPC_MARGIN_MS, teardownRpcDeadline } from './worktree-teardown-deadline'
|
||||
|
||||
/**
|
||||
* Kills every PTY we can prove belongs to `worktreeId`, across all three
|
||||
@@ -95,6 +88,11 @@ export async function killAllProcessesForWorktree(
|
||||
const deadlineError = new Error(
|
||||
`${WORKTREE_TEARDOWN_TIMEOUT_PREFIX} ${worktreeId}. ${WORKTREE_TEARDOWN_FORCE_HINT}`
|
||||
)
|
||||
// FIRST, and before a single PTY sweep starts: a structured agent session is registered on none
|
||||
// of the three surfaces below, so all three answered zero and removal deleted the checkout out
|
||||
// from under a running provider child. Refusing costs nothing when there are none, and the check
|
||||
// is synchronous, so a destructive removal fails fast instead of after the whole sweep budget.
|
||||
const structuredStopped = await sweepStructuredSessions(worktreeId, deps)
|
||||
const sweeps = createWorktreeSweepTracker()
|
||||
const stopAttempts = new Map<string, Promise<boolean>>()
|
||||
const stopPty = (
|
||||
@@ -256,122 +254,49 @@ export async function killAllProcessesForWorktree(
|
||||
}
|
||||
}
|
||||
|
||||
return { runtimeStopped: runtimeResult.stopped, providerStopped, registryStopped }
|
||||
}
|
||||
|
||||
async function sweepProviderByPrefix(
|
||||
worktreeId: string,
|
||||
provider: IPtyProvider,
|
||||
deadline: number,
|
||||
stopPty: (
|
||||
ptyId: string,
|
||||
stop: () => Promise<boolean>
|
||||
) => Promise<{ stopped: boolean; owner: boolean }>,
|
||||
onPtyStopped?: (ptyId: string) => void,
|
||||
failClosed = false
|
||||
): Promise<number> {
|
||||
const prefix = `${worktreeId}@@`
|
||||
// Why (#10252): the cwd fallback only proves ownership when the filesystem path
|
||||
// is the *whole* worktree path. A folder-workspace instance strips its
|
||||
// `::workspace:<uuid>` suffix to a checkout dir shared with sibling instances,
|
||||
// so leave the fallback unset whenever stripping shortened the path — else
|
||||
// deleting one instance would sweep the others.
|
||||
const fullWorktreePath = splitWorktreeId(worktreeId)?.worktreePath
|
||||
const cwdFallbackPath =
|
||||
splitWorktreeIdForFilesystem(worktreeId)?.worktreePath === fullWorktreePath
|
||||
? fullWorktreePath
|
||||
: undefined
|
||||
const rpcDeadline = teardownRpcDeadline(deadline)
|
||||
const sessions = failClosed
|
||||
? await provider.listProcesses({ deadlineMs: rpcDeadline })
|
||||
: await provider.listProcesses({ deadlineMs: rpcDeadline }).catch(() => [])
|
||||
const ownedSessions = sessions.filter((session) => {
|
||||
// Why: older daemon/relay process rows may omit cwd; their established ID
|
||||
// and authoritative worktree ownership must remain usable during teardown.
|
||||
const cwdOwned =
|
||||
cwdFallbackPath !== undefined &&
|
||||
session.worktreeId === undefined &&
|
||||
typeof session.cwd === 'string' &&
|
||||
session.cwd.length > 0 &&
|
||||
isPathInsideOrEqual(cwdFallbackPath, session.cwd)
|
||||
return session.id.startsWith(prefix) || session.worktreeId === worktreeId || cwdOwned
|
||||
})
|
||||
// Why: agent shutdown snapshots coalesce only when requests begin together;
|
||||
// bounded concurrency avoids serial process scans without unbounded fanout.
|
||||
const stopped = await mapWithConcurrency(
|
||||
ownedSessions,
|
||||
WORKTREE_TEARDOWN_CONCURRENCY,
|
||||
async (session) => {
|
||||
if (Date.now() >= deadline) {
|
||||
return 0
|
||||
}
|
||||
const stopResult = await stopPty(session.id, async () => {
|
||||
if (Date.now() >= deadline) {
|
||||
return false
|
||||
}
|
||||
try {
|
||||
await provider.shutdown(session.id, { immediate: true, deadlineMs: rpcDeadline })
|
||||
return Date.now() < deadline
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
})
|
||||
if (stopResult.owner && Date.now() < deadline) {
|
||||
clearStoppedPtyState(session.id, onPtyStopped)
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
)
|
||||
return stopped.reduce<number>((count, value) => count + value, 0)
|
||||
}
|
||||
|
||||
async function sweepRegistryForWorktree(
|
||||
worktreeId: string,
|
||||
localProvider: IPtyProvider,
|
||||
deadline: number,
|
||||
stopPty: (
|
||||
ptyId: string,
|
||||
stop: () => Promise<boolean>
|
||||
) => Promise<{ stopped: boolean; owner: boolean }>,
|
||||
onPtyStopped?: (ptyId: string) => void
|
||||
): Promise<number> {
|
||||
const rpcDeadline = teardownRpcDeadline(deadline)
|
||||
const entries = listRegisteredPtys().filter((r) => r.worktreeId === worktreeId)
|
||||
const stopped = await mapWithConcurrency(
|
||||
entries,
|
||||
WORKTREE_TEARDOWN_CONCURRENCY,
|
||||
async (entry) => {
|
||||
if (Date.now() >= deadline) {
|
||||
return 0
|
||||
}
|
||||
const stopResult = await stopPty(entry.ptyId, async () => {
|
||||
if (Date.now() >= deadline) {
|
||||
return false
|
||||
}
|
||||
try {
|
||||
await localProvider.shutdown(entry.ptyId, { immediate: true, deadlineMs: rpcDeadline })
|
||||
return Date.now() < deadline
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
})
|
||||
if (stopResult.owner && Date.now() < deadline) {
|
||||
clearStoppedPtyState(entry.ptyId, onPtyStopped)
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
)
|
||||
return stopped.reduce<number>((count, value) => count + value, 0)
|
||||
}
|
||||
|
||||
function clearStoppedPtyState(ptyId: string, onPtyStopped?: (ptyId: string) => void): void {
|
||||
try {
|
||||
// Why: daemon shutdown does not always fan a local pty:exit event back
|
||||
// through pty.ts, but removed worktrees must immediately drop memory rows.
|
||||
onPtyStopped?.(ptyId)
|
||||
} catch {
|
||||
/* cleanup is best-effort and must not block git-level removal */
|
||||
return {
|
||||
runtimeStopped: runtimeResult.stopped,
|
||||
providerStopped,
|
||||
registryStopped,
|
||||
...(structuredStopped > 0 ? { structuredStopped } : {})
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The fourth sweep: structured agent sessions bound to this worktree.
|
||||
*
|
||||
* Refuses rather than auto-closing on the ordinary destructive path. `worktree rm` is the verb
|
||||
* that deletes a user's work, and a running agent session is exactly the thing they would want to
|
||||
* be told about before it goes — the same bargain the unstopped-PTY gate already strikes, using
|
||||
* the same `--force` escape hatch. Force closes them properly instead of orphaning a child against
|
||||
* a `cwd` that is about to disappear.
|
||||
*
|
||||
* Only the destructive path (`requirePhysicalStop`) participates: the best-effort callers are
|
||||
* reconciliation sweeps that must never fail a repair, and they delete nothing.
|
||||
*/
|
||||
async function sweepStructuredSessions(
|
||||
worktreeId: string,
|
||||
deps: WorktreeTeardownDeps
|
||||
): Promise<number> {
|
||||
if (!deps.requirePhysicalStop) {
|
||||
return 0
|
||||
}
|
||||
const live = listLiveStructuredSessionsForWorktree(worktreeId)
|
||||
if (live.length === 0) {
|
||||
return 0
|
||||
}
|
||||
if (!deps.allowUnverifiedStop) {
|
||||
throw new Error(
|
||||
`Refusing to remove ${worktreeId}: ${describeLiveStructuredSessions(worktreeId, live)}. ${WORKTREE_TEARDOWN_FORCE_HINT}`
|
||||
)
|
||||
}
|
||||
const { closed, unstopped } = await closeStructuredSessionsForWorktree(worktreeId, deps.runtime)
|
||||
if (unstopped.length > 0) {
|
||||
// Force is the documented escape hatch, so removal continues — but say so, because the child
|
||||
// outliving its `cwd` is the failure this sweep exists to make visible.
|
||||
console.warn(
|
||||
`[worktree-teardown] forcing removal with ${describeLiveStructuredSessions(worktreeId, unstopped)} still attached`
|
||||
)
|
||||
}
|
||||
return closed
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user