diff --git a/desktop/src/stores/chatStore.test.ts b/desktop/src/stores/chatStore.test.ts index 3a15534c..92e27425 100644 --- a/desktop/src/stores/chatStore.test.ts +++ b/desktop/src/stores/chatStore.test.ts @@ -1787,6 +1787,90 @@ describe('chatStore history mapping', () => { .toBeUndefined() }) + it.each(['Bash:0', 'root-shell-task'])( + 'restores a stopped shell on cold load with terminal toolUseId %s', + async (toolUseId) => { + vi.mocked(sessionsApi.getFullHistory).mockResolvedValueOnce({ + messages: [ + { + id: 'root-shell-use', + type: 'assistant', + timestamp: '2026-04-06T00:00:00.000Z', + content: [{ + type: 'tool_use', + id: 'Bash:0', + name: 'Bash', + input: { command: 'bun test', run_in_background: true }, + }], + }, + { + id: 'root-shell-result', + type: 'tool_result', + timestamp: '2026-04-06T00:00:01.000Z', + content: [{ + type: 'tool_result', + tool_use_id: 'Bash:0', + content: 'Command running in background with ID: root-shell-task', + }], + }, + ], + taskNotifications: [ + { + taskId: 'root-shell-task', + toolUseId, + status: 'stopped', + summary: 'Background task stopped', + timestamp: '2026-04-06T00:00:02.000Z', + }, + { + taskId: 'root-shell-task', + toolUseId: 'unrelated-shell-tool', + status: 'failed', + summary: 'A task ID match must not override an unrelated real tool anchor', + timestamp: '2026-04-06T00:00:03.000Z', + }, + { + taskId: 'root-shell-task', + toolUseId: 'root-shell-task', + ownerAgentId: 'child-agent', + status: 'failed', + summary: 'An owned child notification must not replace the root outcome', + timestamp: '2026-04-06T00:00:03.000Z', + }, + { + taskId: 'unjoined-child-task', + toolUseId: 'unjoined-child-task', + status: 'stopped', + summary: 'An unowned task with no root transcript anchor must stay excluded', + }, + ], + }) + useChatStore.setState({ + sessions: { + [TEST_SESSION_ID]: makeSession({ messages: [] }), + }, + }) + + await useChatStore.getState().loadHistory(TEST_SESSION_ID) + + const session = useChatStore.getState().sessions[TEST_SESSION_ID] + expect(session?.backgroundAgentTasks?.['root-shell-task']).toMatchObject({ + taskId: 'root-shell-task', + toolUseId: 'Bash:0', + status: 'stopped', + summary: 'Background task stopped', + }) + expect(session?.agentTaskNotifications?.['Bash:0']).toMatchObject({ + taskId: 'root-shell-task', + toolUseId: 'Bash:0', + status: 'stopped', + }) + expect(session?.agentTaskNotifications?.['root-shell-task']).toBeUndefined() + expect(session?.agentTaskNotifications?.['unrelated-shell-tool']).toBeUndefined() + expect(session?.backgroundAgentTasks?.['unjoined-child-task']).toBeUndefined() + }, + ) + it('does not assign an unjoined child notification to the root run', async () => { vi.mocked(sessionsApi.getFullHistory).mockResolvedValueOnce({ messages: [{ diff --git a/desktop/src/stores/chatStore.ts b/desktop/src/stores/chatStore.ts index e651cd3b..0a2387aa 100644 --- a/desktop/src/stores/chatStore.ts +++ b/desktop/src/stores/chatStore.ts @@ -2333,7 +2333,21 @@ async function fetchAndMapSessionHistory( // jobs and notifications into the parent rail after every reload. const rootRunMessages = messages.filter((message) => !message.parentToolUseId) const rootToolUseIds = transcriptToolUseIds(rootRunMessages) - const rootRunNotifications = (taskNotifications ?? []).filter( + const rootShellTasks = reconstructBackgroundShellTasks(rootRunMessages) + const rootRunNotifications = (taskNotifications ?? []).map((notification) => { + // After a server restart, a late Stop can only persist the task ID as its + // tool anchor. Recover the real anchor from this root run's Bash result so + // the terminal survives cold history loading and links to the right tool. + const rootShellTask = rootShellTasks[notification.taskId] + if ( + !notification.ownerAgentId && + notification.toolUseId === notification.taskId && + rootShellTask?.toolUseId + ) { + return { ...notification, toolUseId: rootShellTask.toolUseId } + } + return notification + }).filter( (notification) => ( !notification.ownerAgentId && ( // workflow_run_id predates owner_agent_id. Keep restoring those diff --git a/src/server/__tests__/backgroundTaskCleanup.test.ts b/src/server/__tests__/backgroundTaskCleanup.test.ts new file mode 100644 index 00000000..239a149c --- /dev/null +++ b/src/server/__tests__/backgroundTaskCleanup.test.ts @@ -0,0 +1,411 @@ +import { afterEach, describe, expect, it, mock, spyOn } from 'bun:test' +import { + __resetWebSocketHandlerStateForTests, + getSessionChatActivityState, + getSessionTurnState, + handleWebSocket, + stopSessionTurn, +} from '../ws/handler.js' +import { activeNonAgentTasks, authoritativeStoppedTaskIds, clearAgentRuntimeState } from '../ws/agentTaskState.js' +import { conversationService } from '../services/conversationService.js' +import { computerUseApprovalService } from '../services/computerUseApprovalService.js' +import { sessionService } from '../services/sessionService.js' +import * as teamPlanRuntime from '../services/teamPlanRuntime.js' +import { migrationMaintenance } from '../migrationMaintenance.js' + +async function flushMicrotasks() { + for (let index = 0; index < 30; index++) await Promise.resolve() +} + +function setup() { + const sessionId = `background-cleanup-${crypto.randomUUID()}` + const clients: any[] = [] + function client() { + const sent: any[] = [] + const ws: any = { + data: { sessionId, channel: 'client', clientKind: 'full', connectedAt: Date.now(), sdkToken: null, serverPort: 0, serverHost: '127.0.0.1' }, + send: mock((payload: string) => sent.push(JSON.parse(payload))), + close: mock(() => {}), + sent, + } + clients.push(ws) + return ws + } + const ws = client() + const callbacks = new Set<(message: any) => void>() + const timers: Array<{ id: number; callback: () => void; delay: number }> = [] + const cancelled = new Set() + let runtimeAlive = true + spyOn(globalThis, 'setTimeout').mockImplementation(((callback: () => void, delay = 0) => { + const id = timers.length + 1 + timers.push({ id, callback, delay }) + return id as any + }) as any) + spyOn(globalThis, 'clearTimeout').mockImplementation((id: any) => { cancelled.add(id) }) + spyOn(conversationService, 'hasSession').mockImplementation(() => runtimeAlive) + const stop = spyOn(conversationService, 'stopSession').mockImplementation(() => { runtimeAlive = false }) + const permissions = spyOn(conversationService, 'getPendingPermissionRequests').mockReturnValue([]) + const computerPermissions = spyOn(computerUseApprovalService, 'getPendingRequests').mockReturnValue([]) + spyOn(conversationService, 'onOutput').mockImplementation((_id, callback) => { callbacks.add(callback) }) + spyOn(conversationService, 'removeOutputCallback').mockImplementation((_id, callback) => { callbacks.delete(callback) }) + spyOn(conversationService, 'sendInterrupt').mockReturnValue(true) + spyOn(conversationService, 'sendMessage').mockResolvedValue(true) + const control = spyOn(conversationService, 'requestControl').mockResolvedValue({}) + const append = spyOn(sessionService, 'appendSessionTaskNotification').mockResolvedValue() + spyOn(sessionService, 'getCustomTitle').mockResolvedValue('Existing title') + const teamWork = spyOn(teamPlanRuntime, 'hasActiveTeamWorkForParent').mockReturnValue(false) + function dispatch(message: any) { for (const callback of [...callbacks]) callback(message) } + function startTask(id = 'bash-1', type = 'local_bash') { + dispatch({ type: 'system', subtype: 'task_started', task_id: id, tool_use_id: `tool-${id}`, task_type: type }) + } + function close() { handleWebSocket.close(ws, 1000, 'test disconnect') } + function ceiling() { return timers.findLast(timer => timer.delay === 31 * 60_000) } + function terminals(socket = ws) { + return socket.sent.filter((message: any) => message.type === 'system_notification' && + message.subtype === 'task_notification' && ['stopped', 'failed'].includes(message.data?.status)) + } + handleWebSocket.open(ws) + return { sessionId, ws, client, timers, cancelled, stop, permissions, computerPermissions, control, append, teamWork, dispatch, startTask, close, ceiling, terminals } +} + +afterEach(() => { + __resetWebSocketHandlerStateForTests() + migrationMaintenance.resetForTests() + mock.restore() +}) + +describe('background task cleanup boundaries', () => { + it('bounds a silent task from disconnect and converges both activity and runtime state', async () => { + const s = setup() + s.startTask() + s.close() + expect(s.ceiling()).toBeDefined() + s.ceiling()!.callback() + await flushMicrotasks() + expect(s.stop).toHaveBeenCalledWith(s.sessionId) + expect(s.append).toHaveBeenCalledWith(s.sessionId, expect.objectContaining({ taskId: 'bash-1', status: 'failed' })) + expect(getSessionChatActivityState(s.sessionId)).toBe('idle') + expect(getSessionTurnState(s.sessionId)).toBe('idle') + expect(activeNonAgentTasks.has(s.sessionId)).toBe(false) + const idleGrace = s.timers.findLast(timer => timer.delay === 30_000) + expect(idleGrace).toBeDefined() + idleGrace!.callback() + expect(authoritativeStoppedTaskIds.has(s.sessionId)).toBe(false) + }) + + it('starts the shell ceiling after the real foreground turn and CLI run both settle', async () => { + const s = setup() + handleWebSocket.message(s.ws, JSON.stringify({ type: 'user_message', content: 'work' })) + await flushMicrotasks() + s.dispatch({ type: 'system', subtype: 'session_state_changed', state: 'running' }) + s.startTask() + s.close() + expect(s.ceiling()).toBeUndefined() + s.dispatch({ type: 'system', subtype: 'session_state_changed', state: 'idle' }) + expect(s.ceiling()).toBeUndefined() + s.dispatch({ type: 'result', subtype: 'success', is_error: false, result: 'done' }) + expect(s.ceiling()).toBeDefined() + expect(s.stop).not.toHaveBeenCalled() + }) + + it('does not reap a foreground run that starts after the shell ceiling was queued', async () => { + const s = setup() + s.startTask() + s.close() + const old = s.ceiling()! + s.dispatch({ type: 'system', subtype: 'session_state_changed', state: 'running' }) + old.callback() + expect(s.stop).not.toHaveBeenCalled() + expect(getSessionTurnState(s.sessionId)).toBe('running') + s.dispatch({ type: 'system', subtype: 'session_state_changed', state: 'idle' }) + expect(s.ceiling()!.id).not.toBe(old.id) + }) + + it('keeps Agent tasks and approved working teams alive alongside a shell', () => { + const s = setup() + s.startTask() + s.startTask('agent-1', 'local_agent') + s.close() + expect(s.ceiling()).toBeUndefined() + s.dispatch({ type: 'system', subtype: 'task_notification', task_id: 'agent-1', status: 'completed' }) + const beforeTeam = s.ceiling()! + expect(beforeTeam).toBeDefined() + s.teamWork.mockReturnValue(true) + beforeTeam.callback() + expect(s.stop).not.toHaveBeenCalled() + }) + + it('starts the shell ceiling after a detached Agent finishes its durable stop', async () => { + const s = setup() + s.startTask() + s.startTask('agent-1', 'local_agent') + await flushMicrotasks() + let saveAgent!: () => void + s.append.mockImplementation((_id, notification) => notification.taskId === 'agent-1' && notification.status === 'stopped' + ? new Promise(resolve => { saveAgent = resolve }) : Promise.resolve()) + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_background_task', taskId: 'agent-1' })) + await flushMicrotasks() + s.close() + expect(s.ceiling()).toBeUndefined() + saveAgent() + await flushMicrotasks() + expect(s.ceiling()).toBeDefined() + expect(s.stop).not.toHaveBeenCalled() + }) + + it('does not use the shell ceiling to shorten a pending tool or computer-use prompt', () => { + const s = setup() + s.startTask() + s.close() + const original = s.ceiling()! + s.permissions.mockReturnValue([{ requestId: 'tool-permission', request: { subtype: 'can_use_tool' } }] as any) + original.callback() + expect(s.stop).not.toHaveBeenCalled() + s.permissions.mockReturnValue([]) + s.dispatch({ type: 'system', subtype: 'session_state_changed', state: 'idle' }) + const next = s.ceiling()! + s.computerPermissions.mockReturnValue([{ requestId: 'computer-permission' }] as any) + next.callback() + expect(s.stop).not.toHaveBeenCalled() + }) + + it('persists shell bookends when an abandoned permission reaches its own cleanup bound', async () => { + const s = setup() + s.startTask() + s.permissions.mockReturnValue([{ requestId: 'abandoned-permission', request: { subtype: 'can_use_tool' } }] as any) + s.close() + // This is the existing permission ceiling, not the background-only ceiling. + const permissionBound = s.timers.find(timer => timer.delay === 31 * 60_000)! + expect(permissionBound).toBeDefined() + permissionBound.callback() + s.permissions.mockReturnValue([]) + await flushMicrotasks() + expect(s.stop).toHaveBeenCalledTimes(1) + expect(s.append).toHaveBeenCalledWith(s.sessionId, expect.objectContaining({ taskId: 'bash-1', status: 'failed' })) + expect(getSessionTurnState(s.sessionId)).toBe('idle') + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).not.toBe(true) + expect(s.timers.findLast(timer => timer.delay === 30_000)).toBeDefined() + }) + + it('does not kill or persist background tasks during migration maintenance', () => { + const s = setup() + s.startTask() + s.close() + const timer = s.ceiling()! + s.append.mockClear() + migrationMaintenance.begin() + timer.callback() + expect(s.stop).not.toHaveBeenCalled() + expect(s.append).not.toHaveBeenCalled() + }) + + it('ignores a queued old ceiling after reconnect and a second disconnect', () => { + const s = setup() + s.startTask() + s.close() + const old = s.ceiling()! + const reconnected = s.client() + handleWebSocket.open(reconnected) + expect(s.cancelled.has(old.id)).toBe(true) + handleWebSocket.close(reconnected, 1000, 'second disconnect') + const current = s.ceiling()! + old.callback() + expect(s.stop).not.toHaveBeenCalled() + current.callback() + expect(s.stop).toHaveBeenCalledTimes(1) + }) + + it('protects a real replacement turn from the previous Stop timeout', async () => { + const s = setup() + s.control.mockImplementation(() => new Promise(() => {})) + s.startTask() + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_generation' })) + const previousStop = s.timers.find(timer => timer.delay === 3_000)! + handleWebSocket.message(s.ws, JSON.stringify({ type: 'user_message', content: 'new work' })) + await flushMicrotasks() + previousStop.callback() + expect(s.stop).not.toHaveBeenCalled() + expect(getSessionTurnState(s.sessionId)).toBe('running') + s.startTask('new-shell') + expect(s.control.mock.calls.filter(call => call[1]?.task_id === 'new-shell')).toHaveLength(0) + }) + + it('stops a late shell once while the original user Stop still owns its scope', async () => { + const s = setup() + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_generation' })) + s.startTask('late-shell') + s.dispatch({ type: 'system', subtype: 'task_notification', task_id: 'late-shell', status: 'running' }) + expect(s.control.mock.calls).toEqual([[s.sessionId, { subtype: 'stop_task', task_id: 'late-shell' }]]) + s.timers.find(timer => timer.delay === 3_000)!.callback() + await flushMicrotasks() + expect(s.stop).toHaveBeenCalledTimes(1) + expect(s.terminals()).toHaveLength(1) + }) + + it('keeps programmatic turn Stop narrow for independent shell tasks', () => { + const s = setup() + s.startTask() + stopSessionTurn(s.sessionId) + s.startTask('later-shell') + expect(s.control).not.toHaveBeenCalled() + }) + + it('persists not_found before dropping tracking or publishing a terminal notification', async () => { + const s = setup() + s.startTask() + await flushMicrotasks() + let save!: () => void + s.append.mockImplementation((_id, notification) => notification.status === 'stopped' + ? new Promise(resolve => { save = resolve }) : Promise.resolve()) + s.control.mockResolvedValue({ reason: 'not_found' }) + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_background_task', taskId: 'bash-1' })) + await flushMicrotasks() + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).toBe(true) + expect(authoritativeStoppedTaskIds.get(s.sessionId)?.has('bash-1')).not.toBe(true) + expect(s.terminals()).toHaveLength(0) + s.dispatch({ type: 'system', subtype: 'task_notification', task_id: 'bash-1', status: 'running' }) + save() + await flushMicrotasks() + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).not.toBe(true) + expect(s.terminals()).toHaveLength(1) + }) + + it('retries failed terminal persistence without reviving or forgetting an exited task', async () => { + const s = setup() + s.startTask() + await flushMicrotasks() + s.append.mockRejectedValue(new Error('disk unavailable')) + s.close() + s.ceiling()!.callback() + await flushMicrotasks() + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).toBe(true) + expect(authoritativeStoppedTaskIds.get(s.sessionId)?.has('bash-1')).not.toBe(true) + expect(getSessionChatActivityState(s.sessionId)).toBe('idle') + // The ordinary cleanup must retain the pending terminal record as well. + s.timers.findLast(timer => timer.delay === 30_000)!.callback() + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).toBe(true) + s.append.mockResolvedValue() + const returning = s.client() + handleWebSocket.open(returning) + await flushMicrotasks() + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).not.toBe(true) + expect(s.terminals(returning)).toHaveLength(1) + }) + + for (const wholeSession of [false, true]) { + for (const response of [{}, { reason: 'not_found' }]) { + it(`ignores a stale ${wholeSession ? 'session' : 'task'} Stop response after its task record is replaced (${JSON.stringify(response)})`, async () => { + const s = setup() + let resolveControl!: (response: any) => void + s.control.mockImplementation(() => new Promise(resolve => { resolveControl = resolve })) + s.startTask() + handleWebSocket.message(s.ws, JSON.stringify(wholeSession + ? { type: 'stop_generation' } + : { type: 'stop_background_task', taskId: 'bash-1' })) + // A runtime replacement releases the old registry; the replacement CLI + // is allowed to reuse a handle while an old request callback is queued. + clearAgentRuntimeState(s.sessionId) + s.startTask() + const replacement = activeNonAgentTasks.get(s.sessionId)!.get('bash-1') + resolveControl(response) + await flushMicrotasks() + expect(activeNonAgentTasks.get(s.sessionId)!.get('bash-1')).toBe(replacement) + expect(s.terminals()).toHaveLength(0) + }) + } + } + + it('treats a successful stop control response as confirmation even without a terminal frame', async () => { + const s = setup() + s.startTask() + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_generation' })) + await flushMicrotasks() + expect(s.terminals().map((message: any) => message.data.status)).toEqual(['stopped']) + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).not.toBe(true) + s.timers.find(timer => timer.delay === 3_000)!.callback() + expect(s.stop).not.toHaveBeenCalled() + }) + + it('bounds a hung terminal write and can retry it after the timeout', async () => { + const s = setup() + s.startTask() + await flushMicrotasks() + s.append.mockImplementation(() => new Promise(() => {})) + s.close() + s.ceiling()!.callback() + for (let attempt = 0; attempt < 3; attempt++) { + const timeout = s.timers.filter(timer => timer.delay === 1_000)[attempt] + expect(timeout).toBeDefined() + timeout.callback() + await flushMicrotasks() + } + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).toBe(true) + expect(authoritativeStoppedTaskIds.get(s.sessionId)?.has('bash-1')).not.toBe(true) + s.append.mockResolvedValue() + s.timers.find(timer => timer.delay === 250)!.callback() + await flushMicrotasks() + expect(activeNonAgentTasks.get(s.sessionId)?.has('bash-1')).not.toBe(true) + expect(s.append.mock.calls.filter(call => call[1].status === 'failed')).toHaveLength(4) + }) + + it('does not infer task termination from a crash merely because Stop was requested', async () => { + const s = setup() + s.control.mockImplementation(() => new Promise(() => {})) + s.startTask() + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_generation' })) + s.stop(s.sessionId) + s.dispatch({ type: 'result', subtype: 'error', is_error: true, result: 'runtime crashed before stop acknowledgement' }) + await flushMicrotasks() + expect(s.terminals().map((message: any) => message.data.status)).toEqual(['failed']) + }) + + it('preserves an acknowledged stop outcome while persistence and a crash race', async () => { + const s = setup() + s.startTask() + await flushMicrotasks() + let save!: () => void + s.append.mockImplementation((_id, notification) => notification.status === 'stopped' + ? new Promise(resolve => { save = resolve }) : Promise.resolve()) + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_generation' })) + await flushMicrotasks() + s.stop(s.sessionId) + s.dispatch({ type: 'result', subtype: 'error', is_error: true, result: 'runtime crashed after confirmed stop' }) + save() + await flushMicrotasks() + expect(s.terminals().map((message: any) => message.data.status)).toEqual(['stopped']) + expect(s.append.mock.calls.filter(call => call[1].status === 'failed')).toHaveLength(0) + }) + + it('records an unacknowledged Windows force-stop as failed rather than claiming the shell died', async () => { + const descriptor = Object.getOwnPropertyDescriptor(process, 'platform')! + Object.defineProperty(process, 'platform', { ...descriptor, value: 'win32' }) + try { + const s = setup() + s.control.mockImplementation(() => new Promise(() => {})) + s.startTask() + handleWebSocket.message(s.ws, JSON.stringify({ type: 'stop_generation' })) + s.timers.find(timer => timer.delay === 3_000)!.callback() + await flushMicrotasks() + s.startTask('late-shell') + await flushMicrotasks() + expect(s.terminals().map((message: any) => message.data.status)).toEqual(['failed', 'failed']) + } finally { + Object.defineProperty(process, 'platform', descriptor) + } + }) + + it('records unrequested runtime loss as failed and closes late task records consistently', async () => { + const s = setup() + s.startTask() + s.stop(s.sessionId) + s.dispatch({ type: 'result', subtype: 'error', is_error: true, result: 'runtime crashed' }) + await flushMicrotasks() + s.startTask('late-shell') + await flushMicrotasks() + expect(s.terminals().map((message: any) => message.data.status)).toEqual(['failed', 'failed']) + expect(s.append).toHaveBeenCalledWith(s.sessionId, expect.objectContaining({ + taskId: 'bash-1', status: 'failed', summary: expect.stringContaining('could not be confirmed'), + })) + expect(getSessionChatActivityState(s.sessionId)).toBe('idle') + }) +}) diff --git a/src/server/__tests__/headlessBackgroundTaskCleanup.test.ts b/src/server/__tests__/headlessBackgroundTaskCleanup.test.ts new file mode 100644 index 00000000..d7d652a5 --- /dev/null +++ b/src/server/__tests__/headlessBackgroundTaskCleanup.test.ts @@ -0,0 +1,109 @@ +import { afterEach, expect, test, mock, spyOn } from 'bun:test' +import { + __resetWebSocketHandlerStateForTests, + ensureCliSessionStartedForControl, + getSessionChatActivityState, + getSessionTurnState, +} from '../ws/handler.js' +import { activeNonAgentTasks } from '../ws/agentTaskState.js' +import { conversationService } from '../services/conversationService.js' +import { computerUseApprovalService } from '../services/computerUseApprovalService.js' +import { sessionService } from '../services/sessionService.js' +import { SettingsService } from '../services/settingsService.js' +import { ProviderService } from '../services/providerService.js' +import { migrationMaintenance } from '../migrationMaintenance.js' + +async function flushMicrotasks() { + for (let index = 0; index < 30; index++) await Promise.resolve() +} + +function setup() { + const sessionId = `headless-background-${crypto.randomUUID()}` + const callbacks = new Set<(message: any) => void>() + const timers: Array<{ callback: () => void; delay: number }> = [] + let runtimeAlive = false + spyOn(conversationService, 'hasSession').mockImplementation(() => runtimeAlive) + const start = spyOn(conversationService, 'startSession').mockImplementation(async () => { runtimeAlive = true }) + const stop = spyOn(conversationService, 'stopSession').mockImplementation(() => { runtimeAlive = false }) + spyOn(conversationService, 'onOutput').mockImplementation((_id, callback) => { callbacks.add(callback) }) + spyOn(conversationService, 'removeOutputCallback').mockImplementation((_id, callback) => { callbacks.delete(callback) }) + spyOn(conversationService, 'getPendingPermissionRequests').mockReturnValue([]) + spyOn(computerUseApprovalService, 'getPendingRequests').mockReturnValue([]) + spyOn(sessionService, 'getSessionWorkDir').mockResolvedValue('/tmp/headless-control-fixture') + spyOn(sessionService, 'getSessionLaunchInfo').mockResolvedValue({ + runtimeModelId: 'fixture-model', runtimeProviderId: 'fixture-provider', permissionMode: 'default', + } as any) + spyOn(ProviderService.prototype, 'listProviders').mockResolvedValue({ + activeId: 'fixture-provider', providers: [{ id: 'fixture-provider' }], + } as any) + spyOn(ProviderService.prototype, 'getProvider').mockResolvedValue(null) + spyOn(SettingsService.prototype, 'getUserSettings').mockResolvedValue({}) + const append = spyOn(sessionService, 'appendSessionTaskNotification').mockResolvedValue() + spyOn(globalThis, 'setTimeout').mockImplementation(((callback: () => void, delay = 0) => { + timers.push({ callback, delay }) + return timers.length as any + }) as any) + spyOn(globalThis, 'clearTimeout').mockImplementation(() => {}) + function dispatch(message: any) { for (const callback of [...callbacks]) callback(message) } + function task(id = 'headless-shell') { + dispatch({ type: 'system', subtype: 'task_started', task_id: id, tool_use_id: `tool-${id}`, task_type: 'local_bash' }) + } + const ensure = () => ensureCliSessionStartedForControl(sessionId, new URL('http://127.0.0.1:12345/api/sessions/control')) + return { sessionId, callbacks, timers, start, stop, append, dispatch, task, ensure } +} + +afterEach(() => { + __resetWebSocketHandlerStateForTests() + migrationMaintenance.resetForTests() + mock.restore() +}) + +test('a control-started session settles a runtime crash without ever opening a client socket', async () => { + const s = setup() + await s.ensure() + expect(s.start).toHaveBeenCalledTimes(1) + s.task() + expect(activeNonAgentTasks.get(s.sessionId)?.has('headless-shell')).toBe(true) + s.stop(s.sessionId) + s.dispatch({ type: 'result', is_error: true, result: 'CLI died' }) + await flushMicrotasks() + expect(s.append).toHaveBeenCalledWith(s.sessionId, expect.objectContaining({ taskId: 'headless-shell', status: 'failed' })) + expect(activeNonAgentTasks.has(s.sessionId)).toBe(false) + expect(getSessionTurnState(s.sessionId)).toBe('idle') + expect(getSessionChatActivityState(s.sessionId)).not.toBe('running') +}) + +test('a control-started silent shell gets a ceiling without a client disconnect event', async () => { + const s = setup() + await s.ensure() + expect(s.timers.some(timer => timer.delay === 30_000)).toBe(false) + s.task() + const ceiling = s.timers.find(timer => timer.delay === 31 * 60_000) + expect(ceiling).toBeDefined() + ceiling!.callback() + await flushMicrotasks() + expect(s.stop).toHaveBeenCalledWith(s.sessionId) + expect(s.append).toHaveBeenCalledWith(s.sessionId, expect.objectContaining({ taskId: 'headless-shell', status: 'failed' })) +}) + +test('initial CLI idle does not arm cleanup ahead of the first headless control request', async () => { + const s = setup() + await s.ensure() + s.dispatch({ type: 'system', subtype: 'session_state_changed', state: 'idle' }) + expect(s.timers.some(timer => timer.delay === 30_000)).toBe(false) + expect(s.stop).not.toHaveBeenCalled() +}) + +test('the independent observer does not persist stale task results after its durable runtime-loss bookend', async () => { + const s = setup() + await s.ensure() + s.task() + s.stop(s.sessionId) + s.dispatch({ type: 'result', is_error: true, result: 'CLI died' }) + await flushMicrotasks() + expect(s.append).toHaveBeenCalledWith(s.sessionId, expect.objectContaining({ taskId: 'headless-shell', status: 'failed' })) + s.append.mockClear() + s.dispatch({ type: 'system', subtype: 'task_notification', task_id: 'headless-shell', tool_use_id: 'tool-headless-shell', status: 'completed' }) + await flushMicrotasks() + expect(s.append).not.toHaveBeenCalled() +}) diff --git a/src/server/__tests__/sessionStateCleanup.test.ts b/src/server/__tests__/sessionStateCleanup.test.ts index 07efd0ed..13c7cb3d 100644 --- a/src/server/__tests__/sessionStateCleanup.test.ts +++ b/src/server/__tests__/sessionStateCleanup.test.ts @@ -45,7 +45,9 @@ const CONTAINERS: Record = { activeNonAgentTasks: { kind: 'cleared' }, activeUserTurns: { kind: 'cleared' }, agentStopRequestedSessions: { kind: 'cleared' }, + nonAgentStopRequestedSessions: { kind: 'cleared' }, authoritativeStoppedTaskIds: { kind: 'cleared' }, + backgroundTaskCleanupTimers: { kind: 'cleared' }, deferredPermissionModes: { kind: 'cleared' }, deferredRuntimeRestarts: { kind: 'cleared' }, interruptedSessionChats: { kind: 'cleared' }, @@ -57,6 +59,7 @@ const CONTAINERS: Record = { prewarmedSessions: { kind: 'cleared' }, rejectedRuntimeConfigs: { kind: 'cleared' }, runtimeExitStoppedSessions: { kind: 'cleared' }, + runtimeExitFailedSessions: { kind: 'cleared' }, runtimeOverrides: { kind: 'cleared' }, runtimeTransitionPromises: { kind: 'cleared' }, sessionDisconnectWatchers: { kind: 'cleared' }, diff --git a/src/server/__tests__/websocket-handler.test.ts b/src/server/__tests__/websocket-handler.test.ts index 0e2e4474..f2b04b7f 100644 --- a/src/server/__tests__/websocket-handler.test.ts +++ b/src/server/__tests__/websocket-handler.test.ts @@ -1331,7 +1331,7 @@ describe('WebSocket handler session isolation', () => { await flushMicrotasks() expect(sendInterrupt).toHaveBeenCalledWith(sessionId) - expect(requestControl).toHaveBeenCalledTimes(4) + expect(requestControl).toHaveBeenCalledTimes(6) for (const taskId of [ 'agent-task-1', 'agent-task-2', @@ -1343,6 +1343,15 @@ describe('WebSocket handler session isolation', () => { task_id: taskId, }, 3_000) } + // A user Stop also reaps the non-Agent tasks sharing the runtime, so their + // shell processes are actually asked to stop rather than only being + // bookended when the runtime happens to exit. + for (const taskId of ['bash-collateral-task', 'provider-analyzer-teammate']) { + expect(requestControl).toHaveBeenCalledWith(sessionId, { + subtype: 'stop_task', + task_id: taskId, + }) + } expect(archiveRemoteSession).toHaveBeenCalledWith( 'remote-session-agent-task-3', { timeoutMs: 1_500 }, @@ -1355,13 +1364,13 @@ describe('WebSocket handler session isolation', () => { expect(append).toHaveBeenCalledWith(sessionId, expect.objectContaining({ taskId: 'bash-collateral-task', toolUseId: 'bash-collateral-tool', - status: 'stopped', + status: 'failed', })) expect(append).toHaveBeenCalledWith(sessionId, expect.objectContaining({ taskId: 'provider-analyzer-teammate', toolUseId: 'provider-analyzer-teammate-tool', ownerAgentId: 'provider-analyzer', - status: 'stopped', + status: 'failed', })) expect(ws.sent.map((payload) => JSON.parse(payload))).toContainEqual({ type: 'system_notification', @@ -1376,7 +1385,7 @@ describe('WebSocket handler session isolation', () => { subtype: 'task_notification', data: expect.objectContaining({ task_id: 'bash-collateral-task', - status: 'stopped', + status: 'failed', }), }) expect(ws.sent.map((payload) => JSON.parse(payload))).toContainEqual({ @@ -1385,7 +1394,7 @@ describe('WebSocket handler session isolation', () => { data: expect.objectContaining({ task_id: 'provider-analyzer-teammate', owner_agent_id: 'provider-analyzer', - status: 'stopped', + status: 'failed', }), }) @@ -1411,19 +1420,19 @@ describe('WebSocket handler session isolation', () => { notification.taskId === 'bash-collateral-task')).toHaveLength(1) expect(append).toHaveBeenCalledWith(sessionId, expect.objectContaining({ taskId: 'late-bash-after-force-kill', - status: 'stopped', + status: 'failed', })) expect(ws.sent.map((payload) => JSON.parse(payload))).toContainEqual({ type: 'system_notification', subtype: 'task_notification', data: expect.objectContaining({ task_id: 'late-bash-after-force-kill', - status: 'stopped', + status: 'failed', }), }) }) - it('does not bulk-stop non-Agent background tasks', async () => { + it('bulk-stops non-Agent background tasks on user Stop', async () => { const sessionId = `stop-agent-filter-${crypto.randomUUID()}` const ws = makeClientSocket(sessionId) const outputCallbacks: Array<(cliMsg: any) => void> = [] @@ -1453,8 +1462,13 @@ describe('WebSocket handler session isolation', () => { handleWebSocket.message(ws, JSON.stringify({ type: 'stop_generation' })) await Promise.resolve() + // A user Stop ends the whole session: the Agent task is stopped through the + // Agent finalization path (with its 3s control timeout), and the non-Agent + // shell/dream tasks are reaped through the same control channel without one. expect(requestControl.mock.calls).toEqual([ [sessionId, { subtype: 'stop_task', task_id: 'agent-task-1' }, 3_000], + [sessionId, { subtype: 'stop_task', task_id: 'bash-task-1' }], + [sessionId, { subtype: 'stop_task', task_id: 'dream-task-1' }], ]) outputCallbacks[0]?.({ type: 'system', @@ -2667,14 +2681,14 @@ describe('WebSocket handler session isolation', () => { expect(append).toHaveBeenCalledWith(sessionId, expect.objectContaining({ taskId: 'bash-clear-failed', toolUseId: 'bash-clear-failed-tool', - status: 'stopped', + status: 'failed', })) expect(first.sent.map((payload) => JSON.parse(payload))).toContainEqual({ type: 'system_notification', subtype: 'task_notification', data: expect.objectContaining({ task_id: 'bash-clear-failed', - status: 'stopped', + status: 'failed', }), }) }) @@ -4766,7 +4780,8 @@ describe('WebSocket handler session isolation', () => { summary: 'Desktop verification passed', }) - expect(setTimeoutSpy).not.toHaveBeenCalled() + expect(setTimeoutSpy).toHaveBeenCalledTimes(1) + expect(setTimeoutSpy.mock.calls[0]?.[1]).toBe(31 * 60_000) expect(stopSession).not.toHaveBeenCalled() outputCallbacks[1]?.({ @@ -4779,11 +4794,11 @@ describe('WebSocket handler session isolation', () => { summary: 'Focused tests passed', }) - expect(setTimeoutSpy).toHaveBeenCalledTimes(1) - expect(setTimeoutSpy.mock.calls[0]?.[1]).toBe(30_000) + expect(setTimeoutSpy).toHaveBeenCalledTimes(2) + expect(setTimeoutSpy.mock.calls[1]?.[1]).toBe(30_000) expect(stopSession).not.toHaveBeenCalled() - const expireIdleGrace = setTimeoutSpy.mock.calls[0]?.[0] as (() => void) | undefined + const expireIdleGrace = setTimeoutSpy.mock.calls[1]?.[0] as (() => void) | undefined expireIdleGrace?.() expect(stopSession).toHaveBeenCalledWith(sessionId) }) @@ -5153,6 +5168,196 @@ describe('WebSocket handler session isolation', () => { activeBackgroundTaskIds: [], }) }) + + it('reaps a non-Agent background task when the user stops a background-only session', async () => { + const sessionId = `stop-background-only-${crypto.randomUUID()}` + const ws = makeClientSocket(sessionId) + const outputCallbacks: Array<(cliMsg: any) => void> = [] + spyOn(globalThis, 'setTimeout').mockImplementation(() => 1 as any) + spyOn(conversationService, 'hasSession').mockReturnValue(true) + spyOn(conversationService, 'onOutput').mockImplementation((_sid, callback) => { + outputCallbacks.push(callback) + }) + spyOn(conversationService, 'removeOutputCallback').mockImplementation(() => {}) + // The CLI reports the task as already evicted, so the Stop converges on a + // terminal bookend rather than waiting for a notification that never comes. + const requestControl = spyOn(conversationService, 'requestControl') + .mockResolvedValue({ reason: 'not_found' }) + + handleWebSocket.open(ws) + outputCallbacks[0]?.({ + type: 'system', + subtype: 'task_started', + task_id: 'orphan-bash-1', + tool_use_id: 'orphan-bash-tool-1', + description: 'Background shell that outlived the turn', + task_type: 'local_bash', + }) + await flushMicrotasks() + ws.sent.length = 0 + + // No foreground turn and no Agent task: before this fix the non-Agent task + // was simply never asked to stop. + handleWebSocket.message(ws, JSON.stringify({ type: 'stop_generation' })) + await flushMicrotasks() + + expect(requestControl).toHaveBeenCalledWith(sessionId, { + subtype: 'stop_task', + task_id: 'orphan-bash-1', + }) + expect(ws.sent.map((payload) => JSON.parse(payload))).toContainEqual({ + type: 'system_notification', + subtype: 'task_notification', + message: 'Background shell that outlived the turn stopped', + data: expect.objectContaining({ + task_id: 'orphan-bash-1', + tool_use_id: 'orphan-bash-tool-1', + status: 'stopped', + }), + }) + }) + + it('bounds a disconnected session kept alive only by a background task', () => { + const sessionId = `background-task-ceiling-${crypto.randomUUID()}` + const ws = makeClientSocket(sessionId) + let nextTimerId = 0 + const timers: Array<{ id: number; callback: () => void; delayMs: number }> = [] + spyOn(globalThis, 'setTimeout').mockImplementation(((callback: () => void, delayMs?: number) => { + const id = ++nextTimerId + timers.push({ id, callback, delayMs: delayMs ?? 0 }) + return id as any + }) as any) + const stopSession = spyOn(conversationService, 'stopSession').mockImplementation(() => {}) + const append = spyOn(sessionService, 'appendSessionTaskNotification').mockResolvedValue() + spyOn(conversationService, 'getPendingPermissionRequests').mockReturnValue([]) + spyOn(conversationService, 'hasSession').mockReturnValue(true) + const outputCallbacks: Array<(cliMsg: any) => void> = [] + spyOn(conversationService, 'onOutput').mockImplementation((_sid, callback) => { + outputCallbacks.push(callback) + }) + spyOn(conversationService, 'removeOutputCallback').mockImplementation(() => {}) + + handleWebSocket.open(ws) + outputCallbacks[0]?.({ + type: 'system', + subtype: 'task_started', + task_id: 'ceiling-bash-1', + tool_use_id: 'ceiling-bash-tool-1', + description: 'Never-ending background shell', + task_type: 'local_bash', + }) + handleWebSocket.close(ws, 1000, 'pet closed while a background task runs') + expect(stopSession).not.toHaveBeenCalled() + + // The watcher sees the background task still running with no client left. + outputCallbacks.at(-1)?.({ + type: 'system', + subtype: 'task_notification', + task_id: 'ceiling-bash-1', + tool_use_id: 'ceiling-bash-tool-1', + task_type: 'local_bash', + status: 'running', + }) + + const ceiling = timers.find((timer) => timer.delayMs === 31 * 60_000) + expect(ceiling).toBeDefined() + ceiling?.callback() + + expect(stopSession).toHaveBeenCalledWith(sessionId) + expect(append).toHaveBeenCalledWith(sessionId, expect.objectContaining({ + taskId: 'ceiling-bash-1', + toolUseId: 'ceiling-bash-tool-1', + status: 'failed', + })) + }) + + it('cancels the background-task ceiling when a client reconnects', () => { + const sessionId = `background-task-ceiling-reconnect-${crypto.randomUUID()}` + const ws = makeClientSocket(sessionId) + let nextTimerId = 0 + const timers: Array<{ id: number; callback: () => void; delayMs: number }> = [] + spyOn(globalThis, 'setTimeout').mockImplementation(((callback: () => void, delayMs?: number) => { + const id = ++nextTimerId + timers.push({ id, callback, delayMs: delayMs ?? 0 }) + return id as any + }) as any) + const clearTimeoutSpy = spyOn(globalThis, 'clearTimeout').mockImplementation(() => {}) + const stopSession = spyOn(conversationService, 'stopSession').mockImplementation(() => {}) + spyOn(conversationService, 'getPendingPermissionRequests').mockReturnValue([]) + spyOn(conversationService, 'hasSession').mockReturnValue(true) + const outputCallbacks: Array<(cliMsg: any) => void> = [] + spyOn(conversationService, 'onOutput').mockImplementation((_sid, callback) => { + outputCallbacks.push(callback) + }) + spyOn(conversationService, 'removeOutputCallback').mockImplementation(() => {}) + + handleWebSocket.open(ws) + outputCallbacks[0]?.({ + type: 'system', + subtype: 'task_started', + task_id: 'ceiling-bash-reconnect-1', + tool_use_id: 'ceiling-bash-reconnect-tool-1', + description: 'Still watched by a returning client', + task_type: 'local_bash', + }) + handleWebSocket.close(ws, 1000, 'pet closed while a background task runs') + outputCallbacks.at(-1)?.({ + type: 'system', + subtype: 'task_notification', + task_id: 'ceiling-bash-reconnect-1', + tool_use_id: 'ceiling-bash-reconnect-tool-1', + task_type: 'local_bash', + status: 'running', + }) + + const ceiling = timers.find((timer) => timer.delayMs === 31 * 60_000) + expect(ceiling).toBeDefined() + + const reconnected = makeClientSocket(sessionId) + handleWebSocket.open(reconnected) + expect(clearTimeoutSpy).toHaveBeenCalledWith(ceiling?.id) + + // Even an already-queued ceiling must not kill a session a client is using. + ceiling?.callback() + expect(stopSession).not.toHaveBeenCalled() + }) + + it('reaps residual background tasks when the runtime is already gone at disconnect', () => { + const sessionId = `dead-runtime-residual-${crypto.randomUUID()}` + const ws = makeClientSocket(sessionId) + spyOn(globalThis, 'setTimeout').mockImplementation(() => 0 as any) + const stopSession = spyOn(conversationService, 'stopSession').mockImplementation(() => {}) + const append = spyOn(sessionService, 'appendSessionTaskNotification').mockResolvedValue() + spyOn(conversationService, 'getPendingPermissionRequests').mockReturnValue([]) + const hasSession = spyOn(conversationService, 'hasSession').mockReturnValue(true) + const outputCallbacks: Array<(cliMsg: any) => void> = [] + spyOn(conversationService, 'onOutput').mockImplementation((_sid, callback) => { + outputCallbacks.push(callback) + }) + spyOn(conversationService, 'removeOutputCallback').mockImplementation(() => {}) + + handleWebSocket.open(ws) + outputCallbacks[0]?.({ + type: 'system', + subtype: 'task_started', + task_id: 'residual-bash-1', + tool_use_id: 'residual-bash-tool-1', + description: 'Task record left behind by a dead CLI', + task_type: 'local_bash', + }) + + // The CLI died before the renderer closed, so no completion event will ever + // arrive to clear this record. + hasSession.mockReturnValue(false) + handleWebSocket.close(ws, 1006, 'renderer closed after the CLI died') + + expect(append).toHaveBeenCalledWith(sessionId, expect.objectContaining({ + taskId: 'residual-bash-1', + toolUseId: 'residual-bash-tool-1', + status: 'failed', + })) + expect(stopSession).not.toHaveBeenCalled() + }) }) describe('prewarm idle timer active-turn guard (issue #865 follow-up)', () => { diff --git a/src/server/ws/agentTaskState.ts b/src/server/ws/agentTaskState.ts index 88698d1e..65c9549e 100644 --- a/src/server/ws/agentTaskState.ts +++ b/src/server/ws/agentTaskState.ts @@ -2,13 +2,13 @@ * Agent and background-task lifecycle state for the session WebSocket handler. * * Moved verbatim out of `handler.ts` as the first cut of that file. This slice was - * chosen because it is provably closed: these functions read and write only the six + * chosen because it is provably closed: these functions read and write only the lifecycle * containers declared here and call nothing outside them, so relocating them cannot * change behavior. The functions that emit stop bookends stayed behind on purpose — * they reach into the broadcast layer and the socket registry, so moving them would * have required injecting dependencies rather than moving code. * - * All six containers are per-session and released by `clearAgentRuntimeState`, which + * All lifecycle containers are per-session and released by `clearAgentRuntimeState`, which * `cleanupSessionRuntimeState` calls. `src/server/__tests__/sessionStateCleanup.test.ts` * follows that closure across this module boundary. * @@ -23,6 +23,15 @@ export type ActiveNonAgentTaskState = { toolUseId: string description?: string ownerAgentId?: string + terminalStatus?: 'stopped' | 'failed' + terminalSummary?: string + terminalMessage?: boolean + stopRequested?: boolean + localStopConfirmed?: boolean + finalizationRetryCount?: number + finalizationRetryTimer?: ReturnType + finalization?: Promise + stopFailureMessage?: string } export type ActiveAgentTaskState = { @@ -54,8 +63,12 @@ export const authoritativeStoppedTaskIds = new Map>() export const agentStopRequestedSessions = new Set() +export const nonAgentStopRequestedSessions = new Set() + export const runtimeExitStoppedSessions = new Set() +export const runtimeExitFailedSessions = new Set() + export type CliBackgroundTaskLifecycle = { taskId: string running: boolean @@ -123,6 +136,8 @@ export function untrackCliBackgroundTask(sessionId: string, taskId: string): voi if (sessionAgentTasks?.size === 0) activeAgentTasks.delete(sessionId) const sessionNonAgentTasks = activeNonAgentTasks.get(sessionId) + const nonAgentTask = sessionNonAgentTasks?.get(taskId) + if (nonAgentTask) clearAgentStopFinalizationRetry(nonAgentTask) sessionNonAgentTasks?.delete(taskId) if (sessionNonAgentTasks?.size === 0) activeNonAgentTasks.delete(sessionId) } @@ -139,16 +154,28 @@ export function clearAgentRuntimeState( ) : new Map() + const retryableNonAgentStops = options?.preserveRetryableStops + ? new Map([...(activeNonAgentTasks.get(sessionId)?.entries() ?? [])] + .filter(([, task]) => task.localStopConfirmed)) + : new Map() for (const task of activeAgentTasks.get(sessionId)?.values() ?? []) { clearAgentStopFinalizationRetry(task) } + for (const task of activeNonAgentTasks.get(sessionId)?.values() ?? []) { + clearAgentStopFinalizationRetry(task) + } activeBackgroundTaskIds.delete(sessionId) activeAgentTasks.delete(sessionId) activeNonAgentTasks.delete(sessionId) authoritativeStoppedTaskIds.delete(sessionId) agentStopRequestedSessions.delete(sessionId) + nonAgentStopRequestedSessions.delete(sessionId) runtimeExitStoppedSessions.delete(sessionId) + runtimeExitFailedSessions.delete(sessionId) + if (retryableNonAgentStops.size > 0) { + activeNonAgentTasks.set(sessionId, retryableNonAgentStops) + } if (retryableStops.size > 0) { activeAgentTasks.set(sessionId, retryableStops) activeBackgroundTaskIds.set(sessionId, new Set(retryableStops.keys())) @@ -171,11 +198,12 @@ export function hasActiveBackgroundTasks(sessionId: string): boolean { const sessionAgentTasks = activeAgentTasks.get(sessionId) return [...taskIds].some((taskId) => { const agentTask = sessionAgentTasks?.get(taskId) - return !agentTask || !agentTask.localStopConfirmed || agentTask.bookendPending + if (agentTask) return !agentTask.localStopConfirmed || agentTask.bookendPending + return !activeNonAgentTasks.get(sessionId)?.get(taskId)?.localStopConfirmed }) } -export function clearAgentStopFinalizationRetry(task: ActiveAgentTaskState): void { +export function clearAgentStopFinalizationRetry(task: { finalizationRetryTimer?: ReturnType }): void { if (task.finalizationRetryTimer === undefined) return clearTimeout(task.finalizationRetryTimer) task.finalizationRetryTimer = undefined diff --git a/src/server/ws/handler.ts b/src/server/ws/handler.ts index a75809f0..a1dfc4ff 100644 --- a/src/server/ws/handler.ts +++ b/src/server/ws/handler.ts @@ -84,7 +84,9 @@ import { activeNonAgentTasks, authoritativeStoppedTaskIds, agentStopRequestedSessions, + nonAgentStopRequestedSessions, runtimeExitStoppedSessions, + runtimeExitFailedSessions, getCliBackgroundTaskLifecycle, isAgentTaskType, untrackCliBackgroundTask, @@ -96,6 +98,7 @@ import { } from './agentTaskState.js' import type { ActiveAgentTaskState, + ActiveNonAgentTaskState, CliBackgroundTaskLifecycle, } from './agentTaskState.js' import { @@ -165,7 +168,15 @@ const sessionSlashCommands = new Map() // The longest automatic question wait is 30 minutes; leave room for its // bounded model request before reclaiming a disconnected CLI. const PENDING_PERMISSION_DISCONNECT_CLEANUP_MS = 31 * 60_000 +// A background shell task may legitimately outlive a disconnected client, but +// nothing else bounds one that never emits a terminal notification — so a +// forgotten run_in_background loop would pin the CLI (and its process group) +// forever. Cap that keep-alive at the same order as the permission bound; +// once it elapses with the client still gone, the shared runtime is stopped and +// terminal bookends are published. +const BACKGROUND_TASK_DISCONNECT_MAX_MS = 31 * 60_000 const sessionCleanupTimers = new Map>() +const backgroundTaskCleanupTimers = new Map>() /** * Per-session removers for the active-work watcher (issue #764). When the last * client disconnects while a turn or background task is still running, we let @@ -304,6 +315,12 @@ function trackCliBackgroundTaskLifecycle( const lifecycle = getCliBackgroundTaskLifecycle(cliMsg) if (!lifecycle) return null + const existingNonAgentTask = activeNonAgentTasks.get(sessionId)?.get(lifecycle.taskId) + if (existingNonAgentTask?.localStopConfirmed) { + void emitAuthoritativeNonAgentStopped(sessionId, existingNonAgentTask) + return { ...lifecycle, running: false, status: 'stopped', suppressForward: true } + } + const existingAgentTask = activeAgentTasks.get(sessionId)?.get(lifecycle.taskId) if ( lifecycle.running && @@ -394,6 +411,11 @@ function trackCliBackgroundTaskLifecycle( return { ...lifecycle, suppressForward: true } } + if (existingNonAgentTask?.stopRequested) { + existingNonAgentTask.localStopConfirmed = true + void emitAuthoritativeNonAgentStopped(sessionId, existingNonAgentTask) + return { ...lifecycle, suppressForward: true } + } untrackCliBackgroundTask(sessionId, lifecycle.taskId) return lifecycle } @@ -617,6 +639,7 @@ export const handleWebSocket = { // Cancel any "let the running turn finish, then clean up" watcher too — // the session is observed again (issue #764). cancelSessionDisconnectWatcher(sessionId) + clearBackgroundTaskDisconnectCeiling(sessionId) addActiveClient(sessionId, ws) if (prewarmPendingSessions.has(sessionId) || prewarmedSessions.has(sessionId)) { @@ -647,6 +670,9 @@ export const handleWebSocket = { turnActive: hasLiveUserTurnForClient(sessionId), }) replayAgentStopFailures(ws, sessionId) + for (const task of activeNonAgentTasks.get(sessionId)?.values() ?? []) { + if (task.localStopConfirmed) void emitAuthoritativeNonAgentStopped(sessionId, task) + } }, message(ws: SessionConnection, rawMessage: string | Buffer) { @@ -758,7 +784,7 @@ export const handleWebSocket = { break case 'stop_generation': - handleStopGeneration(ws) + handleStopGeneration(ws, { reapBackgroundTasks: true }) break case 'stop_background_task': @@ -804,6 +830,23 @@ export const handleWebSocket = { // closed. Defer cleanup until all active work completes, then apply the // idle grace period. Sessions that are already idle go straight to the timer. if (hasActiveSessionWork(sessionId)) { + // If the CLI runtime is already gone but turn/task bookkeeping lingers, + // no completion event will ever arrive to clear it. Publish terminal + // bookends now so a later reconnect does not re-hydrate a ghost "Running". + // A deliberate restart never reaches here: it stops the runtime without + // emitting an error result, and startSession re-creates the session. + if (!conversationService.hasSession(sessionId) && hasTrackedTaskRecords(sessionId)) { + const knownExit = runtimeExitStoppedSessions.has(sessionId) + runtimeExitStoppedSessions.add(sessionId) + if (!knownExit) runtimeExitFailedSessions.add(sessionId) + activeCliRuns.delete(sessionId) + const turn = activeUserTurns.get(sessionId) + if (turn && !sessionStartupPromises.has(sessionId)) clearActiveUserTurn(sessionId, turn) + void emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId) + void emitAuthoritativeStoppedForActiveAgents(sessionId) + scheduleDisconnectCleanup(sessionId) + return + } // A turn blocked on permission cannot finish without user input. Keep the // completion watcher for early cleanup, but also enforce the existing // pending-permission maximum so an abandoned prompt cannot pin the CLI. @@ -868,6 +911,8 @@ async function handleUserMessage( sessionStopRequested.has(sessionId) || agentStopRequestedSessions.has(sessionId) activeTurn.admissionPending = !collaboration activeUserTurns.set(sessionId, activeTurn) + nonAgentStopRequestedSessions.delete(sessionId) + clearBackgroundTaskDisconnectCeiling(sessionId) const content = await resolveSessionReferenceContext(message.content, message.sessionReferences, async id => Boolean(await sessionService.getSessionSummary(id))) @@ -1184,9 +1229,25 @@ function bindSessionTurnObserver(sessionId: string): void { // The observer is independent of renderer subscriptions, so background // sessions keep their task/permission lifecycle when no page is open. if (!hasActiveClients(sessionId)) { - trackCliRunState(sessionId, message) - trackCliBackgroundTaskLifecycle(sessionId, message) - persistThenForwardCliMessage(sessionId, message, () => {}) + const cliRunState = trackCliRunState(sessionId, message) + const taskLifecycle = trackCliBackgroundTaskLifecycle(sessionId, message) + stopLateAgentTaskIfRequested(sessionId, taskLifecycle) + stopLateNonAgentTaskIfRequested(sessionId, taskLifecycle) + closeLateNonAgentTaskAfterRuntimeExit(sessionId, taskLifecycle) + closeStoppedAgentsAfterRuntimeExit(sessionId, message) + if (!sessionStartupPromises.has(sessionId)) { + refreshBackgroundTaskDisconnectCeiling(sessionId) + // Control-only sessions never had a renderer disconnect. Observe their + // completion after work actually starts; an initial idle frame must + // not arm cleanup ahead of a side-question/Agent control request. + if (!sessionDisconnectWatchers.has(sessionId) && + (cliRunState === 'running' || taskLifecycle?.running === true)) { + watchTurnCompletionForCleanup(sessionId) + } + } + if (!taskLifecycle?.suppressForward) { + persistThenForwardCliMessage(sessionId, message, () => {}) + } } queueMicrotask(() => { if (sessionTurnObservers.get(sessionId) === callback) { @@ -1207,6 +1268,7 @@ function clearSessionTurnObserver(sessionId: string): void { function clearActiveUserTurn(sessionId: string, activeTurn: ActiveUserTurnState): void { if (activeUserTurns.get(sessionId) === activeTurn) { activeUserTurns.delete(sessionId) + refreshBackgroundTaskDisconnectCeiling(sessionId) } } @@ -1251,7 +1313,13 @@ function forceStopSharedRuntimeForAgentCancellation(sessionId: string): void { // slash command), otherwise its result can be consumed as the dead turn's. pendingInterruptedTurnResults.delete(sessionId) runtimeExitStoppedSessions.add(sessionId) + // Losing the CLI cannot confirm every shell descendant stopped: POSIX jobs + // can create their own process groups, and Windows has no lifetime guard. + runtimeExitFailedSessions.add(sessionId) conversationService.stopSession(sessionId) + activeCliRuns.delete(sessionId) + clearBackgroundTaskDisconnectCeiling(sessionId) + cancelSessionDisconnectWatcher(sessionId) const stoppedTurn = activeUserTurns.get(sessionId) if ( stoppedTurn?.cancelled && @@ -1259,7 +1327,9 @@ function forceStopSharedRuntimeForAgentCancellation(sessionId: string): void { ) { clearActiveUserTurn(sessionId, stoppedTurn) } - void emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId) + void emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId).then(() => { + scheduleDisconnectedSessionCleanupIfIdle(sessionId) + }) } function consumeInterruptedTurnResult(sessionId: string, cliMsg: any): boolean { @@ -1298,7 +1368,9 @@ function acknowledgeActiveTurnReplay(sessionId: string, cliMsg: any): boolean { pendingInterruptedTurnResults.delete(sessionId) sessionStopRequested.delete(sessionId) agentStopRequestedSessions.delete(sessionId) + nonAgentStopRequestedSessions.delete(sessionId) runtimeExitStoppedSessions.delete(sessionId) + runtimeExitFailedSessions.delete(sessionId) return true } @@ -2082,6 +2154,7 @@ async function restartSessionWithPermissionMode( await conversationService.startSession(sessionId, workDir, sdkUrl, runtimeSettings) if (!agentStopRequestedSessions.has(sessionId)) { runtimeExitStoppedSessions.delete(sessionId) + runtimeExitFailedSessions.delete(sessionId) } await commitConfirmedPermissionMode(sessionId, mode, workDir) @@ -2199,6 +2272,7 @@ async function restartSessionWithRuntimeConfig( const sdkUrl = buildSdkWebSocketUrl(ws, sessionId) await conversationService.startSession(sessionId, workDir, sdkUrl, runtimeSettings) runtimeExitStoppedSessions.delete(sessionId) + runtimeExitFailedSessions.delete(sessionId) broadcastAppliedRuntimeConfig(sessionId) sendMessage(ws, { type: 'status', state: 'idle' }) @@ -2227,11 +2301,23 @@ async function restartSessionWithRuntimeConfig( } } -function handleStopGeneration(ws: SessionConnection) { +function handleStopGeneration( + ws: SessionConnection, + options: { reapBackgroundTasks?: boolean } = {}, +) { const { sessionId } = ws.data emitSessionTurnEvent({ type: 'stopped', sessionId }) const stoppedTurn = activeUserTurns.get(sessionId) const agentTasks = [...(activeAgentTasks.get(sessionId)?.values() ?? [])] + // A user-issued Stop ends the whole session's work. Background shell tasks + // (Bash/PowerShell run_in_background) are tracked as non-Agent tasks, which + // the Agent-stop bookkeeping below never touches — left alone they keep + // running after the user asked everything to stop, and keep the CLI alive. + // Programmatic stops (stopSessionTurn, runtime-config restart) stay narrow. + const backgroundTasks = options.reapBackgroundTasks === true + ? [...(activeNonAgentTasks.get(sessionId)?.values() ?? [])] + : [] + const backgroundTaskIds = backgroundTasks.map((task) => task.taskId) console.log(`[WS] Stop generation requested for session: ${sessionId}`) if (stoppedTurn) { @@ -2273,6 +2359,11 @@ function handleStopGeneration(ws: SessionConnection) { agentTasks.map((task) => requestStopTrackedAgentTask(sessionId, task, ws)), ) + if (options.reapBackgroundTasks) nonAgentStopRequestedSessions.add(sessionId) + for (const taskId of backgroundTaskIds) { + stopTrackedBackgroundTask(sessionId, taskId) + } + if ( stoppedTurn && conversationService.hasSession(sessionId) && @@ -2294,7 +2385,10 @@ function handleStopGeneration(ws: SessionConnection) { conversationService.sendInterrupt(sessionId) } - if ((stoppedTurn || agentTasks.length > 0) && conversationService.hasSession(sessionId)) { + if ( + (stoppedTurn || agentTasks.length > 0 || backgroundTaskIds.length > 0) && + conversationService.hasSession(sessionId) + ) { // Force-kill if still running after 3 seconds setTimeout(() => { const stoppedForegroundStillCurrent = Boolean( @@ -2315,8 +2409,17 @@ function handleStopGeneration(ws: SessionConnection) { [...(activeAgentTasks.get(sessionId)?.values() ?? [])].some( (task) => !task.localStopConfirmed, ) + const currentTurn = activeUserTurns.get(sessionId) + const stoppedBackgroundStillActive = + nonAgentStopRequestedSessions.has(sessionId) && + (!currentTurn || currentTurn === stoppedTurn) && + backgroundTasks.some((task) => + activeNonAgentTasks.get(sessionId)?.get(task.taskId) === task && + !task.localStopConfirmed, + ) if ( - (stoppedForegroundStillCurrent || stoppedAgentsStillActive) && + (stoppedForegroundStillCurrent || stoppedAgentsStillActive || stoppedBackgroundStillActive) && + !migrationMaintenance.isActive && conversationService.hasSession(sessionId) ) { console.log(`[WS] Force-killing CLI subprocess for session: ${sessionId}`) @@ -2359,46 +2462,117 @@ async function requestStopBackgroundTask( return } + const tracked = activeNonAgentTasks.get(sessionId)?.get(taskId) + if (tracked?.localStopConfirmed) { + await emitAuthoritativeNonAgentStopped(sessionId, tracked) + return + } + if (tracked) tracked.stopRequested = true try { const response = await conversationService.requestControl(sessionId, { subtype: 'stop_task', task_id: taskId, }) + if (activeNonAgentTasks.get(sessionId)?.get(taskId) !== tracked) return if (response?.reason === 'not_found') { - convergeEvictedBackgroundTaskStop(sessionId, taskId) + await convergeEvictedBackgroundTaskStop(sessionId, taskId) + } else if (tracked && !tracked.localStopConfirmed && + activeNonAgentTasks.get(sessionId)?.get(taskId) === tracked) { + tracked.localStopConfirmed = true + tracked.terminalStatus = 'stopped' + await emitAuthoritativeNonAgentStopped(sessionId, tracked) } } catch (error) { reportBackgroundTaskStopFailure(sessionId, ws, taskId, error) } } +/** + * Stop a tracked background shell task without an originating socket, for a + * whole-session Stop. Failures are logged rather than answered to one client, + * and a late `not_found` is converged to a terminal bookend just like the + * panel's per-task Stop. + */ +function stopTrackedBackgroundTask(sessionId: string, taskId: string): void { + const task = activeNonAgentTasks.get(sessionId)?.get(taskId) + if (!task || task.stopRequested || task.localStopConfirmed) return + task.stopRequested = true + void migrationMaintenance.track(conversationService + .requestControl(sessionId, { subtype: 'stop_task', task_id: taskId }) + .then((response) => { + if (activeNonAgentTasks.get(sessionId)?.get(taskId) !== task) return + if (response?.reason === 'not_found') { + return convergeEvictedBackgroundTaskStop(sessionId, taskId) + } + if (!task.localStopConfirmed && activeNonAgentTasks.get(sessionId)?.get(taskId) === task) { + task.localStopConfirmed = true + task.terminalStatus = 'stopped' + return emitAuthoritativeNonAgentStopped(sessionId, task) + } + }) + .catch((error) => { + console.warn( + `[WS] Failed to stop background task ${taskId} for session ${sessionId}:`, + error, + ) + })) +} + +function stopLateNonAgentTaskIfRequested( + sessionId: string, + lifecycle: CliBackgroundTaskLifecycle | null, +): void { + if (!lifecycle?.running || isAgentTaskType(lifecycle.taskType) || + !nonAgentStopRequestedSessions.has(sessionId) || runtimeExitStoppedSessions.has(sessionId)) return + const task = activeNonAgentTasks.get(sessionId)?.get(lifecycle.taskId) + if (!task || task.stopRequested) return + if (!conversationService.hasSession(sessionId)) { + runtimeExitStoppedSessions.add(sessionId) + runtimeExitFailedSessions.add(sessionId) + void emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId, 'failed') + return + } + const stoppedTurn = activeUserTurns.get(sessionId) + stopTrackedBackgroundTask(sessionId, lifecycle.taskId) + // A queued shell start can arrive after the original Stop's task snapshot. + // Its fallback belongs to that Stop, never to a newly admitted user turn. + setTimeout(() => { + const currentTurn = activeUserTurns.get(sessionId) + if (migrationMaintenance.isActive || !nonAgentStopRequestedSessions.has(sessionId) || + (currentTurn && currentTurn !== stoppedTurn) || currentTurn?.replacementAfterStop || + activeNonAgentTasks.get(sessionId)?.get(task.taskId) !== task || + task.localStopConfirmed || !conversationService.hasSession(sessionId)) return + forceStopSharedRuntimeForAgentCancellation(sessionId) + void emitAuthoritativeStoppedForActiveAgents(sessionId) + }, 3_000) +} + /** * The CLI evicts a shell task the turn after it terminates (and a process * restart clears the registry outright), so a Stop that lands late is - * answered with `not_found`. That is the stop's goal state, not a failure: - * drop the task from local tracking and send the terminal notification + * answered with `not_found`. Persist that terminal state before dropping + * local tracking and sending the terminal notification * clients need to converge an entry they still show as running. Reporting * `No task found with ID` here only re-arms the stop button for a task that * can never be stopped again. */ -function convergeEvictedBackgroundTaskStop(sessionId: string, taskId: string): void { - const tracked = activeNonAgentTasks.get(sessionId)?.get(taskId) - untrackCliBackgroundTask(sessionId, taskId) - const description = tracked?.description - sendToSession(sessionId, { - type: 'system_notification', - subtype: 'task_notification', - message: description ? `${description} stopped` : 'Background task stopped', - data: { - type: 'system', - subtype: 'task_notification', - task_id: taskId, - tool_use_id: tracked?.toolUseId, - status: 'stopped', - summary: description ? `${description} stopped` : 'Background task stopped', - timestamp: new Date().toISOString(), - }, - }) +function convergeEvictedBackgroundTaskStop(sessionId: string, taskId: string): Promise { + let tracked = activeNonAgentTasks.get(sessionId)?.get(taskId) + if (!tracked) { + tracked = { taskId, toolUseId: taskId } + let tasks = activeNonAgentTasks.get(sessionId) + if (!tasks) { + tasks = new Map() + activeNonAgentTasks.set(sessionId, tasks) + } + tasks.set(taskId, tracked) + } + if (!tracked.localStopConfirmed || !tracked.terminalStatus) { + tracked.terminalStatus = 'stopped' + tracked.terminalMessage = true + } + tracked.localStopConfirmed = true + return emitAuthoritativeNonAgentStopped(sessionId, tracked) } const AGENT_STOP_CONTROL_TIMEOUT_MS = 3_000 @@ -2797,29 +2971,96 @@ function emitAuthoritativeStoppedForActiveAgents(sessionId: string): Promise { +function emitAuthoritativeNonAgentStopped( + sessionId: string, + task: ActiveNonAgentTaskState, +): Promise { + if (activeNonAgentTasks.get(sessionId)?.get(task.taskId) !== task || + sessionClearInProgress.has(sessionId)) return Promise.resolve(false) + if (task.finalization) return task.finalization + clearAgentStopFinalizationRetry(task) + task.localStopConfirmed = true + task.terminalStatus ??= 'stopped' + const status = task.terminalStatus + const cliMsg = { + type: 'system', + subtype: 'task_notification', + task_id: task.taskId, + tool_use_id: task.toolUseId, + ...(task.taskType ? { task_type: task.taskType } : {}), + ...(task.description ? { description: task.description } : {}), + ...(task.ownerAgentId ? { owner_agent_id: task.ownerAgentId } : {}), + status, + ...(task.terminalMessage + ? { message: task.terminalSummary ?? `${task.description ?? 'Background task'} ${status}` } + : {}), + summary: task.terminalSummary ?? `${task.description ?? 'Background task'} ${status}`, + timestamp: new Date().toISOString(), + } + const finalization = migrationMaintenance.track((async () => { + for (let attempt = 0; attempt < AUTHORITATIVE_STOP_PERSIST_ATTEMPTS; attempt++) { + if (activeNonAgentTasks.get(sessionId)?.get(task.taskId) !== task || + sessionClearInProgress.has(sessionId)) return false + try { + await (persistCliTaskNotification(sessionId, cliMsg, { + propagateFailure: true, + timeoutMs: AUTHORITATIVE_STOP_PERSIST_TIMEOUT_MS, + }) ?? Promise.resolve()) + if (activeNonAgentTasks.get(sessionId)?.get(task.taskId) !== task || + sessionClearInProgress.has(sessionId)) return false + task.stopFailureMessage = undefined + markTaskAuthoritativelyStopped(sessionId, task.taskId) + untrackCliBackgroundTask(sessionId, task.taskId) + forwardCliMessageToSessionClients(sessionId, cliMsg) + scheduleDisconnectedSessionCleanupIfIdle(sessionId) + return true + } catch { + // A rejected write is evicted from the cache so each attempt retries it. + } + } + if (activeNonAgentTasks.get(sessionId)?.get(task.taskId) !== task || + sessionClearInProgress.has(sessionId)) return false + task.stopFailureMessage = 'Background task ended, but its terminal state could not be saved' + sendToSession(sessionId, { + type: 'background_task_stop_failed', + taskId: task.taskId, + message: task.stopFailureMessage, + }) + const retryCount = task.finalizationRetryCount ?? 0 + const delay = AGENT_STOP_FINALIZATION_RETRY_DELAYS_MS[retryCount] + if (delay !== undefined && !migrationMaintenance.isActive) { + task.finalizationRetryCount = retryCount + 1 + task.finalizationRetryTimer = setTimeout(() => { + task.finalizationRetryTimer = undefined + if (!migrationMaintenance.isActive) void emitAuthoritativeNonAgentStopped(sessionId, task) + }, delay) + if (typeof task.finalizationRetryTimer === 'object') task.finalizationRetryTimer.unref?.() + } + scheduleDisconnectedSessionCleanupIfIdle(sessionId) + return false + })()) + task.finalization = finalization + void finalization.then(() => { + if (task.finalization === finalization) task.finalization = undefined + }) + return finalization +} + +function emitStoppedForNonAgentTasksAfterRuntimeExit( + sessionId: string, + status: 'stopped' | 'failed' = 'failed', +): Promise { + if (status === 'failed') runtimeExitFailedSessions.add(sessionId) const tasks = [...(activeNonAgentTasks.get(sessionId)?.values() ?? [])] return Promise.all(tasks.map(async (task) => { - if (activeNonAgentTasks.get(sessionId)?.get(task.taskId) !== task) return - // Killing the shared CLI also terminates Bash/Dream/workflow work. Claim - // each task before awaiting persistence so concurrent force-stop paths - // cannot publish duplicate terminal bookends. - markTaskAuthoritativelyStopped(sessionId, task.taskId) - untrackCliBackgroundTask(sessionId, task.taskId) - const cliMsg = { - type: 'system', - subtype: 'task_notification', - task_id: task.taskId, - tool_use_id: task.toolUseId, - ...(task.taskType ? { task_type: task.taskType } : {}), - ...(task.description ? { description: task.description } : {}), - ...(task.ownerAgentId ? { owner_agent_id: task.ownerAgentId } : {}), - status: 'stopped', - summary: `${task.description ?? task.taskId} stopped because the runtime exited`, - timestamp: new Date().toISOString(), + if (!task.localStopConfirmed || !task.terminalStatus) { + task.terminalStatus = status + task.terminalSummary = status === 'failed' + ? `${task.description ?? task.taskId}: runtime connection lost; task completion could not be confirmed` + : `${task.description ?? task.taskId} stopped because the runtime exited` } - await (persistCliTaskNotification(sessionId, cliMsg) ?? Promise.resolve()) - forwardCliMessageToSessionClients(sessionId, cliMsg) + task.localStopConfirmed = true + await emitAuthoritativeNonAgentStopped(sessionId, task) })) } @@ -2854,12 +3095,15 @@ function closeStoppedAgentsAfterRuntimeExit(sessionId: string, cliMsg: any): voi if ( cliMsg?.type === 'result' && cliMsg.is_error && - agentStopRequestedSessions.has(sessionId) && !conversationService.hasSession(sessionId) ) { + // A vanished runtime cannot deliver further task results. Record failed + // shell outcomes instead of claiming their process termination was confirmed; + // even without a user Stop, reconnects must not rehydrate ghost Running tasks. runtimeExitStoppedSessions.add(sessionId) + runtimeExitFailedSessions.add(sessionId) void emitAuthoritativeStoppedForActiveAgents(sessionId) - void emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId) + void emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId, 'failed') } } @@ -3196,6 +3440,7 @@ function cleanupSessionRuntimeState( options?: { preserveRetryableAgentStops?: boolean }, ) { cancelSessionDisconnectWatcher(sessionId) + clearBackgroundTaskDisconnectCeiling(sessionId) clearSessionTurnObserver(sessionId) clearAgentRuntimeState(sessionId, { preserveRetryableStops: options?.preserveRetryableAgentStops, @@ -3363,6 +3608,7 @@ async function ensureCliSessionStarted( await conversationService.startSession(sessionId, workDir, sdkUrl, startupSettings) bindSessionTurnObserver(sessionId) runtimeExitStoppedSessions.delete(sessionId) + runtimeExitFailedSessions.delete(sessionId) })() sessionStartupPromises.set(sessionId, startup) @@ -3411,7 +3657,9 @@ export async function ensureCliSessionStartedForControl( sdkUrl.toString(), { ...runtimeSettings, resumeInterruptedTurn: false }, ) + bindSessionTurnObserver(sessionId) runtimeExitStoppedSessions.delete(sessionId) + runtimeExitFailedSessions.delete(sessionId) })() sessionStartupPromises.set(sessionId, startup) @@ -4145,6 +4393,16 @@ function scheduleDisconnectCleanup(sessionId: string): void { } console.log(`[WS] Session ${sessionId} not reconnected after ${cleanupDelayMs}ms, stopping CLI subprocess`) + if (permissionBoundExpired) { + const turn = activeUserTurns.get(sessionId) + if (turn) { + turn.cancelled = true + turn.replacementAfterStop = false + } + forceStopSharedRuntimeForAgentCancellation(sessionId) + void emitAuthoritativeStoppedForActiveAgents(sessionId) + return + } conversationService.stopSession(sessionId) cleanupSessionRuntimeState(sessionId, { preserveRetryableAgentStops: true }) }, cleanupDelayMs) @@ -4152,6 +4410,7 @@ function scheduleDisconnectCleanup(sessionId: string): void { } function scheduleDisconnectedSessionCleanupIfIdle(sessionId: string): void { + refreshBackgroundTaskDisconnectCeiling(sessionId) if ( hasActiveClients(sessionId) || hasActiveSessionWork(sessionId) @@ -4177,8 +4436,10 @@ function watchTurnCompletionForCleanup(sessionId: string): void { const cliRunState = trackCliRunState(sessionId, cliMsg) const taskLifecycle = trackCliBackgroundTaskLifecycle(sessionId, cliMsg) stopLateAgentTaskIfRequested(sessionId, taskLifecycle) + stopLateNonAgentTaskIfRequested(sessionId, taskLifecycle) closeLateNonAgentTaskAfterRuntimeExit(sessionId, taskLifecycle) closeStoppedAgentsAfterRuntimeExit(sessionId, cliMsg) + refreshBackgroundTaskDisconnectCeiling(sessionId) if ( (cliRunState === 'running' || taskLifecycle?.running) && !hasActiveClients(sessionId) @@ -4191,6 +4452,10 @@ function watchTurnCompletionForCleanup(sessionId: string): void { const cleanupTimer = sessionCleanupTimers.get(sessionId) if (cleanupTimer) clearTimeout(cleanupTimer) sessionCleanupTimers.delete(sessionId) + // Arm the hard ceiling for a non-Agent background task that never + // reports back. Running Agent tasks keep their existing handling: only + // they can prove their own stop through the Agent finalization path. + refreshBackgroundTaskDisconnectCeiling(sessionId) } return } @@ -4211,7 +4476,7 @@ function watchTurnCompletionForCleanup(sessionId: string): void { const backgroundTaskCompleted = taskLifecycle?.running === false if (!foregroundTurnCompleted && !cliRunCompleted && !backgroundTaskCompleted) return if (hasActiveCliRun(sessionId)) return - if (hasActiveBackgroundTasks(sessionId)) return + if (hasActiveBackgroundTasks(sessionId) || hasActiveTeamWorkForParent(sessionId)) return if ( !foregroundTurnCompleted && !cliRunCompleted && @@ -4219,6 +4484,7 @@ function watchTurnCompletionForCleanup(sessionId: string): void { ) return cancelSessionDisconnectWatcher(sessionId) + clearBackgroundTaskDisconnectCeiling(sessionId) // All observed work finished while still disconnected — fall back to the // bounded idle timer rather than stopping the CLI immediately. if (!hasActiveClients(sessionId)) { @@ -4226,6 +4492,7 @@ function watchTurnCompletionForCleanup(sessionId: string): void { } } + refreshBackgroundTaskDisconnectCeiling(sessionId) conversationService.onOutput(sessionId, onComplete) sessionDisconnectWatchers.set(sessionId, () => { conversationService.removeOutputCallback(sessionId, onComplete) @@ -4261,6 +4528,73 @@ function cancelSessionDisconnectWatcher(sessionId: string): void { } } +/** + * Arm the hard ceiling for a disconnected session kept alive only by background + * shell tasks. Idempotent: an already-armed ceiling is left in place so the + * clock measures from the first observation, not from the latest event. + */ +function isDisconnectedBackgroundOnlySession(sessionId: string): boolean { + return !migrationMaintenance.isActive && + !hasActiveClients(sessionId) && + hasTrackedNonAgentTasks(sessionId) && + !hasPendingOrActiveUserTurn(sessionId) && + !hasActiveCliRun(sessionId) && + (activeAgentTasks.get(sessionId)?.size ?? 0) === 0 && + !hasActiveTeamWorkForParent(sessionId) && + conversationService.getPendingPermissionRequests(sessionId).length === 0 && + computerUseApprovalService.getPendingRequests(sessionId).length === 0 +} + +function refreshBackgroundTaskDisconnectCeiling(sessionId: string): void { + if (isDisconnectedBackgroundOnlySession(sessionId)) { + armBackgroundTaskDisconnectCeiling(sessionId) + } else { + clearBackgroundTaskDisconnectCeiling(sessionId) + } +} + +function armBackgroundTaskDisconnectCeiling(sessionId: string): void { + if (backgroundTaskCleanupTimers.has(sessionId)) return + const timer = setTimeout(() => { + if (backgroundTaskCleanupTimers.get(sessionId) !== timer) return + backgroundTaskCleanupTimers.delete(sessionId) + if (!isDisconnectedBackgroundOnlySession(sessionId)) return + console.log( + `[WS] Session ${sessionId} kept alive by background tasks for ${BACKGROUND_TASK_DISCONNECT_MAX_MS}ms without a client; stopping CLI subprocess`, + ) + forceStopSharedRuntimeForAgentCancellation(sessionId) + void emitAuthoritativeStoppedForActiveAgents(sessionId) + }, BACKGROUND_TASK_DISCONNECT_MAX_MS) + backgroundTaskCleanupTimers.set(sessionId, timer) +} + +function clearBackgroundTaskDisconnectCeiling(sessionId: string): void { + const timer = backgroundTaskCleanupTimers.get(sessionId) + if (timer) { + clearTimeout(timer) + backgroundTaskCleanupTimers.delete(sessionId) + } +} + +/** + * Any tracked task record, Agent or not. Used to decide whether a dead runtime + * left behind ghost "Running" entries that need a terminal bookend. + */ +function hasTrackedTaskRecords(sessionId: string): boolean { + return (activeNonAgentTasks.get(sessionId)?.size ?? 0) > 0 || + (activeAgentTasks.get(sessionId)?.size ?? 0) > 0 +} + +/** + * Non-Agent background shell tasks only (Bash/Dream/workflow). Unlike + * {@link hasActiveBackgroundTasks}, this never counts an Agent task, whose + * lifecycle is owned by the Agent finalization path. + */ +function hasTrackedNonAgentTasks(sessionId: string): boolean { + return [...(activeNonAgentTasks.get(sessionId)?.values() ?? [])] + .some((task) => !task.localStopConfirmed) +} + function replayPendingPermissionRequests( ws: SessionConnection, sessionId: string, @@ -4560,6 +4894,7 @@ function bindClientSessionOutput( trackCliRunState(sessionId, cliMsg) const taskLifecycle = trackCliBackgroundTaskLifecycle(sessionId, cliMsg) stopLateAgentTaskIfRequested(sessionId, taskLifecycle) + stopLateNonAgentTaskIfRequested(sessionId, taskLifecycle) closeLateNonAgentTaskAfterRuntimeExit(sessionId, taskLifecycle) closeStoppedAgentsAfterRuntimeExit(sessionId, cliMsg) if (taskLifecycle?.suppressForward) return @@ -5191,6 +5526,7 @@ export function getActiveSessionIds(): string[] { export function __resetWebSocketHandlerStateForTests(): void { for (const timer of sessionCleanupTimers.values()) clearTimeout(timer) + for (const timer of backgroundTaskCleanupTimers.values()) clearTimeout(timer) for (const timer of prewarmIdleTimers.values()) clearTimeout(timer) for (const remove of sessionDisconnectWatchers.values()) remove() for (const tasks of activeAgentTasks.values()) { @@ -5202,6 +5538,7 @@ export function __resetWebSocketHandlerStateForTests(): void { taskNotificationPersistence.clear() sessionTranscriptEpochs.clear() sessionCleanupTimers.clear() + backgroundTaskCleanupTimers.clear() sessionDisconnectWatchers.clear() prewarmPendingSessions.clear() prewarmedSessions.clear() @@ -5213,7 +5550,9 @@ export function __resetWebSocketHandlerStateForTests(): void { activeNonAgentTasks.clear() authoritativeStoppedTaskIds.clear() agentStopRequestedSessions.clear() + nonAgentStopRequestedSessions.clear() runtimeExitStoppedSessions.clear() + runtimeExitFailedSessions.clear() pendingInterruptedTurnResults.clear() sessionClearInProgress.clear() sessionStopRequested.clear() diff --git a/src/utils/Shell.ts b/src/utils/Shell.ts index f2e5fcf4..a098821e 100644 --- a/src/utils/Shell.ts +++ b/src/utils/Shell.ts @@ -1,4 +1,4 @@ -import { execFileSync, spawn } from 'child_process' +import { execFileSync } from 'child_process' import { constants as fsConstants, readFileSync, unlinkSync } from 'fs' import { type FileHandle, mkdir, open, realpath } from 'fs/promises' import memoize from 'lodash-es/memoize.js' @@ -35,6 +35,7 @@ import { getPlatform } from './platform.js' import { SandboxManager } from './sandbox/sandbox-adapter.js' import { invalidateSessionEnvCache } from './sessionEnvironment.js' import { createBashShellProvider } from './shell/bashProvider.js' +import { spawnWithParentProcessGuard } from './shell/parentProcessGuard.js' import { getCachedPowerShellPath } from './shell/powershellDetection.js' import { createPowerShellProvider } from './shell/powershellProvider.js' import type { ShellProvider, ShellType } from './shell/shellProvider.js' @@ -313,7 +314,7 @@ export async function exec( } try { - const childProcess = spawn(spawnBinary, shellArgs, { + const childProcess = spawnWithParentProcessGuard(spawnBinary, shellArgs, { env: { ...subprocessEnv(), SHELL: shellType === 'bash' ? binShell : undefined, diff --git a/src/utils/shell/parentProcessGuard.test.ts b/src/utils/shell/parentProcessGuard.test.ts new file mode 100644 index 00000000..1eb4f41f --- /dev/null +++ b/src/utils/shell/parentProcessGuard.test.ts @@ -0,0 +1,175 @@ +import { afterEach, expect, test } from 'bun:test' +import { spawn, type ChildProcess } from 'child_process' +import { closeSync, mkdtempSync, openSync, readFileSync, rmSync, writeFileSync } from 'fs' +import { tmpdir } from 'os' +import { join } from 'path' +import { spawnWithParentProcessGuard } from './parentProcessGuard.js' + +const directories: string[] = [] +const children: ChildProcess[] = [] +const wait = (ms: number) => new Promise(resolve => setTimeout(resolve, ms)) + +afterEach(() => { + for (const child of children.splice(0)) { + if (!child.pid) continue + try { process.kill(-child.pid, 'SIGKILL') } catch {} + try { child.kill('SIGKILL') } catch {} + } + for (const directory of directories.splice(0)) rmSync(directory, { recursive: true, force: true }) +}) + +function fixtureDirectory(): string { + const directory = mkdtempSync(join(tmpdir(), 'shell-parent-guard-')) + directories.push(directory) + return directory +} + +function exited(child: ChildProcess): Promise { + if (child.exitCode !== null || child.signalCode !== null) return Promise.resolve(child.exitCode) + return new Promise((resolve, reject) => { + child.once('exit', code => resolve(code)) + child.once('error', reject) + }) +} + +async function waitUntil(predicate: () => boolean): Promise { + const deadline = Date.now() + 3_000 + while (!predicate()) { + if (Date.now() >= deadline) throw new Error('Timed out waiting for fixture process') + await wait(10) + } +} + +function ticks(path: string): number { + try { return readFileSync(path, 'utf8').trim().split('\n').length } catch { return 0 } +} + +test.skipIf(process.platform === 'win32')('kills a detached shell and its descendants after the runtime receives SIGKILL', async () => { + const directory = fixtureDirectory() + const output = join(directory, 'ticks') + const pidPath = join(directory, 'child-pid') + const fixture = join(directory, 'runtime.ts') + writeFileSync(fixture, ` + import { openSync, writeFileSync } from 'fs' + import { spawnWithParentProcessGuard } from ${JSON.stringify(import.meta.dir + '/parentProcessGuard.ts')} + const child = spawnWithParentProcessGuard('/bin/sh', ['-c', 'while true; do echo alive; sleep 0.02; done', 'shell'], { + detached: true, + stdio: ['pipe', openSync(${JSON.stringify(output)}, 'a'), 'ignore'], + }) + writeFileSync(${JSON.stringify(pidPath)}, String(child.pid)) + setInterval(() => {}, 1000) + `) + // Write through an inherited file fd, as background Shell.exec does. Pipe + // output would stop when the CLI dies even if its shell kept running. + const runtime = spawn(process.execPath, ['--no-env-file', fixture], { + cwd: directory, + env: { PATH: process.env.PATH, HOME: directory, CLAUDE_CONFIG_DIR: join(directory, 'config') }, + detached: true, + stdio: 'ignore', + }) + children.push(runtime) + let shellPid: number | undefined + try { + await waitUntil(() => ticks(output) >= 3) + shellPid = Number(readFileSync(pidPath, 'utf8')) + runtime.kill('SIGKILL') + await exited(runtime) + expect(runtime.signalCode).toBe('SIGKILL') + await wait(150) + const afterExit = ticks(output) + await wait(150) + expect(ticks(output)).toBe(afterExit) + } finally { + if (shellPid) { + try { process.kill(-shellPid, 'SIGKILL') } catch {} + } + } +}) + +test.skipIf(process.platform === 'win32')('preserves stdin, arguments, bare wait, stdout, stderr and the exit code', async () => { + const child = spawnWithParentProcessGuard('/bin/sh', ['-c', + 'IFS= read -r line; sleep 0.02 & wait; printf "%s|%s|%s\\n" "$line" "$1" "$GUARD_TEST_VALUE"; echo err >&2; exit 7', + 'sandbox-wrapper', 'literal $value `command` "quotes"', + ], { + detached: true, + env: { ...process.env, GUARD_TEST_VALUE: 'env value' }, + stdio: ['pipe', 'pipe', 'pipe'], + }) + children.push(child) + let stdout = '' + let stderr = '' + child.stdout!.on('data', data => { stdout += data }) + child.stderr!.on('data', data => { stderr += data }) + child.stdin!.end('user input\n') + expect(await exited(child)).toBe(7) + expect(stdout).toBe('user input|literal $value `command` "quotes"|env value\n') + expect(stderr).toBe('err\n') +}) + +test.skipIf(process.platform === 'win32')('preserves file output and leaves no process group after normal completion', async () => { + const output = join(fixtureDirectory(), 'output') + const fd = openSync(output, 'a') + let child: ChildProcess + try { + child = spawnWithParentProcessGuard('/bin/sh', ['-c', 'echo stdout; echo stderr >&2; exit 0'], { + detached: true, + stdio: ['pipe', fd, fd], + }) + children.push(child) + } finally { + closeSync(fd) + } + expect(await exited(child)).toBe(0) + expect(readFileSync(output, 'utf8')).toBe('stdout\nstderr\n') + expect(() => process.kill(-child.pid!, 0)).toThrow() +}) + +test.skipIf(process.platform === 'win32')('preserves an existing sandbox wrapper and does not expose the lifetime fd to commands', async () => { + const child = spawnWithParentProcessGuard('/bin/sh', ['-c', + 'exec /bin/sh -c "$1" "$2" "$3"', 'sandbox-wrapper', + 'if (: <&3) 2>/dev/null; then echo leaked; exit 1; fi; printf "%s\\n" "$1"', + 'inner-shell', 'literal $arg `command` "quoted"', + ], { + detached: false, + stdio: ['pipe', 'pipe', 'pipe'], + }) + children.push(child) + let stdout = '' + child.stdout!.on('data', data => { stdout += data }) + expect(await exited(child)).toBe(0) + expect(stdout).toBe('literal $arg `command` "quoted"\n') +}) + +for (const action of ['terminate', 'abort'] as const) { + test.skipIf(process.platform === 'win32')(`reaps the target and guardian when the supervisor is ${action === 'abort' ? 'aborted' : 'terminated'}`, async () => { + const output = join(fixtureDirectory(), 'ticks') + const fd = openSync(output, 'a') + const controller = new AbortController() + let child: ChildProcess + try { + child = spawnWithParentProcessGuard('/bin/sh', ['-c', 'while true; do echo alive; sleep 0.02; done'], { + detached: true, + signal: controller.signal, + stdio: ['pipe', fd, fd], + }) + children.push(child) + } finally { + closeSync(fd) + } + const errors: Error[] = [] + child.on('error', error => { errors.push(error) }) + const exit = new Promise(resolve => child.once('exit', resolve)) + await waitUntil(() => ticks(output) >= 3) + if (action === 'abort') controller.abort() + else child.kill('SIGTERM') + await exit + await waitUntil(() => { + try { process.kill(-child.pid!, 0); return false } catch { return true } + }) + const afterExit = ticks(output) + await wait(75) + expect(ticks(output)).toBe(afterExit) + if (action === 'abort') expect(errors[0]?.name).toBe('AbortError') + else expect(errors).toHaveLength(0) + }) +} diff --git a/src/utils/shell/parentProcessGuard.ts b/src/utils/shell/parentProcessGuard.ts new file mode 100644 index 00000000..44e3c22f --- /dev/null +++ b/src/utils/shell/parentProcessGuard.ts @@ -0,0 +1,51 @@ +import { spawn, type ChildProcess, type SpawnOptions } from 'child_process' +import type { Duplex } from 'stream' + +// Keep the guardian outside the target shell's job table: a command's bare +// `wait` must not wait for our lifetime watcher. Arguments reach the original +// shell unchanged, including login flags and sandbox command wrappers. +const PARENT_PROCESS_GUARD = ` +( + if ! IFS= read -r lifetime <&3; then + kill -KILL "-$$" + fi +) /dev/null 2>&1 & +lifetime_guard=$! +"$@" 3<&- +lifetime_status=$? +kill "$lifetime_guard" 2>/dev/null || : +wait "$lifetime_guard" 2>/dev/null || : +exit "$lifetime_status" +` + +export function spawnWithParentProcessGuard( + command: string, + args: string[], + options: SpawnOptions, +): ChildProcess { + // Windows has no POSIX process group or /bin/sh. Retain its existing spawn + // behavior; callers must not infer that a runtime crash killed its children. + if (process.platform === 'win32') return spawn(command, args, options) + + const stdio = Array.isArray(options.stdio) + ? [...options.stdio] + : Array(3).fill(options.stdio ?? 'pipe') + stdio[3] = 'pipe' + const child = spawn('/bin/sh', [ + '-c', PARENT_PROCESS_GUARD, 'cc-haha-shell-lifetime', command, ...args, + ], { + ...options, + detached: true, + stdio, + }) + + // The CLI owns the pipe's write end. Even SIGKILL closes it in the kernel, + // so EOF reaps this command's group without depending on SDK task_started, + // a model loop, or graceful cleanup. Normal completion kills and waits for + // the guardian in the supervisor before returning the original exit code. + const lifetimePipe = child.stdio[3] as Duplex | null + const release = () => lifetimePipe?.destroy() + child.once('exit', release) + child.once('error', release) + return child +}