diff --git a/packages/agent-core-v2/src/_base/event.ts b/packages/agent-core-v2/src/_base/event.ts index a80483c2c..1c53bb85a 100644 --- a/packages/agent-core-v2/src/_base/event.ts +++ b/packages/agent-core-v2/src/_base/event.ts @@ -112,6 +112,24 @@ export type IWaitUntilData = Omit; export class AsyncEmitter extends Emitter { private _asyncDeliveryQueue?: LinkedList<[(event: T) => void, IWaitUntilData]>; + async fireAsyncConcurrent(data: IWaitUntilData, signal: AbortSignal): Promise { + if (this.isDisposed || this._listeners === undefined || signal.aborted) { + return; + } + const snapshot = Array.from(this._listeners); + await Promise.all( + snapshot.map((entry) => + this.deliverAsync( + (event) => { + entry.listener.call(entry.thisArg, event); + }, + data, + signal, + ), + ), + ); + } + async fireAsync(data: IWaitUntilData, signal: AbortSignal): Promise { if (this.isDisposed || this._listeners === undefined) { return; @@ -129,32 +147,37 @@ export class AsyncEmitter extends Emitter { while (this._asyncDeliveryQueue.size > 0 && !signal.aborted) { const [deliver, eventData] = this._asyncDeliveryQueue.shift()!; - const thenables: Promise[] = []; + await this.deliverAsync(deliver, eventData, signal); + } + } - const event = { - ...eventData, - signal, - waitUntil: (p: Promise): void => { - if (Object.isFrozen(thenables)) { - throw new Error('waitUntil can NOT be called asynchronously'); - } - thenables.push(p); - }, - } as T; - - try { - deliver(event); - } catch (error) { - onUnexpectedError(error); - continue; - } - - void Object.freeze(thenables); - const settled = await Promise.allSettled(thenables); - for (const result of settled) { - if (result.status === 'rejected') { - onUnexpectedError(result.reason); + private async deliverAsync( + deliver: (event: T) => void, + data: IWaitUntilData, + signal: AbortSignal, + ): Promise { + const thenables: Promise[] = []; + const event = { + ...data, + signal, + waitUntil: (p: Promise): void => { + if (Object.isFrozen(thenables)) { + throw new Error('waitUntil can NOT be called asynchronously'); } + thenables.push(p); + }, + } as T; + try { + deliver(event); + } catch (error) { + onUnexpectedError(error); + return; + } + void Object.freeze(thenables); + const settled = await Promise.allSettled(thenables); + for (const result of settled) { + if (result.status === 'rejected') { + onUnexpectedError(result.reason); } } } diff --git a/packages/agent-core-v2/src/app/mcpConfig/configStore.ts b/packages/agent-core-v2/src/app/mcpConfig/configStore.ts index b0f94a169..93b06d722 100644 --- a/packages/agent-core-v2/src/app/mcpConfig/configStore.ts +++ b/packages/agent-core-v2/src/app/mcpConfig/configStore.ts @@ -162,7 +162,7 @@ export class McpConfigStore extends Disposable implements IMcpConfigStore { await this.storage.write(CONFIG_SCOPE, MCP_CONFIG_KEY, textEncoder.encode(text), { atomic: true, }); - await this.writeEmitter.fireAsync({}, NO_ABORT); + await this.writeEmitter.fireAsyncConcurrent({}, NO_ABORT); } } diff --git a/packages/agent-core-v2/test/app/mcpConfig/configStore.test.ts b/packages/agent-core-v2/test/app/mcpConfig/configStore.test.ts index 19948f217..2b4d40b8c 100644 --- a/packages/agent-core-v2/test/app/mcpConfig/configStore.test.ts +++ b/packages/agent-core-v2/test/app/mcpConfig/configStore.test.ts @@ -317,5 +317,33 @@ describe('McpConfigStore', () => { await mutation; expect(completed).toBe(true); }); + + it('starts asynchronous listeners concurrently before waiting for completion', async () => { + let resolveStarted!: () => void; + const started = new Promise((resolve) => { + resolveStarted = resolve; + }); + let release!: () => void; + const gate = new Promise((resolve) => { + release = resolve; + }); + let secondStarted = false; + store.onDidWrite((event) => { + resolveStarted(); + event.waitUntil(gate); + }); + store.onDidWrite(() => { + secondStarted = true; + }); + + const mutation = store.add(stdioServer('alpha')); + await started; + await Promise.resolve(); + const secondStartedBeforeRelease = secondStarted; + release(); + await mutation; + + expect(secondStartedBeforeRelease).toBe(true); + }); }); });