mirror of
https://github.com/NanmiCoder/claude-code-haha.git
synced 2026-10-10 20:03:13 +08:00
fix(runtime): converge background tasks on Stop and runtime loss (#1441)
* fix(server): reap background shell tasks on Stop and after runtime exit A session's background shell task (Bash `run_in_background`, Dream, workflow) is tracked as a non-Agent task, and nothing on the server ever bounded or reaped one. A user Stop only interrupted the foreground turn and the Agent tasks, so the shell kept running after the user had asked everything to stop. A client disconnect let the CLI live forever, because the disconnect watcher waits for `hasActiveBackgroundTasks` to clear and a task that never emits a terminal notification never clears it. A hard runtime death published no terminal state at all, so a later reconnect re-hydrated a ghost "Running" entry that could never reach a terminal state. This closes all three from one file of server bookkeeping: - `handleStopGeneration` now reaps the session's non-Agent background tasks through the same `requestControl` stop path the panel's per-task Stop uses, and the 3-second force-kill guard watches them too. Programmatic stops (turn stop, runtime-config restart) keep the previous narrow behavior. - A new 31-minute ceiling bounds a session that a disconnected client keeps alive through background tasks alone; when it elapses the shared runtime is stopped and terminal bookends are published. A reconnect clears the ceiling. - `close()` and `closeStoppedAgentsAfterRuntimeExit` publish terminal bookends whenever the runtime is already gone but task records remain, instead of leaving them to be re-hydrated as "Running". The new timer registry is released by `cleanupSessionRuntimeState` and classified as `cleared` in the session-state cleanup guard, so a ceiling scheduled for a session cannot outlive that session. 会话的后台 shell 任务(Bash `run_in_background`、Dream、workflow)登记为非 Agent 任务, 服务端一直没有收掉它们的路径:按停止键只断前台轮和 Agent 任务,shell 还在跑;客户端断开后 CLI 因为 `hasActiveBackgroundTasks` 恒为真而一直存活;运行时崩溃又不写终态,重连后又 变成幽灵 Running。这次在 handler 一处收掉这三种: - 按停止键时一并停掉本会话的非 Agent 后台任务,走面板单任务 Stop 同一个 stop 通路, 三秒强杀闸也一并盯住它们;程序化停止(停轮、重启换配置)仍走原来的窄路径。 - 新增 31 分钟硬上限:没有客户端、只靠后台任务续命的会话,到期就停掉共享运行时并写终态; 重连则撤销该计时。 - 运行时已经没了但任务记录还在时,`close()` 与 `closeStoppedAgentsAfterRuntimeExit` 直接写终态,不再留着让重连把 Running 灌回来。 Refs #1132 * fix(runtime): converge background tasks and reap orphaned shells * fix(desktop): restore stopped shell tasks after cold history load --------- Co-authored-by: gugugaga <267102352+omazili-guga@users.noreply.github.com> Co-authored-by: 程序员阿江(Relakkes) <relakkes@gmail.com>
This commit is contained in:
@@ -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<number>()
|
||||
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<void>(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<void>(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<void>(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')
|
||||
})
|
||||
})
|
||||
@@ -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()
|
||||
})
|
||||
@@ -45,7 +45,9 @@ const CONTAINERS: Record<string, Classification> = {
|
||||
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<string, Classification> = {
|
||||
prewarmedSessions: { kind: 'cleared' },
|
||||
rejectedRuntimeConfigs: { kind: 'cleared' },
|
||||
runtimeExitStoppedSessions: { kind: 'cleared' },
|
||||
runtimeExitFailedSessions: { kind: 'cleared' },
|
||||
runtimeOverrides: { kind: 'cleared' },
|
||||
runtimeTransitionPromises: { kind: 'cleared' },
|
||||
sessionDisconnectWatchers: { kind: 'cleared' },
|
||||
|
||||
@@ -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)', () => {
|
||||
|
||||
@@ -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<typeof setTimeout>
|
||||
finalization?: Promise<boolean>
|
||||
stopFailureMessage?: string
|
||||
}
|
||||
|
||||
export type ActiveAgentTaskState = {
|
||||
@@ -54,8 +63,12 @@ export const authoritativeStoppedTaskIds = new Map<string, Set<string>>()
|
||||
|
||||
export const agentStopRequestedSessions = new Set<string>()
|
||||
|
||||
export const nonAgentStopRequestedSessions = new Set<string>()
|
||||
|
||||
export const runtimeExitStoppedSessions = new Set<string>()
|
||||
|
||||
export const runtimeExitFailedSessions = new Set<string>()
|
||||
|
||||
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<string, ActiveAgentTaskState>()
|
||||
|
||||
const retryableNonAgentStops = options?.preserveRetryableStops
|
||||
? new Map([...(activeNonAgentTasks.get(sessionId)?.entries() ?? [])]
|
||||
.filter(([, task]) => task.localStopConfirmed))
|
||||
: new Map<string, ActiveNonAgentTaskState>()
|
||||
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<typeof setTimeout> }): void {
|
||||
if (task.finalizationRetryTimer === undefined) return
|
||||
clearTimeout(task.finalizationRetryTimer)
|
||||
task.finalizationRetryTimer = undefined
|
||||
|
||||
+391
-52
@@ -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<string, SessionSlashCommand[]>()
|
||||
// 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<string, ReturnType<typeof setTimeout>>()
|
||||
const backgroundTaskCleanupTimers = new Map<string, ReturnType<typeof setTimeout>>()
|
||||
/**
|
||||
* 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<boolean> {
|
||||
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<boo
|
||||
}))
|
||||
}
|
||||
|
||||
function emitStoppedForNonAgentTasksAfterRuntimeExit(sessionId: string): Promise<void[]> {
|
||||
function emitAuthoritativeNonAgentStopped(
|
||||
sessionId: string,
|
||||
task: ActiveNonAgentTaskState,
|
||||
): Promise<boolean> {
|
||||
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<void[]> {
|
||||
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()
|
||||
|
||||
+3
-2
@@ -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,
|
||||
|
||||
@@ -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<number | null> {
|
||||
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<void> {
|
||||
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)
|
||||
})
|
||||
}
|
||||
@@ -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 >/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
|
||||
}
|
||||
Reference in New Issue
Block a user