Skip to content

Commit 5bcd3cc

Browse files
committed
fix(agent-core-v2): keep abort listeners under the listener ceiling
1 parent 08edb44 commit 5bcd3cc

5 files changed

Lines changed: 44 additions & 6 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@pymodel/pythinker-code": patch
3+
---
4+
5+
Silence the MaxListenersExceededWarning that could appear during long agent turns with many parallel tool calls.

packages/agent-core-v2/src/agent/loop/loopService.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { randomUUID } from 'node:crypto';
2+
import { EventEmitter } from 'node:events';
23

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

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

86+
const MAX_STEP_SIGNAL_LISTENERS = 64;
87+
8588
export class AgentLoopService extends Disposable implements IAgentLoopService {
8689
declare readonly _serviceBrand: undefined;
8790

@@ -702,6 +705,7 @@ export class AgentLoopService extends Disposable implements IAgentLoopService {
702705
? runtime.turnSignal
703706
: AbortSignal.any([runtime.turnSignal, mutableStep.controller.signal]),
704707
};
708+
EventEmitter.setMaxListeners(MAX_STEP_SIGNAL_LISTENERS, step.signal);
705709
this.materializeBatch(batch);
706710
return { step };
707711
}

packages/agent-core-v2/src/workspace/workspaceFs/internal/fsProcess.ts

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -35,12 +35,16 @@ export async function runCommand(
3535
else signal.addEventListener('abort', onAbort, { once: true });
3636
}
3737

38-
const [stdout, stderr, exitCode] = await Promise.all([
39-
readStream(proc.stdout),
40-
readStream(proc.stderr),
41-
proc.wait().catch(() => -1),
42-
]);
43-
return { exitCode, stdout, stderr };
38+
try {
39+
const [stdout, stderr, exitCode] = await Promise.all([
40+
readStream(proc.stdout),
41+
readStream(proc.stderr),
42+
proc.wait().catch(() => -1),
43+
]);
44+
return { exitCode, stdout, stderr };
45+
} finally {
46+
signal?.removeEventListener('abort', onAbort);
47+
}
4448
}
4549

4650
export function readStream(stream: Readable): Promise<string> {

packages/agent-core-v2/test/agent/loop/loop.test.ts

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import { getMaxListeners } from 'node:events';
2+
13
import { type ToolCall } from '#/kosong/contract/message';
24
import { emptyUsage } from '#/kosong/contract/usage';
35
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
@@ -503,6 +505,22 @@ describe('Agent loop', () => {
503505
);
504506
});
505507

508+
it('raises the abort-listener ceiling on the step signal for parallel tool bursts', async () => {
509+
profile.update({ activeToolNames: [] });
510+
let observed = 0;
511+
loop.hooks.onDidFinishStep.register('test-step-signal-listener-ceiling', async (hookCtx, next) => {
512+
observed = getMaxListeners(hookCtx.signal);
513+
await next();
514+
});
515+
516+
ctx.mockNextResponse({ type: 'text', text: 'answer' });
517+
518+
await ctx.rpc.prompt({ input: [{ type: 'text', text: 'hello' }] });
519+
await ctx.untilTurnEnd();
520+
521+
expect(observed).toBe(64);
522+
});
523+
506524
it('ends the turn when an afterStep hook sets stopTurn even though the model requested tool calls', async () => {
507525
const lookupCall: ToolCall = {
508526
type: 'function',

packages/agent-core-v2/test/workspace/workspaceFs/fsProcess.test.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import { listenerCount, type EventEmitter } from 'node:events';
12
import { Readable, Writable } from 'node:stream';
23

34
import { describe, expect, it } from 'vitest';
@@ -76,4 +77,10 @@ describe('runCommand', () => {
7677
await promise;
7778
expect(killed).toBe(true);
7879
});
80+
81+
it('removes the abort listener once the command completes', async () => {
82+
const controller = new AbortController();
83+
await runCommand(fakeRunner(fakeProcess()), ['echo'], { signal: controller.signal });
84+
expect(listenerCount(controller.signal as unknown as EventEmitter, 'abort')).toBe(0);
85+
});
7986
});

0 commit comments

Comments
 (0)