Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/abort-signal-listener-ceiling.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@moonshot-ai/kimi-code": patch
---

Silence the MaxListenersExceededWarning that could appear during long agent turns with many parallel tool calls.
4 changes: 4 additions & 0 deletions packages/agent-core-v2/src/agent/loop/loopService.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { randomUUID } from 'node:crypto';
import { EventEmitter } from 'node:events';

import { createControlledPromise } from '@antfu/utils';

Expand Down Expand Up @@ -82,6 +83,8 @@ export const loopLastRequestTraceIdKey = defineState<string | undefined>(
);
export const loopDisposingKey = defineState<boolean>('loop.disposing', () => false);

const MAX_STEP_SIGNAL_LISTENERS = 64;

export class AgentLoopService extends Disposable implements IAgentLoopService {
declare readonly _serviceBrand: undefined;

Expand Down Expand Up @@ -702,6 +705,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
? runtime.turnSignal
: AbortSignal.any([runtime.turnSignal, mutableStep.controller.signal]),
};
EventEmitter.setMaxListeners(MAX_STEP_SIGNAL_LISTENERS, step.signal);
Comment thread
tpoisonooo marked this conversation as resolved.
this.materializeBatch(batch);
return { step };
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,12 +35,16 @@ export async function runCommand(
else signal.addEventListener('abort', onAbort, { once: true });
}

const [stdout, stderr, exitCode] = await Promise.all([
readStream(proc.stdout),
readStream(proc.stderr),
proc.wait().catch(() => -1),
]);
return { exitCode, stdout, stderr };
try {
const [stdout, stderr, exitCode] = await Promise.all([
readStream(proc.stdout),
readStream(proc.stderr),
proc.wait().catch(() => -1),
]);
return { exitCode, stdout, stderr };
} finally {
signal?.removeEventListener('abort', onAbort);
}
}

export function readStream(stream: Readable): Promise<string> {
Expand Down
18 changes: 18 additions & 0 deletions packages/agent-core-v2/test/agent/loop/loop.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { getMaxListeners } from 'node:events';

import { type ToolCall } from '#/kosong/contract/message';
import { emptyUsage } from '#/kosong/contract/usage';
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
Expand Down Expand Up @@ -503,6 +505,22 @@ describe('Agent loop', () => {
);
});

it('raises the abort-listener ceiling on the step signal for parallel tool bursts', async () => {
profile.update({ activeToolNames: [] });
let observed = 0;
loop.hooks.onDidFinishStep.register('test-step-signal-listener-ceiling', async (hookCtx, next) => {
observed = getMaxListeners(hookCtx.signal);
await next();
});

ctx.mockNextResponse({ type: 'text', text: 'answer' });

await ctx.rpc.prompt({ input: [{ type: 'text', text: 'hello' }] });
await ctx.untilTurnEnd();

expect(observed).toBe(64);
});

it('ends the turn when an afterStep hook sets stopTurn even though the model requested tool calls', async () => {
const lookupCall: ToolCall = {
type: 'function',
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { listenerCount, type EventEmitter } from 'node:events';
import { Readable, Writable } from 'node:stream';

import { describe, expect, it } from 'vitest';
Expand Down Expand Up @@ -76,4 +77,10 @@ describe('runCommand', () => {
await promise;
expect(killed).toBe(true);
});

it('removes the abort listener once the command completes', async () => {
const controller = new AbortController();
await runCommand(fakeRunner(fakeProcess()), ['echo'], { signal: controller.signal });
expect(listenerCount(controller.signal as unknown as EventEmitter, 'abort')).toBe(0);
});
});
Loading