Compare commits

..

2 Commits

Author SHA1 Message Date
Aiden Cline 8f62645677 fix(opencode): make mcp env handling portable 2026-06-24 16:59:07 -05:00
Aiden Cline d1adfdba16 fix(opencode): restrict local mcp environment 2026-06-24 16:24:48 -05:00
6 changed files with 168 additions and 335 deletions
+78 -95
View File
@@ -117,6 +117,53 @@ type ResourceInfo = Awaited<ReturnType<MCPClient["listResources"]>>["resources"]
type ResourceTemplateInfo = Awaited<ReturnType<MCPClient["listResourceTemplates"]>>["resourceTemplates"][number]
type McpEntry = NonNullable<ConfigV1.Info["mcp"]>[string]
const LOCAL_MCP_INHERITED_ENV = [
"APPDATA",
"COMSPEC",
"HOME",
"HOMEDRIVE",
"HOMEPATH",
"LANG",
"LANGUAGE",
"LC_ADDRESS",
"LC_ALL",
"LC_COLLATE",
"LC_CTYPE",
"LC_IDENTIFICATION",
"LC_MEASUREMENT",
"LC_MESSAGES",
"LC_MONETARY",
"LC_NAME",
"LC_NUMERIC",
"LC_PAPER",
"LC_TELEPHONE",
"LC_TIME",
"LOCALAPPDATA",
"LOGNAME",
"PATH",
"PATHEXT",
"PROCESSOR_ARCHITECTURE",
"PROGRAMDATA",
"PROGRAMFILES",
"PROGRAMFILES(X86)",
"SHELL",
"SYSTEMDRIVE",
"SYSTEMROOT",
"TEMP",
"TERM",
"TMP",
"TMPDIR",
"USER",
"USERNAME",
"USERPROFILE",
"WINDIR",
"XDG_CACHE_HOME",
"XDG_CONFIG_HOME",
"XDG_DATA_HOME",
"XDG_RUNTIME_DIR",
"XDG_STATE_HOME",
] as const
function isMcpConfigured(entry: McpEntry): entry is ConfigMCPV1.Info {
return typeof entry === "object" && entry !== null && "type" in entry
}
@@ -125,6 +172,27 @@ function remoteURL(value: string) {
if (URL.canParse(value)) return new URL(value)
}
function localMcpEnvironment(command: string, environment?: Record<string, string>) {
const inherited = Object.fromEntries(
LOCAL_MCP_INHERITED_ENV.flatMap((key) => {
const value = process.env[key]
if (value === undefined || value.startsWith("()")) return []
return [[key, value] as const]
}),
)
const defaults = {
...inherited,
...(command === "opencode" ? { BUN_BE_BUN: "1" } : {}),
}
if (process.platform !== "win32" || !environment) return { ...defaults, ...environment }
const configured = new Set(Object.keys(environment).map((key) => key.toUpperCase()))
return {
...Object.fromEntries(Object.entries(defaults).filter(([key]) => !configured.has(key.toUpperCase()))),
...environment,
}
}
interface CreateResult {
mcpClient?: MCPClient
status: Status
@@ -146,9 +214,6 @@ interface State {
clients: Record<string, MCPClient>
defs: Record<string, MCPToolDef[]>
instructions: Record<string, string>
generation: Record<string, number>
reconnecting: Record<string, number>
disposed: boolean
}
export interface ServerInstructions {
@@ -341,11 +406,7 @@ export const layer = Layer.effect(
command: cmd,
args,
cwd,
env: {
...process.env,
...(cmd === "opencode" ? { BUN_BE_BUN: "1" } : {}),
...mcp.environment,
},
env: localMcpEnvironment(cmd, mcp.environment),
})
const connectTimeout = mcp.timeout ?? DEFAULT_TIMEOUT
@@ -431,24 +492,16 @@ export const layer = Layer.effect(
Effect.catch(() => Effect.succeed([] as number[])),
)
function watch(
s: State,
name: string,
client: MCPClient,
bridge: EffectBridge.Shape,
mcp: ConfigMCPV1.Info,
generation: number,
) {
function watch(s: State, name: string, client: MCPClient, bridge: EffectBridge.Shape, timeout?: number) {
client.onclose = () => {
if (s.disposed || s.clients[name] !== client || s.generation[name] !== generation) return
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 }).pipe(Effect.ignore)),
Effect.andThen(mcp.type === "remote" ? reconnect(s, name, mcp, generation) : Effect.void),
Effect.andThen(events.publish(ToolsChanged, { server: name })),
Effect.ignore,
),
)
@@ -462,7 +515,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, mcp.timeout))
const listed = await bridge.promise(McpCatalog.defs(client, timeout))
if (!listed) return
if (s.clients[name] !== client || s.status[name]?.status !== "connected") return
@@ -500,9 +553,6 @@ export const layer = Layer.effect(
clients: {},
defs: {},
instructions: {},
generation: {},
reconnecting: {},
disposed: false,
}
yield* Effect.forEach(
@@ -519,14 +569,13 @@ 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, generation)
watch(s, key, result.mcpClient, bridge, mcp.timeout)
}
}),
{ concurrency: "unbounded" },
@@ -534,7 +583,6 @@ export const layer = Layer.effect(
yield* Effect.addFinalizer(() =>
Effect.gen(function* () {
s.disposed = true
const clients = Object.values(s.clients)
s.clients = {}
s.defs = {}
@@ -579,8 +627,7 @@ export const layer = Layer.effect(
client: MCPClient,
listed: MCPToolDef[],
instructions: string | undefined,
mcp: ConfigMCPV1.Info,
generation: number,
timeout?: number,
) {
const bridge = yield* EffectBridge.make()
const previous = s.clients[name]
@@ -589,67 +636,11 @@ 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, mcp, generation)
watch(s, name, client, bridge, timeout)
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
const client = result.mcpClient
if (!ownsGeneration(s, name, generation)) {
return Effect.tryPromise(() => client.close()).pipe(Effect.ignore)
}
return storeClient(s, name, client, 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<CreateResult | undefined> {
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)
@@ -688,14 +679,8 @@ 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)
@@ -703,7 +688,7 @@ export const layer = Layer.effect(
return result.status
}
return yield* storeClient(s, name, result.mcpClient, result.defs!, result.instructions, mcp, generation)
return yield* storeClient(s, name, result.mcpClient, result.defs!, result.instructions, mcp.timeout)
})
const add = Effect.fn("MCP.add")(function* (name: string, mcp: ConfigMCPV1.Info) {
@@ -721,7 +706,6 @@ 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" }
@@ -958,8 +942,7 @@ export const layer = Layer.effect(
const s = yield* InstanceState.get(state)
yield* auth.clearOAuthState(mcpName)
const generation = nextGeneration(s, mcpName)
return yield* storeClient(s, mcpName, client, listed, client.getInstructions()?.trim(), mcpConfig, generation)
return yield* storeClient(s, mcpName, client, listed, client.getInstructions()?.trim(), mcpConfig.timeout)
}
const callbackPromise = McpOAuthCallback.waitForCallback(result.oauthState, mcpName)
@@ -0,0 +1,22 @@
import readline from "node:readline"
import { writeFile } from "node:fs/promises"
await writeFile(process.env.MCP_ENV_OUTPUT, JSON.stringify(process.env))
const lines = readline.createInterface({ input: process.stdin })
lines.on("close", () => process.exit(0))
lines.on("line", (line) => {
const request = JSON.parse(line)
if (request.method !== "initialize") return
process.stdout.write(
`${JSON.stringify({
jsonrpc: "2.0",
id: request.id,
result: {
protocolVersion: request.params?.protocolVersion,
capabilities: {},
serverInfo: { name: "environment-test", version: "1" },
},
})}\n`,
)
})
@@ -1,120 +0,0 @@
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<typeof Bun.spawn>; 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: {} } },
)
@@ -1,94 +0,0 @@
import { LATEST_PROTOCOL_VERSION } from "@modelcontextprotocol/sdk/types.js"
const counts = { initialize: 0, list: 0, call: 0 }
const gates = new Map<string, PromiseWithResolvers<void>>()
const waiters = new Map<string, Array<() => 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<void>((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)
})
@@ -0,0 +1,68 @@
import path from "node:path"
import { expect } from "bun:test"
import { Effect } from "effect"
import { MCP } from "../../src/mcp/index"
import { TestInstance } from "../fixture/fixture"
import { testEffect } from "../lib/effect"
const it = testEffect(MCP.defaultLayer)
const inherited = [
"APPDATA",
"HOME",
"LANG",
"LOCALAPPDATA",
"PATH",
"PATHEXT",
"SYSTEMROOT",
"TEMP",
"TMPDIR",
"USERPROFILE",
] as const
it.instance(
"local subprocess receives only baseline and configured environment",
() =>
Effect.gen(function* () {
const test = yield* TestInstance
const previous = process.env.OPENCODE_MCP_PARENT_SECRET
process.env.OPENCODE_MCP_PARENT_SECRET = "parent-secret"
yield* MCP.Service.use((mcp) =>
Effect.gen(function* () {
const output = path.join(test.directory, "environment.json")
const result = yield* mcp.add("environment", {
type: "local",
command: [process.execPath, path.join(import.meta.dir, "../fixture/mcp-environment.js")],
environment: {
MCP_ENV_OUTPUT: output,
MCP_EXPLICIT_TOKEN: "configured-token",
...(process.platform === "win32" ? { Path: path.dirname(process.execPath) } : {}),
},
})
if (!("environment" in result.status)) throw new Error("Expected MCP status map")
expect(result.status.environment).toEqual({ status: "connected" })
const env = (yield* Effect.promise(() => Bun.file(output).json())) as Record<string, string>
expect(env.OPENCODE_MCP_PARENT_SECRET).toBeUndefined()
expect(env.MCP_EXPLICIT_TOKEN).toBe("configured-token")
inherited.forEach((key) => {
if (process.platform === "win32" && key === "PATH") return
if (process.env[key] !== undefined) expect(env[key]).toBe(process.env[key])
})
if (process.platform === "win32") {
expect(Object.entries(env).find(([key]) => key.toUpperCase() === "PATH")?.[1]).toBe(
path.dirname(process.execPath),
)
}
}).pipe(Effect.ensuring(mcp.disconnect("environment").pipe(Effect.ignore))),
).pipe(
Effect.ensuring(
Effect.sync(() => {
if (previous === undefined) delete process.env.OPENCODE_MCP_PARENT_SECRET
else process.env.OPENCODE_MCP_PARENT_SECRET = previous
}),
),
)
}),
{ config: { mcp: {} } },
)
@@ -1,26 +0,0 @@
import path from "node:path"
import { expect, test } from "bun:test"
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),
])
expect(code, `${stdout}\n${stderr}`).toBe(0)
})