From bd738a0a7cfc995400774c4ed6966aab074fa6ee Mon Sep 17 00:00:00 2001 From: Aiden Cline Date: Wed, 24 Jun 2026 16:28:08 -0500 Subject: [PATCH 1/2] fix(opencode): reconnect closed remote MCP clients --- packages/opencode/src/mcp/index.ts | 98 ++++++++++++-- .../test/fixture/mcp-reconnect-server.ts | 94 +++++++++++++ packages/opencode/test/mcp/reconnect.test.ts | 123 ++++++++++++++++++ 3 files changed, 306 insertions(+), 9 deletions(-) create mode 100644 packages/opencode/test/fixture/mcp-reconnect-server.ts create mode 100644 packages/opencode/test/mcp/reconnect.test.ts diff --git a/packages/opencode/src/mcp/index.ts b/packages/opencode/src/mcp/index.ts index db673244b3..ee485de108 100644 --- a/packages/opencode/src/mcp/index.ts +++ b/packages/opencode/src/mcp/index.ts @@ -146,6 +146,9 @@ interface State { clients: Record defs: Record instructions: Record + generation: Record + reconnecting: Record + disposed: boolean } export interface ServerInstructions { @@ -428,16 +431,24 @@ export const layer = Layer.effect( Effect.catch(() => Effect.succeed([] as number[])), ) - function watch(s: State, name: string, client: MCPClient, bridge: EffectBridge.Shape, timeout?: number) { + function watch( + s: State, + name: string, + client: MCPClient, + bridge: EffectBridge.Shape, + mcp: ConfigMCPV1.Info, + generation: number, + ) { client.onclose = () => { - if (s.clients[name] !== client) return + if (s.disposed || s.clients[name] !== client || s.generation[name] !== generation) 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.andThen(events.publish(ToolsChanged, { server: name }).pipe(Effect.ignore)), + Effect.andThen(mcp.type === "remote" ? reconnect(s, name, mcp, generation) : Effect.void), Effect.ignore, ), ) @@ -451,7 +462,7 @@ export const layer = Layer.effect( client.setNotificationHandler(ToolListChangedNotificationSchema, async () => { if (s.clients[name] !== client || s.status[name]?.status !== "connected") return - const listed = await bridge.promise(McpCatalog.defs(client, timeout)) + const listed = await bridge.promise(McpCatalog.defs(client, mcp.timeout)) if (!listed) return if (s.clients[name] !== client || s.status[name]?.status !== "connected") return @@ -489,6 +500,9 @@ export const layer = Layer.effect( clients: {}, defs: {}, instructions: {}, + generation: {}, + reconnecting: {}, + disposed: false, } yield* Effect.forEach( @@ -505,13 +519,14 @@ export const layer = Layer.effect( return } + const generation = nextGeneration(s, key) const result = yield* create(key, mcp) s.status[key] = result.status if (result.mcpClient) { 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, generation) } }), { concurrency: "unbounded" }, @@ -519,6 +534,7 @@ export const layer = Layer.effect( yield* Effect.addFinalizer(() => Effect.gen(function* () { + s.disposed = true const clients = Object.values(s.clients) s.clients = {} s.defs = {} @@ -563,7 +579,8 @@ export const layer = Layer.effect( client: MCPClient, listed: MCPToolDef[], instructions: string | undefined, - timeout?: number, + mcp: ConfigMCPV1.Info, + generation: number, ) { const bridge = yield* EffectBridge.make() const previous = s.clients[name] @@ -572,11 +589,66 @@ export 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, mcp, generation) if (previous) yield* Effect.tryPromise(() => previous.close()).pipe(Effect.ignore) return s.status[name] }) + const reconnect = Effect.fnUntraced(function* ( + s: State, + name: string, + mcp: ConfigMCPV1.Info & { type: "remote" }, + generation: number, + ) { + if (s.reconnecting[name] === generation) return + s.reconnecting[name] = generation + + yield* reconnectAttempt(s, name, mcp, generation, 0).pipe( + Effect.flatMap((result) => { + if (!result?.mcpClient || !result.defs) return Effect.void + if (!ownsGeneration(s, name, generation)) { + return Effect.tryPromise(() => result.mcpClient.close()).pipe(Effect.ignore) + } + return storeClient(s, name, result.mcpClient, result.defs, result.instructions, mcp, generation).pipe( + Effect.andThen(events.publish(ToolsChanged, { server: name })), + Effect.ignore, + ) + }), + Effect.ensuring( + Effect.sync(() => { + if (s.reconnecting[name] === generation) delete s.reconnecting[name] + }), + ), + ) + }) + + function reconnectAttempt( + s: State, + name: string, + mcp: ConfigMCPV1.Info & { type: "remote" }, + generation: number, + attempt: number, + ): Effect.Effect { + return Effect.gen(function* () { + if (!ownsGeneration(s, name, generation)) return undefined + const result = yield* create(name, mcp) + if (result.mcpClient) return result + if (result.status.status !== "failed" || attempt >= 4) return undefined + yield* Effect.sleep(Math.min(250 * 2 ** attempt, 2_000)) + return yield* reconnectAttempt(s, name, mcp, generation, attempt + 1) + }) + } + + function nextGeneration(s: State, name: string) { + const generation = (s.generation[name] ?? 0) + 1 + s.generation[name] = generation + return generation + } + + function ownsGeneration(s: State, name: string, generation: number) { + return !s.disposed && s.generation[name] === generation + } + const status = Effect.fn("MCP.status")(function* () { const s = yield* InstanceState.get(state) @@ -615,8 +687,14 @@ export const layer = Layer.effect( const createAndStore = Effect.fn("MCP.createAndStore")(function* (name: string, mcp: ConfigMCPV1.Info) { const s = yield* InstanceState.get(state) + const generation = nextGeneration(s, name) const result = yield* create(name, mcp) + if (!ownsGeneration(s, name, generation)) { + const client = result.mcpClient + if (client) yield* Effect.tryPromise(() => client.close()).pipe(Effect.ignore) + return s.status[name] + } s.status[name] = result.status if (!result.mcpClient) { yield* closeClient(s, name) @@ -624,7 +702,7 @@ export 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, generation) }) const add = Effect.fn("MCP.add")(function* (name: string, mcp: ConfigMCPV1.Info) { @@ -642,6 +720,7 @@ export const layer = Layer.effect( const disconnect = Effect.fn("MCP.disconnect")(function* (name: string) { yield* requireMcpConfig(name) const s = yield* InstanceState.get(state) + nextGeneration(s, name) yield* closeClient(s, name) delete s.clients[name] s.status[name] = { status: "disabled" } @@ -878,7 +957,8 @@ export 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) + const generation = nextGeneration(s, mcpName) + return yield* storeClient(s, mcpName, client, listed, client.getInstructions()?.trim(), mcpConfig, generation) } const callbackPromise = McpOAuthCallback.waitForCallback(result.oauthState, mcpName) diff --git a/packages/opencode/test/fixture/mcp-reconnect-server.ts b/packages/opencode/test/fixture/mcp-reconnect-server.ts new file mode 100644 index 0000000000..5637c1e8fa --- /dev/null +++ b/packages/opencode/test/fixture/mcp-reconnect-server.ts @@ -0,0 +1,94 @@ +import { LATEST_PROTOCOL_VERSION } from "@modelcontextprotocol/sdk/types.js" + +const counts = { initialize: 0, list: 0, call: 0 } +const gates = new Map>() +const waiters = new Map void>>() + +function signal(kind: keyof typeof counts) { + for (const [key, resolvers] of waiters) { + const [targetKind, targetCount] = key.split(":") + if (targetKind !== kind || counts[kind] < Number(targetCount)) continue + waiters.delete(key) + resolvers.forEach((resolve) => resolve()) + } +} + +function wait(kind: keyof typeof counts, count: number) { + if (counts[kind] >= count) return Promise.resolve() + return new Promise((resolve) => { + const key = `${kind}:${count}` + waiters.set(key, [...(waiters.get(key) ?? []), resolve]) + }) +} + +const server = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + async fetch(request) { + const url = new URL(request.url) + if (url.pathname === "/control/block") { + gates.set(`${url.searchParams.get("kind")}:${url.searchParams.get("count")}`, Promise.withResolvers()) + return new Response(null, { status: 204 }) + } + if (url.pathname === "/control/release") { + gates.get(`${url.searchParams.get("kind")}:${url.searchParams.get("count")}`)?.resolve() + return new Response(null, { status: 204 }) + } + if (url.pathname === "/control/wait") { + const kind = url.searchParams.get("kind") as keyof typeof counts + await wait(kind, Number(url.searchParams.get("count"))) + return new Response(null, { status: 204 }) + } + if (url.pathname === "/control/state") return Response.json(counts) + if (request.method === "GET") return new Response(null, { status: 405 }) + if (request.method === "DELETE") return new Response(null, { status: 200 }) + + const message = (await request.json()) as { id?: number; method: string } + if (message.method === "initialize") { + counts.initialize++ + signal("initialize") + await gates.get(`initialize:${counts.initialize}`)?.promise + return Response.json( + { + jsonrpc: "2.0", + id: message.id, + result: { + protocolVersion: LATEST_PROTOCOL_VERSION, + capabilities: { tools: {} }, + serverInfo: { name: "reconnect-test", version: "1" }, + }, + }, + { headers: { "mcp-session-id": `session-${counts.initialize}` } }, + ) + } + if (message.method === "notifications/initialized") return new Response(null, { status: 202 }) + if (message.method === "tools/list") { + counts.list++ + signal("list") + return Response.json({ + jsonrpc: "2.0", + id: message.id, + result: { tools: [{ name: "probe", inputSchema: { type: "object", properties: {} } }] }, + }) + } + if (message.method === "tools/call") { + counts.call++ + signal("call") + const call = counts.call + await gates.get(`call:${call}`)?.promise + return Response.json({ + jsonrpc: "2.0", + id: message.id, + result: { content: [{ type: "text", text: `call-${call}-initialize-${counts.initialize}` }] }, + }) + } + return new Response(null, { status: 202 }) + }, +}) + +process.send?.({ url: server.url.href }) + +process.on("SIGTERM", () => { + server.stop(true) + process.exit(0) +}) diff --git a/packages/opencode/test/mcp/reconnect.test.ts b/packages/opencode/test/mcp/reconnect.test.ts new file mode 100644 index 0000000000..d29403e712 --- /dev/null +++ b/packages/opencode/test/mcp/reconnect.test.ts @@ -0,0 +1,123 @@ +import path from "node:path" +import { expect } from "bun:test" +import { Effect, Exit, Fiber } from "effect" +import { MCP } from "../../src/mcp/index" +import { testEffect } from "../lib/effect" + +const it = testEffect(MCP.defaultLayer) + +function server() { + return Effect.acquireRelease( + Effect.promise( + () => + new Promise<{ child: ReturnType; url: string }>((resolve, reject) => { + const child = Bun.spawn( + [process.execPath, path.join(import.meta.dir, "../fixture/mcp-reconnect-server.ts")], + { + cwd: path.join(import.meta.dir, "../.."), + stdout: "inherit", + stderr: "inherit", + ipc(message) { + if ( + typeof message === "object" && + message !== null && + "url" in message && + typeof message.url === "string" + ) { + resolve({ child, url: message.url }) + } + }, + }, + ) + child.exited.then((code) => reject(new Error(`MCP test server exited before readiness with code ${code}`))) + }), + ), + ({ child }) => + Effect.promise(async () => { + child.kill() + await child.exited + }).pipe(Effect.ignore), + ) +} + +function control(url: string, action: "block" | "release" | "wait", kind: string, count: number) { + return Effect.tryPromise(() => fetch(`${url}control/${action}?kind=${kind}&count=${count}`, { method: "POST" })).pipe( + Effect.filterOrFail( + (response) => response.ok, + (response) => new Error(`control request failed: ${response.status}`), + ), + ) +} + +function state(url: string) { + return Effect.promise(() => fetch(`${url}control/state`).then((response) => response.json())) as Effect.Effect<{ + initialize: number + list: number + call: number + }> +} + +it.instance( + "reconnects once without replaying an ambiguous tool call and publishes the replacement", + () => + Effect.scoped( + Effect.gen(function* () { + const fixture = yield* server() + const mcp = yield* MCP.Service + yield* mcp.add("remote", { type: "remote", url: fixture.url, oauth: false }) + yield* control(fixture.url, "block", "call", 1) + + const execute = (yield* mcp.tools()).remote_probe?.execute + if (!execute) return yield* Effect.die("initial tool missing") + const call = yield* Effect.promise(() => execute({}, { toolCallId: "first", messages: [] })).pipe( + Effect.exit, + Effect.forkScoped, + ) + yield* control(fixture.url, "wait", "call", 1) + const original = (yield* mcp.clients()).remote + const transport = original?.transport + if (!transport) return yield* Effect.die("initial client transport missing") + yield* Effect.promise(() => transport.close()) + yield* control(fixture.url, "release", "call", 1) + const callExit = yield* Fiber.await(call) + expect(Exit.isSuccess(callExit) && Exit.isFailure(callExit.value)).toBe(true) + + yield* control(fixture.url, "wait", "list", 2) + const replacement = (yield* mcp.clients()).remote + expect(replacement).toBeDefined() + expect(replacement).not.toBe(original) + const executeLater = (yield* mcp.tools()).remote_probe?.execute + if (!executeLater) return yield* Effect.die("replacement tool missing") + const result = yield* Effect.promise(() => executeLater({}, { toolCallId: "later", messages: [] })) + expect(result).toMatchObject({ content: [{ text: "call-2-initialize-2" }] }) + expect(yield* state(fixture.url)).toEqual({ initialize: 2, list: 2, call: 2 }) + }), + ), + { config: { mcp: {} } }, +) + +it.instance( + "disconnect fences a reconnect that finishes late", + () => + Effect.scoped( + Effect.gen(function* () { + const fixture = yield* server() + const mcp = yield* MCP.Service + yield* mcp.add("remote", { type: "remote", url: fixture.url, oauth: false }) + yield* control(fixture.url, "block", "initialize", 2) + + const transport = (yield* mcp.clients()).remote?.transport + if (!transport) return yield* Effect.die("initial client transport missing") + yield* Effect.promise(() => transport.close()) + yield* control(fixture.url, "wait", "initialize", 2) + yield* mcp.disconnect("remote") + yield* control(fixture.url, "release", "initialize", 2) + yield* control(fixture.url, "wait", "list", 2) + + expect((yield* mcp.status()).remote).toEqual({ status: "disabled" }) + expect((yield* mcp.clients()).remote).toBeUndefined() + expect(yield* state(fixture.url)).toEqual({ initialize: 2, list: 2, call: 0 }) + }), + ), + { config: { mcp: {} } }, +) From 990430fcf8d94e8274466cc4d9e773c4d2e4c25d Mon Sep 17 00:00:00 2001 From: Aiden Cline Date: Wed, 24 Jun 2026 16:59:08 -0500 Subject: [PATCH 2/2] test(opencode): isolate MCP reconnect coverage --- packages/opencode/src/mcp/index.ts | 5 +- .../test/fixture/mcp-reconnect-scenario.ts | 120 +++++++++++++++ packages/opencode/test/mcp/reconnect.test.ts | 141 +++--------------- 3 files changed, 145 insertions(+), 121 deletions(-) create mode 100644 packages/opencode/test/fixture/mcp-reconnect-scenario.ts diff --git a/packages/opencode/src/mcp/index.ts b/packages/opencode/src/mcp/index.ts index ee485de108..c41ef13a95 100644 --- a/packages/opencode/src/mcp/index.ts +++ b/packages/opencode/src/mcp/index.ts @@ -606,10 +606,11 @@ export const layer = Layer.effect( yield* reconnectAttempt(s, name, mcp, generation, 0).pipe( Effect.flatMap((result) => { if (!result?.mcpClient || !result.defs) return Effect.void + const client = result.mcpClient if (!ownsGeneration(s, name, generation)) { - return Effect.tryPromise(() => result.mcpClient.close()).pipe(Effect.ignore) + return Effect.tryPromise(() => client.close()).pipe(Effect.ignore) } - return storeClient(s, name, result.mcpClient, result.defs, result.instructions, mcp, generation).pipe( + return storeClient(s, name, client, result.defs, result.instructions, mcp, generation).pipe( Effect.andThen(events.publish(ToolsChanged, { server: name })), Effect.ignore, ) diff --git a/packages/opencode/test/fixture/mcp-reconnect-scenario.ts b/packages/opencode/test/fixture/mcp-reconnect-scenario.ts new file mode 100644 index 0000000000..e1fda5abeb --- /dev/null +++ b/packages/opencode/test/fixture/mcp-reconnect-scenario.ts @@ -0,0 +1,120 @@ +import path from "node:path" +import { expect } from "bun:test" +import { Effect, Exit, Fiber } from "effect" +import { MCP } from "../../src/mcp/index" +import { testEffect } from "../lib/effect" + +const it = testEffect(MCP.defaultLayer) + +function server() { + return Effect.acquireRelease( + Effect.promise( + () => + new Promise<{ child: ReturnType; url: string }>((resolve, reject) => { + const child = Bun.spawn([process.execPath, path.join(import.meta.dir, "mcp-reconnect-server.ts")], { + cwd: path.join(import.meta.dir, "../.."), + stdout: "inherit", + stderr: "inherit", + ipc(message) { + if ( + typeof message === "object" && + message !== null && + "url" in message && + typeof message.url === "string" + ) { + resolve({ child, url: message.url }) + } + }, + }) + child.exited.then((code) => reject(new Error(`MCP test server exited before readiness with code ${code}`))) + }), + ), + ({ child }) => + Effect.promise(async () => { + child.kill() + await child.exited + }).pipe(Effect.ignore), + ) +} + +function control(url: string, action: "block" | "release" | "wait", kind: string, count: number) { + return Effect.tryPromise(() => fetch(`${url}control/${action}?kind=${kind}&count=${count}`, { method: "POST" })).pipe( + Effect.filterOrFail( + (response) => response.ok, + (response) => new Error(`control request failed: ${response.status}`), + ), + ) +} + +function state(url: string) { + return Effect.promise(() => fetch(`${url}control/state`).then((response) => response.json())) as Effect.Effect<{ + initialize: number + list: number + call: number + }> +} + +it.instance( + "reconnects once without replaying an ambiguous tool call and publishes the replacement", + () => + Effect.scoped( + Effect.gen(function* () { + const fixture = yield* server() + const mcp = yield* MCP.Service + yield* mcp.add("remote", { type: "remote", url: fixture.url, oauth: false }) + yield* control(fixture.url, "block", "call", 1) + + const execute = (yield* mcp.tools()).remote_probe?.execute + if (!execute) return yield* Effect.die("initial tool missing") + const call = yield* Effect.promise(() => execute({}, { toolCallId: "first", messages: [] })).pipe( + Effect.exit, + Effect.forkScoped, + ) + yield* control(fixture.url, "wait", "call", 1) + const original = (yield* mcp.clients()).remote + const transport = original?.transport + if (!transport) return yield* Effect.die("initial client transport missing") + yield* Effect.promise(() => transport.close()) + yield* control(fixture.url, "release", "call", 1) + const callExit = yield* Fiber.await(call) + expect(Exit.isSuccess(callExit) && Exit.isFailure(callExit.value)).toBe(true) + + yield* control(fixture.url, "wait", "list", 2) + const replacement = (yield* mcp.clients()).remote + expect(replacement).toBeDefined() + expect(replacement).not.toBe(original) + const executeLater = (yield* mcp.tools()).remote_probe?.execute + if (!executeLater) return yield* Effect.die("replacement tool missing") + const result = yield* Effect.promise(() => executeLater({}, { toolCallId: "later", messages: [] })) + expect(result).toMatchObject({ content: [{ text: "call-2-initialize-2" }] }) + expect(yield* state(fixture.url)).toEqual({ initialize: 2, list: 2, call: 2 }) + }), + ), + { config: { mcp: {} } }, +) + +it.instance( + "disconnect fences a reconnect that finishes late", + () => + Effect.scoped( + Effect.gen(function* () { + const fixture = yield* server() + const mcp = yield* MCP.Service + yield* mcp.add("remote", { type: "remote", url: fixture.url, oauth: false }) + yield* control(fixture.url, "block", "initialize", 2) + + const transport = (yield* mcp.clients()).remote?.transport + if (!transport) return yield* Effect.die("initial client transport missing") + yield* Effect.promise(() => transport.close()) + yield* control(fixture.url, "wait", "initialize", 2) + yield* mcp.disconnect("remote") + yield* control(fixture.url, "release", "initialize", 2) + yield* control(fixture.url, "wait", "list", 2) + + expect((yield* mcp.status()).remote).toEqual({ status: "disabled" }) + expect((yield* mcp.clients()).remote).toBeUndefined() + expect(yield* state(fixture.url)).toEqual({ initialize: 2, list: 2, call: 0 }) + }), + ), + { config: { mcp: {} } }, +) diff --git a/packages/opencode/test/mcp/reconnect.test.ts b/packages/opencode/test/mcp/reconnect.test.ts index d29403e712..1a9adffe68 100644 --- a/packages/opencode/test/mcp/reconnect.test.ts +++ b/packages/opencode/test/mcp/reconnect.test.ts @@ -1,123 +1,26 @@ import path from "node:path" -import { expect } from "bun:test" -import { Effect, Exit, Fiber } from "effect" -import { MCP } from "../../src/mcp/index" -import { testEffect } from "../lib/effect" +import { expect, test } from "bun:test" -const it = testEffect(MCP.defaultLayer) - -function server() { - return Effect.acquireRelease( - Effect.promise( - () => - new Promise<{ child: ReturnType; url: string }>((resolve, reject) => { - const child = Bun.spawn( - [process.execPath, path.join(import.meta.dir, "../fixture/mcp-reconnect-server.ts")], - { - cwd: path.join(import.meta.dir, "../.."), - stdout: "inherit", - stderr: "inherit", - ipc(message) { - if ( - typeof message === "object" && - message !== null && - "url" in message && - typeof message.url === "string" - ) { - resolve({ child, url: message.url }) - } - }, - }, - ) - child.exited.then((code) => reject(new Error(`MCP test server exited before readiness with code ${code}`))) - }), - ), - ({ child }) => - Effect.promise(async () => { - child.kill() - await child.exited - }).pipe(Effect.ignore), +test("remote MCP reconnect lifecycle", async () => { + const child = Bun.spawn( + [ + process.execPath, + "test", + path.join(import.meta.dir, "../fixture/mcp-reconnect-scenario.ts"), + "--timeout", + "30000", + ], + { + cwd: path.join(import.meta.dir, "../.."), + stdout: "pipe", + stderr: "pipe", + }, ) -} + const [code, stdout, stderr] = await Promise.all([ + child.exited, + Bun.readableStreamToText(child.stdout), + Bun.readableStreamToText(child.stderr), + ]) -function control(url: string, action: "block" | "release" | "wait", kind: string, count: number) { - return Effect.tryPromise(() => fetch(`${url}control/${action}?kind=${kind}&count=${count}`, { method: "POST" })).pipe( - Effect.filterOrFail( - (response) => response.ok, - (response) => new Error(`control request failed: ${response.status}`), - ), - ) -} - -function state(url: string) { - return Effect.promise(() => fetch(`${url}control/state`).then((response) => response.json())) as Effect.Effect<{ - initialize: number - list: number - call: number - }> -} - -it.instance( - "reconnects once without replaying an ambiguous tool call and publishes the replacement", - () => - Effect.scoped( - Effect.gen(function* () { - const fixture = yield* server() - const mcp = yield* MCP.Service - yield* mcp.add("remote", { type: "remote", url: fixture.url, oauth: false }) - yield* control(fixture.url, "block", "call", 1) - - const execute = (yield* mcp.tools()).remote_probe?.execute - if (!execute) return yield* Effect.die("initial tool missing") - const call = yield* Effect.promise(() => execute({}, { toolCallId: "first", messages: [] })).pipe( - Effect.exit, - Effect.forkScoped, - ) - yield* control(fixture.url, "wait", "call", 1) - const original = (yield* mcp.clients()).remote - const transport = original?.transport - if (!transport) return yield* Effect.die("initial client transport missing") - yield* Effect.promise(() => transport.close()) - yield* control(fixture.url, "release", "call", 1) - const callExit = yield* Fiber.await(call) - expect(Exit.isSuccess(callExit) && Exit.isFailure(callExit.value)).toBe(true) - - yield* control(fixture.url, "wait", "list", 2) - const replacement = (yield* mcp.clients()).remote - expect(replacement).toBeDefined() - expect(replacement).not.toBe(original) - const executeLater = (yield* mcp.tools()).remote_probe?.execute - if (!executeLater) return yield* Effect.die("replacement tool missing") - const result = yield* Effect.promise(() => executeLater({}, { toolCallId: "later", messages: [] })) - expect(result).toMatchObject({ content: [{ text: "call-2-initialize-2" }] }) - expect(yield* state(fixture.url)).toEqual({ initialize: 2, list: 2, call: 2 }) - }), - ), - { config: { mcp: {} } }, -) - -it.instance( - "disconnect fences a reconnect that finishes late", - () => - Effect.scoped( - Effect.gen(function* () { - const fixture = yield* server() - const mcp = yield* MCP.Service - yield* mcp.add("remote", { type: "remote", url: fixture.url, oauth: false }) - yield* control(fixture.url, "block", "initialize", 2) - - const transport = (yield* mcp.clients()).remote?.transport - if (!transport) return yield* Effect.die("initial client transport missing") - yield* Effect.promise(() => transport.close()) - yield* control(fixture.url, "wait", "initialize", 2) - yield* mcp.disconnect("remote") - yield* control(fixture.url, "release", "initialize", 2) - yield* control(fixture.url, "wait", "list", 2) - - expect((yield* mcp.status()).remote).toEqual({ status: "disabled" }) - expect((yield* mcp.clients()).remote).toBeUndefined() - expect(yield* state(fixture.url)).toEqual({ initialize: 2, list: 2, call: 0 }) - }), - ), - { config: { mcp: {} } }, -) + expect(code, `${stdout}\n${stderr}`).toBe(0) +})