diff --git a/src/platform/endpoint/node/automodeService.ts b/src/platform/endpoint/node/automodeService.ts index 1341b8d932..e474e5db85 100644 --- a/src/platform/endpoint/node/automodeService.ts +++ b/src/platform/endpoint/node/automodeService.ts @@ -117,6 +117,7 @@ export interface IAutomodeService { export class AutomodeService extends Disposable implements IAutomodeService { readonly _serviceBrand: undefined; private readonly _autoModelCache: Map = new Map(); + private readonly _pendingResolutions: Map> = new Map(); private _reserveTokens: DisposableMap = new DisposableMap(); private readonly _routerDecisionFetcher: RouterDecisionFetcher; @@ -137,6 +138,7 @@ export class AutomodeService extends Disposable implements IAutomodeService { entry.tokenBank.dispose(); } this._autoModelCache.clear(); + this._pendingResolutions.clear(); const keys = Array.from(this._reserveTokens.keys()); this._reserveTokens.clearAndDisposeAll(); for (const location of keys) { @@ -152,6 +154,7 @@ export class AutomodeService extends Disposable implements IAutomodeService { entry.tokenBank.dispose(); } this._autoModelCache.clear(); + this._pendingResolutions.clear(); this._reserveTokens.dispose(); super.dispose(); } @@ -175,6 +178,38 @@ export class AutomodeService extends Disposable implements IAutomodeService { } const conversationId = chatRequest?.sessionResource?.toString() ?? chatRequest?.sessionId ?? 'unknown'; + + // Coalesce concurrent calls for the same conversation so we don't fire + // duplicate router requests or telemetry. Concurrent callers within a + // single turn always share the same prompt, so reusing the result is safe. + return this._singleFlight(conversationId, () => this._resolveAutoModeEndpointCore(chatRequest, conversationId, knownEndpoints)); + } + + /** + * Single-flight coalescer: if a resolution for `key` is already in progress, + * return that pending promise instead of starting a new one. + * + * This function MUST remain synchronous (not `async`). The atomicity of the + * `get -> check -> set` sequence relies on no `await` running between them, and + * keeping the function non-`async` makes that invariant structural rather than + * a comment a future edit might miss. + */ + private _singleFlight(key: string, fn: () => Promise): Promise { + const existing = this._pendingResolutions.get(key); + if (existing) { + return existing; + } + const promise = fn().finally(() => { + // Reference-equality guard: only clear if this entry is still ours. + if (this._pendingResolutions.get(key) === promise) { + this._pendingResolutions.delete(key); + } + }); + this._pendingResolutions.set(key, promise); + return promise; + } + + private async _resolveAutoModeEndpointCore(chatRequest: ChatRequest | undefined, conversationId: string, knownEndpoints: IChatEndpoint[]): Promise { const entry = this._autoModelCache.get(conversationId); const tokenBank = this._acquireTokenBank(entry, chatRequest?.location, conversationId); const token = await tokenBank.getToken(); diff --git a/src/platform/endpoint/node/test/automodeService.spec.ts b/src/platform/endpoint/node/test/automodeService.spec.ts index 3d66f3cece..6239b7a5c3 100644 --- a/src/platform/endpoint/node/test/automodeService.spec.ts +++ b/src/platform/endpoint/node/test/automodeService.spec.ts @@ -916,4 +916,103 @@ describe('AutomodeService', () => { ); }); }); + + describe('concurrent resolution coalescing', () => { + it('should return consistent model selection for concurrent calls with the same conversationId', async () => { + mockApiResponse(['gpt-4o', 'gpt-4o-mini']); + automodeService = createService(); + + const chatRequest: Partial = { + location: ChatLocation.Panel, + prompt: 'hello', + sessionId: 'session-concurrent' + }; + + const [result1, result2, result3] = await Promise.all([ + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + ]); + + expect(result1.model).toBe(result2.model); + expect(result2.model).toBe(result3.model); + expect(result1.modelProvider).toBe(result2.modelProvider); + }); + + it('should make only one CAPI token request for concurrent calls', async () => { + mockApiResponse(['gpt-4o', 'gpt-4o-mini']); + automodeService = createService(); + + const chatRequest: Partial = { + location: ChatLocation.Panel, + prompt: 'hello', + sessionId: 'session-concurrent-token' + }; + + await Promise.all([ + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + ]); + + // Only one CAPI token request for auto models should be made (not two) + const capiCalls = (mockCAPIClientService.makeRequest as ReturnType).mock.calls; + const autoModelsCalls = capiCalls.filter(call => call[1]?.type === RequestType.AutoModels); + expect(autoModelsCalls.length).toBe(1); + }); + + it('should allow a new resolution after the first one completes', async () => { + mockApiResponse(['gpt-4o', 'gpt-4o-mini']); + automodeService = createService(); + + const chatRequest: Partial = { + location: ChatLocation.Panel, + prompt: 'hello', + sessionId: 'session-sequential' + }; + + const result1 = await automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]); + + // Second call after first completes should succeed (hits cache, not pending map) + const result2 = await automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]); + expect(result1.model).toBe(result2.model); + }); + + it('should propagate errors to all concurrent callers', async () => { + (mockCAPIClientService.makeRequest as ReturnType).mockRejectedValue(new Error('token fetch failed')); + automodeService = createService(); + + const chatRequest: Partial = { + location: ChatLocation.Panel, + prompt: 'hello', + sessionId: 'session-error' + }; + + const results = await Promise.allSettled([ + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]), + ]); + + expect(results[0].status).toBe('rejected'); + expect(results[1].status).toBe('rejected'); + }); + + it('should allow retry after a failed concurrent resolution', async () => { + (mockCAPIClientService.makeRequest as ReturnType).mockRejectedValueOnce(new Error('transient failure')); + automodeService = createService(); + + const chatRequest: Partial = { + location: ChatLocation.Panel, + prompt: 'hello', + sessionId: 'session-retry' + }; + + // First call fails + await expect(automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint])).rejects.toThrow(); + + // Pending promise should be cleared; retry should work + mockApiResponse(['gpt-4o', 'gpt-4o-mini']); + const result = await automodeService.resolveAutoModeEndpoint(chatRequest as ChatRequest, [mockChatEndpoint]); + expect(result.model).toBeDefined(); + }); + }); });