diff --git a/src/cli/handlers/orchestration-caller-identity-cli.test.ts b/src/cli/handlers/orchestration-caller-identity-cli.test.ts new file mode 100644 index 00000000000..dc272e4b94b --- /dev/null +++ b/src/cli/handlers/orchestration-caller-identity-cli.test.ts @@ -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) => + ORCHESTRATION_HANDLERS['orchestration dispatch']({ + flags, + client: { call: callMock }, + cwd: '/tmp/repo', + json: true + } as never) + + const invokeDispatchShow = (flags: Map) => + ORCHESTRATION_HANDLERS['orchestration dispatch-show']({ + flags, + client: { call: callMock }, + cwd: '/tmp/repo', + json: true + } as never) + + const invokeRun = (flags: Map) => + 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([ + ['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([ + ['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([ + ['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([ + ['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([['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) => + 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([['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([['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([['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([['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([['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([['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([['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' + }) + }) +}) diff --git a/src/cli/handlers/orchestration-gate-cli.test.ts b/src/cli/handlers/orchestration-gate-cli.test.ts index b4be315c155..a2793ac9323 100644 --- a/src/cli/handlers/orchestration-gate-cli.test.ts +++ b/src/cli/handlers/orchestration-gate-cli.test.ts @@ -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', diff --git a/src/cli/handlers/orchestration-task-create-cli.test.ts b/src/cli/handlers/orchestration-task-create-cli.test.ts index 839ccabdcd6..8b5957782d2 100644 --- a/src/cli/handlers/orchestration-task-create-cli.test.ts +++ b/src/cli/handlers/orchestration-task-create-cli.test.ts @@ -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']({ diff --git a/src/cli/handlers/orchestration.test.ts b/src/cli/handlers/orchestration.test.ts index c43f7b45b62..b6506e3c29c 100644 --- a/src/cli/handlers/orchestration.test.ts +++ b/src/cli/handlers/orchestration.test.ts @@ -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) => - ORCHESTRATION_HANDLERS['orchestration dispatch']({ - flags, - client: { call: callMock }, - cwd: '/tmp/repo', - json: true - } as never) - - const invokeDispatchShow = (flags: Map) => - ORCHESTRATION_HANDLERS['orchestration dispatch-show']({ - flags, - client: { call: callMock }, - cwd: '/tmp/repo', - json: true - } as never) - - const invokeRun = (flags: Map) => - 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([ - ['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([ - ['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([ - ['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([ - ['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([['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) => - 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([['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([['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([['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([['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([['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([['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([['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], diff --git a/src/cli/handlers/orchestration/terminal-identity.ts b/src/cli/handlers/orchestration/terminal-identity.ts index 2eb7be46484..e14c4a1e182 100644 --- a/src/cli/handlers/orchestration/terminal-identity.ts +++ b/src/cli/handlers/orchestration/terminal-identity.ts @@ -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 { + 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 { try { await client.call('terminal.show', { terminal: handle }) return true diff --git a/src/cli/orchestration-structured-sender-identity.test.ts b/src/cli/orchestration-structured-sender-identity.test.ts new file mode 100644 index 00000000000..9b85125d99e --- /dev/null +++ b/src/cli/orchestration-structured-sender-identity.test.ts @@ -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() + }) +}) diff --git a/src/main/runtime/orca-runtime-get-pty-record-for-pane-key.ts b/src/main/runtime/orca-runtime-get-pty-record-for-pane-key.ts index cfe5c623055..6b861cf55c6 100644 --- a/src/main/runtime/orca-runtime-get-pty-record-for-pane-key.ts +++ b/src/main/runtime/orca-runtime-get-pty-record-for-pane-key.ts @@ -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): 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. * diff --git a/src/main/runtime/orca-runtime-stop-requested-pty-ids.ts b/src/main/runtime/orca-runtime-stop-requested-pty-ids.ts index d76b8ad2be3..0c0f3c8bc8c 100644 --- a/src/main/runtime/orca-runtime-stop-requested-pty-ids.ts +++ b/src/main/runtime/orca-runtime-stop-requested-pty-ids.ts @@ -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 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, diff --git a/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts b/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts index a17d4997dcc..afa636915ae 100644 --- a/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts +++ b/src/main/runtime/orchestration/structured-mailbox-pointer-delivery.ts @@ -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, diff --git a/src/main/runtime/orchestration/structured-worker-direct-mailbox-target.test.ts b/src/main/runtime/orchestration/structured-worker-direct-mailbox-target.test.ts index 4c39d840dfd..e5b36cdf67f 100644 --- a/src/main/runtime/orchestration/structured-worker-direct-mailbox-target.test.ts +++ b/src/main/runtime/orchestration/structured-worker-direct-mailbox-target.test.ts @@ -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() + }) +}) diff --git a/src/main/runtime/rpc/methods/orchestration-cli-runtime-boundary.test.ts b/src/main/runtime/rpc/methods/orchestration-cli-runtime-boundary.test.ts index 8ffa47fc34a..1652cc867ad 100644 --- a/src/main/runtime/rpc/methods/orchestration-cli-runtime-boundary.test.ts +++ b/src/main/runtime/rpc/methods/orchestration-cli-runtime-boundary.test.ts @@ -43,8 +43,10 @@ describe('orchestration CLI/runtime boundary', () => { isRemote: false, /** Preserves terminal-handle validation while routing other calls through runtime RPC. */ async call(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 } } diff --git a/src/main/runtime/rpc/methods/orchestration-structured-worker-lifecycle.ts b/src/main/runtime/rpc/methods/orchestration-structured-worker-lifecycle.ts index e82a9698b6d..436750ac991 100644 --- a/src/main/runtime/rpc/methods/orchestration-structured-worker-lifecycle.ts +++ b/src/main/runtime/rpc/methods/orchestration-structured-worker-lifecycle.ts @@ -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 { - 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. */ diff --git a/src/main/runtime/rpc/methods/terminal-manifest-characterization.test.ts b/src/main/runtime/rpc/methods/terminal-manifest-characterization.test.ts index da255066a34..ccdcf5fb7b1 100644 --- a/src/main/runtime/rpc/methods/terminal-manifest-characterization.test.ts +++ b/src/main/runtime/rpc/methods/terminal-manifest-characterization.test.ts @@ -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 { 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() } diff --git a/src/main/runtime/rpc/methods/terminal/terminal-query-methods.ts b/src/main/runtime/rpc/methods/terminal/terminal-query-methods.ts index 96b4455d7dc..82edd55cd79 100644 --- a/src/main/runtime/rpc/methods/terminal/terminal-query-methods.ts +++ b/src/main/runtime/rpc/methods/terminal/terminal-query-methods.ts @@ -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, diff --git a/src/main/runtime/runtime-terminal-agent-presence.ts b/src/main/runtime/runtime-terminal-agent-presence.ts index e87520fccc8..aad71c1a08a 100644 --- a/src/main/runtime/runtime-terminal-agent-presence.ts +++ b/src/main/runtime/runtime-terminal-agent-presence.ts @@ -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 { + // 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) { diff --git a/src/main/runtime/structured-agent-session-close.ts b/src/main/runtime/structured-agent-session-close.ts new file mode 100644 index 00000000000..a7e6af191a2 --- /dev/null +++ b/src/main/runtime/structured-agent-session-close.ts @@ -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 { + 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 } +} diff --git a/src/main/runtime/structured-session-worktree-teardown.test.ts b/src/main/runtime/structured-session-worktree-teardown.test.ts new file mode 100644 index 00000000000..d3af1d0300c --- /dev/null +++ b/src/main/runtime/structured-session-worktree-teardown.test.ts @@ -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 +}): { 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 + }) + }) +}) diff --git a/src/main/runtime/structured-session-worktree-teardown.ts b/src/main/runtime/structured-session-worktree-teardown.ts new file mode 100644 index 00000000000..15801835f89 --- /dev/null +++ b/src/main/runtime/structured-session-worktree-teardown.ts @@ -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 + 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 } +} diff --git a/src/main/runtime/terminal-identity-probe.test.ts b/src/main/runtime/terminal-identity-probe.test.ts new file mode 100644 index 00000000000..cf356bb1bbc --- /dev/null +++ b/src/main/runtime/terminal-identity-probe.test.ts @@ -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') + }) +}) diff --git a/src/main/runtime/terminal-identity-probe.ts b/src/main/runtime/terminal-identity-probe.ts new file mode 100644 index 00000000000..43aca59cb1b --- /dev/null +++ b/src/main/runtime/terminal-identity-probe.ts @@ -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 + } +} diff --git a/src/main/runtime/worktree-pty-surface-sweeps.ts b/src/main/runtime/worktree-pty-surface-sweeps.ts new file mode 100644 index 00000000000..4c663086a39 --- /dev/null +++ b/src/main/runtime/worktree-pty-surface-sweeps.ts @@ -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 +) => Promise<{ stopped: boolean; owner: boolean }> + +export async function sweepProviderByPrefix( + worktreeId: string, + provider: IPtyProvider, + deadline: number, + stopPty: ( + ptyId: string, + stop: () => Promise + ) => Promise<{ stopped: boolean; owner: boolean }>, + onPtyStopped?: (ptyId: string) => void, + failClosed = false +): Promise { + 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:` 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((count, value) => count + value, 0) +} + +export async function sweepRegistryForWorktree( + worktreeId: string, + localProvider: IPtyProvider, + deadline: number, + stopPty: ( + ptyId: string, + stop: () => Promise + ) => Promise<{ stopped: boolean; owner: boolean }>, + onPtyStopped?: (ptyId: string) => void +): Promise { + 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((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 */ + } +} diff --git a/src/main/runtime/worktree-teardown-deadline.ts b/src/main/runtime/worktree-teardown-deadline.ts new file mode 100644 index 00000000000..8dcbfb4b3dd --- /dev/null +++ b/src/main/runtime/worktree-teardown-deadline.ts @@ -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 +} diff --git a/src/main/runtime/worktree-teardown.ts b/src/main/runtime/worktree-teardown.ts index dfe3d4ae5a1..1ce60d3ee65 100644 --- a/src/main/runtime/worktree-teardown.ts +++ b/src/main/runtime/worktree-teardown.ts @@ -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>() 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 - ) => Promise<{ stopped: boolean; owner: boolean }>, - onPtyStopped?: (ptyId: string) => void, - failClosed = false -): Promise { - 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:` 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((count, value) => count + value, 0) -} - -async function sweepRegistryForWorktree( - worktreeId: string, - localProvider: IPtyProvider, - deadline: number, - stopPty: ( - ptyId: string, - stop: () => Promise - ) => Promise<{ stopped: boolean; owner: boolean }>, - onPtyStopped?: (ptyId: string) => void -): Promise { - 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((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 { + 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 +}