Skip to content

Commit 38c7aef

Browse files
committed
fix: keep monitor output flowing until stdout closes and let print mode exit with monitors running
1 parent de6ff0e commit 38c7aef

5 files changed

Lines changed: 91 additions & 2 deletions

File tree

‎.changeset/monitor-tool.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,4 +2,4 @@
22
"@moonshot-ai/kimi-code": minor
33
---
44

5-
Add an experimental Monitor tool that runs a background command and delivers each new output line to the agent as it appears. Enable it with `KIMI_CODE_EXPERIMENTAL_MONITOR=1` or `/experiments`.
5+
Add an experimental Monitor tool, enabled with `KIMI_CODE_EXPERIMENTAL_MONITOR=1` or `/experiments`, that runs a background command and delivers each new output line to the agent as it appears.

‎apps/kimi-code/src/cli/v2/run-v2-print.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,7 @@ import {
7575
shouldEnableTelemetry,
7676
shutdownTelemetry,
7777
} from '@moonshot-ai/kimi-telemetry';
78+
import { isMonitorTaskId } from '@moonshot-ai/agent-core-v2/agent/tools/task/monitor/monitor';
7879
import type { GoalUpdated } from '@moonshot-ai/agent-core-v2/features/goal/goalOps';
7980
import type { TurnEnded } from '@moonshot-ai/agent-core-v2/agent/loop/turnOps';
8081
import type {
@@ -987,7 +988,7 @@ function countPendingBackgroundTasks(session: ISessionScopeHandle): number {
987988
for (const agent of agentManager.list()) {
988989
const handle = agentManager.handleOf(agent.agentId);
989990
if (handle === undefined) continue;
990-
count += handle.accessor.get(IAgentTaskService).list(true).length;
991+
count += handle.accessor.get(IAgentTaskService).list(true).filter((task) => !isMonitorTaskId(task.taskId)).length;
991992
}
992993
return count;
993994
}
@@ -1114,6 +1115,7 @@ async function drainBackgroundTasks(
11141115
if (handle === undefined) continue;
11151116
const taskService = handle.accessor.get(IAgentTaskService);
11161117
for (const task of taskService.list(true)) {
1118+
if (isMonitorTaskId(task.taskId)) continue;
11171119
activeCount++;
11181120
if (seen.has(task.taskId)) continue;
11191121
seen.add(task.taskId);

‎apps/kimi-code/test/cli/v2-run-print.test.ts‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -360,6 +360,35 @@ describe('runV2Print', () => {
360360
expect(app.dispose).toHaveBeenCalled();
361361
});
362362

363+
it.each(['drain', 'steer'])('does not wait for a running monitor in %s mode before shutting down', async (mode) => {
364+
const { app, agentServices, sessionServices, appServices } = makeFakeHarness();
365+
const config = appServices.get(IConfigService) as { get: ReturnType<typeof vi.fn> };
366+
config.get.mockImplementation((section: string) =>
367+
section === 'defaultModel' ? 'k2' : section === 'background' ? { printBackgroundMode: mode } : undefined,
368+
);
369+
const lifecycle = sessionServices.get(IAgentLifecycleService) as { list: ReturnType<typeof vi.fn> };
370+
lifecycle.list.mockReturnValue([{ agentId: 'main' }]);
371+
const taskService = agentServices.get(IAgentTaskService) as {
372+
list: ReturnType<typeof vi.fn>;
373+
stopAllOnExit: ReturnType<typeof vi.fn>;
374+
wait?: ReturnType<typeof vi.fn>;
375+
suppressTerminalNotification?: ReturnType<typeof vi.fn>;
376+
};
377+
taskService.list.mockReturnValue([{ taskId: 'monitor-abc12345', status: 'running' }]);
378+
taskService.wait = vi.fn(() => new Promise(() => {}));
379+
taskService.suppressTerminalNotification = vi.fn(async () => {});
380+
(agentServices.get(IAgentGoalService) as { getGoal: ReturnType<typeof vi.fn> }).getGoal.mockReturnValue({});
381+
mocks.bootstrap.mockReturnValue({ app });
382+
mocks.ensureMainAgent.mockResolvedValue({ agentId: 'main', generation: 1 });
383+
384+
const stderr = writer();
385+
await runV2Print(opts() as never, '1.2.3-test', { stdout: writer(), stderr });
386+
387+
expect(stderr.text()).not.toContain('Warning');
388+
expect(taskService.wait).not.toHaveBeenCalled();
389+
expect(taskService.stopAllOnExit).toHaveBeenCalled();
390+
}, 10_000);
391+
363392
it('passes explicit skill dirs from --skillsDir into bootstrap args', async () => {
364393
const stdout = writer();
365394
const stderr = writer();

‎packages/agent-core-v2/src/agent/tools/os/bash/process-task.ts‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ export class ProcessTask implements AgentTask {
6868
let settlement: AgentTaskSettlement;
6969
try {
7070
const exitCode = await this.proc.wait();
71+
if (this.stdoutEvents) await waitForStreamEndOrAbort(streamDrained, sink.signal);
7172
await waitForStreamDrain(streamDrained);
7273
this.exitCode = exitCode;
7374
settlement = {
@@ -134,6 +135,22 @@ async function waitForStreamDrain(streamDrained: Promise<void>): Promise<void> {
134135
}
135136
}
136137

138+
async function waitForStreamEndOrAbort(streamDrained: Promise<void>, signal: AbortSignal): Promise<void> {
139+
if (signal.aborted) return;
140+
let onAbort: (() => void) | undefined;
141+
try {
142+
await Promise.race([
143+
streamDrained.catch(() => {}),
144+
new Promise<void>((resolve) => {
145+
onAbort = resolve;
146+
signal.addEventListener('abort', onAbort, { once: true });
147+
}),
148+
]);
149+
} finally {
150+
if (onAbort !== undefined) signal.removeEventListener('abort', onAbort);
151+
}
152+
}
153+
137154
async function waitForStreamDrainSettled(streamDrained: Promise<void>): Promise<void> {
138155
try {
139156
await waitForStreamDrain(streamDrained);

‎packages/agent-core-v2/test/agent/task/taskManager.test.ts‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1645,6 +1645,47 @@ describe('AgentTaskService monitor events', () => {
16451645
await fixture.ctx.dispose();
16461646
});
16471647

1648+
it('keeps delivering monitor lines after the shell exits until stdout closes', async () => {
1649+
const fixture = createAgentTaskService();
1650+
const notes = captureNotifications(fixture);
1651+
const { proc, stdout, stderr } = controllableProcess();
1652+
const shellExit = Promise.withResolvers<number>();
1653+
1654+
const taskId = fixture.manager.registerTask(
1655+
new MonitorProcessTask({ ...proc, wait: () => shellExit.promise }, 'tail -F app.log &', 'watch app log'),
1656+
);
1657+
shellExit.resolve(0);
1658+
await new Promise((resolve) => setTimeout(resolve, 500));
1659+
expect(fixture.manager.getTask(taskId)?.status).toBe('running');
1660+
1661+
stdout.end('late line\n');
1662+
stderr.end();
1663+
await waitForTerminal(fixture.manager, taskId);
1664+
await vi.waitFor(() => {
1665+
expect(notes).toHaveLength(2);
1666+
});
1667+
expect(noteText(notes[0]!)).toContain('late line');
1668+
expect(noteText(notes[1]!)).toContain('type="task.completed"');
1669+
await fixture.ctx.dispose();
1670+
});
1671+
1672+
it('stops a monitor whose shell exited while its stdout is still open', async () => {
1673+
const fixture = createAgentTaskService();
1674+
captureNotifications(fixture);
1675+
const { proc } = controllableProcess();
1676+
const shellExit = Promise.withResolvers<number>();
1677+
1678+
const taskId = fixture.manager.registerTask(
1679+
new MonitorProcessTask({ ...proc, wait: () => shellExit.promise }, 'tail -F app.log &', 'watch app log'),
1680+
);
1681+
shellExit.resolve(0);
1682+
await new Promise((resolve) => setTimeout(resolve, 50));
1683+
1684+
expect(await fixture.manager.stop(taskId)).toMatchObject({ status: 'killed' });
1685+
expect(proc.kill).toHaveBeenCalledWith('SIGTERM');
1686+
await fixture.ctx.dispose();
1687+
});
1688+
16481689
it('does not turn stdout of an ordinary background command into events', async () => {
16491690
const fixture = createAgentTaskService();
16501691
const notes = captureNotifications(fixture);

0 commit comments

Comments
 (0)