diff --git a/src/commands/set-context-window/set-context-window.test.ts b/src/commands/set-context-window/set-context-window.test.ts new file mode 100644 index 0000000000..d12c5885f5 --- /dev/null +++ b/src/commands/set-context-window/set-context-window.test.ts @@ -0,0 +1,65 @@ +import { afterEach, beforeEach, expect, test } from 'bun:test' +import { + acquireSharedMutationLock, + releaseSharedMutationLock, +} from '../../test/sharedMutationLock.js' +import { call as clearContextWindow } from '../clear-context-window/clear-context-window.js' +import { call as setContextWindow } from './set-context-window.js' +import { + clearSessionContextWindowOverride, + getContextWindowForModel, +} from '../../utils/context.js' +import { + getAutoCompactThreshold, + getEffectiveContextWindowSize, +} from '../../services/compact/autoCompact.js' + +let hasSharedMutationLock = false +const savedEnv = { + CLAUDE_CODE_AUTO_COMPACT_WINDOW: process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW, + CLAUDE_CODE_MAX_OUTPUT_TOKENS: process.env.CLAUDE_CODE_MAX_OUTPUT_TOKENS, + CLAUDE_AUTOCOMPACT_PCT_OVERRIDE: + process.env.CLAUDE_AUTOCOMPACT_PCT_OVERRIDE, +} + +beforeEach(async () => { + await acquireSharedMutationLock('commands/set-context-window/set-context-window.test.ts') + hasSharedMutationLock = true + process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW = '100000' + process.env.CLAUDE_CODE_MAX_OUTPUT_TOKENS = '20000' + delete process.env.CLAUDE_AUTOCOMPACT_PCT_OVERRIDE + clearSessionContextWindowOverride('claude-sonnet-4') +}) + +afterEach(() => { + clearSessionContextWindowOverride('claude-sonnet-4') + for (const [key, value] of Object.entries(savedEnv)) { + if (value === undefined) { + delete process.env[key] + } else { + process.env[key] = value + } + } + if (hasSharedMutationLock) { + releaseSharedMutationLock() + hasSharedMutationLock = false + } +}) + +test('/set-context-window updates tokens and compaction immediately; clear restores defaults', async () => { + const context = { + options: { mainLoopModel: 'claude-sonnet-4' }, + } as never + + expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(50_000) + const result = await setContextWindow('1000000', context) + expect(result).toMatchObject({ type: 'text', value: expect.stringContaining('1,000,000') }) + expect(getContextWindowForModel('claude-sonnet-4')).toBe(1_000_000) + expect(getEffectiveContextWindowSize('claude-sonnet-4')).toBe(980_000) + expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(950_000) + + await clearContextWindow('claude-sonnet-4', context) + expect(getContextWindowForModel('claude-sonnet-4')).toBe(200_000) + expect(getEffectiveContextWindowSize('claude-sonnet-4')).toBe(80_000) + expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(50_000) +}) diff --git a/src/services/compact/autoCompact.test.ts b/src/services/compact/autoCompact.test.ts index aefa9e9343..bbb09ca0c1 100644 --- a/src/services/compact/autoCompact.test.ts +++ b/src/services/compact/autoCompact.test.ts @@ -122,6 +122,7 @@ beforeEach(async () => { delete process.env.CLAUDE_CODE_MAX_CONTEXT_TOKENS delete process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW delete process.env.CLAUDE_CODE_MAX_OUTPUT_TOKENS + delete process.env.USER_TYPE } catch (error) { releaseSharedMutationLock() hasSharedMutationLock = false @@ -335,6 +336,41 @@ describe('getAutoCompactThreshold', () => { expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(20_000) }) + test('internal context cap keeps precedence over session override and auto-compact cap', async () => { + process.env.USER_TYPE = 'ant' + process.env.CLAUDE_CODE_MAX_CONTEXT_TOKENS = '200000' + process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW = '100000' + realContext.setSessionContextWindowOverride('claude-sonnet-4', 1_000_000) + const { getEffectiveContextWindowSize, getAutoCompactThreshold } = + await importAutoCompact() + + try { + expect(getEffectiveContextWindowSize('claude-sonnet-4')).toBe(80_000) + expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(50_000) + } finally { + realContext.clearSessionContextWindowOverride('claude-sonnet-4') + delete process.env.CLAUDE_CODE_MAX_CONTEXT_TOKENS + } + }) + + test('session context-window override immediately updates auto-compact threshold', async () => { + process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW = '100000' + process.env.CLAUDE_CODE_MAX_OUTPUT_TOKENS = '20000' + delete process.env.CLAUDE_AUTOCOMPACT_PCT_OVERRIDE + const { getAutoCompactThreshold, getEffectiveContextWindowSize } = + await importAutoCompact() + + expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(50_000) + realContext.setSessionContextWindowOverride('claude-sonnet-4', 1_000_000) + + try { + expect(getEffectiveContextWindowSize('claude-sonnet-4')).toBe(980_000) + expect(getAutoCompactThreshold('claude-sonnet-4')).toBe(950_000) + } finally { + realContext.clearSessionContextWindowOverride('claude-sonnet-4') + } + }) + test('keeps compaction and warning thresholds usable across mid-sized windows', async () => { process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW = '64000' const { calculateTokenWarningState, getAutoCompactThreshold } = diff --git a/src/services/compact/autoCompact.ts b/src/services/compact/autoCompact.ts index b185cded64..d93c27a5bc 100644 --- a/src/services/compact/autoCompact.ts +++ b/src/services/compact/autoCompact.ts @@ -5,7 +5,10 @@ import type { QuerySource } from '../../constants/querySource.js' import type { ToolUseContext } from '../../Tool.js' import type { Message } from '../../types/message.js' import { getGlobalConfig } from '../../utils/config.js' -import { getContextWindowForModel } from '../../utils/context.js' +import { + getContextWindowForModel, + getSessionContextWindowOverride, +} from '../../utils/context.js' import { logForDebugging } from '../../utils/debug.js' import { isEnvTruthy } from '../../utils/envUtils.js' import { hasExactErrorMessage } from '../../utils/errors.js' @@ -49,8 +52,21 @@ export function getEffectiveContextWindowSize( runtimeLimits, ) + // An explicit per-model session setting is more specific than the legacy + // auto-compact cap. Otherwise that cap can silently keep a newly configured + // large context window stuck at (for example) 100k for compaction decisions. + const internalContextWindowOverride = + process.env.USER_TYPE === 'ant' + ? parseInt(process.env.CLAUDE_CODE_MAX_CONTEXT_TOKENS ?? '', 10) + : NaN + const internalOverrideTakesPrecedence = + Number.isFinite(internalContextWindowOverride) && + internalContextWindowOverride > 0 + const hasSessionOverride = + getSessionContextWindowOverride(model) !== undefined && + !internalOverrideTakesPrecedence const autoCompactWindow = process.env.CLAUDE_CODE_AUTO_COMPACT_WINDOW - if (autoCompactWindow) { + if (autoCompactWindow && !hasSessionOverride) { const parsed = parseInt(autoCompactWindow, 10) if (!isNaN(parsed) && parsed > 0) { contextWindow = Math.min(contextWindow, parsed) diff --git a/src/services/mcp/client.activity.test.ts b/src/services/mcp/client.activity.test.ts index 5639964d78..0b44feadd3 100644 --- a/src/services/mcp/client.activity.test.ts +++ b/src/services/mcp/client.activity.test.ts @@ -3,8 +3,12 @@ import { ErrorCode, McpError, } from '@modelcontextprotocol/sdk/types.js' +import { AbortError } from '../../utils/errors.js' import { createAssistantMessage } from '../../utils/messages.js' -import { fetchToolsForClient } from './client.js' +import { + callMCPToolWithUrlElicitationRetry, + fetchToolsForClient, +} from './client.js' import type { ConnectedMCPServer } from './types.js' describe('MCP tool activity', () => { @@ -12,6 +16,28 @@ describe('MCP tool activity', () => { vi.useRealTimers() }) + test('URL elicitation cancellation uses the shared AbortError type', async () => { + const abortController = new AbortController() + abortController.abort() + const callToolFn = vi.fn(async () => ({ content: [] })) + + await expect( + callMCPToolWithUrlElicitationRetry({ + client: {} as never, + clientConnection: { + type: 'connected', + name: 'elicitation-cancel-test', + } as never, + tool: 'slow-tool', + args: {}, + signal: abortController.signal, + setAppState: vi.fn(), + callToolFn, + }), + ).rejects.toBeInstanceOf(AbortError) + expect(callToolFn).not.toHaveBeenCalled() + }) + test('emits progress heartbeats while a tool call remains pending', async () => { vi.useFakeTimers() let resolveToolCall: ((result: unknown) => void) | undefined @@ -155,6 +181,79 @@ describe('MCP tool activity', () => { expect(onProgress).toHaveBeenCalledTimes(progressCountAfterCompletion) }) + test('aborts the MCP request when the tool-call timeout expires', async () => { + vi.useFakeTimers() + const originalTimeout = process.env.MCP_TOOL_TIMEOUT + process.env.MCP_TOOL_TIMEOUT = '1000' + let requestSignal: AbortSignal | undefined + const sdkClient = { + request: vi.fn(async () => ({ + tools: [ + { + name: 'slow-tool', + inputSchema: { type: 'object' }, + }, + ], + })), + callTool: vi.fn( + ( + _request: unknown, + _schema: unknown, + options: { signal?: AbortSignal }, + ) => { + requestSignal = options.signal + return new Promise(() => {}) + }, + ), + } + const connection = { + type: 'connected', + name: 'timeout-abort-test', + config: { type: 'sdk' }, + capabilities: { tools: {} }, + client: sdkClient, + } as unknown as ConnectedMCPServer + + try { + const [tool] = await fetchToolsForClient(connection) + expect(tool).toBeDefined() + const parentMessage = createAssistantMessage({ + content: [ + { + type: 'tool_use', + id: 'toolu_timeout_abort', + name: tool!.name, + input: {}, + }, + ], + }) + const callPromise = tool!.call( + {}, + { + abortController: new AbortController(), + setAppState: vi.fn(), + } as never, + undefined as never, + parentMessage, + ) + + await Promise.resolve() + await Promise.resolve() + expect(requestSignal).toBeDefined() + expect(requestSignal?.aborted).toBe(false) + + vi.advanceTimersByTime(1000) + await expect(callPromise).rejects.toThrow('timed out after 1s') + expect(requestSignal?.aborted).toBe(true) + } finally { + if (originalTimeout === undefined) { + delete process.env.MCP_TOOL_TIMEOUT + } else { + process.env.MCP_TOOL_TIMEOUT = originalTimeout + } + } + }) + test('a throwing started callback does not prevent the tool call', async () => { const sdkClient = { request: vi.fn(async () => ({ diff --git a/src/services/mcp/client.test.ts b/src/services/mcp/client.test.ts index 1b58bb2aa1..d0149409d2 100644 --- a/src/services/mcp/client.test.ts +++ b/src/services/mcp/client.test.ts @@ -1,16 +1,20 @@ import assert from 'node:assert/strict' import test from 'node:test' +import { Client } from '@modelcontextprotocol/sdk/client/index.js' import { UnauthorizedError } from '@modelcontextprotocol/sdk/client/auth.js' import { SSEClientTransport } from '@modelcontextprotocol/sdk/client/sse.js' import { StreamableHTTPClientTransport } from '@modelcontextprotocol/sdk/client/streamableHttp.js' +import { CallToolResultSchema } from '@modelcontextprotocol/sdk/types.js' import { appendBoundedMcpStderr, + attachMcpRequestCancellationHandler, buildMcpSseEventSourceHeaders, buildMcpSseRequestHeaders, cleanupFailedConnection, buildMcpStdioCommand, logMcpServerStderr, + wrapFetchWithRequestCancellation, } from './client.js' import { wrapFetchWithStepUpDetection } from './auth.js' import { @@ -206,6 +210,79 @@ test('SSE failed recovery reaches UnauthorizedError without leaking the resource assert.equal(fixture.getRedirectCalls(), 1) }) +test('timed-out Streamable HTTP tool calls cancel the response stream after headers', async () => { + const activeRequests = new Map() + let streamCancelled = false + const server = Bun.serve({ + port: 0, + async fetch(request) { + if (request.method === 'GET') { + return new Response(null, { status: 405 }) + } + const text = await request.text() + if (!text) return new Response(null, { status: 202 }) + const payload = JSON.parse(text) as { id: number; method: string } + if (payload.method === 'initialize') { + return Response.json({ + jsonrpc: '2.0', + id: payload.id, + result: { + protocolVersion: '2025-03-26', + capabilities: {}, + serverInfo: { name: 'timeout-stream-test', version: '1.0.0' }, + }, + }) + } + if (payload.method === 'notifications/initialized') { + return new Response(null, { status: 202 }) + } + if (payload.method === 'tools/call') { + return new Response( + new ReadableStream({ + start() {}, + cancel() { + streamCancelled = true + }, + }), + { + headers: { 'content-type': 'text/event-stream' }, + }, + ) + } + return new Response(null, { status: 202 }) + }, + }) + const transport = new StreamableHTTPClientTransport( + new URL(`http://127.0.0.1:${server.port}/mcp`), + { + fetch: wrapFetchWithRequestCancellation(fetch, activeRequests), + }, + ) + attachMcpRequestCancellationHandler(transport, activeRequests) + const client = new Client({ name: 'timeout-test', version: '1.0.0' }) + + try { + await client.connect(transport) + await assert.rejects( + client.callTool( + { name: 'slow-tool', arguments: {} }, + CallToolResultSchema, + { timeout: 20 }, + ), + ) + + const cancellationDeadline = Date.now() + 1_000 + while (!streamCancelled && Date.now() < cancellationDeadline) { + await new Promise(resolve => setTimeout(resolve, 5)) + } + assert.equal(streamCancelled, true) + assert.equal(activeRequests.size, 0) + } finally { + await client.close() + server.stop(true) + } +}) + test('cleanupFailedConnection awaits transport close before resolving', async () => { let closed = false let resolveClose: (() => void) | undefined diff --git a/src/services/mcp/client.ts b/src/services/mcp/client.ts index c60c850656..eb96d7d9e5 100644 --- a/src/services/mcp/client.ts +++ b/src/services/mcp/client.ts @@ -53,6 +53,7 @@ import { type MCPProgress, MCPTool } from '../../tools/MCPTool/MCPTool.js' import { createMcpAuthTool } from '../../tools/McpAuthTool/McpAuthTool.js' import { ReadMcpResourceTool } from '../../tools/ReadMcpResourceTool/ReadMcpResourceTool.js' import { createAbortController } from '../../utils/abortController.js' +import { createCombinedAbortSignal } from '../../utils/combinedAbortSignal.js' import { AbortError, isAbortError } from '../../utils/errors.js' import { count } from '../../utils/array.js' import { @@ -477,23 +478,72 @@ const MCP_REQUEST_TIMEOUT_MS = 60000 const MCP_STREAMABLE_HTTP_ACCEPT = 'application/json, text/event-stream' /** - * Wraps a fetch function to apply a fresh timeout signal to each request. - * This avoids the bug where a single AbortSignal.timeout() created at connection - * time becomes stale after 60 seconds, causing all subsequent requests to fail - * immediately with "The operation timed out." Uses a 60-second timeout. - * - * Also ensures the Accept header required by the MCP Streamable HTTP spec is - * present on POSTs. The MCP SDK sets this inside StreamableHTTPClientTransport.send(), - * but it is attached to a Headers instance that passes through an object spread here, - * and some runtimes/agents have been observed dropping it before it reaches the wire. - * See https://github.com/anthropics/claude-agent-sdk-typescript/issues/202. - * Normalizing here (the last wrapper before fetch()) guarantees it is sent. - * - * GET requests are excluded from the timeout since, for MCP transports, they are - * long-lived SSE streams meant to stay open indefinitely. (Auth-related GETs use - * a separate fetch wrapper with its own timeout in auth.ts.) + * Keep fetch timeout and cancellation resources alive until a streaming + * response is consumed or cancelled, not merely until response headers arrive. + */ +function wrapResponseBodyWithCleanup( + response: Response, + cleanup: () => void, + signal?: AbortSignal, +): Response { + if (!response.body) { + cleanup() + return response + } + + const reader = response.body.getReader() + let isCleanedUp = false + const finish = () => { + if (isCleanedUp) return + isCleanedUp = true + signal?.removeEventListener('abort', cancelReader) + cleanup() + } + const cancelReader = () => { + void reader.cancel(signal?.reason).catch(() => {}) + finish() + } + signal?.addEventListener('abort', cancelReader, { once: true }) + if (signal?.aborted) { + cancelReader() + } + const body = new ReadableStream({ + async pull(controller) { + try { + const { done, value } = await reader.read() + if (done) { + controller.close() + finish() + } else { + controller.enqueue(value) + } + } catch (error) { + controller.error(error) + finish() + } + }, + async cancel(reason) { + try { + await reader.cancel(reason) + } finally { + finish() + } + }, + }) + + return new Response(body, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }) +} + +/** + * Wrap a fetch function with a fresh request timeout and an MCP-compatible + * Accept header. The request timeout and parent abort listener stay active + * until a streaming response is consumed or cancelled. * - * @param baseFetch - The fetch function to wrap + * GET requests are excluded because MCP GET streams may stay open indefinitely. */ export function wrapFetchWithTimeout(baseFetch: FetchLike): FetchLike { return async (url: string | URL, init?: RequestInit) => { @@ -546,8 +596,8 @@ export function wrapFetchWithTimeout(baseFetch: FetchLike): FetchLike { headers, signal: controller.signal, }) - cleanup() - return response + clearTimeout(timer) + return wrapResponseBodyWithCleanup(response, cleanup, controller.signal) } catch (error) { cleanup() throw error @@ -555,6 +605,96 @@ export function wrapFetchWithTimeout(baseFetch: FetchLike): FetchLike { } } +type McpRequestId = string | number + +/** + * Associate Streamable HTTP POST response bodies with MCP request IDs so + * notifications/cancelled can interrupt a response stream after headers arrive. + */ +export function wrapFetchWithRequestCancellation( + baseFetch: FetchLike, + activeRequests: Map, +): FetchLike { + return async (url, init) => { + let requestId: McpRequestId | undefined + if ( + (init?.method ?? 'GET').toUpperCase() === 'POST' && + typeof init?.body === 'string' + ) { + try { + const request: unknown = JSON.parse(init.body) + if ( + request !== null && + typeof request === 'object' && + 'id' in request && + (typeof request.id === 'string' || typeof request.id === 'number') + ) { + requestId = request.id + } + } catch { + // Leave malformed/non-JSON requests to the underlying fetch unchanged. + } + } + if (requestId === undefined) { + return baseFetch(url, init) + } + + const requestController = createAbortController() + activeRequests.set(requestId, requestController) + const { signal, cleanup } = createCombinedAbortSignal( + init?.signal ?? undefined, + { + signalB: requestController.signal, + trace: { + subsystem: 'mcp_http_request', + controllerRole: 'mcp_http_request', + }, + }, + ) + const finish = () => { + cleanup() + if (activeRequests.get(requestId!) === requestController) { + activeRequests.delete(requestId!) + } + } + try { + const response = await baseFetch(url, { ...init, signal }) + return wrapResponseBodyWithCleanup(response, finish, signal) + } catch (error) { + finish() + throw error + } + } +} + +/** + * Abort the matching HTTP request body when the MCP SDK emits cancellation. + */ +export function attachMcpRequestCancellationHandler( + transport: Transport, + activeRequests: Map, +): void { + const send = transport.send.bind(transport) + transport.send = async (message, options) => { + if ( + 'method' in message && + message.method === 'notifications/cancelled' && + 'params' in message && + message.params !== undefined && + typeof message.params === 'object' && + 'requestId' in message.params + ) { + const requestId = message.params.requestId + if (typeof requestId === 'string' || typeof requestId === 'number') { + activeRequests.get(requestId)?.abort( + new DOMException('MCP request cancelled.', 'AbortError'), + ) + } + } + return send(message, options) + } +} + export function getMcpServerConnectionBatchSize(): number { return parseInt(process.env.MCP_SERVER_CONNECTION_BATCH_SIZE || '', 10) || 3 } @@ -909,17 +1049,21 @@ export const connectToServer = memoize( `Proxy options: ${proxyOptions.dispatcher ? 'custom dispatcher' : 'default'}`, ) + const activeHttpRequests = new Map() const transportOptions: StreamableHTTPClientTransportOptions = { authProvider, // Use fresh timeout per request to avoid stale AbortSignal bug. // Step-up detection wraps innermost so the 403 is seen before the // SDK's handler calls auth() → tokens(). - fetch: wrapFetchWithTimeout( - wrapFetchWithStepUpDetection(createFetchWithInit(), authProvider, { - allowUnauthorizedRefresh, - resourceUrl: serverRef.url, - providerOwnsAuthorization: allowUnauthorizedRefresh, - }), + fetch: wrapFetchWithRequestCancellation( + wrapFetchWithTimeout( + wrapFetchWithStepUpDetection(createFetchWithInit(), authProvider, { + allowUnauthorizedRefresh, + resourceUrl: serverRef.url, + providerOwnsAuthorization: allowUnauthorizedRefresh, + }), + ), + activeHttpRequests, ), requestInit: { ...proxyOptions, @@ -958,6 +1102,7 @@ export const connectToServer = memoize( new URL(serverRef.url), transportOptions, ) + attachMcpRequestCancellationHandler(transport, activeHttpRequests) logMCPDebug(name, `HTTP transport created successfully`) } else if (serverRef.type === 'sdk') { throw new Error('SDK servers should be handled in print.ts') @@ -981,9 +1126,13 @@ export const connectToServer = memoize( const fetchWithAuth = createClaudeAiProxyFetch(globalThis.fetch) const proxyOptions = getProxyFetchOptions() + const activeProxyRequests = new Map() const transportOptions: StreamableHTTPClientTransportOptions = { // Wrap fetchWithAuth with fresh timeout per request - fetch: wrapFetchWithTimeout(fetchWithAuth), + fetch: wrapFetchWithRequestCancellation( + wrapFetchWithTimeout(fetchWithAuth), + activeProxyRequests, + ), requestInit: { ...proxyOptions, headers: { @@ -997,6 +1146,7 @@ export const connectToServer = memoize( new URL(proxyUrl), transportOptions, ) + attachMcpRequestCancellationHandler(transport, activeProxyRequests) logMCPDebug(name, `claude.ai proxy transport created successfully`) } else if ( (serverRef.type === 'stdio' || !serverRef.type) && @@ -3093,7 +3243,7 @@ export async function callMCPToolWithUrlElicitationRetry({ // Check abort signal before each attempt — without this, a cancelled // elicitation retry loop continues spinning until MAX retries if (signal.aborted) { - throw new Error('Tool call aborted during URL elicitation') + throw new AbortError('Tool call aborted during URL elicitation') } try { return await callToolFn({ @@ -3312,27 +3462,33 @@ async function callMCPTool({ tool, ) - // Use Promise.race with our own timeout to handle cases where SDK's - // internal timeout doesn't work (e.g., SSE stream breaks mid-request) + // Use our own timeout to handle cases where SDK's internal timeout + // doesn't work (e.g., SSE stream breaks mid-request). Abort the protocol + // request too; merely racing a timeout leaves the server call running in + // the background and can keep its transport occupied after we return. const timeoutMs = getMcpToolTimeoutMs() + const timeoutController = createAbortController() + const { + signal: requestSignal, + cleanup: cleanupRequestSignal, + } = createCombinedAbortSignal(signal, { + signalB: timeoutController.signal, + trace: { subsystem: 'mcp_tool_call', controllerRole: 'mcp_tool_call' }, + }) let timeoutId: NodeJS.Timeout | undefined const timeoutPromise = new Promise((_, reject) => { - timeoutId = setTimeout( - (reject, name, tool, timeoutMs) => { - reject( - new TelemetrySafeError_I_VERIFIED_THIS_IS_NOT_CODE_OR_FILEPATHS( - `MCP server "${name}" tool "${tool}" timed out after ${Math.floor(timeoutMs / 1000)}s`, - 'MCP tool timeout', - ), - ) - }, - timeoutMs, - reject, - name, - tool, - timeoutMs, - ) + timeoutId = setTimeout(() => { + timeoutController.abort( + new DOMException('The operation timed out.', 'TimeoutError'), + ) + reject( + new TelemetrySafeError_I_VERIFIED_THIS_IS_NOT_CODE_OR_FILEPATHS( + `MCP server "${name}" tool "${tool}" timed out after ${Math.floor(timeoutMs / 1000)}s`, + 'MCP tool timeout', + ), + ) + }, timeoutMs) }) const result = await Promise.race([ @@ -3344,7 +3500,7 @@ async function callMCPTool({ }, CallToolResultSchema, { - signal, + signal: requestSignal, timeout: timeoutMs, onprogress: onProgress ? sdkProgress => { @@ -3366,6 +3522,7 @@ async function callMCPTool({ if (timeoutId) { clearTimeout(timeoutId) } + cleanupRequestSignal() }) if ('isError' in result && result.isError) {