| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560 |
- import path from "node:path"
- import { pathToFileURL } from "node:url"
- import { expect } from "bun:test"
- import { Server } from "@modelcontextprotocol/sdk/server/index.js"
- import { WebStandardStreamableHTTPServerTransport } from "@modelcontextprotocol/sdk/server/webStandardStreamableHttp.js"
- import {
- GetPromptRequestSchema,
- ListPromptsRequestSchema,
- ListResourcesRequestSchema,
- ListResourceTemplatesRequestSchema,
- ListToolsRequestSchema,
- ReadResourceRequestSchema,
- type ServerCapabilities,
- type Tool,
- } from "@modelcontextprotocol/sdk/types.js"
- import { LayerNode } from "@kirincode-ai/core/effect/layer-node"
- import { Cause, Effect, Exit } from "effect"
- import type { MCP as MCPNS } from "../../src/mcp/index"
- import { MCP } from "../../src/mcp/index"
- import { McpOAuthCallback } from "../../src/mcp/oauth-callback"
- import { TestInstance } from "../fixture/fixture"
- import { pollWithTimeout, testEffect } from "../lib/effect"
- const it = testEffect(LayerNode.compile(MCP.node))
- const stdioFixture = path.join(import.meta.dir, "../fixture/mcp-lifecycle-stdio.ts")
- type Page<T> = { items: T[]; nextCursor?: string }
- interface LifecycleServerState {
- tools: Tool[]
- prompts: Array<{ name: string; description?: string }>
- resources: Array<{ name: string; uri: string; description?: string }>
- resourceTemplates: Array<{ name: string; uriTemplate: string; description?: string }>
- toolPages?: Record<string, Page<Tool>>
- promptPages?: Record<string, Page<{ name: string; description?: string }>>
- resourcePages?: Record<string, Page<{ name: string; uri: string; description?: string }>>
- resourceTemplatePages?: Record<string, Page<{ name: string; uriTemplate: string; description?: string }>>
- listToolsError?: string
- requestDelay?: number
- roots?: Array<{ uri: string; name?: string }>
- requests: string[]
- aborted: number
- }
- function lifecycleServer(input?: { capabilities?: ServerCapabilities; instructions?: string; requestRoots?: boolean }) {
- const capabilities = input?.capabilities ?? { tools: {}, prompts: {}, resources: {} }
- return Effect.acquireRelease(
- Effect.promise(async () => {
- const state: LifecycleServerState = {
- tools: [{ name: "test_tool", description: "A test tool", inputSchema: { type: "object", properties: {} } }],
- prompts: [],
- resources: [],
- resourceTemplates: [],
- requests: [],
- aborted: 0,
- }
- const makeProtocol = async () => {
- const protocol = new Server(
- { name: "mcp-lifecycle", version: "1.0.0" },
- { capabilities, instructions: input?.instructions },
- )
- const transport = new WebStandardStreamableHTTPServerTransport({
- sessionIdGenerator: () => crypto.randomUUID(),
- enableJsonResponse: true,
- })
- if (capabilities.tools) {
- protocol.setRequestHandler(ListToolsRequestSchema, (request) => {
- if (state.listToolsError) throw new Error(state.listToolsError)
- const page = state.toolPages?.[request.params?.cursor ?? "initial"]
- return Promise.resolve({ tools: page?.items ?? state.tools, nextCursor: page?.nextCursor })
- })
- }
- if (capabilities.prompts) {
- protocol.setRequestHandler(ListPromptsRequestSchema, (request) => {
- const page = state.promptPages?.[request.params?.cursor ?? "initial"]
- return Promise.resolve({ prompts: page?.items ?? state.prompts, nextCursor: page?.nextCursor })
- })
- protocol.setRequestHandler(GetPromptRequestSchema, async () => {
- if (state.requestDelay) await Bun.sleep(state.requestDelay)
- return { messages: [{ role: "user", content: { type: "text", text: "prompt result" } }] }
- })
- }
- if (capabilities.resources) {
- protocol.setRequestHandler(ListResourcesRequestSchema, (request) => {
- const page = state.resourcePages?.[request.params?.cursor ?? "initial"]
- return Promise.resolve({ resources: page?.items ?? state.resources, nextCursor: page?.nextCursor })
- })
- protocol.setRequestHandler(ListResourceTemplatesRequestSchema, (request) => {
- const page = state.resourceTemplatePages?.[request.params?.cursor ?? "initial"]
- return Promise.resolve({
- resourceTemplates: page?.items ?? state.resourceTemplates,
- nextCursor: page?.nextCursor,
- })
- })
- protocol.setRequestHandler(ReadResourceRequestSchema, async (request) => {
- if (state.requestDelay) await Bun.sleep(state.requestDelay)
- return { contents: [{ uri: request.params.uri, text: "resource result" }] }
- })
- }
- protocol.oninitialized = () => {
- if (!input?.requestRoots) return
- if (!protocol.getClientCapabilities()?.roots) return
- void Bun.sleep(25)
- .then(() => protocol.listRoots())
- .then((result) => {
- state.roots = result.roots
- })
- .catch(() => {})
- }
- await protocol.connect(transport)
- return { protocol, transport }
- }
- let current = await makeProtocol()
- const http = Bun.serve({
- port: 0,
- fetch(request) {
- state.requests.push(request.method)
- request.signal.addEventListener("abort", () => state.aborted++)
- return current.transport.handleRequest(request)
- },
- })
- return {
- state,
- url: http.url.toString(),
- sendToolListChanged: () => current.protocol.sendToolListChanged(),
- restart: async () => {
- current = await makeProtocol()
- },
- close: async () => {
- await current.protocol.close().catch(() => {})
- http.stop(true)
- },
- }
- }),
- (server) => Effect.promise(server.close),
- )
- }
- function hangingLifecycleServer() {
- return Effect.acquireRelease(
- Effect.promise(async () => {
- const protocol = new Server({ name: "mcp-lifecycle-hanging", version: "1.0.0" }, { capabilities: { tools: {} } })
- protocol.setRequestHandler(ListToolsRequestSchema, () => Promise.resolve({ tools: [] }))
- const transport = new WebStandardStreamableHTTPServerTransport({
- sessionIdGenerator: () => crypto.randomUUID(),
- enableJsonResponse: true,
- })
- await protocol.connect(transport)
- const requests: string[] = []
- let aborted = 0
- const http = Bun.serve({
- port: 0,
- fetch(request) {
- requests.push(request.method)
- request.signal.addEventListener("abort", () => aborted++)
- return new Promise<Response>(() => {})
- },
- })
- return {
- requests,
- aborted: () => aborted,
- url: http.url.toString(),
- close: async () => {
- await protocol.close().catch(() => {})
- http.stop(true)
- },
- }
- }),
- (server) => Effect.promise(server.close),
- )
- }
- function statusName(status: Record<string, MCPNS.Status> | MCPNS.Status, server: string) {
- if ("status" in status) return status.status
- return status[server]?.status
- }
- const remote = (url: string, timeout?: number) => ({ type: "remote" as const, url, oauth: false as const, timeout })
- it.instance("advertises and lists the instance directory as its root", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer({ requestRoots: true })
- const mcp = yield* MCP.Service
- const test = yield* TestInstance
- yield* mcp.add("roots", remote(server.url))
- const roots = yield* pollWithTimeout(
- Effect.sync(() => server.state.roots),
- "server did not receive roots",
- )
- expect(roots).toEqual([{ uri: pathToFileURL(test.directory).href }])
- }),
- )
- it.instance(
- "local mcp cwd resolves relative paths against instance directory",
- () =>
- Effect.gen(function* () {
- const mcp = yield* MCP.Service
- const test = yield* TestInstance
- yield* mcp.add("rel-cwd", {
- type: "local",
- command: [process.execPath, stdioFixture],
- cwd: "plugins/sub",
- })
- expect((yield* mcp.tools())["rel-cwd_current_directory"]?.def.description).toBe(
- path.resolve(test.directory, "plugins/sub"),
- )
- }),
- { init: (directory) => Effect.promise(() => Bun.$`mkdir -p ${path.join(directory, "plugins/sub")}`.quiet()) },
- )
- it.instance("tools() reuses cached definitions until a protocol notification", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer({ capabilities: { tools: { listChanged: true } } })
- const mcp = yield* MCP.Service
- yield* mcp.add("cache-server", remote(server.url))
- server.state.tools = [{ name: "next_tool", inputSchema: { type: "object", properties: {} } }]
- expect(Object.keys(yield* mcp.tools())).toEqual(["cache-server_test_tool"])
- yield* Effect.promise(server.sendToolListChanged)
- yield* pollWithTimeout(
- Effect.gen(function* () {
- const keys = Object.keys(yield* mcp.tools())
- return keys.includes("cache-server_next_tool") ? keys : undefined
- }),
- "tool cache did not refresh",
- )
- expect(Object.keys(yield* mcp.tools())).toEqual(["cache-server_next_tool"])
- }),
- )
- it.instance("instructions() returns non-empty connected server instructions with tool names", () =>
- Effect.gen(function* () {
- const guide = yield* lifecycleServer({ instructions: "Use lookup before mutate." })
- const blank = yield* lifecycleServer({ instructions: " " })
- const mcp = yield* MCP.Service
- yield* mcp.add("guide-server", remote(guide.url))
- yield* mcp.add("blank-server", remote(blank.url))
- expect(yield* mcp.instructions()).toEqual([
- { name: "guide-server", instructions: "Use lookup before mutate.", tools: ["guide-server_test_tool"] },
- ])
- yield* mcp.disconnect("guide-server")
- expect(yield* mcp.instructions()).toEqual([])
- }),
- )
- it.instance("follows cursors when listing tools, prompts, resources, and templates", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer()
- server.state.toolPages = {
- initial: { items: [{ name: "tool-one", inputSchema: { type: "object" } }], nextCursor: "tools-2" },
- "tools-2": { items: [{ name: "tool-two", inputSchema: { type: "object" } }] },
- }
- server.state.promptPages = {
- initial: { items: [{ name: "prompt-one" }], nextCursor: "prompts-2" },
- "prompts-2": { items: [{ name: "prompt-two" }] },
- }
- server.state.resourcePages = {
- initial: { items: [{ name: "resource-one", uri: "test://one" }], nextCursor: "resources-2" },
- "resources-2": { items: [{ name: "resource-two", uri: "test://two" }] },
- }
- server.state.resourceTemplatePages = {
- initial: { items: [{ name: "template-one", uriTemplate: "test://one/{id}" }], nextCursor: "templates-2" },
- "templates-2": { items: [{ name: "template-two", uriTemplate: "test://two/{id}" }] },
- }
- const mcp = yield* MCP.Service
- yield* mcp.add("paged-server", remote(server.url))
- expect(Object.keys(yield* mcp.tools())).toEqual(["paged-server_tool-one", "paged-server_tool-two"])
- expect(Object.keys(yield* mcp.prompts())).toEqual(["paged-server:prompt-one", "paged-server:prompt-two"])
- expect(Object.keys(yield* mcp.resources())).toEqual(["paged-server:test://one", "paged-server:test://two"])
- expect(Object.keys(yield* mcp.resourceTemplates())).toEqual([
- "paged-server:test://one/{id}",
- "paged-server:test://two/{id}",
- ])
- }),
- )
- it.instance("accepts empty cursors and rejects repeated cursors", () =>
- Effect.gen(function* () {
- const empty = yield* lifecycleServer({ capabilities: { prompts: {} } })
- empty.state.promptPages = {
- initial: { items: [{ name: "prompt-one" }], nextCursor: "" },
- "": { items: [{ name: "prompt-two" }] },
- }
- const looping = yield* lifecycleServer({ capabilities: { tools: {} } })
- looping.state.toolPages = {
- initial: { items: [], nextCursor: "repeat" },
- repeat: { items: [], nextCursor: "repeat" },
- }
- const mcp = yield* MCP.Service
- yield* mcp.add("empty-cursor", remote(empty.url))
- const result = yield* mcp.add("looping-cursor", remote(looping.url))
- expect(Object.keys(yield* mcp.prompts())).toEqual(["empty-cursor:prompt-one", "empty-cursor:prompt-two"])
- expect(statusName(result.status, "looping-cursor")).toBe("failed")
- }),
- )
- it.instance("disconnect removes protocol data and reconnect establishes a new session", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer()
- const mcp = yield* MCP.Service
- yield* mcp.add("reconnect-server", remote(server.url))
- expect((yield* mcp.status())["reconnect-server"]?.status).toBe("connected")
- yield* mcp.disconnect("reconnect-server")
- expect((yield* mcp.status())["reconnect-server"]?.status).toBe("disabled")
- expect(Object.keys(yield* mcp.tools())).toEqual([])
- yield* pollWithTimeout(
- Effect.sync(() => (server.state.aborted > 0 ? server.state.aborted : undefined)),
- "disconnected HTTP session was not aborted",
- )
- yield* Effect.promise(server.restart)
- yield* mcp.connect("reconnect-server")
- expect((yield* mcp.status())["reconnect-server"]?.status).toBe("connected")
- expect(Object.keys(yield* mcp.tools())).toEqual(["reconnect-server_test_tool"])
- }),
- )
- it.instance("add() closes the old protocol session when replacing a server", () =>
- Effect.gen(function* () {
- const first = yield* lifecycleServer()
- const second = yield* lifecycleServer()
- const mcp = yield* MCP.Service
- yield* mcp.add("replace-server", remote(first.url))
- yield* mcp.add("replace-server", remote(second.url))
- yield* pollWithTimeout(
- Effect.sync(() => (first.state.aborted > 0 ? first.state.aborted : undefined)),
- "replaced HTTP session was not aborted",
- )
- expect(second.state.aborted).toBe(0)
- expect(Object.keys(yield* mcp.tools())).toEqual(["replace-server_test_tool"])
- }),
- )
- it.instance("one failed server does not affect another connected server", () =>
- Effect.gen(function* () {
- const good = yield* lifecycleServer()
- good.state.tools = [{ name: "good_tool", inputSchema: { type: "object" } }]
- const bad = yield* lifecycleServer()
- bad.state.listToolsError = "listTools failed"
- const mcp = yield* MCP.Service
- yield* mcp.add("good-server", remote(good.url))
- yield* mcp.add("bad-server", remote(bad.url))
- expect((yield* mcp.status())["good-server"]?.status).toBe("connected")
- expect((yield* mcp.status())["bad-server"]?.status).toBe("failed")
- expect(Object.keys(yield* mcp.tools())).toEqual(["good-server_good_tool"])
- }),
- )
- it.instance("falls back when output schema refs fail SDK tool discovery", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer({ capabilities: { tools: {} } })
- server.state.tools = [
- {
- name: "render_screen",
- inputSchema: { type: "object", properties: { prompt: { type: "string" } }, required: ["prompt"] },
- outputSchema: { type: "object", properties: { screen: { $ref: "#/$defs/ScreenInstance" } } },
- },
- ]
- const mcp = yield* MCP.Service
- const result = yield* mcp.add("schema-server", remote(server.url))
- expect(statusName(result.status, "schema-server")).toBe("connected")
- expect(Object.keys(yield* mcp.tools())).toEqual(["schema-server_render_screen"])
- }),
- )
- it.instance("does not fall back for protocol tool discovery errors", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer({ capabilities: { tools: {} } })
- server.state.listToolsError = "transport closed"
- const mcp = yield* MCP.Service
- const result = yield* mcp.add("broken-server", remote(server.url))
- expect(statusName(result.status, "broken-server")).toBe("failed")
- }),
- )
- it.instance("disabled server is marked disabled without opening a protocol session", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer()
- const mcp = yield* MCP.Service
- yield* mcp.add("disabled-server", { ...remote(server.url), enabled: false })
- expect((yield* mcp.status())["disabled-server"]?.status).toBe("disabled")
- expect(server.state.requests).toEqual([])
- }),
- )
- it.instance("returns prompts and URI-keyed resources from connected servers", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer()
- server.state.prompts = [{ name: "my-prompt", description: "A test prompt" }]
- server.state.resources = [
- { name: "same-name", uri: "file:///test.txt" },
- { name: "same-name", uri: "ui://component-state" },
- ]
- const mcp = yield* MCP.Service
- yield* mcp.add("content-server", remote(server.url))
- expect(Object.keys(yield* mcp.prompts())).toEqual(["content-server:my-prompt"])
- expect(Object.keys(yield* mcp.resources())).toEqual([
- "content-server:file:///test.txt",
- "content-server:ui://component-state",
- ])
- yield* mcp.disconnect("content-server")
- expect(yield* mcp.prompts()).toEqual({})
- }),
- )
- it.instance("uses per-server timeouts for prompt and resource requests", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer()
- server.state.requestDelay = 200
- const mcp = yield* MCP.Service
- yield* mcp.add("timeout-server", remote(server.url, 50))
- expect(yield* mcp.getPrompt("timeout-server", "test")).toBeUndefined()
- expect(yield* mcp.readResource("timeout-server", "test://resource")).toBeUndefined()
- }),
- )
- it.instance("connects resource-only, prompt-only, and tools-only servers", () =>
- Effect.gen(function* () {
- const resources = yield* lifecycleServer({ capabilities: { resources: {} } })
- resources.state.resources = [{ name: "docs", uri: "docs://readme" }]
- const prompts = yield* lifecycleServer({ capabilities: { prompts: {} } })
- prompts.state.prompts = [{ name: "review" }]
- const tools = yield* lifecycleServer({ capabilities: { tools: {} } })
- const mcp = yield* MCP.Service
- yield* mcp.add("resource-only", remote(resources.url))
- yield* mcp.add("prompt-only", remote(prompts.url))
- yield* mcp.add("tools-only", remote(tools.url))
- expect(Object.keys(yield* mcp.tools())).toEqual(["tools-only_test_tool"])
- expect(Object.keys(yield* mcp.prompts())).toEqual(["prompt-only:review"])
- expect(Object.keys(yield* mcp.resources())).toEqual(["resource-only:docs://readme"])
- }),
- )
- it.instance("connect and disconnect fail for unknown servers", () =>
- Effect.gen(function* () {
- const mcp = yield* MCP.Service
- for (const operation of [mcp.connect("missing"), mcp.disconnect("missing")]) {
- const exit = yield* operation.pipe(Effect.exit)
- expect(Exit.isFailure(exit)).toBe(true)
- if (Exit.isFailure(exit)) {
- expect(Cause.squash(exit.cause)).toMatchObject({ _tag: "MCP.NotFoundError", name: "missing" })
- }
- }
- expect(yield* mcp.status()).toEqual({})
- expect(yield* mcp.tools()).toEqual({})
- }),
- )
- it.instance("unavailable remote server is marked failed without tools", () =>
- Effect.gen(function* () {
- const server = yield* Effect.acquireRelease(
- Effect.sync(() => Bun.serve({ port: 0, fetch: () => new Response("unavailable", { status: 503 }) })),
- (http) => Effect.promise(() => http.stop(true)),
- )
- const mcp = yield* MCP.Service
- yield* mcp.add("unavailable", remote(server.url.toString(), 500))
- expect((yield* mcp.status()).unavailable?.status).toBe("failed")
- expect(yield* mcp.tools()).toEqual({})
- }),
- )
- it.instance("tools() prefixes sanitized server and tool names", () =>
- Effect.gen(function* () {
- const server = yield* lifecycleServer({ capabilities: { tools: {} } })
- server.state.tools = [
- { name: "tool-a", inputSchema: { type: "object" } },
- { name: "tool.b", inputSchema: { type: "object" } },
- ]
- const mcp = yield* MCP.Service
- yield* mcp.add("my.special-server", remote(server.url))
- expect(Object.keys(yield* mcp.tools())).toEqual(["my_special-server_tool-a", "my_special-server_tool_b"])
- }),
- )
- it.instance("local stdio timeout terminates the real server process", () =>
- Effect.gen(function* () {
- const test = yield* TestInstance
- const pidFile = path.join(test.directory, "mcp.pid")
- const mcp = yield* MCP.Service
- const result = yield* mcp.add("hanging-stdio", {
- type: "local",
- command: [process.execPath, stdioFixture, "--hang"],
- environment: { MCP_LIFECYCLE_PID_FILE: pidFile },
- timeout: 100,
- })
- expect(statusName(result.status, "hanging-stdio")).toBe("failed")
- const pid = yield* pollWithTimeout(
- Effect.promise(async () => {
- const file = Bun.file(pidFile)
- return (await file.exists()) ? Number(await file.text()) : undefined
- }),
- "stdio fixture did not publish its pid",
- )
- yield* pollWithTimeout(
- Effect.sync(() => {
- try {
- process.kill(pid, 0)
- return undefined
- } catch {
- return true
- }
- }),
- "stdio fixture process was not terminated",
- )
- }),
- )
- it.instance("remote timeout aborts both real HTTP transport attempts", () =>
- Effect.gen(function* () {
- const server = yield* hangingLifecycleServer()
- const mcp = yield* MCP.Service
- const result = yield* mcp.add("hanging-remote", remote(server.url, 100))
- expect(statusName(result.status, "hanging-remote")).toBe("failed")
- yield* pollWithTimeout(
- Effect.sync(() => (server.aborted() >= 2 ? server.aborted() : undefined)),
- "remote transport requests were not aborted",
- )
- expect(server.requests).toEqual(["POST", "GET"])
- }),
- )
- it.live("McpOAuthCallback.cancelPending rejects the pending callback", () =>
- Effect.acquireUseRelease(
- Effect.sync(() => McpOAuthCallback.waitForCallback("abc123hexstate", "my-mcp-server")),
- (callback) =>
- Effect.gen(function* () {
- McpOAuthCallback.cancelPending("my-mcp-server")
- const exit = yield* Effect.tryPromise({
- try: () => callback,
- catch: (error) => (error instanceof Error ? error : new Error(String(error))),
- }).pipe(Effect.exit)
- expect(Exit.isFailure(exit)).toBe(true)
- }),
- () => Effect.promise(() => McpOAuthCallback.stop()).pipe(Effect.ignore),
- ),
- )
|