diff --git a/.changeset/abort-signal-listener-ceiling.md b/.changeset/abort-signal-listener-ceiling.md new file mode 100644 index 00000000000..cc0a62ca2e8 --- /dev/null +++ b/.changeset/abort-signal-listener-ceiling.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Silence the MaxListenersExceededWarning that could appear during long agent turns with many parallel tool calls. diff --git a/packages/agent-core-v2/src/agent/loop/loopService.ts b/packages/agent-core-v2/src/agent/loop/loopService.ts index 87e8c7f6854..d0a7e312a08 100644 --- a/packages/agent-core-v2/src/agent/loop/loopService.ts +++ b/packages/agent-core-v2/src/agent/loop/loopService.ts @@ -1,4 +1,5 @@ import { randomUUID } from 'node:crypto'; +import { EventEmitter } from 'node:events'; import { createControlledPromise } from '@antfu/utils'; @@ -82,6 +83,8 @@ export const loopLastRequestTraceIdKey = defineState( ); export const loopDisposingKey = defineState('loop.disposing', () => false); +const MAX_STEP_SIGNAL_LISTENERS = 64; + export class AgentLoopService extends Disposable implements IAgentLoopService { declare readonly _serviceBrand: undefined; @@ -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); this.materializeBatch(batch); return { step }; } diff --git a/packages/agent-core-v2/src/workspace/workspaceFs/internal/fsProcess.ts b/packages/agent-core-v2/src/workspace/workspaceFs/internal/fsProcess.ts index 669f33c8049..c51f795e818 100644 --- a/packages/agent-core-v2/src/workspace/workspaceFs/internal/fsProcess.ts +++ b/packages/agent-core-v2/src/workspace/workspaceFs/internal/fsProcess.ts @@ -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 { diff --git a/packages/agent-core-v2/test/agent/loop/loop.test.ts b/packages/agent-core-v2/test/agent/loop/loop.test.ts index 5fd1356b53c..f123de319e4 100644 --- a/packages/agent-core-v2/test/agent/loop/loop.test.ts +++ b/packages/agent-core-v2/test/agent/loop/loop.test.ts @@ -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'; @@ -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', diff --git a/packages/agent-core-v2/test/workspace/workspaceFs/fsProcess.test.ts b/packages/agent-core-v2/test/workspace/workspaceFs/fsProcess.test.ts index 370ad7636cb..329cc934cbf 100644 --- a/packages/agent-core-v2/test/workspace/workspaceFs/fsProcess.test.ts +++ b/packages/agent-core-v2/test/workspace/workspaceFs/fsProcess.test.ts @@ -1,3 +1,4 @@ +import { listenerCount, type EventEmitter } from 'node:events'; import { Readable, Writable } from 'node:stream'; import { describe, expect, it } from 'vitest'; @@ -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); + }); });