Download packages/agent-core-v2/test/agent/task/taskManager.test.ts from SaylorTwift/kimi-code: direct link, hf CLI and curl.
- Browser
- Download file 46 kB
-
https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/test/agent/task/taskManager.test.ts
- Command line
-
hf download hf://SaylorTwift/kimi-code/packages/agent-core-v2/test/agent/task/taskManager.test.ts
-
curl -L -o taskManager.test.ts https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/test/agent/task/taskManager.test.ts
46 kB
| import { mkdtemp, rm } from 'node:fs/promises'; | |
| import { tmpdir } from 'node:os'; | |
| import { PassThrough, Readable } from 'node:stream'; | |
| import type { Writable } from 'node:stream'; | |
| import { join } from 'pathe'; | |
| import type { IHostProcess } from '#/os/interface/hostProcess'; | |
| import { afterEach, describe, expect, it, vi } from 'vitest'; | |
| import { | |
| IAgentTaskService, | |
| type AgentTaskInfo, | |
| } from '#/agent/task/task'; | |
| import { | |
| SubagentTask, | |
| type SubagentHandle, | |
| } from '#/agent/tools/agent/subagent-task'; | |
| import { ProcessTask } from '#/agent/tools/os/bash/process-task'; | |
| import { isUserCancellation, userCancellationReason } from '#/_base/utils/abort'; | |
| import { ISessionMetadata } from '#/session/sessionMetadata/sessionMetadata'; | |
| import { | |
| configServices, | |
| createTestAgent, | |
| homeDirServices, | |
| type TestAgentContext, | |
| type TestAgentServiceOverride, | |
| } from '../../harness'; | |
| import { | |
| createAgentTaskPersistence, | |
| type TaskServiceTestManager, | |
| } from './stubs'; | |
| const MiB = 1024 * 1024; | |
| const LIMIT_BYTES = 16 * MiB; | |
| interface TaskServiceFixture { | |
| ctx: TestAgentContext; | |
| manager: TaskServiceTestManager; | |
| persistence?: ReturnType<typeof createAgentTaskPersistence>; | |
| } | |
| function createAgentTaskService(options: { | |
| sessionDir?: string; | |
| maxRunningTasks?: number; | |
| } = {}): TaskServiceFixture { | |
| const persistence = | |
| options.sessionDir === undefined | |
| ? undefined | |
| : createAgentTaskPersistence(options.sessionDir); | |
| const overrides: TestAgentServiceOverride[] = []; | |
| if (options.sessionDir !== undefined) { | |
| overrides.push(homeDirServices(options.sessionDir)); | |
| } | |
| const maxRunningTasks = options.maxRunningTasks; | |
| if (maxRunningTasks !== undefined) { | |
| overrides.push(configServices(() => ({ | |
| providers: {}, | |
| task: { maxRunningTasks }, | |
| }))); | |
| } | |
| const ctx = createTestAgent(...overrides); | |
| return { | |
| ctx, | |
| manager: ctx.get(IAgentTaskService) as TaskServiceTestManager, | |
| persistence, | |
| }; | |
| } | |
| function registerProcess( | |
| manager: IAgentTaskService, | |
| proc: IHostProcess, | |
| command: string, | |
| description: string, | |
| ): string { | |
| return manager.registerTask(new ProcessTask(proc, command, description)); | |
| } | |
| function agentTask( | |
| completion: Promise<{ result: string }>, | |
| description: string, | |
| options: { | |
| readonly agentId?: string; | |
| readonly subagentType?: string; | |
| readonly parentToolCallId?: string; | |
| readonly abortController?: AbortController; | |
| readonly timeoutMs?: number; | |
| } = {}, | |
| ): SubagentTask { | |
| const handle: SubagentHandle = { | |
| agentId: options.agentId ?? 'agent-child', | |
| profileName: options.subagentType ?? 'coder', | |
| parentToolCallId: options.parentToolCallId, | |
| completion, | |
| }; | |
| const task = new SubagentTask( | |
| handle, | |
| description, | |
| options.abortController ?? new AbortController(), | |
| ); | |
| if (options.timeoutMs !== undefined) { | |
| Object.defineProperty(task, 'timeoutMs', { | |
| value: options.timeoutMs, | |
| enumerable: true, | |
| }); | |
| } | |
| return task; | |
| } | |
| async function waitForTerminal( | |
| manager: IAgentTaskService, | |
| taskId: string, | |
| timeoutMs = 30_000, | |
| ): Promise<AgentTaskInfo | undefined> { | |
| const deadline = Date.now() + timeoutMs; | |
| while (Date.now() <= deadline) { | |
| const info = await manager.wait(taskId, 5); | |
| if ( | |
| info?.status === 'completed' || | |
| info?.status === 'failed' || | |
| info?.status === 'timed_out' || | |
| info?.status === 'killed' || | |
| info?.status === 'lost' | |
| ) { | |
| return info; | |
| } | |
| await new Promise((resolve) => setTimeout(resolve, 1)); | |
| } | |
| return manager.getTask(taskId); | |
| } | |
| async function waitForOutput( | |
| manager: IAgentTaskService, | |
| taskId: string, | |
| expected: string, | |
| ): Promise<void> { | |
| for (let i = 0; i < 20; i++) { | |
| const output = await manager.readOutput(taskId); | |
| if (output.includes(expected)) return; | |
| await new Promise((resolve) => setTimeout(resolve, 5)); | |
| } | |
| throw new Error(`Timed out waiting for output: ${expected}`); | |
| } | |
| function immediateProcess(exitCode: number, stdoutText = ''): IHostProcess { | |
| return { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout: Readable.from(stdoutText ? [stdoutText] : []), | |
| stderr: Readable.from([]), | |
| pid: 10000 + exitCode, | |
| exitCode, | |
| wait: vi.fn().mockResolvedValue(exitCode) as IHostProcess['wait'], | |
| kill: vi.fn().mockResolvedValue(undefined) as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| } | |
| function rejectedProcess(error: Error): IHostProcess { | |
| return { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout: Readable.from([]), | |
| stderr: Readable.from([]), | |
| pid: 99999, | |
| exitCode: null, | |
| wait: vi.fn().mockRejectedValue(error) as IHostProcess['wait'], | |
| kill: vi.fn().mockResolvedValue(undefined) as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| } | |
| function processWithStdoutError(message = 'stdout read failed'): IHostProcess { | |
| const stdout = new PassThrough(); | |
| return { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout, | |
| stderr: Readable.from([]), | |
| pid: 99998, | |
| exitCode: 0, | |
| wait: vi.fn(async () => { | |
| stdout.destroy(new Error(message)); | |
| return 0; | |
| }) as IHostProcess['wait'], | |
| kill: vi.fn().mockResolvedValue(undefined) as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| } | |
| function processWithStdoutErrorBeforeWait(message = 'stdout read failed'): { | |
| proc: IHostProcess; | |
| failStdout: () => void; | |
| resolveWait: (exitCode: number) => void; | |
| } { | |
| const stdout = new PassThrough(); | |
| let currentExitCode: number | null = null; | |
| let resolveWait: (n: number) => void = () => {}; | |
| const waitPromise = new Promise<number>((resolve) => { | |
| resolveWait = resolve; | |
| }); | |
| return { | |
| proc: { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout, | |
| stderr: Readable.from([]), | |
| pid: 99997, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: vi.fn(() => waitPromise) as IHostProcess['wait'], | |
| kill: vi.fn().mockResolvedValue(undefined) as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }, | |
| failStdout: () => { | |
| stdout.destroy(new Error(message)); | |
| }, | |
| resolveWait: (exitCode) => { | |
| currentExitCode = exitCode; | |
| resolveWait(exitCode); | |
| }, | |
| }; | |
| } | |
| function pendingProcess(exitOnKill = 143): { | |
| proc: IHostProcess; | |
| killSpy: ReturnType<typeof vi.fn>; | |
| } { | |
| let resolveWait: (n: number) => void = () => {}; | |
| const waitPromise = new Promise<number>((resolve) => { | |
| resolveWait = resolve; | |
| }); | |
| let currentExitCode: number | null = null; | |
| const killSpy = vi.fn(async () => { | |
| if (currentExitCode !== null) return; | |
| currentExitCode = exitOnKill; | |
| resolveWait(exitOnKill); | |
| }); | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout: Readable.from([]), | |
| stderr: Readable.from([]), | |
| pid: 54321, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: () => waitPromise, | |
| kill: killSpy as unknown as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| return { proc, killSpy }; | |
| } | |
| function streamingProcess(chunks: string[]): { | |
| proc: IHostProcess; | |
| killSpy: ReturnType<typeof vi.fn>; | |
| } { | |
| const stdout = Readable.from(chunks); | |
| const stderr = Readable.from([]); | |
| let currentExitCode: number | null = null; | |
| let resolveWait: (code: number) => void = () => {}; | |
| const waitPromise = new Promise<number>((resolve) => { | |
| resolveWait = resolve; | |
| }); | |
| stdout.on('end', () => { | |
| currentExitCode = 0; | |
| resolveWait(0); | |
| }); | |
| const killSpy = vi.fn(async (signal: NodeJS.Signals) => { | |
| if (currentExitCode !== null) return; | |
| currentExitCode = signal === 'SIGKILL' ? 137 : 143; | |
| stdout.destroy(); | |
| resolveWait(currentExitCode); | |
| }); | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout, | |
| stderr, | |
| pid: 54325, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: () => waitPromise, | |
| kill: killSpy as unknown as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| return { proc, killSpy }; | |
| } | |
| function sigtermIgnoringProcess(chunks: string[]): { | |
| proc: IHostProcess; | |
| killSpy: ReturnType<typeof vi.fn>; | |
| } { | |
| const stdout = Readable.from(chunks); | |
| const stderr = Readable.from([]); | |
| let currentExitCode: number | null = null; | |
| let resolveWait: (code: number) => void = () => {}; | |
| const waitPromise = new Promise<number>((resolve) => { | |
| resolveWait = resolve; | |
| }); | |
| stdout.on('end', () => { | |
| currentExitCode = 0; | |
| resolveWait(0); | |
| }); | |
| const killSpy = vi.fn(async (signal: NodeJS.Signals) => { | |
| if (signal !== 'SIGKILL' || currentExitCode !== null) return; | |
| currentExitCode = 137; | |
| stdout.destroy(); | |
| resolveWait(137); | |
| }); | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout, | |
| stderr, | |
| pid: 54326, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: () => waitPromise, | |
| kill: killSpy as unknown as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| return { proc, killSpy }; | |
| } | |
| function manuallyResolvedProcess(): { | |
| proc: IHostProcess; | |
| killSpy: ReturnType<typeof vi.fn>; | |
| resolve: (exitCode: number) => void; | |
| } { | |
| let resolveWait: (n: number) => void = () => {}; | |
| const waitPromise = new Promise<number>((resolve) => { | |
| resolveWait = resolve; | |
| }); | |
| let currentExitCode: number | null = null; | |
| const killSpy = vi.fn().mockResolvedValue(undefined); | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout: Readable.from([]), | |
| stderr: Readable.from([]), | |
| pid: 54324, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: () => waitPromise, | |
| kill: killSpy as unknown as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| return { | |
| proc, | |
| killSpy, | |
| resolve: (exitCode) => { | |
| if (currentExitCode !== null) return; | |
| currentExitCode = exitCode; | |
| resolveWait(exitCode); | |
| }, | |
| }; | |
| } | |
| function processWithVisibleExitCodeBeforeWait(exitCode = 143): { | |
| proc: IHostProcess; | |
| markExited: () => void; | |
| } { | |
| let currentExitCode: number | null = null; | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout: Readable.from([]), | |
| stderr: Readable.from([]), | |
| pid: 54322, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: () => new Promise<number>(() => {}), | |
| kill: vi.fn().mockResolvedValue(undefined) as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| return { | |
| proc, | |
| markExited: () => { | |
| currentExitCode = exitCode; | |
| }, | |
| }; | |
| } | |
| describe('AgentTaskService', () => { | |
| afterEach(() => { | |
| vi.useRealTimers(); | |
| }); | |
| it('registers process tasks and exposes process metadata', () => { | |
| const { manager } = createAgentTaskService(); | |
| const proc = immediateProcess(0); | |
| const taskId = registerProcess(manager, proc, 'echo hello', 'test echo'); | |
| expect(taskId).toMatch(/^bash-[0-9a-z]{8}$/); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| taskId, | |
| kind: 'process', | |
| command: 'echo hello', | |
| description: 'test echo', | |
| pid: proc.pid, | |
| status: 'running', | |
| }); | |
| }); | |
| it('registers agent tasks and exposes agent metadata', () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = manager.registerTask( | |
| agentTask(new Promise(() => {}), 'investigate bug', { | |
| agentId: 'agent-child', | |
| subagentType: 'coder', | |
| parentToolCallId: 'call-parent-1', | |
| }), | |
| ); | |
| expect(taskId).toMatch(/^agent-[0-9a-z]{8}$/); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| taskId, | |
| kind: 'agent', | |
| description: 'investigate bug', | |
| agentId: 'agent-child', | |
| subagentType: 'coder', | |
| parentToolCallId: 'call-parent-1', | |
| status: 'running', | |
| }); | |
| }); | |
| it('tracks foreground tasks and releases their waiter when detached', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = manager.registerTask( | |
| agentTask(new Promise(() => {}), 'foreground agent'), | |
| { detached: false }, | |
| ); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| detached: false, | |
| }); | |
| const waiting = manager.waitForForegroundRelease(taskId); | |
| await Promise.resolve(); | |
| expect(manager.detach(taskId)).toMatchObject({ | |
| taskId, | |
| detached: true, | |
| }); | |
| await expect(waiting).resolves.toBe('detached'); | |
| }); | |
| it('releases foreground waiters when a foreground task completes', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = manager.registerTask( | |
| agentTask(Promise.resolve({ result: 'done' }), 'foreground agent'), | |
| { detached: false }, | |
| ); | |
| await expect(manager.waitForForegroundRelease(taskId)).resolves.toBe('terminal'); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| detached: false, | |
| status: 'completed', | |
| }); | |
| }); | |
| it('stops foreground tasks from their register-time signal', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(); | |
| const controller = new AbortController(); | |
| const taskId = manager.registerTask( | |
| new ProcessTask(proc, 'sleep 10', 'foreground process'), | |
| { | |
| detached: false, | |
| signal: controller.signal, | |
| }, | |
| ); | |
| const waiting = manager.waitForForegroundRelease(taskId); | |
| controller.abort(); | |
| await expect(waiting).resolves.toBe('terminal'); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| status: 'killed', | |
| stopReason: 'Aborted by the user', | |
| }); | |
| }); | |
| it('keeps a detached process task running when the register-time signal aborts', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(); | |
| const controller = new AbortController(); | |
| const taskId = manager.registerTask( | |
| new ProcessTask(proc, 'sleep 10', 'foreground process'), | |
| { | |
| detached: false, | |
| signal: controller.signal, | |
| }, | |
| ); | |
| const waiting = manager.waitForForegroundRelease(taskId); | |
| expect(manager.detach(taskId)).toMatchObject({ detached: true }); | |
| controller.abort(); | |
| await expect(waiting).resolves.toBe('detached'); | |
| expect(killSpy).not.toHaveBeenCalled(); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| status: 'running', | |
| detached: true, | |
| }); | |
| }); | |
| it('forwards foreground signal abort reasons to agent task controllers', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const foregroundController = new AbortController(); | |
| const subagentController = new AbortController(); | |
| const completion = new Promise<{ result: string }>((_resolve, reject) => { | |
| subagentController.signal.addEventListener( | |
| 'abort', | |
| () => { | |
| reject(subagentController.signal.reason); | |
| }, | |
| { once: true }, | |
| ); | |
| }); | |
| const taskId = manager.registerTask( | |
| agentTask(completion, 'foreground agent', { abortController: subagentController }), | |
| { | |
| detached: false, | |
| signal: foregroundController.signal, | |
| }, | |
| ); | |
| foregroundController.abort(userCancellationReason()); | |
| const info = await manager.wait(taskId); | |
| expect(info).toMatchObject({ | |
| status: 'killed', | |
| stopReason: 'Aborted by the user', | |
| }); | |
| expect(isUserCancellation(subagentController.signal.reason)).toBe(true); | |
| }); | |
| it('does not forward register-time signal aborts to a detached agent task', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const foregroundController = new AbortController(); | |
| const subagentController = new AbortController(); | |
| const taskId = manager.registerTask( | |
| agentTask(new Promise(() => {}), 'foreground agent', { | |
| abortController: subagentController, | |
| }), | |
| { | |
| detached: false, | |
| signal: foregroundController.signal, | |
| }, | |
| ); | |
| expect(manager.detach(taskId)).toMatchObject({ detached: true }); | |
| foregroundController.abort(userCancellationReason()); | |
| expect(subagentController.signal.aborted).toBe(false); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| status: 'running', | |
| detached: true, | |
| }); | |
| }); | |
| it('does not count foreground tasks against the detached task limit', () => { | |
| const { manager } = createAgentTaskService({ maxRunningTasks: 1 }); | |
| manager.registerTask(agentTask(new Promise(() => {}), 'foreground agent'), { | |
| detached: false, | |
| }); | |
| manager.registerTask(agentTask(new Promise(() => {}), 'background agent')); | |
| expect(() => { | |
| manager.registerTask(agentTask(new Promise(() => {}), 'second background')); | |
| }).toThrow('Too many background tasks are already running.'); | |
| }); | |
| it('does not count foreground tasks detached later against the detached task limit', () => { | |
| const { manager } = createAgentTaskService({ maxRunningTasks: 1 }); | |
| const taskId = manager.registerTask( | |
| agentTask(new Promise(() => {}), 'foreground agent'), | |
| { detached: false }, | |
| ); | |
| manager.detach(taskId); | |
| manager.registerTask(agentTask(new Promise(() => {}), 'background agent')); | |
| expect(() => { | |
| manager.registerTask(agentTask(new Promise(() => {}), 'second background')); | |
| }).toThrow('Too many background tasks are already running.'); | |
| }); | |
| it('lists active tasks by default', () => { | |
| const { manager } = createAgentTaskService(); | |
| registerProcess(manager, pendingProcess().proc, 'sleep 60', 'task 1'); | |
| registerProcess(manager, pendingProcess().proc, 'sleep 60', 'task 2'); | |
| expect(manager.list()).toHaveLength(2); | |
| }); | |
| it('excludes terminal detached tasks from active listings and includes them in all-task listings', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess(manager, immediateProcess(0), 'echo done', 'done'); | |
| await manager.wait(taskId); | |
| expect(manager.list(true)).toEqual([]); | |
| expect(manager.list(false)).toEqual([ | |
| expect.objectContaining({ | |
| taskId, | |
| kind: 'process', | |
| status: 'completed', | |
| exitCode: 0, | |
| }), | |
| ]); | |
| }); | |
| it('honours the list limit parameter', () => { | |
| const { manager } = createAgentTaskService(); | |
| const first = registerProcess(manager, pendingProcess().proc, 'sleep 1', 'one'); | |
| const second = registerProcess(manager, pendingProcess().proc, 'sleep 2', 'two'); | |
| expect(manager.list(true, 1)).toEqual([ | |
| expect.objectContaining({ taskId: first }), | |
| ]); | |
| expect(manager.list(true, 1)).not.toEqual([ | |
| expect.objectContaining({ taskId: second }), | |
| ]); | |
| }); | |
| it('lists running tasks synchronously without waiting for task completion', () => { | |
| vi.useFakeTimers(); | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess(manager, pendingProcess().proc, 'sleep 60', 'running list'); | |
| const tasks = manager.list(true); | |
| expect(tasks).toEqual([ | |
| expect.objectContaining({ | |
| taskId, | |
| status: 'running', | |
| description: 'running list', | |
| }), | |
| ]); | |
| }); | |
| it('rejects new tasks when maxRunningTasks is reached', () => { | |
| const { manager } = createAgentTaskService({ maxRunningTasks: 1 }); | |
| registerProcess(manager, pendingProcess().proc, 'sleep 60', 'first task'); | |
| expect(() => { | |
| registerProcess(manager, pendingProcess().proc, 'sleep 60', 'second task'); | |
| }).toThrow('Too many background tasks are already running.'); | |
| expect(() => { | |
| manager.registerTask(agentTask(new Promise(() => {}), 'agent task')); | |
| }).toThrow('Too many background tasks are already running.'); | |
| }); | |
| it('captures process output', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess( | |
| manager, | |
| immediateProcess(0, 'captured output\n'), | |
| 'echo captured output', | |
| 'capture test', | |
| ); | |
| await waitForOutput(manager, taskId, 'captured output'); | |
| expect(await manager.readOutput(taskId)).toContain('captured output'); | |
| }); | |
| it('terminates a foreground process task that exceeds the output limit', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const chunks = Array.from({ length: 20 }, () => 'x'.repeat(MiB)); | |
| const { proc, killSpy } = streamingProcess(chunks); | |
| let forwardedChars = 0; | |
| const onOutput = vi.fn((_kind: 'stdout' | 'stderr', text: string) => { | |
| forwardedChars += text.length; | |
| }); | |
| const taskId = manager.registerTask( | |
| new ProcessTask( | |
| proc, | |
| 'b3sum --length 18446744073709551615', | |
| 'hash', | |
| onOutput, | |
| ), | |
| { | |
| detached: false, | |
| signal: new AbortController().signal, | |
| timeoutMs: 60_000, | |
| }, | |
| ); | |
| const info = await waitForTerminal(manager, taskId); | |
| expect(info).toMatchObject({ status: 'killed' }); | |
| expect(info?.stopReason ?? '').toMatch(/output limit/i); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(forwardedChars).toBeLessThanOrEqual(LIMIT_BYTES); | |
| }); | |
| it('also terminates a detached process task that exceeds the output limit', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const chunks = Array.from({ length: 20 }, () => 'x'.repeat(MiB)); | |
| const { proc, killSpy } = streamingProcess(chunks); | |
| const taskId = manager.registerTask(new ProcessTask(proc, 'producer', 'bg'), { | |
| detached: true, | |
| timeoutMs: 60_000, | |
| }); | |
| const info = await waitForTerminal(manager, taskId); | |
| expect(info).toMatchObject({ status: 'killed' }); | |
| expect(info?.stopReason ?? '').toMatch(/output limit/i); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| }); | |
| it('stops appending persisted foreground output once the output limit trips', async () => { | |
| const sessionDir = await mkdtemp(join(tmpdir(), 'kimi-bg-limit-fg-')); | |
| try { | |
| const { manager } = createAgentTaskService({ sessionDir }); | |
| const chunks = Array.from({ length: 20 }, () => 'x'.repeat(MiB)); | |
| const { proc } = sigtermIgnoringProcess(chunks); | |
| const taskId = manager.registerTask( | |
| new ProcessTask(proc, 'runaway', 'hash', () => {}), | |
| { | |
| detached: false, | |
| signal: new AbortController().signal, | |
| timeoutMs: 60_000, | |
| }, | |
| ); | |
| const info = await waitForTerminal(manager, taskId); | |
| const output = await manager.getOutputSnapshot(taskId, 1); | |
| expect(info).toMatchObject({ status: 'killed' }); | |
| expect(output.outputSizeBytes).toBeLessThanOrEqual(LIMIT_BYTES); | |
| } finally { | |
| await rm(sessionDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); | |
| } | |
| }); | |
| it('stops appending persisted output once the output limit trips for a detached process task', async () => { | |
| const sessionDir = await mkdtemp(join(tmpdir(), 'kimi-bg-limit-bg-')); | |
| try { | |
| const { manager } = createAgentTaskService({ sessionDir }); | |
| const chunks = Array.from({ length: 20 }, () => 'x'.repeat(MiB)); | |
| const { proc } = sigtermIgnoringProcess(chunks); | |
| const taskId = manager.registerTask( | |
| new ProcessTask(proc, 'runaway', 'background runaway', () => {}), | |
| { | |
| detached: true, | |
| timeoutMs: 60_000, | |
| }, | |
| ); | |
| const info = await waitForTerminal(manager, taskId); | |
| const output = await manager.getOutputSnapshot(taskId, 1); | |
| expect(info).toMatchObject({ status: 'killed' }); | |
| expect(info?.stopReason ?? '').toMatch(/output limit/i); | |
| expect(output.outputSizeBytes).toBeLessThanOrEqual(LIMIT_BYTES); | |
| } finally { | |
| await rm(sessionDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); | |
| } | |
| }); | |
| it('does not cap a detached subagent result larger than the process output limit', async () => { | |
| const sessionDir = await mkdtemp(join(tmpdir(), 'kimi-bg-limit-agent-')); | |
| try { | |
| const { manager } = createAgentTaskService({ sessionDir }); | |
| const result = 'y'.repeat(20 * MiB); | |
| const taskId = manager.registerTask( | |
| agentTask(Promise.resolve({ result }), 'big subagent result'), | |
| { detached: true, timeoutMs: 60_000 }, | |
| ); | |
| const info = await waitForTerminal(manager, taskId); | |
| const output = await manager.getOutputSnapshot(taskId, 1); | |
| expect(info).toMatchObject({ status: 'completed' }); | |
| expect(output.outputSizeBytes).toBe(Buffer.byteLength(result)); | |
| } finally { | |
| await rm(sessionDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); | |
| } | |
| }); | |
| it('fails process tasks when output capture errors after successful exit', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess( | |
| manager, | |
| processWithStdoutError(), | |
| 'ssh example.test', | |
| 'stream error test', | |
| ); | |
| await expect(manager.wait(taskId)).resolves.toMatchObject({ | |
| kind: 'process', | |
| status: 'failed', | |
| exitCode: 0, | |
| stopReason: 'stdout read failed', | |
| }); | |
| }); | |
| it('fails the process task once wait settles after an earlier stream error', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, failStdout, resolveWait } = processWithStdoutErrorBeforeWait(); | |
| const taskId = registerProcess( | |
| manager, | |
| proc, | |
| 'ssh example.test', | |
| 'stream error before wait test', | |
| ); | |
| await Promise.resolve(); | |
| failStdout(); | |
| await Promise.resolve(); | |
| expect(await manager.wait(taskId, 0)).toMatchObject({ | |
| kind: 'process', | |
| status: 'running', | |
| exitCode: null, | |
| }); | |
| resolveWait(0); | |
| await expect(manager.wait(taskId)).resolves.toMatchObject({ | |
| kind: 'process', | |
| status: 'failed', | |
| exitCode: 0, | |
| stopReason: 'stdout read failed', | |
| }); | |
| }); | |
| it('disposes process resources after a process task completes', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const dispose = vi.fn(); | |
| const proc = { | |
| ...immediateProcess(0, 'hello'), | |
| dispose, | |
| } as unknown as IHostProcess; | |
| const taskId = registerProcess(manager, proc, 'echo hello', 'test echo'); | |
| await waitForTerminal(manager, taskId); | |
| await vi.waitFor(() => { | |
| expect(dispose).toHaveBeenCalledTimes(1); | |
| }); | |
| }); | |
| it('transitions process status from exit code', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const successId = registerProcess(manager, immediateProcess(0), 'echo done', 'ok'); | |
| const failureId = registerProcess(manager, immediateProcess(42), 'exit 42', 'fail'); | |
| expect(await manager.wait(successId)).toMatchObject({ | |
| kind: 'process', | |
| status: 'completed', | |
| exitCode: 0, | |
| }); | |
| expect(await manager.wait(failureId)).toMatchObject({ | |
| kind: 'process', | |
| status: 'failed', | |
| exitCode: 42, | |
| }); | |
| }); | |
| it('records failed runtime when proc.wait rejects', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess( | |
| manager, | |
| rejectedProcess(new Error('launch failed')), | |
| '/bogus/cmd', | |
| 'broken launch', | |
| ); | |
| const info = await manager.wait(taskId); | |
| expect(info).toMatchObject({ | |
| status: 'failed', | |
| stopReason: 'launch failed', | |
| }); | |
| expect(info?.endedAt).not.toBeNull(); | |
| }); | |
| it('does not finalize from a visible process exit code before wait settles', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, markExited } = processWithVisibleExitCodeBeforeWait(143); | |
| const taskId = registerProcess(manager, proc, 'sleep 60', 'external kill test'); | |
| markExited(); | |
| expect(manager.getTask(taskId)).toMatchObject({ | |
| kind: 'process', | |
| status: 'running', | |
| exitCode: null, | |
| endedAt: null, | |
| }); | |
| expect(await manager.wait(taskId, 1)).toMatchObject({ | |
| kind: 'process', | |
| status: 'running', | |
| exitCode: null, | |
| }); | |
| }); | |
| it('stop kills a running process and records the stop reason', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(143); | |
| const taskId = registerProcess(manager, proc, 'sleep 60', 'kill test'); | |
| const result = await manager.stop(taskId, 'user requested'); | |
| expect(result).toMatchObject({ | |
| status: 'killed', | |
| stopReason: 'user requested', | |
| exitCode: 143, | |
| }); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| }); | |
| it('includes stopReason for stopped tasks in all-task listings', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess(manager, pendingProcess().proc, 'sleep 60', 'stop reason'); | |
| await manager.stop(taskId, 'superseded by newer task'); | |
| expect(manager.list(false)).toEqual([ | |
| expect.objectContaining({ | |
| taskId, | |
| status: 'killed', | |
| stopReason: 'superseded by newer task', | |
| }), | |
| ]); | |
| }); | |
| it('disposes process resources after a stopped process task settles', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(143); | |
| const dispose = vi.fn(); | |
| const disposableProc = { | |
| ...proc, | |
| dispose, | |
| } as unknown as IHostProcess; | |
| const taskId = registerProcess(manager, disposableProc, 'sleep 60', 'kill test'); | |
| await manager.stop(taskId, 'user requested'); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(dispose).toHaveBeenCalledTimes(1); | |
| }); | |
| it('stop normalizes blank reasons', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, resolve } = manuallyResolvedProcess(); | |
| const taskId = registerProcess(manager, proc, 'sleep 60', 'blank reason test'); | |
| const stopPromise = manager.stop(taskId, ' '); | |
| resolve(0); | |
| const result = await stopPromise; | |
| expect(result).toMatchObject({ status: 'killed' }); | |
| expect(result?.stopReason).toBeUndefined(); | |
| }); | |
| it('stop keeps graceful process shutdown classified as killed', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy, resolve } = manuallyResolvedProcess(); | |
| const taskId = registerProcess(manager, proc, 'sleep 60', 'process race test'); | |
| const stopPromise = manager.stop(taskId, 'user requested'); | |
| resolve(0); | |
| const result = await stopPromise; | |
| expect(result).toMatchObject({ | |
| status: 'killed', | |
| stopReason: 'user requested', | |
| exitCode: 0, | |
| }); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(killSpy).not.toHaveBeenCalledWith('SIGKILL'); | |
| }); | |
| function sigtermOnlyKillProcess(pid: number): { | |
| proc: IHostProcess; | |
| killSpy: ReturnType<typeof vi.fn>; | |
| } { | |
| const stdout = new PassThrough(); | |
| let currentExitCode: number | null = null; | |
| let resolveWait: (code: number) => void = () => {}; | |
| const waitPromise = new Promise<number>((resolve) => { | |
| resolveWait = resolve; | |
| }); | |
| const killSpy = vi.fn(async (signal: NodeJS.Signals) => { | |
| if (currentExitCode !== null) return; | |
| if (signal !== 'SIGKILL') return; | |
| currentExitCode = 137; | |
| stdout.destroy(); | |
| resolveWait(137); | |
| }); | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout, | |
| stderr: Readable.from([]), | |
| pid, | |
| get exitCode(): number | null { | |
| return currentExitCode; | |
| }, | |
| wait: () => waitPromise, | |
| kill: killSpy as unknown as IHostProcess['kill'], | |
| dispose: vi.fn().mockResolvedValue(undefined) as IHostProcess['dispose'], | |
| }; | |
| return { proc, killSpy }; | |
| } | |
| it('escalates a wall-clock timeout to SIGKILL when the process ignores SIGTERM', async () => { | |
| vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }); | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = sigtermOnlyKillProcess(54327); | |
| const taskId = manager.registerTask(new ProcessTask(proc, 'runaway', 'timeout sigkill'), { | |
| timeoutMs: 1, | |
| }); | |
| const terminal = manager.wait(taskId); | |
| await vi.advanceTimersByTimeAsync(1); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(killSpy).not.toHaveBeenCalledWith('SIGKILL'); | |
| await vi.advanceTimersByTimeAsync(5_000); | |
| const info = await terminal; | |
| expect(info?.status).toBe('timed_out'); | |
| expect(killSpy).toHaveBeenCalledWith('SIGKILL'); | |
| }); | |
| it('reports timed_out when a timed-out process exits to SIGTERM within the grace window', async () => { | |
| vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }); | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(); | |
| const taskId = manager.registerTask(new ProcessTask(proc, 'sleep 60', 'timeout graceful'), { | |
| timeoutMs: 1, | |
| }); | |
| const terminal = manager.wait(taskId); | |
| await vi.advanceTimersByTimeAsync(1); | |
| const info = await terminal; | |
| expect(info?.status).toBe('timed_out'); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(killSpy).not.toHaveBeenCalledWith('SIGKILL'); | |
| }); | |
| it('applies the SIGTERM grace + SIGKILL escalation to a detachTimeout deadline', async () => { | |
| vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }); | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = sigtermOnlyKillProcess(54328); | |
| const taskId = manager.registerTask(new ProcessTask(proc, 'runaway', 'detach timeout'), { | |
| detached: false, | |
| detachTimeoutMs: 1, | |
| }); | |
| manager.detach(taskId); | |
| const terminal = manager.wait(taskId); | |
| await vi.advanceTimersByTimeAsync(1); | |
| await vi.advanceTimersByTimeAsync(5_000); | |
| const info = await terminal; | |
| expect(info?.status).toBe('timed_out'); | |
| expect(killSpy).toHaveBeenCalledWith('SIGTERM'); | |
| expect(killSpy).toHaveBeenCalledWith('SIGKILL'); | |
| }); | |
| it('auto-backgrounds a foreground task instead of killing it when its deadline fires', async () => { | |
| vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }); | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(); | |
| const taskId = manager.registerTask(new ProcessTask(proc, 'sleep 60', 'auto background'), { | |
| detached: false, | |
| timeoutMs: 1_000, | |
| detachTimeoutMs: 5_000, | |
| autoBackgroundOnTimeout: true, | |
| }); | |
| const waiting = manager.waitForForegroundRelease(taskId); | |
| await vi.advanceTimersByTimeAsync(1_000); | |
| await expect(waiting).resolves.toBe('timeout_detached'); | |
| expect(killSpy).not.toHaveBeenCalled(); | |
| expect(manager.getTask(taskId)).toMatchObject({ status: 'running', detached: true }); | |
| await vi.advanceTimersByTimeAsync(1_000); | |
| expect(manager.getTask(taskId)?.status).toBe('running'); | |
| await vi.advanceTimersByTimeAsync(4_000); | |
| expect(manager.getTask(taskId)?.status).toBe('timed_out'); | |
| expect(killSpy).toHaveBeenCalled(); | |
| }); | |
| it('kills a foreground task on timeout when auto-background is not enabled', async () => { | |
| vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }); | |
| const { manager } = createAgentTaskService(); | |
| const { proc, killSpy } = pendingProcess(); | |
| const taskId = manager.registerTask(new ProcessTask(proc, 'sleep 60', 'plain timeout'), { | |
| detached: false, | |
| timeoutMs: 1_000, | |
| detachTimeoutMs: 5_000, | |
| }); | |
| const waiting = manager.waitForForegroundRelease(taskId); | |
| await vi.advanceTimersByTimeAsync(1_000); | |
| await expect(waiting).resolves.toBe('terminal'); | |
| expect(killSpy).toHaveBeenCalled(); | |
| expect(manager.getTask(taskId)?.status).toBe('timed_out'); | |
| }); | |
| it('persists graceful process shutdown as killed when stop was requested', async () => { | |
| const sessionDir = await mkdtemp(join(tmpdir(), 'kimi-bg-stop-race-')); | |
| try { | |
| const writer = createAgentTaskService({ sessionDir }).manager; | |
| const { proc, resolve } = manuallyResolvedProcess(); | |
| const taskId = registerProcess(writer, proc, 'sleep 60', 'persisted race'); | |
| const stopPromise = writer.stop(taskId, 'user requested'); | |
| resolve(0); | |
| await stopPromise; | |
| const reader = createAgentTaskService({ sessionDir }).manager; | |
| await reader.loadFromDisk(); | |
| expect(reader.getTask(taskId)).toMatchObject({ | |
| kind: 'process', | |
| status: 'killed', | |
| exitCode: 0, | |
| stopReason: 'user requested', | |
| }); | |
| } finally { | |
| await rm(sessionDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); | |
| } | |
| }); | |
| it('stop preserves agent completion when it wins the stop race', async () => { | |
| const { manager } = createAgentTaskService(); | |
| let resolveCompletion!: (value: { result: string }) => void; | |
| const completion = new Promise<{ result: string }>((resolve) => { | |
| resolveCompletion = resolve; | |
| }); | |
| const controller = new AbortController(); | |
| const abort = vi.spyOn(controller, 'abort'); | |
| const taskId = manager.registerTask( | |
| agentTask(completion, 'agent race test', { abortController: controller }), | |
| ); | |
| const stopPromise = manager.stop(taskId, 'user requested'); | |
| resolveCompletion({ result: 'finished naturally' }); | |
| const result = await stopPromise; | |
| expect(result).toMatchObject({ status: 'completed' }); | |
| expect(result?.stopReason).toBeUndefined(); | |
| expect(await manager.readOutput(taskId)).toContain('finished naturally'); | |
| expect(abort).toHaveBeenCalled(); | |
| }); | |
| it('stop preserves agent failure when a non-abort rejection wins', async () => { | |
| const { manager } = createAgentTaskService(); | |
| let rejectCompletion!: (error: Error) => void; | |
| const completion = new Promise<{ result: string }>((_resolve, reject) => { | |
| rejectCompletion = reject; | |
| }); | |
| const controller = new AbortController(); | |
| const abort = vi.spyOn(controller, 'abort'); | |
| const taskId = manager.registerTask( | |
| agentTask(completion, 'agent failure race test', { abortController: controller }), | |
| ); | |
| const stopPromise = manager.stop(taskId, 'user requested'); | |
| rejectCompletion(new Error('model failed')); | |
| const result = await stopPromise; | |
| expect(result).toMatchObject({ | |
| status: 'failed', | |
| stopReason: 'model failed', | |
| }); | |
| expect(abort).toHaveBeenCalled(); | |
| }); | |
| it('stop marks agent task killed when abort rejection wins', async () => { | |
| const { manager } = createAgentTaskService(); | |
| let rejectCompletion!: (error: Error) => void; | |
| const completion = new Promise<{ result: string }>((_resolve, reject) => { | |
| rejectCompletion = reject; | |
| }); | |
| const abortError = new Error('The operation was aborted.'); | |
| abortError.name = 'AbortError'; | |
| const controller = new AbortController(); | |
| const abort = vi.spyOn(controller, 'abort').mockImplementation((reason?: unknown) => { | |
| AbortController.prototype.abort.call(controller, reason); | |
| rejectCompletion(abortError); | |
| }); | |
| const taskId = manager.registerTask( | |
| agentTask(completion, 'agent abort test', { abortController: controller }), | |
| ); | |
| const result = await manager.stop(taskId, 'user requested'); | |
| expect(result).toMatchObject({ | |
| status: 'killed', | |
| stopReason: 'user requested', | |
| }); | |
| expect(abort).toHaveBeenCalled(); | |
| }); | |
| it('stop finalizes a never-settling agent task after the grace window', async () => { | |
| vi.useFakeTimers(); | |
| const { manager } = createAgentTaskService(); | |
| const controller = new AbortController(); | |
| const abort = vi.spyOn(controller, 'abort'); | |
| const taskId = manager.registerTask( | |
| agentTask(new Promise(() => {}), 'hung agent task', { abortController: controller }), | |
| ); | |
| const stopPromise = manager.stop(taskId, 'user requested'); | |
| await Promise.resolve(); | |
| await vi.advanceTimersByTimeAsync(5_000); | |
| const stopped = await stopPromise; | |
| expect(stopped).toMatchObject({ | |
| status: 'killed', | |
| stopReason: 'user requested', | |
| }); | |
| expect(abort).toHaveBeenCalled(); | |
| }); | |
| it('wait resolves on completion and returns the current snapshot on timeout', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const completedId = registerProcess(manager, immediateProcess(0), 'echo fast', 'wait test'); | |
| expect(await manager.wait(completedId, 5_000)).toMatchObject({ status: 'completed' }); | |
| const runningId = registerProcess(manager, pendingProcess().proc, 'sleep 60', 'timeout'); | |
| expect(await manager.wait(runningId, 0)).toMatchObject({ status: 'running' }); | |
| }); | |
| it('rejects a cancelled wait without stopping the running task', async () => { | |
| const { ctx, manager } = createAgentTaskService(); | |
| const taskId = registerProcess( | |
| manager, | |
| pendingProcess().proc, | |
| 'sleep 60', | |
| 'cancelled wait', | |
| ); | |
| const controller = new AbortController(); | |
| const waiting = manager.wait(taskId, 60_000, controller.signal); | |
| const reason = userCancellationReason(); | |
| controller.abort(reason); | |
| await expect(waiting).rejects.toBe(reason); | |
| expect(manager.getTask(taskId)).toMatchObject({ status: 'running' }); | |
| await manager.stop(taskId, 'test cleanup'); | |
| await ctx.dispose(); | |
| }); | |
| it('wait with a zero timeout returns the immediate snapshot before next-tick completion', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const proc = manuallyResolvedProcess(); | |
| const taskId = registerProcess( | |
| manager, | |
| proc.proc, | |
| 'sleep 0', | |
| 'next-tick completion', | |
| ); | |
| await Promise.resolve(); | |
| setTimeout(() => { | |
| proc.resolve(0); | |
| }, 0); | |
| expect(await manager.wait(taskId, 0)).toMatchObject({ | |
| status: 'running', | |
| exitCode: null, | |
| }); | |
| await expect(manager.wait(taskId)).resolves.toMatchObject({ | |
| status: 'completed', | |
| exitCode: 0, | |
| }); | |
| }); | |
| it('clears task deadline timers when completion wins the race', async () => { | |
| vi.useFakeTimers(); | |
| const { manager } = createAgentTaskService(); | |
| const baselineTimerCount = vi.getTimerCount(); | |
| const taskId = manager.registerTask( | |
| agentTask(Promise.resolve({ result: 'done' }), 'fast deadline task', { | |
| timeoutMs: 60_000, | |
| }), | |
| ); | |
| await expect(manager.wait(taskId, 60_000)).resolves.toMatchObject({ status: 'completed' }); | |
| expect(vi.getTimerCount()).toBeLessThanOrEqual(baselineTimerCount); | |
| }); | |
| it('returns undefined or empty output for unknown task ids', async () => { | |
| const { manager } = createAgentTaskService(); | |
| expect(manager.getTask('bash-nonexist')).toBeUndefined(); | |
| expect(await manager.readOutput('bash-nonexist')).toBe(''); | |
| expect(await manager.stop('bash-nonexist')).toBeUndefined(); | |
| }); | |
| it('stop returns terminal info for an already-exited task', async () => { | |
| const { manager } = createAgentTaskService(); | |
| const taskId = registerProcess(manager, immediateProcess(0), 'echo done', 'already done'); | |
| await manager.wait(taskId); | |
| expect(await manager.stop(taskId, 'too late')).toMatchObject({ | |
| status: 'completed', | |
| stopReason: undefined, | |
| }); | |
| }); | |
| it('getTask on an unknown id does not create persisted state', async () => { | |
| const sessionDir = await mkdtemp(join(tmpdir(), 'kimi-bg-mgr-missing-')); | |
| try { | |
| const { ctx, manager, persistence } = createAgentTaskService({ sessionDir }); | |
| expect(manager.getTask('bash-bogusss0')).toBeUndefined(); | |
| expect(await persistence!.listTasks()).toEqual([]); | |
| await ctx.get(ISessionMetadata).ready; | |
| } finally { | |
| await rm(sessionDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); | |
| } | |
| }); | |
| it('launches a real process and waits to completion', async () => { | |
| const { spawn } = await import('node:child_process'); | |
| const { manager } = createAgentTaskService(); | |
| const child = spawn( | |
| process.execPath, | |
| ['-e', "process.stdout.write('bg-ok\\n')"], | |
| { stdio: 'pipe' }, | |
| ); | |
| const proc: IHostProcess = { | |
| _serviceBrand: undefined, | |
| stdin: { write: vi.fn(), end: vi.fn() } as unknown as Writable, | |
| stdout: child.stdout, | |
| stderr: child.stderr, | |
| pid: child.pid ?? 0, | |
| get exitCode(): number | null { | |
| return child.exitCode; | |
| }, | |
| wait: () => | |
| new Promise<number>((resolve) => { | |
| child.on('exit', (code) => { | |
| resolve(code ?? 0); | |
| }); | |
| }), | |
| kill: vi.fn(async (signal?: NodeJS.Signals) => { | |
| child.kill(signal ?? 'SIGTERM'); | |
| }) as unknown as IHostProcess['kill'], | |
| dispose: vi.fn(async () => { | |
| child.stdin?.destroy(); | |
| child.stdout?.destroy(); | |
| child.stderr?.destroy(); | |
| }) as IHostProcess['dispose'], | |
| }; | |
| const taskId = registerProcess(manager, proc, 'node -e <stdout bg-ok>', 'real worker'); | |
| const info = await manager.wait(taskId, 10_000); | |
| expect(info).toMatchObject({ kind: 'process', status: 'completed', exitCode: 0 }); | |
| expect(await manager.readOutput(taskId)).toContain('bg-ok'); | |
| }, 15_000); | |
| }); | |