Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 30 additions & 24 deletions packages/core/src/mcp/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -531,35 +531,35 @@ export const layer = (options?: Options) =>
yield* bus.publish(McpEvent.StatusChanged, { server: name })
})

// New servers connect in the background so one slow server cannot block Location startup or
// the plugin activation that triggered this reconcile; readers wait on the startup latch.
const addServer = Effect.fnUntraced(function* (name: ServerName, serverConfig: Mcp.ServerConfig) {
const entry: ServerEntry = {
config: serverConfig,
status: { status: "pending" },
startup: Latch.makeUnsafe(),
}
entries.set(name, entry)
yield* register(name, entry)
if (serverConfig.disabled) {
entry.status = { status: "disabled" }
entry.startup.openUnsafe()
yield* bus.publish(McpEvent.StatusChanged, { server: name })
return
}
// A later reconcile can replace or remove this entry before the start acquires the lock.
fork(
Effect.suspend(() => (entries.get(name) === entry ? startServer(name, entry) : entry.startup.open)).pipe(
locks.withLock(name),
),
)
})

let applied: Map<ServerName, Mcp.ServerConfig> | undefined
const overrides = new Map<ServerName, Mcp.ServerConfig | false>()
const reconcileLock = Semaphore.makeUnsafe(1)
const reconcile = Effect.fnUntraced(function* () {
const servers = state.get().servers
if (!applied && entries.size === 0) {
for (const [name, server] of servers) {
entries.set(name, {
config: server,
status: { status: "pending" },
startup: Latch.makeUnsafe(),
})
}
yield* Effect.forEach(entries, ([name, entry]) => register(name, entry), { discard: true })
applied = servers

// Initial connections stay asynchronous so one slow server does not block Location startup.
for (const [name, entry] of entries) {
if (entry.config.disabled) {
entry.status = { status: "disabled" }
entry.startup.openUnsafe()
yield* bus.publish(McpEvent.StatusChanged, { server: name })
continue
}
fork(startServer(name, entry).pipe(locks.withLock(name)))
}
return
}

const names = new Set([...(applied?.keys() ?? []), ...servers.keys()])
for (const name of names) {
const previous = applied?.get(name)
Expand All @@ -569,6 +569,10 @@ export const layer = (options?: Options) =>
yield* removeServer(name).pipe(locks.withLock(name))
continue
}
if (!entries.has(name)) {
yield* addServer(name, updated).pipe(locks.withLock(name))
continue
}
yield* replaceServer(name, updated).pipe(locks.withLock(name))
}
applied = servers
Expand Down Expand Up @@ -636,6 +640,8 @@ export const layer = (options?: Options) =>
const name = ServerName.make(server)
overrides.set(name, config)
yield* state.reload()
// Reconcile starts new servers in the background; an explicit add reports once this one settles.
yield* entries.get(name)?.startup.await ?? Effect.void
}),
connect: Effect.fn("MCP.connect")(function* (server) {
const name = ServerName.make(server)
Expand Down
59 changes: 53 additions & 6 deletions packages/core/test/mcp.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import path from "node:path"
import fs from "node:fs/promises"
import os from "node:os"
import { describe, expect, test } from "bun:test"
import { Client, InMemoryTransport, StreamableHTTPClientTransport } from "@modelcontextprotocol/client"
import {
Expand Down Expand Up @@ -1728,9 +1729,8 @@ testEffect(resourceMcpLayer(new ConfigMCP.Local({ type: "local", command: ["unus
command: [process.execPath, path.join(import.meta.dir, "fixture/mcp-output-schema.ts")],
})
})
expect((yield* service.servers()).find((server) => server.name === "dynamic")?.status.status).toBe(
"connected",
)
// A newly added server starts in the background.
expect((yield* settled(service, "dynamic"))?.status).toBe("connected")
expect(yield* service.tools()).toHaveLength(2)

yield* service.transform((editor) => editor.update("dynamic", (server) => (server.codemode = false)))
Expand All @@ -1757,9 +1757,7 @@ testEffect(resourceMcpLayer(new ConfigMCP.Local({ type: "local", command: ["unus
expect(yield* service.tools()).toEqual([])

yield* removed.dispose
expect((yield* service.servers()).find((server) => server.name === "dynamic")?.status.status).toBe(
"connected",
)
expect((yield* settled(service, "dynamic"))?.status).toBe("connected")
expect(yield* service.tools()).toHaveLength(2)
}),
)
Expand Down Expand Up @@ -1866,6 +1864,55 @@ testEffect(Layer.empty).live("batches MCP transforms without connecting intermed
}),
)

testEffect(Layer.empty).live("does not wait for MCP servers added after the first reconcile to start", () =>
Effect.gen(function* () {
const service = yield* Mcp.Service
// The layer already reconciled its configured server, so this is a later addition (e.g. from a plugin).
yield* service
.transform((editor) =>
editor.set(
"hanging",
new ConfigMCP.Local({ type: "local", command: [process.execPath, "-e", "setInterval(() => {}, 1000)"] }),
),
)
.pipe(Effect.timeout("2 seconds"))
expect((yield* service.servers()).find((server) => server.name === "hanging")?.status).toEqual({
status: "pending",
})
}).pipe(
Effect.provide(resourceMcpLayer(new ConfigMCP.Local({ type: "local", command: ["unused"], disabled: true }))),
),
)

testEffect(Layer.empty).live("does not start an MCP server removed before its background start", () =>
Effect.gen(function* () {
const service = yield* Mcp.Service
const marker = path.join(yield* Effect.promise(() => fs.mkdtemp(path.join(os.tmpdir(), "mcp-stale-"))), "pids")
for (let i = 0; i < 20; i++) {
const added = yield* service.transform((editor) =>
editor.set(
"stale",
new ConfigMCP.Local({
type: "local",
command: [
process.execPath,
"-e",
`require("fs").appendFileSync(${JSON.stringify(marker)}, process.pid + "\\n"); setInterval(() => {}, 1000)`,
],
}),
),
)
yield* added.dispose
}
yield* Effect.sleep("500 millis")
const spawned = yield* Effect.promise(() => fs.readFile(marker, "utf8").catch(() => ""))
expect(spawned).toBe("")
expect((yield* service.servers()).some((server) => server.name === "stale")).toBe(false)
}).pipe(
Effect.provide(resourceMcpLayer(new ConfigMCP.Local({ type: "local", command: ["unused"], disabled: true }))),
),
)

test("reconciles only changed MCP server config", async () => {
await Effect.runPromise(
Effect.scoped(
Expand Down
Loading