From 3f14feb2a75f3f05f248041b6b203e701f77367f Mon Sep 17 00:00:00 2001 From: kelvinwww <1970138194@qq.com> Date: Thu, 20 Aug 2026 11:15:17 +0800 Subject: [PATCH] =?UTF-8?q?fix(mcp):=20=E4=BF=AE=E5=A4=8D=E8=BF=9C?= =?UTF-8?q?=E7=A8=8B=20MCP=20=E8=BF=9E=E6=8E=A5=E4=B8=A2=E5=A4=B1=E5=90=8E?= =?UTF-8?q?=E7=9A=84=E8=87=AA=E5=8A=A8=E9=87=8D=E8=BF=9E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/opencode/src/mcp/index.ts | 246 ++++++++++++++++--- packages/opencode/test/mcp/lifecycle.test.ts | 243 ++++++++++++++++++ 2 files changed, 454 insertions(+), 35 deletions(-) diff --git a/packages/opencode/src/mcp/index.ts b/packages/opencode/src/mcp/index.ts index 05f12fa2ee45..63cc44e4c794 100644 --- a/packages/opencode/src/mcp/index.ts +++ b/packages/opencode/src/mcp/index.ts @@ -26,7 +26,7 @@ import { McpOAuthCallback } from "./oauth-callback" import { McpAuth } from "./auth" import { EventV2Bridge } from "@/event-v2-bridge" import { TuiEvent } from "@/server/tui-event" -import { Cause, Effect, Exit, Layer, Context, Schema, Stream } from "effect" +import { Cause, Effect, Exit, Fiber, Layer, Context, Schema, Stream } from "effect" import { EffectBridge } from "@/effect/bridge" import { InstanceState } from "@/effect/instance-state" import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process" @@ -36,6 +36,7 @@ import { McpEvent } from "@opencode-ai/schema/mcp-event" import { McpBrowser } from "./browser" const DEFAULT_TIMEOUT = 30_000 +const RECONNECT_DELAYS = [1_000, 2_000, 4_000, 8_000, 16_000, 30_000] as const const CLIENT_OPTIONS = { capabilities: { // https://github.com/anomalyco/opencode/issues/11948 @@ -145,6 +146,14 @@ interface State { clients: Record defs: Record instructions: Record + reconnecting: Map + reconnectFibers: Map +} + +interface ReconnectWorker { + token: symbol + client: MCPClient + fiber: Fiber.Fiber } export interface ServerInstructions { @@ -439,19 +448,49 @@ const layer = Layer.effect( Effect.catch(() => Effect.succeed([] as number[])), ) - function watch(s: State, name: string, client: MCPClient, bridge: EffectBridge.Shape, timeout?: number) { + /** + * Invalidates a current remote connection and schedules its replacement. + * Transport errors and close events share this path so a broken client is + * removed before its close callback can fire again. + */ + function connectionLost( + s: State, + name: string, + client: MCPClient, + bridge: EffectBridge.Shape, + error: string, + reconnectable: boolean, + ) { + if (s.clients[name] !== client) return + delete s.clients[name] + delete s.defs[name] + delete s.instructions[name] + s.status[name] = { status: "failed", error } + if (reconnectable) startReconnect(s, name, client, bridge) + bridge.fork( + Effect.gen(function* () { + yield* Effect.logWarning("MCP connection lost", { server: name, error }) + yield* events.publish(ToolsChanged, { server: name }).pipe(Effect.ignore) + }).pipe(Effect.ignore), + ) + } + + function watch( + s: State, + name: string, + client: MCPClient, + bridge: EffectBridge.Shape, + timeout?: number, + reconnectable = false, + ) { client.onclose = () => { - if (s.clients[name] !== client) return - delete s.clients[name] - delete s.defs[name] - delete s.instructions[name] - s.status[name] = { status: "failed", error: "Connection closed" } - bridge.fork( - Effect.logWarning("MCP connection closed", { server: name }).pipe( - Effect.andThen(events.publish(ToolsChanged, { server: name })), - Effect.ignore, - ), - ) + connectionLost(s, name, client, bridge, "Connection closed", reconnectable) + } + + if (reconnectable) { + client.onerror = (error) => { + connectionLost(s, name, client, bridge, error instanceof Error ? error.message : String(error), true) + } } client.setNotificationHandler(LoggingMessageNotificationSchema, (notification) => @@ -471,6 +510,40 @@ const layer = Layer.effect( }) } + /** + * Releases reconnect ownership only if it still belongs to the given worker. + */ + function releaseReconnect(s: State, name: string, token: symbol) { + if (s.reconnecting.get(name) !== token) return + s.reconnecting.delete(name) + const worker = s.reconnectFibers.get(name) + if (worker?.token === token) s.reconnectFibers.delete(name) + } + + /** + * Starts one reconnect worker per source client generation. + */ + function startReconnect(s: State, name: string, client: MCPClient, bridge: EffectBridge.Shape) { + const current = s.reconnectFibers.get(name) + if (current?.client === client) return + const token = Symbol() + s.reconnecting.set(name, token) + const fiber = bridge.fork( + Effect.gen(function* () { + yield* Effect.tryPromise(() => client.close()).pipe(Effect.ignore) + yield* reconnect(s, name, client, token) + }).pipe( + Effect.ensuring( + Effect.sync(() => { + releaseReconnect(s, name, token) + }), + ), + Effect.ignore, + ), + ) + s.reconnectFibers.set(name, { token, client, fiber }) + } + function serverLog(name: string, params: LoggingMessageNotification["params"]) { const fields = { server: name, logger: params.logger, level: params.level, data: params.data } switch (params.level) { @@ -500,6 +573,8 @@ const layer = Layer.effect( clients: {}, defs: {}, instructions: {}, + reconnecting: new Map(), + reconnectFibers: new Map(), } yield* Effect.forEach( @@ -522,7 +597,7 @@ const layer = Layer.effect( s.clients[key] = result.mcpClient s.defs[key] = result.defs! if (result.instructions) s.instructions[key] = result.instructions - watch(s, key, result.mcpClient, bridge, mcp.timeout) + watch(s, key, result.mcpClient, bridge, mcp.timeout, mcp.type === "remote") } }), { concurrency: "unbounded" }, @@ -530,6 +605,11 @@ const layer = Layer.effect( yield* Effect.addFinalizer(() => Effect.gen(function* () { + const reconnectFibers = Array.from(s.reconnectFibers.values()).map((worker) => worker.fiber) + s.reconnectFibers.clear() + s.reconnecting.clear() + yield* Fiber.interruptAll(reconnectFibers) + const clients = Object.values(s.clients) s.clients = {} s.defs = {} @@ -559,13 +639,46 @@ const layer = Layer.effect( }), ) + /** + * Cancels the in-flight reconnect worker for a server, if one exists. + */ + function cancelReconnect(s: State, name: string) { + const token = s.reconnecting.get(name) + if (token === undefined) return Effect.void + const worker = s.reconnectFibers.get(name) + releaseReconnect(s, name, token) + if (!worker || worker.token !== token) return Effect.void + return Fiber.interrupt(worker.fiber) + } + + const getMcpConfig = Effect.fnUntraced(function* (mcpName: string) { + const s = yield* InstanceState.get(state) + if (s.config[mcpName]) return s.config[mcpName] + + const cfg = yield* cfgSvc.get() + const mcpConfig = cfg.mcp?.[mcpName] + if (!mcpConfig || !isMcpConfigured(mcpConfig)) return undefined + return mcpConfig + }) + + const requireMcpConfig = Effect.fnUntraced(function* (mcpName: string) { + const mcpConfig = yield* getMcpConfig(mcpName) + if (!mcpConfig) return yield* new NotFoundError({ name: mcpName }) + return mcpConfig + }) + + /** + * Removes a client before closing its transport so the close callback is + * recognized as intentional and cannot start another reconnect worker. + */ function closeClient(s: State, name: string) { const client = s.clients[name] delete s.clients[name] delete s.defs[name] delete s.instructions[name] - if (!client) return Effect.void - return Effect.tryPromise(() => client.close()).pipe(Effect.ignore) + return cancelReconnect(s, name).pipe( + Effect.andThen(client ? Effect.tryPromise(() => client.close()).pipe(Effect.ignore) : Effect.void), + ) } const storeClient = Effect.fnUntraced(function* ( @@ -575,6 +688,7 @@ const layer = Layer.effect( listed: MCPToolDef[], instructions: string | undefined, timeout?: number, + reconnectable = false, ) { const bridge = yield* EffectBridge.make() const previous = s.clients[name] @@ -583,7 +697,7 @@ const layer = Layer.effect( s.defs[name] = listed if (instructions) s.instructions[name] = instructions else delete s.instructions[name] - watch(s, name, client, bridge, timeout) + watch(s, name, client, bridge, timeout, reconnectable) if (previous) yield* Effect.tryPromise(() => previous.close()).pipe(Effect.ignore) return s.status[name] }) @@ -626,6 +740,7 @@ const layer = Layer.effect( const createAndStore = Effect.fn("MCP.createAndStore")(function* (name: string, mcp: ConfigMCPV1.Info) { const s = yield* InstanceState.get(state) + yield* cancelReconnect(s, name) const result = yield* create(name, mcp) s.status[name] = result.status @@ -635,7 +750,83 @@ const layer = Layer.effect( return result.status } - return yield* storeClient(s, name, result.mcpClient, result.defs!, result.instructions, mcp.timeout) + return yield* storeClient( + s, + name, + result.mcpClient, + result.defs!, + result.instructions, + mcp.timeout, + mcp.type === "remote", + ) + }) + + /** + * Reconnects an unexpectedly closed remote MCP client until it recovers or + * configuration/user activity makes the retry no longer applicable. + */ + const reconnect = Effect.fn("MCP.reconnect")(function* (s: State, name: string, client: MCPClient, token: symbol) { + let attempt = 0 + + while (true) { + if (s.reconnecting.get(name) !== token) return + const mcp = yield* getMcpConfig(name) + if (!mcp || mcp.enabled === false || mcp.type !== "remote") { + if (mcp?.enabled === false) s.status[name] = { status: "disabled" } + return + } + + const currentBeforeWait = s.clients[name] + if (currentBeforeWait !== undefined) { + if (currentBeforeWait !== client) return + return + } + + attempt++ + yield* Effect.logInfo("reconnecting MCP server", { server: name, attempt }) + yield* Effect.sleep(RECONNECT_DELAYS[Math.min(attempt - 1, RECONNECT_DELAYS.length - 1)]) + if (s.reconnecting.get(name) !== token) return + + const latest = yield* getMcpConfig(name) + if (!latest || latest.enabled === false || latest.type !== "remote") { + if (latest?.enabled === false) s.status[name] = { status: "disabled" } + return + } + + const currentAfterWait = s.clients[name] + if (currentAfterWait !== undefined) { + if (currentAfterWait !== client) return + return + } + + const result = yield* create(name, latest) + const reconnectClient = result.mcpClient + if (s.reconnecting.get(name) !== token) { + if (reconnectClient) yield* Effect.tryPromise(() => reconnectClient.close()).pipe(Effect.ignore) + return + } + s.status[name] = result.status + + if (result.status.status === "needs_auth" || result.status.status === "needs_client_registration") return + if (!reconnectClient) { + yield* Effect.logWarning("MCP reconnect failed", { + server: name, + attempt, + error: result.status.status === "failed" ? result.status.error : result.status.status, + }) + continue + } + + if (s.clients[name] !== undefined) { + yield* Effect.tryPromise(() => reconnectClient.close()).pipe(Effect.ignore) + return + } + + yield* storeClient(s, name, reconnectClient, result.defs!, result.instructions, latest.timeout, true) + yield* events.publish(ToolsChanged, { server: name }).pipe(Effect.ignore) + yield* Effect.logInfo("MCP reconnected", { server: name }) + return + } }) const add = Effect.fn("MCP.add")(function* (name: string, mcp: ConfigMCPV1.Info) { @@ -787,22 +978,6 @@ const layer = Layer.effect( ) }) - const getMcpConfig = Effect.fnUntraced(function* (mcpName: string) { - const s = yield* InstanceState.get(state) - if (s.config[mcpName]) return s.config[mcpName] - - const cfg = yield* cfgSvc.get() - const mcpConfig = cfg.mcp?.[mcpName] - if (!mcpConfig || !isMcpConfigured(mcpConfig)) return undefined - return mcpConfig - }) - - const requireMcpConfig = Effect.fnUntraced(function* (mcpName: string) { - const mcpConfig = yield* getMcpConfig(mcpName) - if (!mcpConfig) return yield* new NotFoundError({ name: mcpName }) - return mcpConfig - }) - const startAuth = Effect.fn("MCP.startAuth")(function* (mcpName: string) { const mcpConfig = yield* requireMcpConfig(mcpName) if (mcpConfig.type !== "remote") throw new Error(`MCP server ${mcpName} is not a remote server`) @@ -892,7 +1067,8 @@ const layer = Layer.effect( const s = yield* InstanceState.get(state) yield* auth.clearOAuthState(mcpName) - return yield* storeClient(s, mcpName, client, listed, client.getInstructions()?.trim(), mcpConfig.timeout) + yield* cancelReconnect(s, mcpName) + return yield* storeClient(s, mcpName, client, listed, client.getInstructions()?.trim(), mcpConfig.timeout, true) } const callbackPromise = McpOAuthCallback.waitForCallback(result.oauthState, mcpName) diff --git a/packages/opencode/test/mcp/lifecycle.test.ts b/packages/opencode/test/mcp/lifecycle.test.ts index 80c8fd22f886..c5bccf33204d 100644 --- a/packages/opencode/test/mcp/lifecycle.test.ts +++ b/packages/opencode/test/mcp/lifecycle.test.ts @@ -40,6 +40,7 @@ interface LifecycleServerState { roots?: Array<{ uri: string; name?: string }> requests: string[] aborted: number + connectionFailures: number } function lifecycleServer(input?: { capabilities?: ServerCapabilities; instructions?: string; requestRoots?: boolean }) { @@ -53,6 +54,7 @@ function lifecycleServer(input?: { capabilities?: ServerCapabilities; instructio resourceTemplates: [], requests: [], aborted: 0, + connectionFailures: 0, } const makeProtocol = async () => { @@ -120,6 +122,10 @@ function lifecycleServer(input?: { capabilities?: ServerCapabilities; instructio fetch(request) { state.requests.push(request.method) request.signal.addEventListener("abort", () => state.aborted++) + if (state.connectionFailures > 0) { + state.connectionFailures-- + return new Response("unavailable", { status: 503 }) + } return current.transport.handleRequest(request) }, }) @@ -327,6 +333,243 @@ it.instance("disconnect removes protocol data and reconnect establishes a new se }), ) +it.instance("unexpected remote disconnect reconnects and restores cached tools", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("unexpected-close", remote(server.url)) + const previous = (yield* mcp.clients())["unexpected-close"] + + yield* Effect.promise(server.restart) + yield* Effect.sync(() => { + previous?.onclose?.() + previous?.onclose?.() + }) + + expect((yield* mcp.status())["unexpected-close"]?.status).toBe("failed") + yield* pollWithTimeout( + Effect.gen(function* () { + const client = (yield* mcp.clients())["unexpected-close"] + return client && client !== previous && (yield* mcp.status())["unexpected-close"]?.status === "connected" + ? client + : undefined + }), + "remote MCP did not reconnect", + "3 seconds", + ) + expect(Object.keys(yield* mcp.tools())).toEqual(["unexpected-close_test_tool"]) + }), +) + +it.instance("transport errors invalidate the current remote client without requiring onclose", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("transport-error", remote(server.url)) + const previous = (yield* mcp.clients())["transport-error"] + + yield* Effect.promise(server.restart) + yield* Effect.sync(() => previous?.onerror?.(new Error("fetch failed: getaddrinfo ENOTFOUND mcp.test"))) + + expect((yield* mcp.clients())["transport-error"]).toBeUndefined() + expect((yield* mcp.status())["transport-error"]).toEqual({ + status: "failed", + error: "fetch failed: getaddrinfo ENOTFOUND mcp.test", + }) + yield* pollWithTimeout( + Effect.gen(function* () { + const client = (yield* mcp.clients())["transport-error"] + return client && client !== previous && (yield* mcp.status())["transport-error"]?.status === "connected" + ? client + : undefined + }), + "remote MCP did not reconnect after a transport error", + "3 seconds", + ) + }), +) + +it.instance("newly reconnected remote clients can reconnect immediately", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("reconnect-generation", remote(server.url)) + const first = (yield* mcp.clients())["reconnect-generation"] + + if (!first) throw new Error("initial remote MCP client was not created") + yield* Effect.promise(server.restart) + yield* Effect.sync(() => first.onerror?.(new Error("fetch failed"))) + + const second = yield* pollWithTimeout( + Effect.gen(function* () { + const client = (yield* mcp.clients())["reconnect-generation"] + return client && client !== first && (yield* mcp.status())["reconnect-generation"]?.status === "connected" + ? client + : undefined + }), + "first remote MCP reconnect did not complete", + "3 seconds", + ) + if (!second) throw new Error("first remote MCP reconnect did not create a client") + + yield* Effect.sync(() => second.onerror?.(new Error("fetch failed again"))) + yield* Effect.promise(server.restart) + + const third = yield* pollWithTimeout( + Effect.gen(function* () { + const client = (yield* mcp.clients())["reconnect-generation"] + return client && + client !== first && + client !== second && + (yield* mcp.status())["reconnect-generation"]?.status === "connected" + ? client + : undefined + }), + "second remote MCP reconnect did not complete", + "3 seconds", + ) + expect(third).toBeDefined() + expect((yield* mcp.status())["reconnect-generation"]?.status).toBe("connected") + expect((yield* mcp.clients())["reconnect-generation"]).not.toBe(first) + expect((yield* mcp.clients())["reconnect-generation"]).not.toBe(second) + expect(Object.keys(yield* mcp.tools())).toEqual(["reconnect-generation_test_tool"]) + }), +) + +it.instance("manual disconnect cancels a newly restored remote client", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("disconnect-restored", remote(server.url)) + const previous = (yield* mcp.clients())["disconnect-restored"] + + yield* Effect.promise(server.restart) + yield* Effect.sync(() => previous?.onerror?.(new Error("fetch failed"))) + yield* pollWithTimeout( + Effect.gen(function* () { + const client = (yield* mcp.clients())["disconnect-restored"] + return client && client !== previous && (yield* mcp.status())["disconnect-restored"]?.status === "connected" + ? client + : undefined + }), + "restored remote MCP did not connect", + "3 seconds", + ) + + const requests = server.state.requests.length + yield* mcp.disconnect("disconnect-restored") + yield* Effect.sleep("1100 millis") + + expect((yield* mcp.status())["disconnect-restored"]?.status).toBe("disabled") + expect((yield* mcp.clients())["disconnect-restored"]).toBeUndefined() + expect(server.state.requests.length).toBe(requests) + }), +) + +it.instance("JSON-RPC application errors do not invalidate a connected remote client", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("application-error", remote(server.url)) + const client = (yield* mcp.clients())["application-error"] + server.state.listToolsError = "application failure" + + const exit = yield* Effect.tryPromise({ + try: () => client.listTools(), + catch: (error) => (error instanceof Error ? error : new Error(String(error))), + }).pipe(Effect.exit) + + expect(Exit.isFailure(exit)).toBe(true) + expect((yield* mcp.status())["application-error"]?.status).toBe("connected") + expect((yield* mcp.clients())["application-error"]).toBe(client) + }), +) + +it.instance("remote reconnect retries transient failures until the network recovers", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("transient-close", remote(server.url)) + const previous = (yield* mcp.clients())["transient-close"] + + server.state.connectionFailures = 4 + yield* Effect.promise(server.restart) + yield* Effect.sync(() => { + previous?.onerror?.(new Error("fetch failed")) + previous?.onclose?.() + }) + + yield* pollWithTimeout( + Effect.gen(function* () { + const client = (yield* mcp.clients())["transient-close"] + return client && client !== previous && (yield* mcp.status())["transient-close"]?.status === "connected" + ? client + : undefined + }), + "remote MCP did not recover after transient failures", + "10 seconds", + ) + expect(server.state.connectionFailures).toBe(0) + expect(Object.keys(yield* mcp.tools())).toEqual(["transient-close_test_tool"]) + }), +) + +it.instance("manual disconnect cancels an unexpected remote reconnect", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("manual-disconnect", remote(server.url)) + const previous = (yield* mcp.clients())["manual-disconnect"] + const requests = server.state.requests.length + + yield* Effect.sync(() => previous?.onclose?.()) + yield* mcp.disconnect("manual-disconnect") + yield* Effect.sleep("1100 millis") + + expect((yield* mcp.status())["manual-disconnect"]?.status).toBe("disabled") + expect(server.state.requests.length).toBe(requests) + }), +) + +it.instance("disabled remote configuration stops a pending reconnect", () => + Effect.gen(function* () { + const server = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("disabled-after-close", remote(server.url)) + const previous = (yield* mcp.clients())["disabled-after-close"] + const requests = server.state.requests.length + + yield* Effect.sync(() => previous?.onclose?.()) + yield* mcp.add("disabled-after-close", { ...remote(server.url), enabled: false }) + yield* Effect.sleep("1100 millis") + + expect((yield* mcp.status())["disabled-after-close"]?.status).toBe("disabled") + expect(server.state.requests.length).toBe(requests) + }), +) + +it.instance("replaced remote clients cannot start a reconnect for the old client", () => + Effect.gen(function* () { + const first = yield* lifecycleServer() + const second = yield* lifecycleServer() + const mcp = yield* MCP.Service + yield* mcp.add("replaced-close", remote(first.url)) + const previous = (yield* mcp.clients())["replaced-close"] + yield* mcp.add("replaced-close", remote(second.url)) + const requests = second.state.requests.length + + yield* Effect.sync(() => { + previous?.onerror?.(new Error("fetch failed")) + previous?.onclose?.() + }) + yield* Effect.sleep("1100 millis") + + expect((yield* mcp.status())["replaced-close"]?.status).toBe("connected") + expect((yield* mcp.clients())["replaced-close"]).not.toBe(previous) + expect(second.state.requests.length).toBe(requests) + }), +) + it.instance("add() closes the old protocol session when replacing a server", () => Effect.gen(function* () { const first = yield* lifecycleServer()