diff --git a/.changeset/mcp-tool-call-reconnect.md b/.changeset/mcp-tool-call-reconnect.md new file mode 100644 index 00000000000..74ba19eeaaf --- /dev/null +++ b/.changeset/mcp-tool-call-reconnect.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Reconnect a dropped MCP server connection automatically when one of its tools is called, and retry the call once. diff --git a/packages/agent-core-v2/src/agent/mcp/client-http.ts b/packages/agent-core-v2/src/agent/mcp/client-http.ts index fbeab86e794..367a32834d4 100644 --- a/packages/agent-core-v2/src/agent/mcp/client-http.ts +++ b/packages/agent-core-v2/src/agent/mcp/client-http.ts @@ -7,6 +7,7 @@ import { buildRequestOptions, KIMI_MCP_CLIENT_NAME, KIMI_MCP_CLIENT_VERSION, + MCP_LIVENESS_PROBE_TIMEOUT_MS, toMcpToolDefinition, toMcpToolResult, type UnexpectedCloseListener, @@ -103,6 +104,10 @@ export class HttpMcpClient implements MCPClient { return toMcpToolResult(result); } + async ping(signal?: AbortSignal): Promise { + await this.client.ping(buildRequestOptions(MCP_LIVENESS_PROBE_TIMEOUT_MS, signal)); + } + private async closeStartedClient(): Promise { if (!this.started) return; this.started = false; diff --git a/packages/agent-core-v2/src/agent/mcp/client-shared.ts b/packages/agent-core-v2/src/agent/mcp/client-shared.ts index f878aae9b60..cd51273e96c 100644 --- a/packages/agent-core-v2/src/agent/mcp/client-shared.ts +++ b/packages/agent-core-v2/src/agent/mcp/client-shared.ts @@ -1,6 +1,7 @@ import { getCoreVersion } from '#/_base/version'; +import { ErrorCode, McpError } from '@modelcontextprotocol/sdk/types.js'; -import type { MCPToolDefinition, MCPToolResult } from './types'; +import type { MCPClient, MCPToolDefinition, MCPToolResult } from './types'; export const KIMI_MCP_CLIENT_NAME = 'kimi-code'; export const KIMI_MCP_CLIENT_VERSION = getCoreVersion(); @@ -12,6 +13,63 @@ export interface UnexpectedCloseReason { export type UnexpectedCloseListener = (reason: UnexpectedCloseReason) => void; +export function isMcpConnectionClosedError(error: unknown): boolean { + return ( + error instanceof Error && + (error as Error & { readonly code?: unknown }).code === ErrorCode.ConnectionClosed + ); +} + +export function isMcpTransportFailure(error: unknown): boolean { + if (!(error instanceof Error)) return false; + if (isMcpConnectionClosedError(error)) return true; + return !(error instanceof McpError); +} + +/** + * Timeout for the liveness probe sent after an ambiguous tool-call failure. + * Kept short: the probe runs on an already-failed call, so it must not add + * anywhere near a tool-call timeout to the turn. + */ +export const MCP_LIVENESS_PROBE_TIMEOUT_MS = 5_000; + +/** + * True when the error is a client-side validation failure of an otherwise + * well-formed JSON-RPC response: the SDK rejects with a `ZodError` when the + * result of `tools/call` does not match `CallToolResultSchema` + * (shared/protocol.js rejects with `parseResult.error`). The server did + * answer, so reconnecting is pointless — but the error is not an `McpError`, + * so `isMcpTransportFailure` alone cannot tell it apart from a dead + * transport. Matched by name because the repo carries more than one zod + * copy, which makes `instanceof` unreliable. + */ +export function isMcpMalformedResultError(error: unknown): boolean { + return error instanceof Error && error.name === 'ZodError'; +} + +/** + * Probes whether the client's transport is still usable by sending a ping. + * A server that answers in any way — including `MethodNotFound`, a JSON-RPC + * error, or an unparseable result — counts as alive; only errors that prove + * the bytes never made a round trip (closed connection, fetch failures) or + * a probe that itself timed out (alive socket, unresponsive server) count + * as dead. Never rejects; an abort surfaces as a dead verdict and is the + * caller's job to detect via the signal. + */ +export async function probeMcpLiveness(client: MCPClient, signal: AbortSignal): Promise { + try { + await client.ping(signal); + return true; + } catch (error) { + if (isMcpConnectionClosedError(error)) return false; + if (isMcpMalformedResultError(error)) return true; + if (error instanceof McpError) { + return (error as Error & { readonly code?: unknown }).code !== ErrorCode.RequestTimeout; + } + return false; + } +} + export interface McpRequestOptions { readonly timeout?: number; readonly signal?: AbortSignal; diff --git a/packages/agent-core-v2/src/agent/mcp/client-sse.ts b/packages/agent-core-v2/src/agent/mcp/client-sse.ts index ee758e91fb8..e1ba7daaaf3 100644 --- a/packages/agent-core-v2/src/agent/mcp/client-sse.ts +++ b/packages/agent-core-v2/src/agent/mcp/client-sse.ts @@ -7,6 +7,7 @@ import { buildRequestOptions, KIMI_MCP_CLIENT_NAME, KIMI_MCP_CLIENT_VERSION, + MCP_LIVENESS_PROBE_TIMEOUT_MS, toMcpToolDefinition, toMcpToolResult, type UnexpectedCloseListener, @@ -103,6 +104,10 @@ export class SseMcpClient implements MCPClient { return toMcpToolResult(result); } + async ping(signal?: AbortSignal): Promise { + await this.client.ping(buildRequestOptions(MCP_LIVENESS_PROBE_TIMEOUT_MS, signal)); + } + private async closeStartedClient(): Promise { if (!this.started) return; this.started = false; diff --git a/packages/agent-core-v2/src/agent/mcp/client-stdio.ts b/packages/agent-core-v2/src/agent/mcp/client-stdio.ts index e5acf2c5242..bfe726ee57f 100644 --- a/packages/agent-core-v2/src/agent/mcp/client-stdio.ts +++ b/packages/agent-core-v2/src/agent/mcp/client-stdio.ts @@ -9,6 +9,7 @@ import { buildRequestOptions, KIMI_MCP_CLIENT_NAME, KIMI_MCP_CLIENT_VERSION, + MCP_LIVENESS_PROBE_TIMEOUT_MS, toMcpToolDefinition, toMcpToolResult, type UnexpectedCloseListener, @@ -115,6 +116,10 @@ export class StdioMcpClient implements MCPClient { return toMcpToolResult(result); } + async ping(signal?: AbortSignal): Promise { + await this.client.ping(buildRequestOptions(MCP_LIVENESS_PROBE_TIMEOUT_MS, signal)); + } + private async closeStartedClient(): Promise { if (!this.started) return; this.started = false; diff --git a/packages/agent-core-v2/src/agent/mcp/connection-manager.ts b/packages/agent-core-v2/src/agent/mcp/connection-manager.ts index 6ad83a386ee..4b8f70b4734 100644 --- a/packages/agent-core-v2/src/agent/mcp/connection-manager.ts +++ b/packages/agent-core-v2/src/agent/mcp/connection-manager.ts @@ -68,6 +68,7 @@ export interface McpConnectionManagerOptions { export class McpConnectionManager { private readonly entries = new Map(); private readonly listeners = new Set(); + private readonly inFlightReconnects = new Map>(); private initialLoad: Promise = Promise.resolve(); private initialLoadAttemptId = 0; private initialLoadStartedAt: number | undefined; @@ -232,6 +233,18 @@ export class McpConnectionManager { await this.connectOne(entry, attemptId); } + reconnectAndJoin(name: string): Promise { + const existing = this.inFlightReconnects.get(name); + if (existing !== undefined) return existing; + const work = this.reconnect(name).finally(() => { + if (this.inFlightReconnects.get(name) === work) { + this.inFlightReconnects.delete(name); + } + }); + this.inFlightReconnects.set(name, work); + return work; + } + async shutdown(): Promise { const entries = Array.from(this.entries.values()); this.entries.clear(); diff --git a/packages/agent-core-v2/src/agent/mcp/mcpService.ts b/packages/agent-core-v2/src/agent/mcp/mcpService.ts index b4d13c00d7f..22ed4826e55 100644 --- a/packages/agent-core-v2/src/agent/mcp/mcpService.ts +++ b/packages/agent-core-v2/src/agent/mcp/mcpService.ts @@ -7,6 +7,7 @@ import type { Tool as KosongTool } from '#/app/llmProtocol/tool'; import { Disposable, type IDisposable } from "#/_base/di/lifecycle"; import type { KimiErrorPayload } from '#/_base/errors/serialize'; import { ErrorCodes, makeErrorPayload } from "#/errors"; +import { abortable } from '#/_base/utils/abort'; import { IEventBus } from '#/app/event/eventBus'; import { ITelemetryService } from '#/app/telemetry/telemetry'; import { sessionMediaOriginalsDir } from '#/agent/media/image-originals'; @@ -130,6 +131,26 @@ export class AgentMcpService extends Disposable implements IAgentMcpService { signal?.throwIfAborted(); } + private reconnectForToolCall( + serverName: string, + staleClient: MCPClient, + signal?: AbortSignal, + ): Promise { + const work = this.joinHealedOrReconnect(serverName, staleClient); + return signal === undefined ? work : abortable(work, signal); + } + + private async joinHealedOrReconnect( + serverName: string, + staleClient: MCPClient, + ): Promise { + const healed = this.resolved(serverName)?.client; + if (healed !== undefined && healed !== staleClient) return healed; + await this.sessionMcp.connectionManager().reconnectAndJoin(serverName); + const current = this.resolved(serverName)?.client; + return current !== undefined && current !== staleClient ? current : undefined; + } + onStatusChange(listener: Parameters[0]) { const unsubscribe = this.sessionMcp.connectionManager().onStatusChange(listener); return { @@ -167,16 +188,15 @@ export class AgentMcpService extends Disposable implements IAgentMcpService { this.registerNeedsAuthMcpServer(entry); return; } - if (entry.status === 'failed') { - this.unregisterMcpServer(entry.name); - this.eventBus.publish({ - type: 'tool.list.updated', - reason: 'mcp.failed', - serverName: entry.name, - }); + if (entry.status === 'failed' || entry.status === 'pending') { + // Keep the server's tools registered while it is down or reconnecting. + // The captured client is closed, so the next call fails fast at the + // transport layer and the tool adapter's reconnect-and-retry path heals + // the connection — a dropped server surfaces as a slow call instead of + // "tool not found" for the rest of the session. return; } - if (entry.status === 'disabled' || entry.status === 'pending') { + if (entry.status === 'disabled') { const removed = this.unregisterMcpServer(entry.name); if (removed) { this.eventBus.publish({ @@ -267,6 +287,7 @@ export class AgentMcpService extends Disposable implements IAgentMcpService { createMcpTool(qualified, tool, client, { originalsDir: sessionMediaOriginalsDir(this.sessionContext.sessionDir), telemetry: this.telemetry, + reconnect: (signal) => this.reconnectForToolCall(serverName, client, signal), }), { source: 'mcp' }, ), diff --git a/packages/agent-core-v2/src/agent/mcp/tools/mcp.ts b/packages/agent-core-v2/src/agent/mcp/tools/mcp.ts index 4184d03f6bf..ac2fc97ca32 100644 --- a/packages/agent-core-v2/src/agent/mcp/tools/mcp.ts +++ b/packages/agent-core-v2/src/agent/mcp/tools/mcp.ts @@ -3,19 +3,45 @@ * * Each tool exposed by a connected MCP server is adapted into an * `ExecutableTool` whose `resolveExecution` forwards the call to the client - * and normalizes the result. + * and normalizes the result. When a call fails, the adapter picks one of + * three recoveries based on why it failed: + * + * - The server answered (a JSON-RPC error, or a response that failed + * client-side schema validation) → the error is rethrown; reconnecting + * would not change the answer. + * - The failure is ambiguous (a raw fetch/socket error) → the client is + * probed with a ping: alive means a transient blip and the call is + * retried once in place; dead means the transport is gone. + * - The transport is provably dead (the SDK fired `onclose`, or the probe + * failed) → the server is reconnected once through `options.reconnect` + * and the call retried on the fresh client, so a dropped connection + * surfaces as a slow call instead of a failed turn. + * + * Retries are at-least-once: if the transport died after the server + * processed the call but before the response arrived, the retry may + * duplicate side effects. There is no protocol-level dedup across + * reconnects, so this trade-off is accepted deliberately. */ import type { Tool as KosongTool } from '#/app/llmProtocol/tool'; import type { ITelemetryService } from '#/app/telemetry/telemetry'; +import { toErrorMessage } from '#/errors'; +import { isAbortError } from '#/_base/utils/abort'; -import type { ExecutableTool, ExecutableToolResult } from '#/tool/toolContract'; +import type { ExecutableTool, ExecutableToolContext, ExecutableToolResult } from '#/tool/toolContract'; import { mcpResultToExecutableOutput } from '#/agent/mcp/output'; -import type { MCPClient } from '#/agent/mcp/types'; +import type { MCPClient, MCPToolResult } from '#/agent/mcp/types'; +import { + isMcpConnectionClosedError, + isMcpMalformedResultError, + isMcpTransportFailure, + probeMcpLiveness, +} from '#/agent/mcp/client-shared'; interface McpToolOptions { readonly originalsDir?: string; readonly telemetry?: ITelemetryService; + readonly reconnect?: (signal?: AbortSignal) => Promise; } export function createMcpTool( @@ -24,6 +50,8 @@ export function createMcpTool( client: MCPClient, options: McpToolOptions = {}, ): ExecutableTool { + const callTool = (activeClient: MCPClient, args: unknown, signal: AbortSignal) => + activeClient.callTool(tool.name, (args ?? {}) as Record, signal); return { name: qualifiedName, description: tool.description, @@ -31,11 +59,12 @@ export function createMcpTool( resolveExecution: (args) => ({ approvalRule: qualifiedName, execute: async (context) => { - const result = await client.callTool( - tool.name, - (args ?? {}) as Record, - context.signal, - ); + let result; + try { + result = await callTool(client, args, context.signal); + } catch (error) { + result = await retryAfterReconnect(error, client, args, context, options, callTool); + } return normalizeMcpToolResult( await mcpResultToExecutableOutput(result, qualifiedName, { originalsDir: options.originalsDir, @@ -47,6 +76,70 @@ export function createMcpTool( }; } +async function retryAfterReconnect( + error: unknown, + client: MCPClient, + args: unknown, + context: Pick, + options: McpToolOptions, + callTool: (client: MCPClient, args: unknown, signal: AbortSignal) => Promise, +): Promise { + const reconnect = options.reconnect; + // Errors that can never be fixed by a retry: user cancellation, and the + // server having answered — a JSON-RPC error (`McpError`, including a tool + // call timeout) or a malformed result that failed schema validation. + const isUnrecoverable = (e: unknown): boolean => + context.signal.aborted || + isAbortError(e) || + !isMcpTransportFailure(e) || + isMcpMalformedResultError(e); + if (reconnect === undefined || isUnrecoverable(error)) { + throw error; + } + + // A ConnectionClosed error is a measured death (the SDK already fired + // `onclose` and rejected every pending request), so it goes straight to + // reconnect. Anything else is ambiguous about whether the transport + // still works — probe it instead of guessing from the error's type. + let failure = error; + if (!isMcpConnectionClosedError(failure)) { + const alive = await probeMcpLiveness(client, context.signal); + context.signal.throwIfAborted(); + if (alive) { + // The transport is fine and the failure was transient: retry once in + // place instead of paying a full reconnect for a network blip. If the + // transport dies between probe and retry, fall through to reconnect — + // still capped at one reconnect per call. + try { + return await callTool(client, args, context.signal); + } catch (retryError) { + if (isUnrecoverable(retryError)) { + throw retryError; + } + failure = retryError; + } + } + } + + context.onUpdate?.({ kind: 'status', text: 'MCP connection lost — reconnecting…' }); + let freshClient: MCPClient | undefined; + try { + freshClient = await reconnect(context.signal); + } catch (reconnectError) { + if (context.signal.aborted || isAbortError(reconnectError)) { + throw reconnectError; + } + throw new Error( + `${toErrorMessage(failure)} (reconnecting the MCP server also failed: ${toErrorMessage(reconnectError)})`, + { cause: reconnectError }, + ); + } + if (freshClient === undefined) { + throw failure; + } + return callTool(freshClient, args, context.signal); +} + function normalizeMcpToolResult(result: { readonly output: ExecutableToolResult['output']; readonly isError: boolean; diff --git a/packages/agent-core-v2/src/agent/mcp/types.ts b/packages/agent-core-v2/src/agent/mcp/types.ts index 33e2b46887a..ff783c6b6a7 100644 --- a/packages/agent-core-v2/src/agent/mcp/types.ts +++ b/packages/agent-core-v2/src/agent/mcp/types.ts @@ -49,6 +49,13 @@ export interface MCPClient { args: Record, signal?: AbortSignal, ): Promise; + /** + * Sends a protocol-level `ping` with a short built-in timeout, so a hung + * server rejects instead of blocking. Used to probe liveness after an + * ambiguous call failure; a server that answers in any way — even with + * `MethodNotFound` — proves the transport is usable. + */ + ping(signal?: AbortSignal): Promise; } export function assertMcpInputSchema( diff --git a/packages/agent-core-v2/test/agent/mcp/connection-manager.test.ts b/packages/agent-core-v2/test/agent/mcp/connection-manager.test.ts index d37bec9d97c..ec317c429ee 100644 --- a/packages/agent-core-v2/test/agent/mcp/connection-manager.test.ts +++ b/packages/agent-core-v2/test/agent/mcp/connection-manager.test.ts @@ -254,6 +254,51 @@ describe('McpConnectionManager', () => { } }); + it('reconnectAndJoin joins an in-flight reconnect instead of starting a second one', async () => { + const cm = new McpConnectionManager(); + const seen: Array<{ name: string; status: McpServerEntry['status'] }> = []; + cm.onStatusChange((entry) => { + seen.push({ name: entry.name, status: entry.status }); + }); + const delayedMockServer = `setTimeout(() => import(${JSON.stringify( + pathToFileURL(stdioFixture).href, + )}), 250)`; + + try { + await cm.connectAll({ + slow: { + transport: 'stdio', + command: process.execPath, + args: ['-e', delayedMockServer], + startupTimeoutMs: 5_000, + }, + }); + seen.length = 0; + + await Promise.all([cm.reconnectAndJoin('slow'), cm.reconnectAndJoin('slow')]); + + expect(cm.get('slow')?.status).toBe('connected'); + expect(seen.filter((event) => event.name === 'slow').map((event) => event.status)).toEqual([ + 'pending', + 'connected', + ]); + } finally { + await cm.shutdown(); + } + }, 20000); + + it('reconnectAndJoin rejects for unknown servers', async () => { + const cm = new McpConnectionManager(); + try { + await expect(cm.reconnectAndJoin('nope')).rejects.toBeInstanceOf(Error2); + await expect(cm.reconnectAndJoin('nope')).rejects.toMatchObject({ + code: 'mcp.server_not_found', + }); + } finally { + await cm.shutdown(); + } + }); + it('shutdown clears entries and is idempotent', async () => { const cm = new McpConnectionManager(); await cm.connectAll({ alpha: stdioConfig() }); diff --git a/packages/agent-core-v2/test/agent/mcp/mcp.test.ts b/packages/agent-core-v2/test/agent/mcp/mcp.test.ts index 3f62a94cd6f..44b0211df7d 100644 --- a/packages/agent-core-v2/test/agent/mcp/mcp.test.ts +++ b/packages/agent-core-v2/test/agent/mcp/mcp.test.ts @@ -1,12 +1,14 @@ import type { ContentPart } from '#/app/llmProtocol/message'; import type { Tool as KosongTool } from '#/app/llmProtocol/tool'; import { Jimp } from 'jimp'; +import { CallToolResultSchema, ErrorCode, McpError } from '@modelcontextprotocol/sdk/types.js'; import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { SyncDescriptor } from '#/_base/di/descriptors'; import { DisposableStore, toDisposable } from '#/_base/di/lifecycle'; import { TestInstantiationService } from '#/_base/di/test'; import { Event } from '#/_base/event'; +import { abortError } from '#/_base/utils/abort'; import { type DomainEvent, IEventBus } from '#/app/event/eventBus'; import { ITelemetryService } from '#/app/telemetry/telemetry'; import type { McpConnectionManager, McpServerEntry } from '#/agent/mcp/connection-manager'; @@ -61,6 +63,7 @@ class FakeMcpManager { } resolved(name: string): ResolvedServer | undefined { + if (this.entries.get(name)?.status !== 'connected') return undefined; return this.resolvedEntries.get(name); } @@ -68,7 +71,25 @@ class FakeMcpManager { return name === 'needs-auth' ? 'https://example.com/mcp' : undefined; } - async reconnect(): Promise {} + reconnectHandler: (name: string) => Promise = async () => {}; + + async reconnect(name: string): Promise { + await this.reconnectHandler(name); + } + + private readonly inFlightReconnects = new Map>(); + + reconnectAndJoin(name: string): Promise { + const existing = this.inFlightReconnects.get(name); + if (existing !== undefined) return existing; + const work = this.reconnect(name).finally(() => { + if (this.inFlightReconnects.get(name) === work) { + this.inFlightReconnects.delete(name); + } + }); + this.inFlightReconnects.set(name, work); + return work; + } async waitForInitialLoad(): Promise {} @@ -136,6 +157,14 @@ class FakeMcpManager { this.emit(entry); } + pending(name: string): void { + const current = this.entries.get(name); + if (current === undefined) return; + const entry: McpServerEntry = { ...current, status: 'pending', toolCount: 0 }; + this.entries.set(name, entry); + this.emit(entry); + } + disconnect(name: string): void { const current = this.entries.get(name); if (current === undefined) return; @@ -360,6 +389,505 @@ describe('AgentMcpService', () => { expect(result.output).toBe('hello world'); }); + function throwingClient( + base: MCPClient = fakeMcpClient(), + onCall?: () => void, + makeError: () => Error = () => new McpError(ErrorCode.ConnectionClosed, 'Connection closed'), + ): MCPClient { + return { + listTools: () => base.listTools(), + async callTool() { + onCall?.(); + throw makeError(); + }, + async ping() { + throw makeError(); + }, + }; + } + + function countingClient(base: MCPClient, counter: { calls: number }): MCPClient { + return { + listTools: () => base.listTools(), + callTool: (name, args, signal) => { + counter.calls += 1; + return base.callTool(name, args, signal); + }, + ping: (signal) => base.ping(signal), + }; + } + + function deferred(): { + readonly promise: Promise; + readonly resolve: (value: T | PromiseLike) => void; + } { + let resolvePromise!: (value: T | PromiseLike) => void; + const promise = new Promise((resolve) => { + resolvePromise = resolve; + }); + return { promise, resolve: resolvePromise }; + } + + it('reconnects the server and retries the call once when the transport dies', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient(fakeMcpClient(), () => manager.fail('s')); + const freshCounter = { calls: 0 }; + const freshClient = countingClient(fakeMcpClient(), freshCounter); + let reconnects = 0; + manager.reconnectHandler = async (name) => { + reconnects += 1; + manager.setResolved(name, freshClient, await discoverTools(freshClient)); + manager.connect(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + const result = await executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-reconnect', + args: { text: 'hello again' }, + signal: new AbortController().signal, + }); + + expect(result.isError).toBeUndefined(); + expect(result.output).toBe('hello again'); + expect(freshCounter.calls).toBe(1); + expect(reconnects).toBe(1); + }); + + it('heals a server that died between turns when its tool is called again', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient(fakeMcpClient()); + const freshCounter = { calls: 0 }; + const freshClient = countingClient(fakeMcpClient(), freshCounter); + let reconnects = 0; + manager.reconnectHandler = async (name) => { + reconnects += 1; + manager.setResolved(name, freshClient, await discoverTools(freshClient)); + manager.connect(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + // The connection drops while no call is in flight: the manager marks the + // server failed. The tools must stay registered so the next call reaches + // the adapter and its reconnect-and-retry path instead of failing with + // "tool not found". + manager.fail('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + expect(echo).toBeDefined(); + const result = await executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-between-turns', + args: { text: 'back from the dead' }, + signal: new AbortController().signal, + }); + + expect(result.isError).toBeUndefined(); + expect(result.output).toBe('back from the dead'); + expect(freshCounter.calls).toBe(1); + expect(reconnects).toBe(1); + }); + + it('returns a non-transport MCP error without reconnecting the server', async () => { + const manager = new FakeMcpManager(); + const base = fakeMcpClient(); + const client: MCPClient = { + listTools: () => base.listTools(), + async callTool() { + throw new McpError(ErrorCode.InvalidParams, 'Invalid tool arguments'); + }, + ping: () => base.ping(), + }; + let reconnects = 0; + manager.reconnectHandler = async () => { + reconnects += 1; + }; + manager.setResolved('s', client, await discoverTools(client)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + await expect( + executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-non-transport-error', + args: { text: 'hi' }, + signal: new AbortController().signal, + }), + ).rejects.toThrow('Invalid tool arguments'); + expect(reconnects).toBe(0); + }); + + it('rethrows the original error when the server does not come back', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient(fakeMcpClient(), () => manager.fail('s')); + manager.reconnectHandler = async (name) => { + manager.fail(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + await expect( + executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-still-dead', + args: { text: 'hi' }, + signal: new AbortController().signal, + }), + ).rejects.toThrow('Connection closed'); + // The tools stay registered after the failed reconnect so a later call + // can try healing the server again instead of hitting "tool not found". + expect(ix.get(IAgentToolRegistryService).list().filter((tool) => tool.source === 'mcp')).toHaveLength(2); + }); + + it('reports both errors when the reconnect attempt itself fails', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient(fakeMcpClient(), () => manager.fail('s')); + manager.reconnectHandler = async () => { + throw new Error('spawn failed'); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + await expect( + executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-reconnect-fails', + args: { text: 'hi' }, + signal: new AbortController().signal, + }), + ).rejects.toThrow(/Connection closed .*spawn failed/); + }); + + it('does not reconnect when the call was aborted', async () => { + const manager = new FakeMcpManager(); + const base = fakeMcpClient(); + const abortingClient: MCPClient = { + listTools: () => base.listTools(), + async callTool() { + throw abortError('This operation was aborted'); + }, + ping: () => base.ping(), + }; + let reconnects = 0; + manager.reconnectHandler = async () => { + reconnects += 1; + }; + manager.setResolved('s', abortingClient, await discoverTools(abortingClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + await expect( + executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-aborted', + args: { text: 'hi' }, + signal: new AbortController().signal, + }), + ).rejects.toThrow('This operation was aborted'); + expect(reconnects).toBe(0); + }); + + it('dedupes concurrent reconnects from parallel failing tool calls', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient(fakeMcpClient(), () => manager.fail('s')); + const freshClient = fakeMcpClient(); + let reconnects = 0; + manager.reconnectHandler = async (name) => { + reconnects += 1; + manager.setResolved(name, freshClient, await discoverTools(freshClient)); + manager.connect(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const registry = ix.get(IAgentToolRegistryService); + const echo = registry.resolve('mcp__s__echo'); + const noop = registry.resolve('mcp__s__noop'); + const [echoResult, noopResult] = await Promise.all([ + executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-par-1', + args: { text: 'one' }, + signal: new AbortController().signal, + }), + executeTool(noop!, { + turnId: 1, + toolCallId: 'tc-par-2', + args: {}, + signal: new AbortController().signal, + }), + ]); + + expect(echoResult.output).toBe('one'); + expect(noopResult.output).toBe('ok'); + expect(reconnects).toBe(1); + }); + + it('keeps the shared reconnect alive when one parallel call is aborted', async () => { + const manager = new FakeMcpManager(); + const reconnectStarted = deferred(); + const reconnectReleased = deferred(); + const deadClient = throwingClient(fakeMcpClient(), () => manager.fail('s')); + const freshClient = fakeMcpClient(); + let reconnects = 0; + manager.reconnectHandler = async (name) => { + reconnects += 1; + reconnectStarted.resolve(); + await reconnectReleased.promise; + manager.setResolved(name, freshClient, await discoverTools(freshClient)); + manager.connect(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const registry = ix.get(IAgentToolRegistryService); + const echo = registry.resolve('mcp__s__echo'); + const noop = registry.resolve('mcp__s__noop'); + const firstController = new AbortController(); + const firstCall = executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-par-abort-1', + args: { text: 'one' }, + signal: firstController.signal, + }); + const secondCall = executeTool(noop!, { + turnId: 1, + toolCallId: 'tc-par-abort-2', + args: {}, + signal: new AbortController().signal, + }); + + await reconnectStarted.promise; + firstController.abort(new Error('cancelled by test')); + await expect(firstCall).rejects.toThrow('cancelled by test'); + + reconnectReleased.resolve(); + await expect(secondCall).resolves.toMatchObject({ output: 'ok' }); + expect(reconnects).toBe(1); + }); + + it('reconnects and retries when the call fails with a raw transport error the manager did not observe', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient( + fakeMcpClient(), + undefined, + () => new TypeError('fetch failed'), + ); + const freshClient = fakeMcpClient(); + let reconnects = 0; + manager.reconnectHandler = async (name) => { + reconnects += 1; + manager.setResolved(name, freshClient, await discoverTools(freshClient)); + manager.connect(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + const result = await executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-raw-transport', + args: { text: 'hello again' }, + signal: new AbortController().signal, + }); + + expect(result.isError).toBeUndefined(); + expect(result.output).toBe('hello again'); + expect(reconnects).toBe(1); + }); + + it('retries on the healed client without reconnecting again when the server already came back', async () => { + const manager = new FakeMcpManager(); + const deadClient = throwingClient(fakeMcpClient(), undefined, () => new Error('Not connected')); + const freshClient = fakeMcpClient(); + let reconnects = 0; + manager.reconnectHandler = async () => { + reconnects += 1; + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const registry = ix.get(IAgentToolRegistryService); + const staleEcho = registry.resolve('mcp__s__echo'); + + // Resolve the stale tool first, then heal the server the way a parallel + // call's reconnect would: the resolved entry swaps to a fresh client and + // the registry re-seeds, leaving `staleEcho` bound to the dead client. + manager.setResolved('s', freshClient, await discoverTools(freshClient)); + manager.connect('s'); + + const result = await executeTool(staleEcho!, { + turnId: 1, + toolCallId: 'tc-healed', + args: { text: 'late call' }, + signal: new AbortController().signal, + }); + + expect(result.isError).toBeUndefined(); + expect(result.output).toBe('late call'); + expect(reconnects).toBe(0); + }); + + it('rethrows a malformed tool result without reconnecting or retrying when the server answered', async () => { + const manager = new FakeMcpManager(); + const base = fakeMcpClient(); + const malformed = CallToolResultSchema.safeParse({ content: [{ text: 'missing type' }] }); + if (malformed.success) throw new Error('expected the fixture result to fail validation'); + let calls = 0; + const client: MCPClient = { + listTools: () => base.listTools(), + ping: () => base.ping(), + async callTool() { + calls += 1; + throw malformed.error; + }, + }; + let reconnects = 0; + manager.reconnectHandler = async () => { + reconnects += 1; + }; + manager.setResolved('s', client, await discoverTools(client)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + await expect( + executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-malformed-result', + args: { text: 'hi' }, + signal: new AbortController().signal, + }), + ).rejects.toBe(malformed.error); + expect(calls).toBe(1); + expect(reconnects).toBe(0); + }); + + it('retries a transient transport failure in place without reconnecting', async () => { + const manager = new FakeMcpManager(); + const base = fakeMcpClient(); + let calls = 0; + const flakyClient: MCPClient = { + listTools: () => base.listTools(), + ping: () => base.ping(), + callTool: (name, args, signal) => { + calls += 1; + if (calls === 1) return Promise.reject(new TypeError('fetch failed')); + return base.callTool(name, args, signal); + }, + }; + let reconnects = 0; + manager.reconnectHandler = async () => { + reconnects += 1; + }; + manager.setResolved('s', flakyClient, await discoverTools(flakyClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + const result = await executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-transient', + args: { text: 'hello again' }, + signal: new AbortController().signal, + }); + + expect(result.isError).toBeUndefined(); + expect(result.output).toBe('hello again'); + expect(calls).toBe(2); + expect(reconnects).toBe(0); + }); + + it('reconnects when the transport failure persists past a successful probe', async () => { + const manager = new FakeMcpManager(); + const base = fakeMcpClient(); + let calls = 0; + const deadClient: MCPClient = { + listTools: () => base.listTools(), + ping: () => base.ping(), + async callTool() { + calls += 1; + throw new TypeError('fetch failed'); + }, + }; + const freshClient = fakeMcpClient(); + let reconnects = 0; + manager.reconnectHandler = async (name) => { + reconnects += 1; + manager.setResolved(name, freshClient, await discoverTools(freshClient)); + manager.connect(name); + }; + manager.setResolved('s', deadClient, await discoverTools(deadClient)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + const result = await executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-persistent-transport', + args: { text: 'hello again' }, + signal: new AbortController().signal, + }); + + expect(result.isError).toBeUndefined(); + expect(result.output).toBe('hello again'); + expect(calls).toBe(2); + expect(reconnects).toBe(1); + }); + + it('abandons the retry when the call is aborted during the liveness probe', async () => { + const manager = new FakeMcpManager(); + const base = fakeMcpClient(); + const probeStarted = deferred(); + const releaseProbe = deferred(); + const client: MCPClient = { + listTools: () => base.listTools(), + async ping() { + probeStarted.resolve(); + await releaseProbe.promise; + }, + async callTool() { + throw new TypeError('fetch failed'); + }, + }; + let reconnects = 0; + manager.reconnectHandler = async () => { + reconnects += 1; + }; + manager.setResolved('s', client, await discoverTools(client)); + createService(manager); + manager.connect('s'); + + const echo = ix.get(IAgentToolRegistryService).resolve('mcp__s__echo'); + const controller = new AbortController(); + const call = executeTool(echo!, { + turnId: 1, + toolCallId: 'tc-abort-during-probe', + args: { text: 'hi' }, + signal: controller.signal, + }); + await probeStarted.promise; + controller.abort(new Error('cancelled by test')); + releaseProbe.resolve(); + await expect(call).rejects.toThrow('cancelled by test'); + expect(reconnects).toBe(0); + }); + it('truncates oversized MCP text output through the wrapped tool path', async () => { const manager = new FakeMcpManager(); const client: MCPClient = { @@ -378,6 +906,7 @@ describe('AgentMcpService', () => { isError: false, }; }, + async ping() {}, }; manager.setResolved('s', client, await discoverTools(client)); createService(manager); @@ -413,6 +942,7 @@ describe('AgentMcpService', () => { isError: false, }; }, + async ping() {}, }; manager.setResolved('s', client, await discoverTools(client)); createService(manager); @@ -459,6 +989,7 @@ describe('AgentMcpService', () => { isError: false, }; }, + async ping() {}, }; manager.setResolved('s', client, await discoverTools(client)); createService(manager); @@ -506,6 +1037,7 @@ describe('AgentMcpService', () => { receivedSignal = signal; return { content: [{ type: 'text', text: String(args['text']) }], isError: false }; }, + async ping() {}, }; manager.setResolved('s', client, await discoverTools(client)); createService(manager); @@ -545,7 +1077,7 @@ describe('AgentMcpService', () => { ]); }); - it('emits tool.list.updated(mcp.failed) when a connected server fails', async () => { + it('keeps tools registered when a connected server fails so later calls can heal', async () => { const manager = new FakeMcpManager(); const client = fakeMcpClient(); manager.setResolved('s', client, await discoverTools(client)); @@ -554,16 +1086,33 @@ describe('AgentMcpService', () => { manager.connect('s'); manager.fail('s'); - expect(ix.get(IAgentToolRegistryService).list().filter((tool) => tool.source === 'mcp')).toEqual([]); + expect(ix.get(IAgentToolRegistryService).list().filter((tool) => tool.source === 'mcp')).toHaveLength(2); + expect(events).not.toContainEqual( + expect.objectContaining({ type: 'tool.list.updated', reason: 'mcp.failed' }), + ); expect(events).toContainEqual( expect.objectContaining({ - type: 'tool.list.updated', - reason: 'mcp.failed', - serverName: 's', + type: 'mcp.server.status', + server: expect.objectContaining({ name: 's', status: 'failed' }), }), ); }); + it('keeps tools registered while the server is reconnecting', async () => { + const manager = new FakeMcpManager(); + const client = fakeMcpClient(); + manager.setResolved('s', client, await discoverTools(client)); + createService(manager); + + manager.connect('s'); + manager.pending('s'); + + expect(ix.get(IAgentToolRegistryService).list().filter((tool) => tool.source === 'mcp')).toHaveLength(2); + expect(events).not.toContainEqual( + expect.objectContaining({ type: 'tool.list.updated', reason: 'mcp.disconnected' }), + ); + }); + const RAW_QUERY: MCPToolDefinition = { name: 'query_range', description: 'Query a metrics range', diff --git a/packages/agent-core-v2/test/agent/mcp/output.test.ts b/packages/agent-core-v2/test/agent/mcp/output.test.ts index cf5ca3fb75d..bc226bc39ed 100644 --- a/packages/agent-core-v2/test/agent/mcp/output.test.ts +++ b/packages/agent-core-v2/test/agent/mcp/output.test.ts @@ -536,6 +536,7 @@ describe('createMcpTool', () => { async callTool() { return { content: [{ type: 'text', text: 'ok' }], isError: false }; }, + async ping() {}, } satisfies MCPClient; const tool = createMcpTool( 'mcp__server__tool', diff --git a/packages/agent-core-v2/test/agent/mcp/stubs.ts b/packages/agent-core-v2/test/agent/mcp/stubs.ts index ad04fbf6167..af76e429e98 100644 --- a/packages/agent-core-v2/test/agent/mcp/stubs.ts +++ b/packages/agent-core-v2/test/agent/mcp/stubs.ts @@ -77,6 +77,7 @@ export function fakeMcpClient( } return { content: [{ type: 'text', text: 'ok' }], isError: false }; }, + async ping() {}, }; }