mirror of
https://github.com/anomalyco/opencode.git
synced 2026-09-30 20:47:39 +00:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6b61e05776 | ||
|
|
704ef1c693 | ||
|
|
2f0a2b607d |
@@ -1,5 +1,6 @@
|
||||
import { isDeepStrictEqual } from "node:util"
|
||||
import {
|
||||
isMcpServerNotFoundError,
|
||||
isSessionNotFoundError,
|
||||
type CommandInfo,
|
||||
type ModelRef,
|
||||
@@ -280,6 +281,7 @@ export function make(input: {
|
||||
if (!isSessionNotFoundError(error)) throw error
|
||||
})
|
||||
await turn?.turn.catch(() => {})
|
||||
await releaseMcpServers(input.client, registeredMcp, params.sessionId)
|
||||
detach(params.sessionId)
|
||||
return {}
|
||||
},
|
||||
@@ -487,21 +489,25 @@ async function registerMcpServers(
|
||||
const current = registered.get(session.id) ?? new Set<string>()
|
||||
registered.set(session.id, current)
|
||||
await Promise.all(
|
||||
servers.flatMap((server) => {
|
||||
const config = mcpConfig(server)
|
||||
const key = `${server.name}:${stableStringify(config)}`
|
||||
if (current.has(key)) return []
|
||||
current.add(key)
|
||||
return [
|
||||
client.mcp.add({ server: server.name, location: { directory: session.cwd }, config }).catch((error) => {
|
||||
current.delete(key)
|
||||
throw error
|
||||
}),
|
||||
]
|
||||
servers.map(async (server) => {
|
||||
await client.session.mcp.add({ sessionID: session.id, server: server.name, config: mcpConfig(server) })
|
||||
current.add(server.name)
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
async function releaseMcpServers(client: OpenCodeClient, registered: Map<string, Set<string>>, sessionID: string) {
|
||||
const servers = Array.from(registered.get(sessionID) ?? [])
|
||||
registered.delete(sessionID)
|
||||
await Promise.all(
|
||||
servers.map((server) =>
|
||||
client.session.mcp.remove({ sessionID, server }).catch((error) => {
|
||||
if (!isSessionNotFoundError(error) && !isMcpServerNotFoundError(error)) throw error
|
||||
}),
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
function mcpConfig(server: McpServer) {
|
||||
if ("type" in server) {
|
||||
if (server.type === "acp") throw new Error("MCP-over-ACP is not supported")
|
||||
@@ -519,15 +525,6 @@ function mcpConfig(server: McpServer) {
|
||||
}
|
||||
}
|
||||
|
||||
function stableStringify(value: unknown): string {
|
||||
if (Array.isArray(value)) return `[${value.map(stableStringify).join(",")}]`
|
||||
if (!value || typeof value !== "object") return JSON.stringify(value)
|
||||
return `{${Object.entries(value)
|
||||
.toSorted(([a], [b]) => a.localeCompare(b))
|
||||
.map(([key, item]) => `${JSON.stringify(key)}:${stableStringify(item)}`)
|
||||
.join(",")}}`
|
||||
}
|
||||
|
||||
async function sendUsageUpdate(
|
||||
client: OpenCodeClient,
|
||||
connection: ACPConnection.Connection,
|
||||
|
||||
@@ -317,7 +317,7 @@ describe("acp service directory behavior", () => {
|
||||
expect(invalidConfig).toMatchObject({ _tag: "ACPInvalidConfigOptionError" })
|
||||
})
|
||||
|
||||
test("converts MCP configs and deduplicates registrations per session and config", async () => {
|
||||
test("registers MCP configs for the session on every attach", async () => {
|
||||
const local: McpServer = {
|
||||
name: "tools",
|
||||
command: "bun",
|
||||
@@ -332,7 +332,7 @@ describe("acp service directory behavior", () => {
|
||||
headers: [{ name: "Authorization", value: "Bearer x" }],
|
||||
}
|
||||
let created = 0
|
||||
const mcp = "/api/experimental/mcp/"
|
||||
const mcp = /^\/api\/experimental\/session\/(ses_\d+)\/mcp\/(\w+)$/
|
||||
await using fixture = makeACPFixture({
|
||||
fetch(request) {
|
||||
if (request.method === "POST" && request.path === "/api/session") {
|
||||
@@ -342,21 +342,35 @@ describe("acp service directory behavior", () => {
|
||||
if (request.method === "GET" && request.path === "/api/session/ses_1") {
|
||||
return Response.json({ data: makeSession("ses_1") })
|
||||
}
|
||||
if (request.method === "PUT" && request.path.startsWith(mcp)) {
|
||||
if (request.method === "PUT" && mcp.test(request.path)) {
|
||||
return new Response(null, { status: 204 })
|
||||
}
|
||||
return undefined
|
||||
},
|
||||
})
|
||||
|
||||
await fixture.service.newSession({ cwd: "/workspace", mcpServers: [local, local, remote] })
|
||||
await fixture.service.newSession({ cwd: "/workspace", mcpServers: [local, remote] })
|
||||
await fixture.service.resumeSession({ cwd: "/workspace", sessionId: "ses_1", mcpServers: [local, remote] })
|
||||
await fixture.service.resumeSession({ cwd: "/workspace", sessionId: "ses_1", mcpServers: [changed] })
|
||||
await fixture.service.newSession({ cwd: "/workspace", mcpServers: [local] })
|
||||
|
||||
const adds = fixture.requests.filter((request) => request.method === "PUT" && request.path.startsWith(mcp))
|
||||
expect(adds).toHaveLength(4)
|
||||
expect(adds.filter((request) => request.path === `${mcp}tools`).map((request) => request.body)).toEqual([
|
||||
const adds = fixture.requests.filter((request) => request.method === "PUT" && mcp.test(request.path))
|
||||
expect(adds.map((request) => request.path.match(mcp)?.slice(1))).toEqual([
|
||||
["ses_1", "tools"],
|
||||
["ses_1", "docs"],
|
||||
["ses_1", "tools"],
|
||||
["ses_1", "docs"],
|
||||
["ses_1", "tools"],
|
||||
["ses_2", "tools"],
|
||||
])
|
||||
expect(adds.filter((request) => request.path.endsWith("/mcp/tools")).map((request) => request.body)).toEqual([
|
||||
{
|
||||
config: {
|
||||
type: "local",
|
||||
command: ["bun", "server.ts"],
|
||||
environment: { TOKEN: "x" },
|
||||
},
|
||||
},
|
||||
{
|
||||
config: {
|
||||
type: "local",
|
||||
@@ -379,7 +393,7 @@ describe("acp service directory behavior", () => {
|
||||
},
|
||||
},
|
||||
])
|
||||
expect(adds.find((request) => request.path === `${mcp}docs`)?.body).toEqual({
|
||||
expect(adds.find((request) => request.path.endsWith("/mcp/docs"))?.body).toEqual({
|
||||
config: {
|
||||
type: "remote",
|
||||
url: "https://example.com/mcp",
|
||||
@@ -387,12 +401,7 @@ describe("acp service directory behavior", () => {
|
||||
oauth: false,
|
||||
},
|
||||
})
|
||||
expect(adds.map((request) => request.query)).toEqual([
|
||||
{ "location[directory]": "/workspace" },
|
||||
{ "location[directory]": "/workspace" },
|
||||
{ "location[directory]": "/workspace" },
|
||||
{ "location[directory]": "/workspace" },
|
||||
])
|
||||
expect(adds.map((request) => request.query)).toEqual([{}, {}, {}, {}, {}, {}])
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
@@ -295,6 +295,46 @@ describe("acp service lifecycle", () => {
|
||||
.catch((error: unknown) => error)
|
||||
expect(missing).toMatchObject({ _tag: "ACPSessionNotFoundError", sessionId: session.sessionId })
|
||||
})
|
||||
|
||||
test("releases session MCP servers on close and leaves deletion cleanup to the server", async () => {
|
||||
let created = 0
|
||||
await using fixture = makeACPFixture({
|
||||
fetch(request) {
|
||||
if (request.method === "POST" && request.path === "/api/session") {
|
||||
created++
|
||||
return Response.json({ data: makeSession(`ses_${created}`) })
|
||||
}
|
||||
if (request.method === "GET" && request.path === "/api/session/ses_1") {
|
||||
return Response.json({ data: makeSession("ses_1") })
|
||||
}
|
||||
if (request.path.startsWith("/api/experimental/session/")) return new Response(null, { status: 204 })
|
||||
if (request.method === "POST" && request.path.endsWith("/interrupt"))
|
||||
return Response.json({ interrupted: false })
|
||||
if (request.method === "DELETE" && request.path === "/api/session/ses_2") {
|
||||
return new Response(null, { status: 204 })
|
||||
}
|
||||
return undefined
|
||||
},
|
||||
})
|
||||
const server = { name: "ctx", command: "bun", args: ["ctx.ts"], env: [] }
|
||||
const mcpRequests = () =>
|
||||
fixture.requests
|
||||
.filter((request) => request.path.startsWith("/api/experimental/session/"))
|
||||
.map((request) => `${request.method} ${request.path}`)
|
||||
|
||||
const closed = await fixture.service.newSession({ cwd: "/workspace", mcpServers: [server] })
|
||||
const deleted = await fixture.service.newSession({ cwd: "/workspace", mcpServers: [server] })
|
||||
await fixture.service.closeSession({ sessionId: closed.sessionId })
|
||||
await fixture.service.resumeSession({ cwd: "/workspace", sessionId: closed.sessionId, mcpServers: [server] })
|
||||
await fixture.service.deleteSession({ sessionId: deleted.sessionId })
|
||||
|
||||
expect(mcpRequests()).toEqual([
|
||||
"PUT /api/experimental/session/ses_1/mcp/ctx",
|
||||
"PUT /api/experimental/session/ses_2/mcp/ctx",
|
||||
"DELETE /api/experimental/session/ses_1/mcp/ctx",
|
||||
"PUT /api/experimental/session/ses_1/mcp/ctx",
|
||||
])
|
||||
})
|
||||
})
|
||||
|
||||
function currentValue(result: { readonly configOptions?: readonly SessionConfigOption[] | null }, id: string) {
|
||||
|
||||
@@ -24,7 +24,7 @@ describe("acp service", () => {
|
||||
if (url.pathname === "/api/command")
|
||||
return Response.json({ location, data: [{ name: "review", template: "" }] })
|
||||
if (url.pathname === "/api/session" && request.method === "POST") return Response.json({ data: session })
|
||||
if (url.pathname === "/api/experimental/mcp/docs" && request.method === "PUT")
|
||||
if (url.pathname === "/api/experimental/session/ses_acp/mcp/docs" && request.method === "PUT")
|
||||
return new Response(null, { status: 204 })
|
||||
return new Response(null, { status: 404 })
|
||||
},
|
||||
@@ -54,7 +54,7 @@ describe("acp service", () => {
|
||||
expect(result.configOptions?.map((option) => option.id)).toEqual(["model", "effort", "mode"])
|
||||
expect(requests).toContainEqual({
|
||||
method: "PUT",
|
||||
path: "/api/experimental/mcp/docs",
|
||||
path: "/api/experimental/session/ses_acp/mcp/docs",
|
||||
body: {
|
||||
config: { type: "local", command: ["bun", "docs.ts"], environment: { TOKEN: "x" } },
|
||||
},
|
||||
|
||||
@@ -19,13 +19,13 @@ import type { Skill } from "@opencode/schema/skill"
|
||||
import type { FileDiff } from "@opencode/schema/file-diff"
|
||||
import type { InstructionEntry } from "@opencode/schema/instruction-entry"
|
||||
import type { Schema } from "effect"
|
||||
import type { Mcp } from "@opencode/schema/mcp"
|
||||
import type { Event } from "@opencode/schema/event"
|
||||
import type { EventLog } from "@opencode/schema/event-log"
|
||||
import type { Shell } from "@opencode/schema/shell"
|
||||
import type { Provider } from "@opencode/schema/provider"
|
||||
import type { Form } from "@opencode/schema/form"
|
||||
import type { Integration } from "@opencode/schema/integration"
|
||||
import type { Mcp } from "@opencode/schema/mcp"
|
||||
import type { Credential } from "@opencode/schema/credential"
|
||||
import type { PermissionSaved } from "@opencode/schema/permission-saved"
|
||||
import type { FileSystem } from "@opencode/schema/filesystem"
|
||||
@@ -411,6 +411,20 @@ export type SessionInstructionsEntryRemoveOperation<E = never> = (
|
||||
input: SessionInstructionsEntryRemoveInput,
|
||||
) => Effect.Effect<SessionInstructionsEntryRemoveOutput, E>
|
||||
|
||||
export type SessionMcpAddInput = {
|
||||
readonly sessionID: Session.ID
|
||||
readonly server: string
|
||||
readonly config: Mcp.LocalConfig | Mcp.RemoteConfig
|
||||
}
|
||||
export type SessionMcpAddOutput = void
|
||||
export type SessionMcpAddOperation<E = never> = (input: SessionMcpAddInput) => Effect.Effect<SessionMcpAddOutput, E>
|
||||
|
||||
export type SessionMcpRemoveInput = { readonly sessionID: Session.ID; readonly server: string }
|
||||
export type SessionMcpRemoveOutput = void
|
||||
export type SessionMcpRemoveOperation<E = never> = (
|
||||
input: SessionMcpRemoveInput,
|
||||
) => Effect.Effect<SessionMcpRemoveOutput, E>
|
||||
|
||||
export type SessionGenerateInput = { readonly sessionID: Session.ID; readonly prompt: string }
|
||||
export type SessionGenerateOutput = { readonly text: string }
|
||||
export type SessionGenerateOperation<E = never> = (
|
||||
@@ -1477,6 +1491,7 @@ export interface SessionApi<E = never> {
|
||||
readonly remove: SessionInstructionsEntryRemoveOperation<E>
|
||||
}
|
||||
}
|
||||
readonly mcp: { readonly add: SessionMcpAddOperation<E>; readonly remove: SessionMcpRemoveOperation<E> }
|
||||
readonly generate: SessionGenerateOperation<E>
|
||||
readonly log: SessionLogOperation<E>
|
||||
readonly interrupt: SessionInterruptOperation<E>
|
||||
|
||||
@@ -83,6 +83,10 @@ import type {
|
||||
SessionInstructionsEntryPutOutput,
|
||||
SessionInstructionsEntryRemoveInput,
|
||||
SessionInstructionsEntryRemoveOutput,
|
||||
SessionMcpAddInput,
|
||||
SessionMcpAddOutput,
|
||||
SessionMcpRemoveInput,
|
||||
SessionMcpRemoveOutput,
|
||||
SessionGenerateInput,
|
||||
SessionGenerateOutput,
|
||||
SessionLogInput,
|
||||
@@ -657,6 +661,21 @@ const EndpointSessionInstructionsEntryRemove =
|
||||
),
|
||||
)
|
||||
|
||||
const EndpointSessionMcpAdd = (raw: RawClient["server.session"]) => (input: SessionMcpAddInput) =>
|
||||
preserveEffect<SessionMcpAddOutput>()(
|
||||
raw["session.mcp.add"]({
|
||||
params: { sessionID: input["sessionID"], server: input["server"] },
|
||||
payload: { config: input["config"] },
|
||||
}).pipe(Effect.mapError(mapClientError)),
|
||||
)
|
||||
|
||||
const EndpointSessionMcpRemove = (raw: RawClient["server.session"]) => (input: SessionMcpRemoveInput) =>
|
||||
preserveEffect<SessionMcpRemoveOutput>()(
|
||||
raw["session.mcp.remove"]({ params: { sessionID: input["sessionID"], server: input["server"] } }).pipe(
|
||||
Effect.mapError(mapClientError),
|
||||
),
|
||||
)
|
||||
|
||||
const EndpointSessionGenerate = (raw: RawClient["server.session"]) => (input: SessionGenerateInput) =>
|
||||
preserveEffect<SessionGenerateOutput>()(
|
||||
raw["session.generate"]({ params: { sessionID: input["sessionID"] }, payload: { prompt: input["prompt"] } }).pipe(
|
||||
@@ -796,6 +815,7 @@ const adaptGroupSession = (raw: RawClient["server.session"]) => ({
|
||||
remove: EndpointSessionInstructionsEntryRemove(raw),
|
||||
},
|
||||
},
|
||||
mcp: { add: EndpointSessionMcpAdd(raw), remove: EndpointSessionMcpRemove(raw) },
|
||||
generate: EndpointSessionGenerate(raw),
|
||||
log: EndpointSessionLog(raw),
|
||||
interrupt: EndpointSessionInterrupt(raw),
|
||||
|
||||
@@ -77,6 +77,10 @@ import type {
|
||||
SessionInstructionsEntryPutOutput,
|
||||
SessionInstructionsEntryRemoveInput,
|
||||
SessionInstructionsEntryRemoveOutput,
|
||||
SessionMcpAddInput,
|
||||
SessionMcpAddOutput,
|
||||
SessionMcpRemoveInput,
|
||||
SessionMcpRemoveOutput,
|
||||
SessionGenerateInput,
|
||||
SessionGenerateOutput,
|
||||
SessionLogInput,
|
||||
@@ -946,6 +950,31 @@ export function make(options: ClientOptions) {
|
||||
),
|
||||
},
|
||||
},
|
||||
mcp: {
|
||||
add: (input: SessionMcpAddInput, requestOptions?: RequestOptions) =>
|
||||
request<SessionMcpAddOutput>(
|
||||
{
|
||||
method: "PUT",
|
||||
path: `/api/experimental/session/${encodeURIComponent(input.sessionID)}/mcp/${encodeURIComponent(input.server)}`,
|
||||
body: { config: input["config"] },
|
||||
successStatus: 204,
|
||||
declaredStatuses: [400, 401, 404],
|
||||
empty: true,
|
||||
},
|
||||
requestOptions,
|
||||
),
|
||||
remove: (input: SessionMcpRemoveInput, requestOptions?: RequestOptions) =>
|
||||
request<SessionMcpRemoveOutput>(
|
||||
{
|
||||
method: "DELETE",
|
||||
path: `/api/experimental/session/${encodeURIComponent(input.sessionID)}/mcp/${encodeURIComponent(input.server)}`,
|
||||
successStatus: 204,
|
||||
declaredStatuses: [400, 401, 404],
|
||||
empty: true,
|
||||
},
|
||||
requestOptions,
|
||||
),
|
||||
},
|
||||
generate: (input: SessionGenerateInput, requestOptions?: RequestOptions) =>
|
||||
request<{ readonly data: SessionGenerateOutput }>(
|
||||
{
|
||||
|
||||
@@ -2618,6 +2618,14 @@ export const isInstructionEntryValueTooLargeError = (value: unknown): value is I
|
||||
"_tag" in value &&
|
||||
value["_tag"] === "InstructionEntryValueTooLargeError"
|
||||
|
||||
export type McpServerNotFoundError = {
|
||||
readonly _tag: "McpServerNotFoundError"
|
||||
readonly server: string
|
||||
readonly message: string
|
||||
}
|
||||
export const isMcpServerNotFoundError = (value: unknown): value is McpServerNotFoundError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "McpServerNotFoundError"
|
||||
|
||||
export type FormNotFoundError = { readonly _tag: "FormNotFoundError"; readonly id: string; readonly message: string }
|
||||
export const isFormNotFoundError = (value: unknown): value is FormNotFoundError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "FormNotFoundError"
|
||||
@@ -2672,14 +2680,6 @@ export type IntegrationMethodNotFoundError = {
|
||||
export const isIntegrationMethodNotFoundError = (value: unknown): value is IntegrationMethodNotFoundError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "IntegrationMethodNotFoundError"
|
||||
|
||||
export type McpServerNotFoundError = {
|
||||
readonly _tag: "McpServerNotFoundError"
|
||||
readonly server: string
|
||||
readonly message: string
|
||||
}
|
||||
export const isMcpServerNotFoundError = (value: unknown): value is McpServerNotFoundError =>
|
||||
typeof value === "object" && value !== null && "_tag" in value && value["_tag"] === "McpServerNotFoundError"
|
||||
|
||||
export type ProjectNotFoundError = {
|
||||
readonly _tag: "ProjectNotFoundError"
|
||||
readonly projectID: string
|
||||
@@ -4556,6 +4556,52 @@ export type SessionInstructionsEntryRemoveInput = {
|
||||
|
||||
export type SessionInstructionsEntryRemoveOutput = void
|
||||
|
||||
export type SessionMcpAddInput = {
|
||||
readonly sessionID: { readonly sessionID: string; readonly server: string }["sessionID"]
|
||||
readonly server: { readonly sessionID: string; readonly server: string }["server"]
|
||||
readonly config: {
|
||||
readonly config:
|
||||
| {
|
||||
readonly type: "local"
|
||||
readonly command: ReadonlyArray<string>
|
||||
readonly cwd?: string
|
||||
readonly environment?: { readonly [x: string]: string }
|
||||
readonly disabled?: boolean
|
||||
readonly codemode?: boolean
|
||||
readonly timeout?: { readonly startup?: number; readonly catalog?: number; readonly execution?: number }
|
||||
readonly protocol?: "legacy" | "auto" | "2026-07-28"
|
||||
}
|
||||
| {
|
||||
readonly type: "remote"
|
||||
readonly url: string
|
||||
readonly headers?: { readonly [x: string]: string }
|
||||
readonly oauth?:
|
||||
| {
|
||||
readonly client_id?: string
|
||||
readonly client_secret?: string
|
||||
readonly scope?: string
|
||||
readonly callback_port?: number
|
||||
readonly redirect_uri?: string
|
||||
readonly auth_server_metadata_url?: string
|
||||
}
|
||||
| false
|
||||
readonly disabled?: boolean
|
||||
readonly codemode?: boolean
|
||||
readonly timeout?: { readonly startup?: number; readonly catalog?: number; readonly execution?: number }
|
||||
readonly protocol?: "legacy" | "auto" | "2026-07-28"
|
||||
}
|
||||
}["config"]
|
||||
}
|
||||
|
||||
export type SessionMcpAddOutput = void
|
||||
|
||||
export type SessionMcpRemoveInput = {
|
||||
readonly sessionID: { readonly sessionID: string; readonly server: string }["sessionID"]
|
||||
readonly server: { readonly sessionID: string; readonly server: string }["server"]
|
||||
}
|
||||
|
||||
export type SessionMcpRemoveOutput = void
|
||||
|
||||
export type SessionGenerateInput = {
|
||||
readonly sessionID: { readonly sessionID: string }["sessionID"]
|
||||
readonly prompt: { readonly prompt: string }["prompt"]
|
||||
|
||||
@@ -22,6 +22,7 @@ import { LocationLifecycle } from "./location-lifecycle.js"
|
||||
import { FileAccess } from "./file-access.js"
|
||||
import { ModelResolver } from "./model-resolver.js"
|
||||
import { Mcp } from "./mcp/index.js"
|
||||
import { McpSession } from "./mcp/session.js"
|
||||
import { Permission } from "./permission.js"
|
||||
import { Plugin } from "./plugin.js"
|
||||
import { PluginHooks } from "./plugin/hooks.js"
|
||||
@@ -88,6 +89,7 @@ const nodes = [
|
||||
FileMutation.node,
|
||||
Formatter.node,
|
||||
Mcp.node,
|
||||
McpSession.node,
|
||||
Permission.node,
|
||||
Tool.node,
|
||||
ToolOutput.node,
|
||||
|
||||
@@ -0,0 +1,145 @@
|
||||
export * as McpElicitation from "./elicitation.js"
|
||||
|
||||
import { Effect } from "effect"
|
||||
import { waitForAbort } from "@opencode/util/process"
|
||||
import { Form } from "../form.js"
|
||||
import type { McpClient } from "./client.js"
|
||||
|
||||
const URL_FIELD_KEY = "elicitation"
|
||||
|
||||
/** Answers one connection's elicitation requests with forms owned by `sessionID`. */
|
||||
export const handler = (forms: Form.Interface, sessionID: string): McpClient.ElicitationHandler => {
|
||||
// Legacy era only: pending URL-mode elicitation forms, settled by notifications/elicitation/complete.
|
||||
const pending = new Map<string, Form.ID>()
|
||||
return {
|
||||
create: (request) =>
|
||||
Effect.gen(function* () {
|
||||
if (request.params.mode === "url") {
|
||||
const formID = Form.ID.create()
|
||||
// Legacy only: 2026-07-28 has no elicitationId and no completion notification, so the form
|
||||
// settles when the user confirms and the SDK retries the tool call itself.
|
||||
const elicitationID: string | undefined = request.params.elicitationId
|
||||
if (elicitationID !== undefined) pending.set(elicitationID, formID)
|
||||
return yield* forms
|
||||
.ask({
|
||||
id: formID,
|
||||
sessionID,
|
||||
title: `${request.server} is requesting input`,
|
||||
metadata: {
|
||||
kind: "mcp-elicitation",
|
||||
server: request.server,
|
||||
...(elicitationID === undefined ? {} : { elicitationID }),
|
||||
message: request.params.message,
|
||||
},
|
||||
fields: [{ key: URL_FIELD_KEY, type: "external", url: request.params.url }],
|
||||
})
|
||||
.pipe(
|
||||
Effect.raceFirst(waitForAbort(request.signal)),
|
||||
Effect.ensuring(Effect.sync(() => elicitationID !== undefined && pending.delete(elicitationID))),
|
||||
Effect.map(
|
||||
(state): McpClient.ElicitationResult => ({
|
||||
action: state.status === "answered" ? "accept" : "cancel",
|
||||
}),
|
||||
),
|
||||
)
|
||||
}
|
||||
const params = request.params
|
||||
const [field, ...fields] = Object.entries(params.requestedSchema.properties).map(([key, property]) =>
|
||||
toField(key, property, params.requestedSchema.required?.includes(key) === true),
|
||||
)
|
||||
if (!field) return { action: "accept", content: {} }
|
||||
return yield* forms
|
||||
.ask({
|
||||
sessionID,
|
||||
title: `${request.server} is requesting input`,
|
||||
metadata: { kind: "mcp-elicitation", server: request.server, message: params.message },
|
||||
fields: [field, ...fields],
|
||||
})
|
||||
.pipe(
|
||||
Effect.raceFirst(waitForAbort(request.signal)),
|
||||
Effect.map((state): McpClient.ElicitationResult => {
|
||||
if (state.status !== "answered") return { action: "cancel" }
|
||||
return {
|
||||
action: "accept",
|
||||
content: Object.fromEntries(
|
||||
Object.entries(state.answer).map(
|
||||
([key, value]): [string, NonNullable<McpClient.ElicitationResult["content"]>[string]] =>
|
||||
typeof value === "object" ? [key, Array.from(value)] : [key, value],
|
||||
),
|
||||
),
|
||||
}
|
||||
}),
|
||||
)
|
||||
}),
|
||||
complete: (request) =>
|
||||
Effect.gen(function* () {
|
||||
const formID = pending.get(request.elicitationID)
|
||||
if (!formID) return
|
||||
yield* forms.reply({ id: formID, answer: { [URL_FIELD_KEY]: true } }).pipe(Effect.ignore)
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
// Schema `optional` strips undefined-valued properties on encode, so fields can assign
|
||||
// optional properties directly instead of conditionally spreading them.
|
||||
function toField(key: string, property: ElicitationProperty, required: boolean): Form.Field {
|
||||
// Some servers emit machine titles like "string with format email"; prefer description/key over those.
|
||||
const machineTitle = /^(boolean|string|number|integer|array|object)(\s+with\b.*|\s+in\b.*)?$/i
|
||||
const title =
|
||||
property.title && !machineTitle.test(property.title.trim()) ? property.title : (property.description ?? key)
|
||||
const base = {
|
||||
key,
|
||||
title,
|
||||
description: property.description === title ? undefined : property.description,
|
||||
required: required || undefined,
|
||||
}
|
||||
switch (property.type) {
|
||||
case "boolean":
|
||||
return { ...base, type: "boolean", default: property.default }
|
||||
case "number":
|
||||
case "integer":
|
||||
return {
|
||||
...base,
|
||||
type: property.type,
|
||||
minimum: property.minimum,
|
||||
maximum: property.maximum,
|
||||
default: property.default,
|
||||
}
|
||||
case "array":
|
||||
return {
|
||||
...base,
|
||||
type: "multiselect",
|
||||
options:
|
||||
"anyOf" in property.items
|
||||
? property.items.anyOf.map((option) => ({ value: option.const, label: option.title }))
|
||||
: property.items.enum.map((value) => ({ value, label: value })),
|
||||
custom: false,
|
||||
minItems: property.minItems,
|
||||
maxItems: property.maxItems,
|
||||
default: property.default,
|
||||
}
|
||||
case "string": {
|
||||
const options =
|
||||
"oneOf" in property
|
||||
? property.oneOf.map((option) => ({ value: option.const, label: option.title }))
|
||||
: "enum" in property
|
||||
? property.enum.map((value, index) => ({
|
||||
value,
|
||||
label: ("enumNames" in property ? property.enumNames?.[index] : undefined) ?? value,
|
||||
}))
|
||||
: undefined
|
||||
return {
|
||||
...base,
|
||||
type: "string",
|
||||
format: "format" in property ? property.format : undefined,
|
||||
minLength: "minLength" in property ? property.minLength : undefined,
|
||||
maxLength: "maxLength" in property ? property.maxLength : undefined,
|
||||
default: property.default,
|
||||
options,
|
||||
custom: options ? false : undefined,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type ElicitationProperty = McpClient.ElicitationFormParams["requestedSchema"]["properties"][string]
|
||||
+340
-425
@@ -15,9 +15,9 @@ import { Form } from "../form.js"
|
||||
import { Integration } from "../integration.js"
|
||||
import { KeyedMutex } from "../effect/keyed-mutex.js"
|
||||
import { Location } from "../location.js"
|
||||
import { waitForAbort } from "@opencode/util/process"
|
||||
import { State } from "../state.js"
|
||||
import type { McpClient } from "./client.js"
|
||||
import { McpElicitation } from "./elicitation.js"
|
||||
|
||||
export const ServerName = Schema.String.pipe(Schema.brand("MCP.ServerName"))
|
||||
export type ServerName = typeof ServerName.Type
|
||||
@@ -25,6 +25,8 @@ export const PromptsChanged = ephemeral({ type: "mcp.prompts.changed", schema: {
|
||||
|
||||
export const Status = Mcp.Status
|
||||
export type Status = Mcp.Status
|
||||
export const ServerConfig = Mcp.ServerConfig
|
||||
export type ServerConfig = Mcp.ServerConfig
|
||||
export type ServerInfo = Mcp.Server
|
||||
|
||||
export interface ServerInstructions {
|
||||
@@ -62,7 +64,11 @@ export class ToolCallError extends Schema.TaggedError<ToolCallError>()("MCP.Tool
|
||||
message: Schema.String,
|
||||
}) {}
|
||||
|
||||
type ServerEntry = {
|
||||
export type Owner = { readonly type: "location" } | { readonly type: "session"; readonly id: Session.ID }
|
||||
export type Lock = <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E, R>
|
||||
|
||||
export type Entry<O extends Owner = Owner> = {
|
||||
readonly owner: O
|
||||
readonly config: Mcp.ServerConfig
|
||||
status: Status
|
||||
readonly startup: Latch.Latch
|
||||
@@ -75,11 +81,13 @@ type ServerEntry = {
|
||||
registration?: State.Registration
|
||||
}
|
||||
|
||||
// MCP elicitations are Location-scoped, not Session-scoped: the server cannot attribute them to a
|
||||
// persisted session row, so their forms are owned by this opaque sentinel session identifier.
|
||||
const LOCATION = { type: "location" } as const
|
||||
type LocationEntry = Entry<typeof LOCATION>
|
||||
|
||||
// Location servers cannot attribute elicitations to a persisted session row, so their forms are owned
|
||||
// by this opaque sentinel session identifier. Session-scoped servers ask on behalf of their owner.
|
||||
const GLOBAL_ELICITATION_SESSION_ID = "global"
|
||||
const URL_ELICITATION_FIELD_KEY = "elicitation"
|
||||
// Connections remain Location-scoped, but shared remote endpoints should not receive concurrent startup bursts.
|
||||
// Shared remote endpoints should not receive concurrent startup bursts, whichever registry owns them.
|
||||
const endpointLoads = KeyedMutex.makeUnsafe<string>()
|
||||
|
||||
type Data = {
|
||||
@@ -95,7 +103,7 @@ export type Editor = {
|
||||
remove: (server: ServerName | string) => void
|
||||
}
|
||||
|
||||
const cloneConfig = (config: Mcp.ServerConfig) => structuredClone(config) as Types.DeepMutable<Mcp.ServerConfig>
|
||||
export const cloneConfig = (config: Mcp.ServerConfig) => structuredClone(config) as Types.DeepMutable<Mcp.ServerConfig>
|
||||
|
||||
export interface Interface extends State.Transformable<Editor> {
|
||||
readonly servers: () => Effect.Effect<ServerInfo[]>
|
||||
@@ -137,30 +145,324 @@ export const Options = Schema.Struct({
|
||||
})
|
||||
export type Options = typeof Options.Type
|
||||
|
||||
/**
|
||||
* Connection lifecycle shared by the Location registry and the Session-scoped registry. Each registry
|
||||
* builds its own runtime and supplies the per-server lock; only Location entries publish MCP events.
|
||||
*/
|
||||
export const makeRuntime = Effect.fnUntraced(function* <O extends Owner>(input: {
|
||||
readonly options?: Options
|
||||
readonly lock: (name: ServerName, owner: O) => Lock
|
||||
}) {
|
||||
const location = yield* Location.Service
|
||||
const environment = yield* Environment.Service
|
||||
const bus = yield* Bus.Service
|
||||
const forms = yield* Form.Service
|
||||
const credentials = yield* Credential.Service
|
||||
const root = yield* Effect.scope
|
||||
const fork = yield* FiberSet.makeRuntime<never, void, never>()
|
||||
// Lifecycle operations close scopes while holding the server lock, firing onClose, so anything
|
||||
// taking it from a connection callback must stay forked.
|
||||
const lock = (name: ServerName, entry: Entry<O>) => input.lock(name, entry.owner)
|
||||
const notify = (
|
||||
entry: Entry<O>,
|
||||
definition:
|
||||
| typeof McpEvent.StatusChanged
|
||||
| typeof McpEvent.ToolsChanged
|
||||
| typeof McpEvent.ResourcesChanged
|
||||
| typeof PromptsChanged,
|
||||
name: ServerName,
|
||||
) => (entry.owner.type === "location" ? bus.publish(definition, { server: name }).pipe(Effect.asVoid) : Effect.void)
|
||||
|
||||
const connectProvider = Effect.fnUntraced(function* (entry: Entry<O>) {
|
||||
if (entry.config.type !== "remote" || !entry.integrationID) return undefined
|
||||
const { McpOAuth } = yield* Effect.promise(() => import("./oauth.js"))
|
||||
return yield* McpOAuth.connectProvider({ config: entry.config, integrationID: entry.integrationID }).pipe(
|
||||
Effect.provideService(Credential.Service, credentials),
|
||||
)
|
||||
})
|
||||
|
||||
const toTool = (server: ServerName, entry: Entry<O>, tool: McpClient.Tool): Tool => ({
|
||||
...tool,
|
||||
server,
|
||||
...(entry.config.codemode === undefined ? {} : { codemode: entry.config.codemode }),
|
||||
})
|
||||
|
||||
const refreshTools = (name: ServerName, entry: Entry<O>, connection: McpClient.Connection) =>
|
||||
connection.tools().pipe(
|
||||
Effect.map((tools) => {
|
||||
entry.tools = tools.map((tool) => toTool(name, entry, tool))
|
||||
}),
|
||||
)
|
||||
|
||||
const refreshPrompts = (name: ServerName, entry: Entry<O>, connection: McpClient.Connection) =>
|
||||
connection.prompts().pipe(
|
||||
Effect.orElseSucceed(() => []),
|
||||
Effect.map((prompts) => {
|
||||
entry.prompts = prompts.map((prompt): Prompt => ({ ...prompt, server: name }))
|
||||
}),
|
||||
Effect.andThen(notify(entry, PromptsChanged, name)),
|
||||
)
|
||||
|
||||
// Runs a connection callback under the server lock, dropping it if the connection is no longer
|
||||
// the entry's live client, so late SDK callbacks cannot commit obsolete state.
|
||||
const whenLive =
|
||||
(name: ServerName, entry: Entry<O>, connection: McpClient.Connection) =>
|
||||
<E>(effect: Effect.Effect<void, E>) =>
|
||||
fork(
|
||||
Effect.suspend(() => (entry.client === connection ? effect : Effect.void)).pipe(
|
||||
lock(name, entry),
|
||||
Effect.ignore,
|
||||
),
|
||||
)
|
||||
|
||||
// Re-establishes a server whose HTTP session the server dropped, unless another path already
|
||||
// replaced the connection while this one waited for the lock.
|
||||
const recover = (name: ServerName, entry: Entry<O>, connection: McpClient.Connection) =>
|
||||
Effect.gen(function* () {
|
||||
if (entry.client !== connection) return
|
||||
yield* Effect.logInfo("mcp session expired, reconnecting", { server: name })
|
||||
yield* stopServer(name, entry)
|
||||
yield* startServer(name, entry)
|
||||
}).pipe(lock(name, entry))
|
||||
|
||||
// Runs a request against the live connection and, if that request observed a session expiry,
|
||||
// reconnects and runs it once more against the replacement. Any other failure passes through.
|
||||
const recovering = <A, E extends Error>(
|
||||
name: ServerName,
|
||||
entry: Entry<O>,
|
||||
connection: McpClient.Connection,
|
||||
run: (connection: McpClient.Connection) => Effect.Effect<A, E>,
|
||||
) =>
|
||||
run(connection).pipe(
|
||||
Effect.catchIf(
|
||||
// The client module is loaded lazily, so match the tagged error by tag rather than class.
|
||||
(error) => "_tag" in error && error._tag === "MCP.SessionExpiredError",
|
||||
(error) =>
|
||||
recover(name, entry, connection).pipe(
|
||||
Effect.flatMap(() => (entry.client ? run(entry.client) : Effect.fail(error))),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
const loadCatalog = (name: ServerName, entry: Entry<O>, connection: McpClient.Connection) =>
|
||||
recovering(name, entry, connection, (connection) =>
|
||||
Effect.all(
|
||||
{
|
||||
resources: connection.resources(),
|
||||
// Some servers declare resources without implementing template listing.
|
||||
templates: connection.resourceTemplates().pipe(Effect.orElseSucceed(() => [])),
|
||||
},
|
||||
{ concurrency: "unbounded" },
|
||||
),
|
||||
).pipe(
|
||||
Effect.map((catalog) =>
|
||||
ResourceCatalog.make({
|
||||
resources: catalog.resources.map((resource) =>
|
||||
Resource.make({
|
||||
server: name,
|
||||
name: resource.name,
|
||||
uri: resource.uri,
|
||||
description: resource.description,
|
||||
mimeType: resource.mimeType,
|
||||
}),
|
||||
),
|
||||
templates: catalog.templates.map((template) =>
|
||||
ResourceTemplate.make({
|
||||
server: name,
|
||||
name: template.name,
|
||||
uriTemplate: template.uriTemplate,
|
||||
description: template.description,
|
||||
mimeType: template.mimeType,
|
||||
}),
|
||||
),
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
const watch = (name: ServerName, entry: Entry<O>, connection: McpClient.Connection) => {
|
||||
const live = whenLive(name, entry, connection)
|
||||
connection.onClose(() =>
|
||||
live(
|
||||
Effect.gen(function* () {
|
||||
entry.status = { status: "failed", error: "Connection closed" }
|
||||
yield* stopServer(name, entry)
|
||||
yield* notify(entry, McpEvent.StatusChanged, name)
|
||||
}),
|
||||
),
|
||||
)
|
||||
// Background refreshes (list-changed) can observe the expiry too; they do not retry, so
|
||||
// reconnect here for them. Foreground calls reconnect through `recovering`.
|
||||
connection.onSessionExpired(() => fork(recover(name, entry, connection)))
|
||||
connection.onToolsChanged(() =>
|
||||
live(refreshTools(name, entry, connection).pipe(Effect.andThen(notify(entry, McpEvent.ToolsChanged, name)))),
|
||||
)
|
||||
connection.onPromptsChanged(() => live(refreshPrompts(name, entry, connection)))
|
||||
connection.onResourcesChanged(() => live(notify(entry, McpEvent.ResourcesChanged, name)))
|
||||
}
|
||||
|
||||
const startServer = (name: ServerName, entry: Entry<O>) =>
|
||||
Effect.gen(function* () {
|
||||
// Announce the handshake so connect() and credential reconnects don't show a stale
|
||||
// disabled/failed status for the duration of the connection attempt.
|
||||
entry.status = { status: "pending" }
|
||||
yield* notify(entry, McpEvent.StatusChanged, name)
|
||||
const scope = yield* Scope.fork(root)
|
||||
entry.scope = scope
|
||||
const authProvider = yield* connectProvider(entry)
|
||||
const { McpClient } = yield* Effect.promise(() => import("./client.js"))
|
||||
// List tools as part of connect so a failure here marks the server failed rather than
|
||||
// leaving it connected with a silently empty tool list and no path to recover.
|
||||
const load = McpClient.connect(
|
||||
name,
|
||||
entry.config,
|
||||
location.directory,
|
||||
authProvider,
|
||||
McpElicitation.handler(forms, entry.owner.type === "session" ? entry.owner.id : GLOBAL_ELICITATION_SESSION_ID),
|
||||
input.options?.clientInfo,
|
||||
).pipe(
|
||||
Effect.flatMap((connection) => connection.tools().pipe(Effect.map((tools) => ({ connection, tools })))),
|
||||
// A stdio server is spawned on this location's execution plane, not the host's.
|
||||
Effect.provideService(Environment.Service, environment),
|
||||
Scope.provide(scope),
|
||||
)
|
||||
const result = yield* (
|
||||
entry.config.type === "remote" ? endpointLoads.withLock(entry.config.url)(load) : load
|
||||
).pipe(Effect.exit)
|
||||
if (Exit.isSuccess(result)) {
|
||||
entry.client = result.value.connection
|
||||
entry.tools = result.value.tools.map((tool) => toTool(name, entry, tool))
|
||||
entry.prompts = []
|
||||
entry.status = { status: "connected" }
|
||||
watch(name, entry, result.value.connection)
|
||||
yield* Effect.logInfo("mcp connected", { server: name, tools: entry.tools.length })
|
||||
// The tool registry reads on this event; a late-connecting server has no other way to appear.
|
||||
yield* notify(entry, McpEvent.ToolsChanged, name)
|
||||
yield* notify(entry, McpEvent.ResourcesChanged, name)
|
||||
yield* notify(entry, McpEvent.StatusChanged, name)
|
||||
whenLive(name, entry, result.value.connection)(refreshPrompts(name, entry, result.value.connection))
|
||||
return
|
||||
}
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
entry.scope = undefined
|
||||
const error = Cause.squash(result.cause)
|
||||
entry.status =
|
||||
error instanceof McpClient.NeedsAuthError
|
||||
? { status: "needs_auth", error: error.message }
|
||||
: { status: "failed", error: error instanceof Error ? error.message : String(error) }
|
||||
yield* Effect.logWarning("mcp connect failed", { server: name, status: entry.status })
|
||||
yield* notify(entry, McpEvent.StatusChanged, name)
|
||||
}).pipe(
|
||||
Effect.ensuring(entry.startup.open),
|
||||
Effect.annotateLogs({ server: name, directory: location.directory, connectionID: crypto.randomUUID() }),
|
||||
)
|
||||
|
||||
const stopServer = Effect.fnUntraced(function* (name: ServerName, entry: Entry<O>) {
|
||||
const scope = entry.scope
|
||||
if (!scope) return
|
||||
entry.scope = undefined
|
||||
entry.client = undefined
|
||||
entry.tools = undefined
|
||||
entry.prompts = undefined
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
yield* notify(entry, McpEvent.ToolsChanged, name)
|
||||
yield* notify(entry, McpEvent.ResourcesChanged, name)
|
||||
yield* notify(entry, PromptsChanged, name)
|
||||
})
|
||||
|
||||
// Brings a new entry online, or marks it disabled, and settles its startup latch even when
|
||||
// `prepare` fails or the caller is interrupted, so readers cannot hang.
|
||||
const open = (name: ServerName, entry: Entry<O>, prepare: Effect.Effect<void> = Effect.void) =>
|
||||
Effect.gen(function* () {
|
||||
yield* prepare
|
||||
if (entry.config.disabled) {
|
||||
entry.status = { status: "disabled" }
|
||||
yield* notify(entry, McpEvent.StatusChanged, name)
|
||||
return
|
||||
}
|
||||
yield* startServer(name, entry)
|
||||
}).pipe(Effect.ensuring(entry.startup.open))
|
||||
|
||||
const settled = (entry: Entry<O>) => entry.startup.await.pipe(Effect.map(() => entry.client))
|
||||
|
||||
return {
|
||||
fork,
|
||||
start: startServer,
|
||||
stop: stopServer,
|
||||
open,
|
||||
callTool: Effect.fnUntraced(function* (
|
||||
name: ServerName,
|
||||
entry: Entry<O>,
|
||||
call: { readonly name: string; readonly args?: Record<string, unknown>; readonly sessionID?: Session.ID },
|
||||
) {
|
||||
const client = yield* settled(entry)
|
||||
if (!client)
|
||||
return yield* new ToolCallError({ server: name, tool: call.name, message: "MCP server is not connected" })
|
||||
const result = yield* recovering(name, entry, client, (connection) =>
|
||||
connection.callTool({ name: call.name, args: call.args, sessionID: call.sessionID }),
|
||||
).pipe(Effect.mapError((error) => new ToolCallError({ server: name, tool: call.name, message: error.message })))
|
||||
return { ...result, server: name, tool: call.name }
|
||||
}),
|
||||
prompt: Effect.fnUntraced(function* (
|
||||
name: ServerName,
|
||||
entry: Entry<O>,
|
||||
prompt: { readonly name: string; readonly args?: Record<string, string> },
|
||||
) {
|
||||
const client = yield* settled(entry)
|
||||
if (!client) return undefined
|
||||
const result = yield* recovering(name, entry, client, (connection) =>
|
||||
connection.prompt({ name: prompt.name, args: prompt.args }),
|
||||
).pipe(Effect.orElseSucceed(() => undefined))
|
||||
if (!result) return undefined
|
||||
return { ...result, server: name, name: prompt.name }
|
||||
}),
|
||||
// Reports what is connected now without waiting for servers still starting.
|
||||
catalog: (name: ServerName, entry: Entry<O>) =>
|
||||
entry.client
|
||||
? loadCatalog(name, entry, entry.client).pipe(Effect.orElseSucceed(() => emptyCatalog))
|
||||
: Effect.succeed(emptyCatalog),
|
||||
resources: Effect.fnUntraced(function* (name: ServerName, entry: Entry<O>) {
|
||||
const client = yield* settled(entry)
|
||||
if (!client) return emptyCatalog
|
||||
return mergeCatalogs([yield* loadCatalog(name, entry, client)])
|
||||
}),
|
||||
readResource: Effect.fnUntraced(function* (name: ServerName, entry: Entry<O>, uri: string) {
|
||||
const client = yield* settled(entry)
|
||||
if (!client) return undefined
|
||||
const result = yield* recovering(name, entry, client, (connection) => connection.readResource({ uri }))
|
||||
if (!result) return undefined
|
||||
return ResourceContent.make({
|
||||
server: name,
|
||||
uri,
|
||||
contents: result.contents.map((part) =>
|
||||
"text" in part
|
||||
? { type: "text", uri: part.uri, text: part.text, mimeType: part.mimeType }
|
||||
: { type: "blob", uri: part.uri, blob: part.blob, mimeType: part.mimeType },
|
||||
),
|
||||
})
|
||||
}),
|
||||
}
|
||||
})
|
||||
|
||||
const emptyCatalog = ResourceCatalog.make({ resources: [], templates: [] })
|
||||
|
||||
export const layer = (options?: Options) =>
|
||||
Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const location = yield* Location.Service
|
||||
const environment = yield* Environment.Service
|
||||
const bus = yield* Bus.Service
|
||||
const forms = yield* Form.Service
|
||||
const integration = yield* Integration.Service
|
||||
const credentials = yield* Credential.Service
|
||||
const root = yield* Effect.scope
|
||||
const fork = yield* FiberSet.makeRuntime<never, void, never>()
|
||||
|
||||
const entries = new Map<ServerName, ServerEntry>()
|
||||
// Serializes lifecycle operations per server. Anything taking this lock from a connection
|
||||
// callback must stay forked: lifecycle operations close scopes while holding it, firing onClose.
|
||||
const entries = new Map<ServerName, LocationEntry>()
|
||||
const locks = KeyedMutex.makeUnsafe<ServerName>()
|
||||
// Legacy era only: pending URL-mode elicitation forms, settled by notifications/elicitation/complete.
|
||||
const urlElicitations = new Map<string, Form.ID>()
|
||||
const runtime = yield* makeRuntime<typeof LOCATION>({ options, lock: (name) => locks.withLock(name) })
|
||||
const fork = runtime.fork
|
||||
|
||||
// Register every remote server as an OAuth integration so credentials live in the global store
|
||||
// rather than in committed config. Servers that connect anonymously simply never use the method.
|
||||
const owned = new Set<Integration.ID>()
|
||||
const register = Effect.fnUntraced(function* (name: ServerName, entry: ServerEntry) {
|
||||
const register = Effect.fnUntraced(function* (name: ServerName, entry: LocationEntry) {
|
||||
if (entry.config.type !== "remote" || entry.config.oauth === false) return
|
||||
const remote = entry.config
|
||||
// Key identity on name + url, not url alone: two configs for the same url under different names are
|
||||
@@ -207,281 +509,8 @@ export const layer = (options?: Options) =>
|
||||
return { name, entry }
|
||||
})
|
||||
|
||||
const connectProvider = Effect.fnUntraced(function* (entry: ServerEntry) {
|
||||
if (entry.config.type !== "remote" || !entry.integrationID) return undefined
|
||||
const { McpOAuth } = yield* Effect.promise(() => import("./oauth.js"))
|
||||
return yield* McpOAuth.connectProvider({ config: entry.config, integrationID: entry.integrationID }).pipe(
|
||||
Effect.provideService(Credential.Service, credentials),
|
||||
)
|
||||
})
|
||||
|
||||
const elicitation = {
|
||||
create: (input: {
|
||||
readonly server: string
|
||||
readonly params: McpClient.ElicitationParams
|
||||
readonly signal: AbortSignal
|
||||
}) =>
|
||||
Effect.gen(function* () {
|
||||
if (input.params.mode === "url") {
|
||||
const formID = Form.ID.create()
|
||||
// Legacy only: 2026-07-28 has no elicitationId and no completion notification, so the form
|
||||
// settles when the user confirms and the SDK retries the tool call itself.
|
||||
const elicitationID: string | undefined = input.params.elicitationId
|
||||
const key = elicitationID === undefined ? undefined : input.server + "\u0000" + elicitationID
|
||||
if (key) urlElicitations.set(key, formID)
|
||||
return yield* forms
|
||||
.ask({
|
||||
id: formID,
|
||||
sessionID: GLOBAL_ELICITATION_SESSION_ID,
|
||||
title: `${input.server} is requesting input`,
|
||||
metadata: {
|
||||
kind: "mcp-elicitation",
|
||||
server: input.server,
|
||||
...(elicitationID === undefined ? {} : { elicitationID }),
|
||||
message: input.params.message,
|
||||
},
|
||||
fields: [{ key: URL_ELICITATION_FIELD_KEY, type: "external", url: input.params.url }],
|
||||
})
|
||||
.pipe(
|
||||
Effect.raceFirst(waitForAbort(input.signal)),
|
||||
Effect.ensuring(Effect.sync(() => key && urlElicitations.delete(key))),
|
||||
Effect.map(
|
||||
(state): McpClient.ElicitationResult => ({
|
||||
action: state.status === "answered" ? "accept" : "cancel",
|
||||
}),
|
||||
),
|
||||
)
|
||||
}
|
||||
const params = input.params
|
||||
const [field, ...fields] = Object.entries(params.requestedSchema.properties).map(([key, property]) =>
|
||||
toElicitationField(key, property, params.requestedSchema.required?.includes(key) === true),
|
||||
)
|
||||
if (!field) return { action: "accept", content: {} }
|
||||
return yield* forms
|
||||
.ask({
|
||||
sessionID: GLOBAL_ELICITATION_SESSION_ID,
|
||||
title: `${input.server} is requesting input`,
|
||||
metadata: { kind: "mcp-elicitation", server: input.server, message: params.message },
|
||||
fields: [field, ...fields],
|
||||
})
|
||||
.pipe(
|
||||
Effect.raceFirst(waitForAbort(input.signal)),
|
||||
Effect.map((state): McpClient.ElicitationResult => {
|
||||
if (state.status !== "answered") return { action: "cancel" }
|
||||
return {
|
||||
action: "accept",
|
||||
content: Object.fromEntries(
|
||||
Object.entries(state.answer).map(
|
||||
([key, value]): [string, NonNullable<McpClient.ElicitationResult["content"]>[string]] =>
|
||||
typeof value === "object" ? [key, Array.from(value)] : [key, value],
|
||||
),
|
||||
),
|
||||
}
|
||||
}),
|
||||
)
|
||||
}),
|
||||
complete: (input: { readonly server: string; readonly elicitationID: string }) =>
|
||||
Effect.gen(function* () {
|
||||
const formID = urlElicitations.get(input.server + "\u0000" + input.elicitationID)
|
||||
if (!formID) return
|
||||
yield* forms.reply({ id: formID, answer: { [URL_ELICITATION_FIELD_KEY]: true } }).pipe(Effect.ignore)
|
||||
}),
|
||||
} satisfies McpClient.ElicitationHandler
|
||||
|
||||
const toTool = (server: ServerName, entry: ServerEntry, tool: McpClient.Tool): Tool => ({
|
||||
...tool,
|
||||
server,
|
||||
...(entry.config.codemode === undefined ? {} : { codemode: entry.config.codemode }),
|
||||
})
|
||||
|
||||
const refreshTools = (name: ServerName, entry: ServerEntry, connection: McpClient.Connection) =>
|
||||
connection.tools().pipe(
|
||||
Effect.map((tools) => {
|
||||
entry.tools = tools.map((tool) => toTool(name, entry, tool))
|
||||
}),
|
||||
)
|
||||
|
||||
const refreshPrompts = (name: ServerName, entry: ServerEntry, connection: McpClient.Connection) =>
|
||||
connection.prompts().pipe(
|
||||
Effect.orElseSucceed(() => []),
|
||||
Effect.map((prompts) => {
|
||||
entry.prompts = prompts.map((prompt): Prompt => ({ ...prompt, server: name }))
|
||||
}),
|
||||
Effect.andThen(bus.publish(PromptsChanged, { server: name })),
|
||||
)
|
||||
|
||||
// Runs a connection callback under the server lock, dropping it if the connection is no longer
|
||||
// the entry's live client, so late SDK callbacks cannot commit obsolete state.
|
||||
const whenLive =
|
||||
(name: ServerName, entry: ServerEntry, connection: McpClient.Connection) =>
|
||||
<E>(effect: Effect.Effect<void, E>) =>
|
||||
fork(
|
||||
Effect.suspend(() => (entry.client === connection ? effect : Effect.void)).pipe(
|
||||
locks.withLock(name),
|
||||
Effect.ignore,
|
||||
),
|
||||
)
|
||||
|
||||
// Re-establishes a server whose HTTP session the server dropped, unless another path already
|
||||
// replaced the connection while this one waited for the lock.
|
||||
const recover = (name: ServerName, entry: ServerEntry, connection: McpClient.Connection) =>
|
||||
Effect.gen(function* () {
|
||||
if (entry.client !== connection) return
|
||||
yield* Effect.logInfo("mcp session expired, reconnecting", { server: name })
|
||||
yield* stopServer(name, entry)
|
||||
yield* startServer(name, entry)
|
||||
}).pipe(locks.withLock(name))
|
||||
|
||||
// Runs a request against the live connection and, if that request observed a session expiry,
|
||||
// reconnects and runs it once more against the replacement. Any other failure passes through.
|
||||
const recovering = <A, E extends Error>(
|
||||
name: ServerName,
|
||||
entry: ServerEntry,
|
||||
connection: McpClient.Connection,
|
||||
run: (connection: McpClient.Connection) => Effect.Effect<A, E>,
|
||||
) =>
|
||||
run(connection).pipe(
|
||||
Effect.catchIf(
|
||||
// The client module is loaded lazily, so match the tagged error by tag rather than class.
|
||||
(error) => "_tag" in error && error._tag === "MCP.SessionExpiredError",
|
||||
(error) =>
|
||||
recover(name, entry, connection).pipe(
|
||||
Effect.flatMap(() => (entry.client ? run(entry.client) : Effect.fail(error))),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
const loadCatalog = (name: ServerName, entry: ServerEntry, connection: McpClient.Connection) =>
|
||||
recovering(name, entry, connection, (connection) =>
|
||||
Effect.all(
|
||||
{
|
||||
resources: connection.resources(),
|
||||
// Some servers declare resources without implementing template listing.
|
||||
templates: connection.resourceTemplates().pipe(Effect.orElseSucceed(() => [])),
|
||||
},
|
||||
{ concurrency: "unbounded" },
|
||||
),
|
||||
).pipe(
|
||||
Effect.map((catalog) =>
|
||||
ResourceCatalog.make({
|
||||
resources: catalog.resources.map((resource) =>
|
||||
Resource.make({
|
||||
server: name,
|
||||
name: resource.name,
|
||||
uri: resource.uri,
|
||||
description: resource.description,
|
||||
mimeType: resource.mimeType,
|
||||
}),
|
||||
),
|
||||
templates: catalog.templates.map((template) =>
|
||||
ResourceTemplate.make({
|
||||
server: name,
|
||||
name: template.name,
|
||||
uriTemplate: template.uriTemplate,
|
||||
description: template.description,
|
||||
mimeType: template.mimeType,
|
||||
}),
|
||||
),
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
const watch = (name: ServerName, entry: ServerEntry, connection: McpClient.Connection) => {
|
||||
const live = whenLive(name, entry, connection)
|
||||
connection.onClose(() =>
|
||||
live(
|
||||
Effect.gen(function* () {
|
||||
entry.status = { status: "failed", error: "Connection closed" }
|
||||
yield* stopServer(name, entry)
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
}),
|
||||
),
|
||||
)
|
||||
// Background refreshes (list-changed) can observe the expiry too; they do not retry, so
|
||||
// reconnect here for them. Foreground calls reconnect through `recovering`.
|
||||
connection.onSessionExpired(() => fork(recover(name, entry, connection)))
|
||||
connection.onToolsChanged(() =>
|
||||
live(
|
||||
refreshTools(name, entry, connection).pipe(
|
||||
Effect.andThen(bus.publish(McpEvent.ToolsChanged, { server: name })),
|
||||
),
|
||||
),
|
||||
)
|
||||
connection.onPromptsChanged(() => live(refreshPrompts(name, entry, connection)))
|
||||
connection.onResourcesChanged(() => live(bus.publish(McpEvent.ResourcesChanged, { server: name })))
|
||||
}
|
||||
|
||||
const startServer = (name: ServerName, entry: ServerEntry) =>
|
||||
Effect.gen(function* () {
|
||||
// Announce the handshake so connect() and credential reconnects don't show a stale
|
||||
// disabled/failed status for the duration of the connection attempt.
|
||||
entry.status = { status: "pending" }
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
const scope = yield* Scope.fork(root)
|
||||
entry.scope = scope
|
||||
const authProvider = yield* connectProvider(entry)
|
||||
const { McpClient } = yield* Effect.promise(() => import("./client.js"))
|
||||
// List tools as part of connect so a failure here marks the server failed rather than
|
||||
// leaving it connected with a silently empty tool list and no path to recover.
|
||||
const load = McpClient.connect(
|
||||
name,
|
||||
entry.config,
|
||||
location.directory,
|
||||
authProvider,
|
||||
elicitation,
|
||||
options?.clientInfo,
|
||||
).pipe(
|
||||
Effect.flatMap((connection) => connection.tools().pipe(Effect.map((tools) => ({ connection, tools })))),
|
||||
// A stdio server is spawned on this location's execution plane, not the host's.
|
||||
Effect.provideService(Environment.Service, environment),
|
||||
Scope.provide(scope),
|
||||
)
|
||||
const result = yield* (
|
||||
entry.config.type === "remote" ? endpointLoads.withLock(entry.config.url)(load) : load
|
||||
).pipe(Effect.exit)
|
||||
if (Exit.isSuccess(result)) {
|
||||
entry.client = result.value.connection
|
||||
entry.tools = result.value.tools.map((tool) => toTool(name, entry, tool))
|
||||
entry.prompts = []
|
||||
entry.status = { status: "connected" }
|
||||
watch(name, entry, result.value.connection)
|
||||
yield* Effect.logInfo("mcp connected", { server: name, tools: entry.tools.length })
|
||||
// The tool registry reads on this event; a late-connecting server has no other way to appear.
|
||||
yield* bus.publish(McpEvent.ToolsChanged, { server: name })
|
||||
yield* bus.publish(McpEvent.ResourcesChanged, { server: name })
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
whenLive(name, entry, result.value.connection)(refreshPrompts(name, entry, result.value.connection))
|
||||
return
|
||||
}
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
entry.scope = undefined
|
||||
const error = Cause.squash(result.cause)
|
||||
entry.status =
|
||||
error instanceof McpClient.NeedsAuthError
|
||||
? { status: "needs_auth", error: error.message }
|
||||
: { status: "failed", error: error instanceof Error ? error.message : String(error) }
|
||||
yield* Effect.logWarning("mcp connect failed", { server: name, status: entry.status })
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
}).pipe(
|
||||
Effect.ensuring(entry.startup.open),
|
||||
Effect.annotateLogs({ server: name, directory: location.directory, connectionID: crypto.randomUUID() }),
|
||||
)
|
||||
|
||||
const stopServer = Effect.fnUntraced(function* (name: ServerName, entry: ServerEntry) {
|
||||
const scope = entry.scope
|
||||
if (!scope) return
|
||||
entry.scope = undefined
|
||||
entry.client = undefined
|
||||
entry.tools = undefined
|
||||
entry.prompts = undefined
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
yield* bus.publish(McpEvent.ToolsChanged, { server: name })
|
||||
yield* bus.publish(McpEvent.ResourcesChanged, { server: name })
|
||||
yield* bus.publish(PromptsChanged, { server: name })
|
||||
})
|
||||
|
||||
const disposeServer = Effect.fnUntraced(function* (name: ServerName, entry: ServerEntry) {
|
||||
yield* stopServer(name, entry)
|
||||
const disposeServer = Effect.fnUntraced(function* (name: ServerName, entry: LocationEntry) {
|
||||
yield* runtime.stop(name, entry)
|
||||
if (entry.integrationID) owned.delete(entry.integrationID)
|
||||
if (entry.registration) yield* entry.registration.dispose
|
||||
})
|
||||
@@ -489,24 +518,14 @@ export const layer = (options?: Options) =>
|
||||
const replaceServer = Effect.fnUntraced(function* (name: ServerName, serverConfig: Mcp.ServerConfig) {
|
||||
const previous = entries.get(name)
|
||||
if (previous) yield* disposeServer(name, previous)
|
||||
const entry: ServerEntry = {
|
||||
const entry: LocationEntry = {
|
||||
owner: LOCATION,
|
||||
config: serverConfig,
|
||||
status: { status: "pending" },
|
||||
startup: Latch.makeUnsafe(),
|
||||
}
|
||||
entries.set(name, entry)
|
||||
yield* Effect.gen(function* () {
|
||||
yield* register(name, entry)
|
||||
if (serverConfig.disabled) {
|
||||
entry.status = { status: "disabled" }
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
return
|
||||
}
|
||||
yield* startServer(name, entry)
|
||||
}).pipe(
|
||||
// Settle startup even when registration fails or replacement is interrupted, so readers cannot hang.
|
||||
Effect.ensuring(entry.startup.open),
|
||||
)
|
||||
yield* runtime.open(name, entry, register(name, entry))
|
||||
})
|
||||
|
||||
const removeServer = Effect.fnUntraced(function* (name: ServerName) {
|
||||
@@ -526,6 +545,7 @@ export const layer = (options?: Options) =>
|
||||
if (!applied && entries.size === 0) {
|
||||
for (const [name, server] of servers) {
|
||||
entries.set(name, {
|
||||
owner: LOCATION,
|
||||
config: server,
|
||||
status: { status: "pending" },
|
||||
startup: Latch.makeUnsafe(),
|
||||
@@ -542,7 +562,7 @@ export const layer = (options?: Options) =>
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
continue
|
||||
}
|
||||
fork(startServer(name, entry).pipe(locks.withLock(name)))
|
||||
fork(runtime.start(name, entry).pipe(locks.withLock(name)))
|
||||
}
|
||||
return
|
||||
}
|
||||
@@ -573,8 +593,8 @@ export const layer = (options?: Options) =>
|
||||
const entry = entries.get(name)
|
||||
if (!entry || entry.integrationID !== integrationID) return
|
||||
if (entry.status.status === "disabled") return
|
||||
yield* stopServer(name, entry)
|
||||
yield* startServer(name, entry)
|
||||
yield* runtime.stop(name, entry)
|
||||
yield* runtime.start(name, entry)
|
||||
}).pipe(locks.withLock(name))
|
||||
})
|
||||
fork(
|
||||
@@ -628,15 +648,15 @@ export const layer = (options?: Options) =>
|
||||
const name = ServerName.make(server)
|
||||
yield* Effect.gen(function* () {
|
||||
const target = yield* requireServer(name)
|
||||
yield* stopServer(name, target.entry)
|
||||
yield* startServer(name, target.entry)
|
||||
yield* runtime.stop(name, target.entry)
|
||||
yield* runtime.start(name, target.entry)
|
||||
}).pipe(locks.withLock(name))
|
||||
}),
|
||||
disconnect: Effect.fn("MCP.disconnect")(function* (server) {
|
||||
const name = ServerName.make(server)
|
||||
yield* Effect.gen(function* () {
|
||||
const target = yield* requireServer(name)
|
||||
yield* stopServer(name, target.entry)
|
||||
yield* runtime.stop(name, target.entry)
|
||||
target.entry.status = { status: "disabled" }
|
||||
yield* bus.publish(McpEvent.StatusChanged, { server: name })
|
||||
}).pipe(locks.withLock(name))
|
||||
@@ -655,21 +675,7 @@ export const layer = (options?: Options) =>
|
||||
}),
|
||||
callTool: Effect.fn("MCP.callTool")(function* (input) {
|
||||
const target = yield* requireServer(input.server)
|
||||
yield* target.entry.startup.await
|
||||
if (!target.entry.client)
|
||||
return yield* new ToolCallError({
|
||||
server: target.name,
|
||||
tool: input.name,
|
||||
message: "MCP server is not connected",
|
||||
})
|
||||
const result = yield* recovering(target.name, target.entry, target.entry.client, (connection) =>
|
||||
connection.callTool({ name: input.name, args: input.args, sessionID: input.sessionID }),
|
||||
).pipe(
|
||||
Effect.mapError(
|
||||
(error) => new ToolCallError({ server: target.name, tool: input.name, message: error.message }),
|
||||
),
|
||||
)
|
||||
return { ...result, server: target.name, tool: input.name }
|
||||
return yield* runtime.callTool(target.name, target.entry, input)
|
||||
}),
|
||||
instructions: Effect.fn("MCP.instructions")(function* () {
|
||||
return Array.from(entries)
|
||||
@@ -687,55 +693,28 @@ export const layer = (options?: Options) =>
|
||||
}),
|
||||
prompt: Effect.fn("MCP.prompt")(function* (input) {
|
||||
const target = yield* requireServer(input.server)
|
||||
yield* target.entry.startup.await
|
||||
if (!target.entry.client) return undefined
|
||||
const result = yield* recovering(target.name, target.entry, target.entry.client, (connection) =>
|
||||
connection.prompt({ name: input.name, args: input.args }),
|
||||
).pipe(Effect.orElseSucceed(() => undefined))
|
||||
if (!result) return undefined
|
||||
return { ...result, server: target.name, name: input.name }
|
||||
return yield* runtime.prompt(target.name, target.entry, input)
|
||||
}),
|
||||
resourceCatalog: Effect.fn("MCP.resourceCatalog")(function* () {
|
||||
const empty = ResourceCatalog.make({ resources: [], templates: [] })
|
||||
const catalogs = yield* Effect.forEach(
|
||||
Array.from(entries),
|
||||
([name, entry]) =>
|
||||
entry.client
|
||||
? loadCatalog(name, entry, entry.client).pipe(Effect.orElseSucceed(() => empty))
|
||||
: Effect.succeed(empty),
|
||||
{ concurrency: "unbounded" },
|
||||
return mergeCatalogs(
|
||||
yield* Effect.forEach(entries, ([name, entry]) => runtime.catalog(name, entry), {
|
||||
concurrency: "unbounded",
|
||||
}),
|
||||
)
|
||||
return mergeCatalogs(catalogs)
|
||||
}),
|
||||
resources: Effect.fn("MCP.resources")(function* (input) {
|
||||
const target = yield* requireServer(input.server)
|
||||
yield* target.entry.startup.await
|
||||
if (!target.entry.client) return ResourceCatalog.make({ resources: [], templates: [] })
|
||||
return mergeCatalogs([yield* loadCatalog(target.name, target.entry, target.entry.client)])
|
||||
return yield* runtime.resources(target.name, target.entry)
|
||||
}),
|
||||
readResource: Effect.fn("MCP.readResource")(function* (input) {
|
||||
const target = yield* requireServer(input.server)
|
||||
yield* target.entry.startup.await
|
||||
if (!target.entry.client) return undefined
|
||||
const result = yield* recovering(target.name, target.entry, target.entry.client, (connection) =>
|
||||
connection.readResource({ uri: input.uri }),
|
||||
)
|
||||
if (!result) return undefined
|
||||
return ResourceContent.make({
|
||||
server: target.name,
|
||||
uri: input.uri,
|
||||
contents: result.contents.map((part) =>
|
||||
"text" in part
|
||||
? { type: "text", uri: part.uri, text: part.text, mimeType: part.mimeType }
|
||||
: { type: "blob", uri: part.uri, blob: part.blob, mimeType: part.mimeType },
|
||||
),
|
||||
})
|
||||
return yield* runtime.readResource(target.name, target.entry, input.uri)
|
||||
}),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
function mergeCatalogs(catalogs: ReadonlyArray<ResourceCatalog>) {
|
||||
export function mergeCatalogs(catalogs: ReadonlyArray<ResourceCatalog>) {
|
||||
return ResourceCatalog.make({
|
||||
resources: catalogs
|
||||
.flatMap((catalog) => catalog.resources)
|
||||
@@ -762,67 +741,3 @@ export function configured(options?: Options) {
|
||||
}
|
||||
|
||||
export const node = configured()
|
||||
|
||||
// Schema `optional` strips undefined-valued properties on encode, so fields can assign
|
||||
// optional properties directly instead of conditionally spreading them.
|
||||
function toElicitationField(key: string, property: ElicitationProperty, required: boolean): Form.Field {
|
||||
// Some servers emit machine titles like "string with format email"; prefer description/key over those.
|
||||
const machineTitle = /^(boolean|string|number|integer|array|object)(\s+with\b.*|\s+in\b.*)?$/i
|
||||
const title =
|
||||
property.title && !machineTitle.test(property.title.trim()) ? property.title : (property.description ?? key)
|
||||
const base = {
|
||||
key,
|
||||
title,
|
||||
description: property.description === title ? undefined : property.description,
|
||||
required: required || undefined,
|
||||
}
|
||||
switch (property.type) {
|
||||
case "boolean":
|
||||
return { ...base, type: "boolean", default: property.default }
|
||||
case "number":
|
||||
case "integer":
|
||||
return {
|
||||
...base,
|
||||
type: property.type,
|
||||
minimum: property.minimum,
|
||||
maximum: property.maximum,
|
||||
default: property.default,
|
||||
}
|
||||
case "array":
|
||||
return {
|
||||
...base,
|
||||
type: "multiselect",
|
||||
options:
|
||||
"anyOf" in property.items
|
||||
? property.items.anyOf.map((option) => ({ value: option.const, label: option.title }))
|
||||
: property.items.enum.map((value) => ({ value, label: value })),
|
||||
custom: false,
|
||||
minItems: property.minItems,
|
||||
maxItems: property.maxItems,
|
||||
default: property.default,
|
||||
}
|
||||
case "string": {
|
||||
const options =
|
||||
"oneOf" in property
|
||||
? property.oneOf.map((option) => ({ value: option.const, label: option.title }))
|
||||
: "enum" in property
|
||||
? property.enum.map((value, index) => ({
|
||||
value,
|
||||
label: ("enumNames" in property ? property.enumNames?.[index] : undefined) ?? value,
|
||||
}))
|
||||
: undefined
|
||||
return {
|
||||
...base,
|
||||
type: "string",
|
||||
format: "format" in property ? property.format : undefined,
|
||||
minLength: "minLength" in property ? property.minLength : undefined,
|
||||
maxLength: "maxLength" in property ? property.maxLength : undefined,
|
||||
default: property.default,
|
||||
options,
|
||||
custom: options ? false : undefined,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type ElicitationProperty = McpClient.ElicitationFormParams["requestedSchema"]["properties"][string]
|
||||
|
||||
@@ -4,7 +4,7 @@ import { makeLocationNode } from "@opencode/util/effect/app-node"
|
||||
import { Context, Effect, Layer, Schema } from "effect"
|
||||
import { Permission } from "../permission.js"
|
||||
import { McpTool } from "../tool/mcp.js"
|
||||
import { Mcp } from "./index.js"
|
||||
import type { McpSession } from "./session.js"
|
||||
import { Instructions } from "../instructions/index.js"
|
||||
|
||||
const Summary = Schema.Struct({
|
||||
@@ -55,57 +55,53 @@ const update = (previous: ReadonlyArray<Summary>, current: ReadonlyArray<Summary
|
||||
|
||||
export interface Interface {
|
||||
/** Lists server instructions reachable under the given ruleset; callers pass the merged agent and Session permissions. */
|
||||
readonly load: (permissions: Permission.Ruleset) => Effect.Effect<Instructions.List>
|
||||
readonly load: (
|
||||
permissions: Permission.Ruleset,
|
||||
view: Pick<McpSession.View, "instructions" | "tools" | "owned">,
|
||||
) => Effect.Effect<Instructions.List>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/McpInstructions") {}
|
||||
|
||||
export const layer = Layer.effect(
|
||||
export const layer = Layer.succeed(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const mcp = yield* Mcp.Service
|
||||
|
||||
return Service.of({
|
||||
load: Effect.fn("McpInstructions.load")(function* (permissions) {
|
||||
const source = (value: ReadonlyArray<Summary> | Instructions.Removed) =>
|
||||
Instructions.make<ReadonlyArray<Summary>>({
|
||||
key: Instructions.Key.make("core/mcp-guidance"),
|
||||
codec: Schema.toCodecJson(Schema.Array(Summary)),
|
||||
read: Effect.succeed(value),
|
||||
render: {
|
||||
initial: render,
|
||||
changed: update,
|
||||
removed: () => "MCP server instructions are no longer available.",
|
||||
},
|
||||
})
|
||||
const [instructions, tools] = yield* Effect.all([mcp.instructions(), mcp.tools()], {
|
||||
concurrency: "unbounded",
|
||||
Service.of({
|
||||
load: Effect.fn("McpInstructions.load")(function* (permissions, view) {
|
||||
const source = (value: ReadonlyArray<Summary> | Instructions.Removed) =>
|
||||
Instructions.make<ReadonlyArray<Summary>>({
|
||||
key: Instructions.Key.make("core/mcp-guidance"),
|
||||
codec: Schema.toCodecJson(Schema.Array(Summary)),
|
||||
read: Effect.succeed(value),
|
||||
render: {
|
||||
initial: render,
|
||||
changed: update,
|
||||
removed: () => "MCP server instructions are no longer available.",
|
||||
},
|
||||
})
|
||||
const canExecute = Permission.evaluate("execute", "*", permissions).effect !== "deny"
|
||||
// Instructions are useful only when this Session can reach at least one server tool.
|
||||
const visible = instructions
|
||||
.flatMap((item) => {
|
||||
const owned = tools.filter((tool) => tool.server === item.server)
|
||||
const codemode = owned[0]?.codemode !== false
|
||||
if (codemode && !canExecute) return []
|
||||
if (
|
||||
!owned.some(
|
||||
(tool) =>
|
||||
Permission.evaluate(McpTool.name(tool.server, tool.name), "*", permissions).effect !== "deny",
|
||||
)
|
||||
const tools = [...view.tools, ...view.owned.map((owned) => owned.tool)]
|
||||
const canExecute = Permission.evaluate("execute", "*", permissions).effect !== "deny"
|
||||
// Instructions are useful only when this Session can reach at least one server tool.
|
||||
const visible = view.instructions
|
||||
.flatMap((item) => {
|
||||
const owned = tools.filter((tool) => tool.server === item.server)
|
||||
const codemode = owned[0]?.codemode !== false
|
||||
if (codemode && !canExecute) return []
|
||||
if (
|
||||
!owned.some(
|
||||
(tool) => Permission.evaluate(McpTool.name(tool.server, tool.name), "*", permissions).effect !== "deny",
|
||||
)
|
||||
return []
|
||||
return [
|
||||
codemode
|
||||
? { server: item.server, instructions: item.instructions }
|
||||
: { server: item.server, instructions: item.instructions, codemode: false as const },
|
||||
]
|
||||
})
|
||||
.toSorted((a, b) => a.server.localeCompare(b.server))
|
||||
return source(visible.length === 0 ? Instructions.removed : visible)
|
||||
}),
|
||||
})
|
||||
)
|
||||
return []
|
||||
return [
|
||||
codemode
|
||||
? { server: item.server, instructions: item.instructions }
|
||||
: { server: item.server, instructions: item.instructions, codemode: false as const },
|
||||
]
|
||||
})
|
||||
.toSorted((a, b) => a.server.localeCompare(b.server))
|
||||
return source(visible.length === 0 ? Instructions.removed : visible)
|
||||
}),
|
||||
}),
|
||||
)
|
||||
|
||||
export const node = makeLocationNode({ service: Service, layer, deps: [Mcp.node] })
|
||||
export const node = makeLocationNode({ service: Service, layer, deps: [] })
|
||||
|
||||
@@ -0,0 +1,229 @@
|
||||
export * as McpSession from "./session.js"
|
||||
|
||||
import type { Session } from "@opencode/schema/session"
|
||||
import { SessionEvent } from "@opencode/schema/session-event"
|
||||
import { isDeepStrictEqual } from "node:util"
|
||||
import { Context, Effect, Latch, Layer, Stream } from "effect"
|
||||
import { makeLocationNode } from "@opencode/util/effect/app-node"
|
||||
import { Bus } from "../bus.js"
|
||||
import { Credential } from "../credential.js"
|
||||
import { KeyedMutex } from "../effect/keyed-mutex.js"
|
||||
import { Environment } from "../environment/index.js"
|
||||
import { Form } from "../form.js"
|
||||
import { Location } from "../location.js"
|
||||
import { SessionStore } from "../session/store.js"
|
||||
import { Mcp } from "./index.js"
|
||||
|
||||
type Owner = Extract<Mcp.Owner, { readonly type: "session" }>
|
||||
type Entry = Mcp.Entry<Owner>
|
||||
type Registry = { readonly servers: Map<Mcp.ServerName, Entry>; readonly locks: KeyedMutex.KeyedMutex<Mcp.ServerName> }
|
||||
|
||||
export interface OwnedTool {
|
||||
readonly tool: Mcp.Tool
|
||||
readonly call: (input: {
|
||||
readonly args?: Record<string, unknown>
|
||||
readonly sessionID: Session.ID
|
||||
}) => Effect.Effect<Mcp.ToolResult, Mcp.ToolCallError>
|
||||
}
|
||||
|
||||
/**
|
||||
* MCP capabilities visible to one Session: Location servers plus the servers registered for the Session
|
||||
* or one of its ancestors, nearest owner first. A visible Session-scoped server shadows the Location
|
||||
* server of the same name.
|
||||
*/
|
||||
export interface View {
|
||||
readonly shadowed: ReadonlySet<string>
|
||||
readonly servers: ReadonlyArray<Mcp.ServerInfo>
|
||||
/** Tools of the visible Location servers. */
|
||||
readonly tools: ReadonlyArray<Mcp.Tool>
|
||||
readonly owned: ReadonlyArray<OwnedTool>
|
||||
readonly instructions: ReadonlyArray<Mcp.ServerInstructions>
|
||||
readonly resourceCatalog: Effect.Effect<Mcp.ResourceCatalog>
|
||||
readonly resources: (server: string) => Effect.Effect<Mcp.ResourceCatalog, Error>
|
||||
readonly readResource: (input: {
|
||||
readonly server: string
|
||||
readonly uri: string
|
||||
}) => Effect.Effect<Mcp.ResourceContent | undefined, Error>
|
||||
}
|
||||
|
||||
/** Process-local MCP servers owned by a Session, released on removal or when the Session is deleted or moved away. */
|
||||
export interface Interface {
|
||||
/** Re-adding an identical config is a no-op; a different config replaces only this Session's server. */
|
||||
readonly add: (sessionID: Session.ID, server: string, config: Mcp.ServerConfig) => Effect.Effect<void>
|
||||
readonly remove: (sessionID: Session.ID, server: string) => Effect.Effect<void, Mcp.NotFoundError>
|
||||
readonly view: (sessionID: Session.ID) => Effect.Effect<View>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/McpSession") {}
|
||||
|
||||
export const layer = (options?: Mcp.Options) =>
|
||||
Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const mcp = yield* Mcp.Service
|
||||
const store = yield* SessionStore.Service
|
||||
const bus = yield* Bus.Service
|
||||
const location = yield* Location.Service
|
||||
const registries = new Map<Session.ID, Registry>()
|
||||
|
||||
const withRegistry = <A, E, R>(
|
||||
sessionID: Session.ID,
|
||||
name: Mcp.ServerName,
|
||||
effect: (servers: Map<Mcp.ServerName, Entry>) => Effect.Effect<A, E, R>,
|
||||
) =>
|
||||
Effect.suspend(() => {
|
||||
const registry = registries.get(sessionID) ?? { servers: new Map(), locks: KeyedMutex.makeUnsafe() }
|
||||
registries.set(sessionID, registry)
|
||||
return registry.locks
|
||||
.withLock(name)(effect(registry.servers))
|
||||
.pipe(
|
||||
// Forget the Session once its last server is gone and no operation still waits on its locks.
|
||||
Effect.ensuring(
|
||||
Effect.gen(function* () {
|
||||
if (registry.servers.size > 0 || (yield* registry.locks.size) > 0) return
|
||||
if (registries.get(sessionID) === registry) registries.delete(sessionID)
|
||||
}),
|
||||
),
|
||||
)
|
||||
})
|
||||
|
||||
const runtime = yield* Mcp.makeRuntime<Owner>({
|
||||
options,
|
||||
lock: (name, owner) => (effect) => withRegistry(owner.id, name, () => effect),
|
||||
})
|
||||
|
||||
const remove = (sessionID: Session.ID, server: string) => {
|
||||
const name = Mcp.ServerName.make(server)
|
||||
return withRegistry(sessionID, name, (servers) =>
|
||||
Effect.gen(function* () {
|
||||
const entry = servers.get(name)
|
||||
if (!entry) return yield* new Mcp.NotFoundError({ server: name })
|
||||
yield* runtime.stop(name, entry)
|
||||
servers.delete(name)
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
const visible = Effect.fnUntraced(function* (sessionID: Session.ID) {
|
||||
const scoped = new Map<Mcp.ServerName, Entry>()
|
||||
if (registries.size === 0) return scoped
|
||||
let current: Session.ID | undefined = sessionID
|
||||
while (current) {
|
||||
for (const [name, entry] of registries.get(current)?.servers ?? [])
|
||||
if (!scoped.has(name)) scoped.set(name, entry)
|
||||
current = (yield* store.get(current))?.parentID
|
||||
}
|
||||
return scoped
|
||||
})
|
||||
|
||||
yield* bus.subscribe([SessionEvent.Deleted, SessionEvent.Moved]).pipe(
|
||||
Stream.filter(
|
||||
(event) =>
|
||||
event.type === "session.deleted" ||
|
||||
event.data.location.directory !== location.directory ||
|
||||
event.data.location.workspaceID !== location.workspaceID,
|
||||
),
|
||||
Stream.runForEach((event) =>
|
||||
Effect.forEach(
|
||||
Array.from(registries.get(event.data.sessionID)?.servers.keys() ?? []),
|
||||
(name) => remove(event.data.sessionID, name).pipe(Effect.ignore),
|
||||
{ concurrency: "unbounded", discard: true },
|
||||
),
|
||||
),
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
|
||||
return Service.of({
|
||||
add: Effect.fn("McpSession.add")(function* (sessionID, server, serverConfig) {
|
||||
const name = Mcp.ServerName.make(server)
|
||||
const config = Mcp.cloneConfig(serverConfig)
|
||||
yield* withRegistry(sessionID, name, (servers) =>
|
||||
Effect.gen(function* () {
|
||||
const previous = servers.get(name)
|
||||
if (previous && isDeepStrictEqual(previous.config, config)) return
|
||||
if (previous) yield* runtime.stop(name, previous)
|
||||
const entry: Entry = {
|
||||
owner: { type: "session", id: sessionID },
|
||||
config,
|
||||
status: { status: "pending" },
|
||||
startup: Latch.makeUnsafe(),
|
||||
}
|
||||
servers.set(name, entry)
|
||||
// Session-scoped servers are not registered as OAuth integrations: integration IDs derive from
|
||||
// name and URL, so owners sharing one would share and prematurely dispose its registration.
|
||||
yield* runtime.open(name, entry)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
remove: Effect.fn("McpSession.remove")(remove),
|
||||
view: Effect.fn("McpSession.view")(function* (sessionID) {
|
||||
const scoped = yield* visible(sessionID)
|
||||
const local = yield* Effect.all({
|
||||
servers: mcp.servers(),
|
||||
tools: mcp.tools(),
|
||||
instructions: mcp.instructions(),
|
||||
})
|
||||
const shadowed = new Set<string>(scoped.keys())
|
||||
const owned = Array.from(scoped)
|
||||
return {
|
||||
shadowed,
|
||||
servers: [
|
||||
...local.servers.filter((server) => !shadowed.has(server.name)),
|
||||
...owned.map(([name, entry]): Mcp.ServerInfo => ({ name, status: entry.status })),
|
||||
].toSorted((a, b) => a.name.localeCompare(b.name)),
|
||||
tools: local.tools.filter((tool) => !shadowed.has(tool.server)),
|
||||
owned: owned
|
||||
.flatMap(([name, entry]) =>
|
||||
(entry.tools ?? []).map(
|
||||
(tool): OwnedTool => ({
|
||||
tool,
|
||||
call: (input) => runtime.callTool(name, entry, { ...input, name: tool.name }),
|
||||
}),
|
||||
),
|
||||
)
|
||||
.toSorted((a, b) => a.tool.server.localeCompare(b.tool.server) || a.tool.name.localeCompare(b.tool.name)),
|
||||
instructions: [
|
||||
...local.instructions.filter((item) => !shadowed.has(item.server)),
|
||||
...owned.flatMap(([server, entry]) =>
|
||||
entry.client?.instructions ? [{ server, instructions: entry.client.instructions }] : [],
|
||||
),
|
||||
].toSorted((a, b) => a.server.localeCompare(b.server)),
|
||||
resourceCatalog: Effect.all(
|
||||
[mcp.resourceCatalog(), ...owned.map(([name, entry]) => runtime.catalog(name, entry))],
|
||||
{ concurrency: "unbounded" },
|
||||
).pipe(
|
||||
Effect.map(([catalog, ...catalogs]) =>
|
||||
Mcp.mergeCatalogs([
|
||||
{
|
||||
resources: catalog.resources.filter((resource) => !shadowed.has(resource.server)),
|
||||
templates: catalog.templates.filter((template) => !shadowed.has(template.server)),
|
||||
},
|
||||
...catalogs,
|
||||
]),
|
||||
),
|
||||
),
|
||||
resources: (server) => {
|
||||
const name = Mcp.ServerName.make(server)
|
||||
const entry = scoped.get(name)
|
||||
return entry ? runtime.resources(name, entry) : mcp.resources({ server })
|
||||
},
|
||||
readResource: (input) => {
|
||||
const name = Mcp.ServerName.make(input.server)
|
||||
const entry = scoped.get(name)
|
||||
return entry ? runtime.readResource(name, entry, input.uri) : mcp.readResource(input)
|
||||
},
|
||||
}
|
||||
}),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
export function configured(options?: Mcp.Options) {
|
||||
return makeLocationNode({
|
||||
service: Service,
|
||||
layer: layer(options),
|
||||
deps: [Mcp.node, SessionStore.node, Location.node, Environment.node, Bus.node, Form.node, Credential.node],
|
||||
})
|
||||
}
|
||||
|
||||
export const node = configured()
|
||||
@@ -53,6 +53,7 @@ import { Location } from "../location.js"
|
||||
import { ManagedPolicy } from "../managed-policy.js"
|
||||
import { ModelsDev } from "../models-dev.js"
|
||||
import { Mcp } from "../mcp/index.js"
|
||||
import { McpSession } from "../mcp/session.js"
|
||||
import { Npm } from "@opencode/util/npm"
|
||||
import { Permission } from "../permission.js"
|
||||
import { Reference } from "../reference.js"
|
||||
@@ -132,6 +133,7 @@ const services = [
|
||||
ManagedPolicy.Service,
|
||||
ModelsDev.Service,
|
||||
Mcp.Service,
|
||||
McpSession.Service,
|
||||
Npm.Service,
|
||||
Permission.Service,
|
||||
Form.Service,
|
||||
@@ -185,6 +187,7 @@ export const requirements = LayerNode.group([
|
||||
ManagedPolicy.node,
|
||||
ModelsDev.node,
|
||||
Mcp.node,
|
||||
McpSession.node,
|
||||
Npm.node,
|
||||
Permission.node,
|
||||
Form.node,
|
||||
|
||||
@@ -12,6 +12,7 @@ import { Instructions } from "../instructions/index.js"
|
||||
import { InstructionBuiltIns } from "../instructions/builtins.js"
|
||||
import { Location } from "../location.js"
|
||||
import { McpInstructions } from "../mcp/instructions.js"
|
||||
import { McpSession } from "../mcp/session.js"
|
||||
import { McpTool } from "../tool/mcp.js"
|
||||
import { ReferenceInstructions } from "../reference/instructions.js"
|
||||
import { SkillInstructions } from "../skill/instructions.js"
|
||||
@@ -82,6 +83,7 @@ const layer = Layer.effect(
|
||||
const entries = yield* InstructionEntry.Service
|
||||
const location = yield* Location.Service
|
||||
const mcpInstructions = yield* McpInstructions.Service
|
||||
const mcpSessions = yield* McpSession.Service
|
||||
const mcpTools = yield* McpTool.Service
|
||||
const models = yield* SessionRunnerModel.Service
|
||||
const request = yield* SessionModelRequest.Service
|
||||
@@ -129,14 +131,15 @@ const layer = Layer.effect(
|
||||
if (!agent.info) return yield* new AgentNotFoundError({ sessionID: session.id, agent: session.agent ?? agent.id })
|
||||
// Session permissions narrow discovery the same way they narrow the tool snapshot.
|
||||
const permissions = Permission.merge(agent.info.permissions, session.permissions ?? [])
|
||||
const mcp = yield* mcpSessions.view(sessionID)
|
||||
const loaded = yield* Effect.all(
|
||||
{
|
||||
tools: registry.snapshot(permissions),
|
||||
tools: mcpTools.overlay(mcp).pipe(Effect.flatMap((overlay) => registry.snapshot(permissions, overlay))),
|
||||
builtins: builtins.load(),
|
||||
discovery: discovery.load(),
|
||||
skills: skillInstructions.load(permissions),
|
||||
references: referenceInstructions.load(),
|
||||
mcp: mcpInstructions.load(permissions),
|
||||
mcp: mcpInstructions.load(permissions, mcp),
|
||||
entries: entries.load(sessionID),
|
||||
},
|
||||
{ concurrency: "unbounded" },
|
||||
@@ -197,6 +200,7 @@ export const node = makeLocationNode({
|
||||
InstructionEntry.node,
|
||||
Location.node,
|
||||
McpInstructions.node,
|
||||
McpSession.node,
|
||||
McpTool.node,
|
||||
ReferenceInstructions.node,
|
||||
SessionRunnerModel.node,
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
export * as Tool from "./tool.js"
|
||||
export { CallID, Content, Error, FileContent, TextContent } from "@opencode/schema/tool"
|
||||
export type { Context, Metadata, Namespace, Options, Result } from "@opencode/schema/tool"
|
||||
export type { Context, Info, Metadata, Namespace, Options, Result } from "@opencode/schema/tool"
|
||||
|
||||
import { ToolDefinition, type ToolCall } from "@opencode/ai"
|
||||
import { Tool } from "@opencode/schema/tool"
|
||||
@@ -38,9 +38,16 @@ type Data = {
|
||||
errors: { kind: "tool" | "namespace"; name: string; namespace?: string; error: RegistrationError }[]
|
||||
}
|
||||
|
||||
/** Session-specific tools layered over the Location registry for one snapshot. */
|
||||
export interface Overlay {
|
||||
readonly tools: ReadonlyArray<Tool.Info>
|
||||
/** Effective names of registered tools the overlay hides. */
|
||||
readonly hidden: ReadonlySet<string>
|
||||
}
|
||||
|
||||
export interface Interface extends State.Transformable<Editor> {
|
||||
readonly list: () => Effect.Effect<ReadonlyArray<Tool.Info & { readonly id: string }>>
|
||||
readonly snapshot: (permissions?: Permission.Ruleset) => Effect.Effect<Snapshot>
|
||||
readonly snapshot: (permissions?: Permission.Ruleset, overlay?: Overlay) => Effect.Effect<Snapshot>
|
||||
}
|
||||
|
||||
/** A local execution result after hooks and content normalization. */
|
||||
@@ -155,7 +162,9 @@ const layer = Layer.effect(
|
||||
}
|
||||
})
|
||||
|
||||
let catalog: { data: Data; names: string; value: CodeModeCatalog.Inventory } | undefined
|
||||
let catalog:
|
||||
| { data: Data; names: string; overlay: ReadonlyArray<Tool.Info>; value: CodeModeCatalog.Inventory }
|
||||
| undefined
|
||||
const state = State.create<Data, Editor>({
|
||||
name: "tool",
|
||||
initial: () => ({
|
||||
@@ -222,12 +231,19 @@ const layer = Layer.effect(
|
||||
transform: state.transform,
|
||||
reload: state.reload,
|
||||
list: () => Effect.sync(() => Array.from(state.get().tools.values())),
|
||||
snapshot: Effect.fn("Tool.snapshot")((permissions) =>
|
||||
snapshot: Effect.fn("Tool.snapshot")((permissions, overlay) =>
|
||||
Effect.sync(() => {
|
||||
const data = state.get()
|
||||
const active = new Map<string, Tool.Info>()
|
||||
const rules = permissions ?? []
|
||||
for (const [name, tool] of data.tools) {
|
||||
const overlaid = overlay?.tools.filter((tool) => !registrationError(tool)) ?? []
|
||||
const registered = overlay
|
||||
? [
|
||||
...Array.from(data.tools).filter(([name]) => !overlay.hidden.has(name)),
|
||||
...overlaid.map((tool) => [effectiveName(tool), tool] as const),
|
||||
]
|
||||
: data.tools
|
||||
for (const [name, tool] of registered) {
|
||||
if (whollyDisabled(tool.options?.permission ?? name, rules)) continue
|
||||
active.set(name, tool)
|
||||
}
|
||||
@@ -248,10 +264,13 @@ const layer = Layer.effect(
|
||||
// definitions/executors fresh, but share the much larger rendered catalog across steps.
|
||||
const codeModeCatalog = !codeModeEnabled
|
||||
? undefined
|
||||
: catalog?.data === data && catalog.names === names
|
||||
: catalog?.data === data &&
|
||||
catalog.names === names &&
|
||||
catalog.overlay.length === overlaid.length &&
|
||||
catalog.overlay.every((tool, index) => tool === overlaid[index])
|
||||
? catalog.value
|
||||
: CodeModeTool.catalog(codeModeInventory)
|
||||
if (codeModeCatalog) catalog = { data, names, value: codeModeCatalog }
|
||||
if (codeModeCatalog) catalog = { data, names, overlay: overlaid, value: codeModeCatalog }
|
||||
return {
|
||||
...(codeModeCatalog === undefined ? {} : { codeModeCatalog }),
|
||||
definitions: [
|
||||
|
||||
@@ -43,6 +43,8 @@ The service uses shared `State` to replay synchronous transforms in registration
|
||||
|
||||
MCP owns one stable tool transform that reads its latest discovered tools. Tool-list changes update that source and reload the tool state instead of re-registering at the end of the transform order. MCP refresh therefore preserves the precedence of later plugin overrides.
|
||||
|
||||
Session-scoped MCP servers never enter the registry. `McpTool.overlay(view)` turns a Session's `McpSession.View` into its owned tools plus the effective names of the Location MCP tools whose servers the view shadows, and `Tool.Service.snapshot(permissions, overlay)` layers them over the registry for that snapshot only. Registry transforms, including plugin overrides, do not apply to overlay tools.
|
||||
|
||||
Type safety ends at registration. The registry validates model input and declared output at runtime and should not carry producer schema generics through storage or execution.
|
||||
|
||||
`Tool.Service` is Location-scoped. Do not make the registry process-global or construct a separate application-tool service for each Location.
|
||||
@@ -61,4 +63,4 @@ Producer capture remains local to producers. Shell stores combined process outpu
|
||||
|
||||
## Current Gaps
|
||||
|
||||
- Future Session-scoped registrations still need an explicit canonical registration design.
|
||||
- Session-scoped registrations beyond MCP servers still need an explicit canonical registration design.
|
||||
|
||||
+107
-79
@@ -2,6 +2,8 @@ export * as McpTool from "./mcp.js"
|
||||
|
||||
import { ToolFailure } from "@opencode/ai"
|
||||
import { McpEvent } from "@opencode/schema/mcp-event"
|
||||
import type { Session } from "@opencode/schema/session"
|
||||
import type { McpSession } from "../mcp/session.js"
|
||||
import { Context, Effect, Fiber, type JsonSchema, Layer, PubSub, Semaphore, Stream } from "effect"
|
||||
import { makeLocationNode } from "@opencode/util/effect/app-node"
|
||||
import { Bus } from "../bus.js"
|
||||
@@ -19,6 +21,8 @@ export const name = (server: string, tool: string) => `${namespace(server)}_${to
|
||||
export interface Interface {
|
||||
/** Wait for the initial MCP tool registration to settle. */
|
||||
readonly flush: Effect.Effect<void>
|
||||
/** Session-owned server tools from a Session's MCP view, hiding the Location server tools they shadow. */
|
||||
readonly overlay: (view: Pick<McpSession.View, "shadowed" | "owned">) => Effect.Effect<Tool.Overlay | undefined>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/McpTool") {}
|
||||
@@ -33,90 +37,102 @@ export const layer = Layer.effect(
|
||||
const lock = Semaphore.makeUnsafe(1)
|
||||
let discovered: Mcp.Tool[] = []
|
||||
|
||||
// Keyed by the MCP tool object so unchanged tools keep their identity across snapshots. A tool object
|
||||
// belongs to exactly one server entry, so the first call bound to it stays valid.
|
||||
const infos = new WeakMap<Mcp.Tool, Tool.Info>()
|
||||
const info = (
|
||||
tool: Mcp.Tool,
|
||||
call: (input: {
|
||||
readonly args: Record<string, unknown>
|
||||
readonly sessionID: Session.ID
|
||||
}) => Effect.Effect<Mcp.ToolResult, Mcp.NotFoundError | Mcp.ToolCallError>,
|
||||
) => {
|
||||
const cached = infos.get(tool)
|
||||
if (cached) return cached
|
||||
const created: Tool.Info = {
|
||||
name: tool.name,
|
||||
options: { namespace: namespace(tool.server), codemode: tool.codemode !== false },
|
||||
description: tool.description ?? "",
|
||||
input: (tool.inputSchema ?? { type: "object", properties: {} }) as JsonSchema.JsonSchema,
|
||||
output: (tool.outputSchema ?? {}) as JsonSchema.JsonSchema,
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
yield* permission.assert({
|
||||
action: name(tool.server, tool.name),
|
||||
resources: ["*"],
|
||||
save: ["*"],
|
||||
metadata: {},
|
||||
sessionID: context.sessionID,
|
||||
agent: context.agent,
|
||||
source: {
|
||||
type: "tool",
|
||||
messageID: context.messageID,
|
||||
id: context.id,
|
||||
},
|
||||
})
|
||||
const result = yield* call({
|
||||
args: (input ?? {}) as Record<string, unknown>,
|
||||
sessionID: context.sessionID,
|
||||
}).pipe(
|
||||
Effect.catchTags({
|
||||
"MCP.NotFoundError": (error) =>
|
||||
new ToolFailure({ message: `MCP server "${error.server}" is not available` }),
|
||||
"MCP.ToolCallError": (error) => new ToolFailure({ message: error.message }),
|
||||
}),
|
||||
)
|
||||
if (result.isError)
|
||||
return yield* new ToolFailure({
|
||||
message:
|
||||
result.content
|
||||
.flatMap((part) => (part.type === "text" ? [part.text] : []))
|
||||
.join("\n")
|
||||
.trim() || "MCP tool returned an error",
|
||||
})
|
||||
const content = result.content.map((part) =>
|
||||
part.type === "text"
|
||||
? { type: "text" as const, text: part.text }
|
||||
: {
|
||||
type: "file" as const,
|
||||
uri: `data:${part.mimeType};base64,${part.data}`,
|
||||
mime: part.mimeType,
|
||||
},
|
||||
)
|
||||
const text = content.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n")
|
||||
const output = () => {
|
||||
if (result.structured !== undefined) return result.structured
|
||||
if (text === "") return null
|
||||
// Agents assume JSON returned as text is already an object, so parse it when the server declares no schema.
|
||||
if (tool.outputSchema === undefined && (text.startsWith("{") || text.startsWith("["))) {
|
||||
try {
|
||||
return JSON.parse(text)
|
||||
} catch {}
|
||||
}
|
||||
return text
|
||||
}
|
||||
return {
|
||||
output: output(),
|
||||
...(content.length === 0 ? {} : { content }),
|
||||
}
|
||||
}).pipe(
|
||||
Effect.mapError((error) =>
|
||||
error instanceof ToolFailure
|
||||
? error
|
||||
: new ToolFailure({ message: `Unable to execute ${name(tool.server, tool.name)}` }),
|
||||
),
|
||||
),
|
||||
}
|
||||
infos.set(tool, created)
|
||||
return created
|
||||
}
|
||||
|
||||
// Register once after initial discovery; only subsequent updates need a debounced reload.
|
||||
const initial = yield* lock
|
||||
.withPermit(
|
||||
Effect.gen(function* () {
|
||||
discovered = yield* mcp.tools()
|
||||
yield* tools.transform((editor) => {
|
||||
for (const tool of discovered) {
|
||||
editor.add({
|
||||
name: tool.name,
|
||||
options: { namespace: namespace(tool.server), codemode: tool.codemode !== false },
|
||||
description: tool.description ?? "",
|
||||
input: (tool.inputSchema ?? { type: "object", properties: {} }) as JsonSchema.JsonSchema,
|
||||
output: (tool.outputSchema ?? {}) as JsonSchema.JsonSchema,
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
yield* permission.assert({
|
||||
action: name(tool.server, tool.name),
|
||||
resources: ["*"],
|
||||
save: ["*"],
|
||||
metadata: {},
|
||||
sessionID: context.sessionID,
|
||||
agent: context.agent,
|
||||
source: {
|
||||
type: "tool",
|
||||
messageID: context.messageID,
|
||||
id: context.id,
|
||||
},
|
||||
})
|
||||
const result = yield* mcp
|
||||
.callTool({
|
||||
server: tool.server,
|
||||
name: tool.name,
|
||||
args: (input ?? {}) as Record<string, unknown>,
|
||||
sessionID: context.sessionID,
|
||||
})
|
||||
.pipe(
|
||||
Effect.catchTags({
|
||||
"MCP.NotFoundError": (error) =>
|
||||
new ToolFailure({ message: `MCP server "${error.server}" is not available` }),
|
||||
"MCP.ToolCallError": (error) => new ToolFailure({ message: error.message }),
|
||||
}),
|
||||
)
|
||||
if (result.isError)
|
||||
return yield* new ToolFailure({
|
||||
message:
|
||||
result.content
|
||||
.flatMap((part) => (part.type === "text" ? [part.text] : []))
|
||||
.join("\n")
|
||||
.trim() || "MCP tool returned an error",
|
||||
})
|
||||
const content = result.content.map((part) =>
|
||||
part.type === "text"
|
||||
? { type: "text" as const, text: part.text }
|
||||
: {
|
||||
type: "file" as const,
|
||||
uri: `data:${part.mimeType};base64,${part.data}`,
|
||||
mime: part.mimeType,
|
||||
},
|
||||
)
|
||||
const text = content.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n")
|
||||
const output = () => {
|
||||
if (result.structured !== undefined) return result.structured
|
||||
if (text === "") return null
|
||||
// Agents assume JSON returned as text is already an object, so parse it when the server declares no schema.
|
||||
if (tool.outputSchema === undefined && (text.startsWith("{") || text.startsWith("["))) {
|
||||
try {
|
||||
return JSON.parse(text)
|
||||
} catch {}
|
||||
}
|
||||
return text
|
||||
}
|
||||
return {
|
||||
output: output(),
|
||||
...(content.length === 0 ? {} : { content }),
|
||||
}
|
||||
}).pipe(
|
||||
Effect.mapError((error) =>
|
||||
error instanceof ToolFailure
|
||||
? error
|
||||
: new ToolFailure({ message: `Unable to execute ${name(tool.server, tool.name)}` }),
|
||||
),
|
||||
),
|
||||
})
|
||||
}
|
||||
for (const tool of discovered)
|
||||
editor.add(info(tool, (input) => mcp.callTool({ ...input, server: tool.server, name: tool.name })))
|
||||
})
|
||||
}),
|
||||
)
|
||||
@@ -141,7 +157,19 @@ export const layer = Layer.effect(
|
||||
Stream.runForEach(() => reconcile),
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
return Service.of({ flush: Effect.asVoid(Fiber.await(initial)) })
|
||||
return Service.of({
|
||||
flush: Effect.asVoid(Fiber.await(initial)),
|
||||
overlay: (view) =>
|
||||
lock.withPermit(
|
||||
Effect.sync(() => {
|
||||
const hidden = new Set(
|
||||
discovered.filter((tool) => view.shadowed.has(tool.server)).map((tool) => name(tool.server, tool.name)),
|
||||
)
|
||||
if (view.owned.length === 0 && hidden.size === 0) return undefined
|
||||
return { tools: view.owned.map((owned) => info(owned.tool, owned.call)), hidden }
|
||||
}),
|
||||
),
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -4,12 +4,13 @@ import { ToolFailure } from "@opencode/ai"
|
||||
import type { Context } from "@opencode/plugin/effect/plugin"
|
||||
import { Effect, Schema } from "effect"
|
||||
import { Mcp } from "../../mcp/index.js"
|
||||
import { McpSession } from "../../mcp/session.js"
|
||||
import { Permission } from "../../permission.js"
|
||||
|
||||
export const Plugin = {
|
||||
id: "opencode.tools.mcp-resources",
|
||||
effect: Effect.fn("McpResourceTools.Plugin")(function* (ctx: Context) {
|
||||
const mcp = yield* Mcp.Service
|
||||
const mcp = yield* McpSession.Service
|
||||
const permission = yield* Permission.Service
|
||||
|
||||
yield* ctx.tool
|
||||
@@ -30,9 +31,9 @@ export const Plugin = {
|
||||
}),
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
const view = yield* mcp.view(context.sessionID)
|
||||
// Listing every server checks each server name so per-server rules still apply.
|
||||
const servers =
|
||||
input.server === undefined ? (yield* mcp.servers()).map((server) => server.name) : [input.server]
|
||||
const servers = input.server === undefined ? view.servers.map((server) => server.name) : [input.server]
|
||||
yield* permission.assert({
|
||||
action: "opencode_list_mcp_resources",
|
||||
resources: servers,
|
||||
@@ -42,8 +43,8 @@ export const Plugin = {
|
||||
agent: context.agent,
|
||||
source: { type: "tool", messageID: context.messageID, id: context.id },
|
||||
})
|
||||
if (input.server === undefined) return { output: yield* mcp.resourceCatalog() }
|
||||
return { output: yield* mcp.resources({ server: input.server }) }
|
||||
if (input.server === undefined) return { output: yield* view.resourceCatalog }
|
||||
return { output: yield* view.resources(input.server) }
|
||||
}).pipe(Effect.mapError((error) => new ToolFailure({ message: error.message, error }))),
|
||||
})
|
||||
editor.add({
|
||||
@@ -71,7 +72,8 @@ export const Plugin = {
|
||||
agent: context.agent,
|
||||
source: { type: "tool", messageID: context.messageID, id: context.id },
|
||||
})
|
||||
const resource = yield* mcp.readResource(input)
|
||||
const view = yield* mcp.view(context.sessionID)
|
||||
const resource = yield* view.readResource(input)
|
||||
if (!resource)
|
||||
return yield* new ToolFailure({
|
||||
message: `MCP server "${input.server}" is not connected or does not expose resources`,
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
import { Server } from "@modelcontextprotocol/server"
|
||||
import { StdioServerTransport } from "@modelcontextprotocol/server/stdio"
|
||||
|
||||
const [identity = "unknown", pidFile] = process.argv.slice(2)
|
||||
if (pidFile) await Bun.write(pidFile, String(process.pid))
|
||||
|
||||
const server = new Server({ name: `identity-${identity}`, version: "1.0.0" }, { capabilities: { tools: {} } })
|
||||
|
||||
server.setRequestHandler("tools/list", () =>
|
||||
Promise.resolve({
|
||||
tools: [
|
||||
{ name: "whoami", inputSchema: { type: "object" } },
|
||||
{ name: `only_${identity}`, inputSchema: { type: "object" } },
|
||||
],
|
||||
}),
|
||||
)
|
||||
|
||||
server.setRequestHandler("tools/call", () => Promise.resolve({ content: [{ type: "text", text: identity }] }))
|
||||
|
||||
await server.connect(new StdioServerTransport())
|
||||
@@ -1,6 +1,5 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Effect, Layer } from "effect"
|
||||
import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder"
|
||||
import { Effect } from "effect"
|
||||
import { Mcp } from "@opencode/core/mcp/index"
|
||||
import { McpInstructions } from "@opencode/core/mcp/instructions"
|
||||
import { McpTool } from "@opencode/core/tool/mcp"
|
||||
@@ -14,25 +13,28 @@ const schema = { type: "object" as const }
|
||||
const tool = (server: string, name = "search") =>
|
||||
({ server: Mcp.ServerName.make(server), name, inputSchema: schema }) satisfies Mcp.Tool
|
||||
|
||||
const layer = (catalog: () => Mcp.ServerInstructions[], tools: () => Mcp.Tool[]) =>
|
||||
AppNodeBuilder.build(McpInstructions.node, [
|
||||
Mcp.node.replace(
|
||||
Layer.mock(Mcp.Service, {
|
||||
instructions: () => Effect.succeed(catalog()),
|
||||
tools: () => Effect.succeed(tools()),
|
||||
}),
|
||||
),
|
||||
])
|
||||
const view = (instructions: Mcp.ServerInstructions[], tools: Mcp.Tool[]) => ({ instructions, tools, owned: [] })
|
||||
|
||||
describe("McpInstructions", () => {
|
||||
it.effect("renders instructions for servers with at least one permitted tool", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* McpInstructions.Service
|
||||
const generation = yield* service
|
||||
.load([
|
||||
{ action: McpTool.name("alpha", "restricted"), resource: "*", effect: "deny" },
|
||||
{ action: McpTool.name("hidden", "search"), resource: "*", effect: "deny" },
|
||||
])
|
||||
.load(
|
||||
[
|
||||
{ action: McpTool.name("alpha", "restricted"), resource: "*", effect: "deny" },
|
||||
{ action: McpTool.name("hidden", "search"), resource: "*", effect: "deny" },
|
||||
],
|
||||
view(
|
||||
[
|
||||
instructions("beta", "Beta instructions"),
|
||||
instructions("unused", "No tools"),
|
||||
instructions("hidden", "Denied tool"),
|
||||
instructions("alpha", "Alpha line one\nAlpha line two"),
|
||||
],
|
||||
[tool("alpha"), tool("alpha", "restricted"), tool("beta"), tool("hidden")],
|
||||
),
|
||||
)
|
||||
.pipe(Effect.flatMap(readInitial))
|
||||
|
||||
expect(generation.text).toBe(
|
||||
@@ -50,44 +52,31 @@ describe("McpInstructions", () => {
|
||||
"</mcp_instructions>",
|
||||
].join("\n"),
|
||||
)
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
layer(
|
||||
() => [
|
||||
instructions("beta", "Beta instructions"),
|
||||
instructions("unused", "No tools"),
|
||||
instructions("hidden", "Denied tool"),
|
||||
instructions("alpha", "Alpha line one\nAlpha line two"),
|
||||
],
|
||||
() => [tool("alpha"), tool("alpha", "restricted"), tool("beta"), tool("hidden")],
|
||||
),
|
||||
),
|
||||
),
|
||||
}).pipe(Effect.provide(McpInstructions.layer)),
|
||||
)
|
||||
|
||||
it.effect("omits instructions when the agent cannot use execute", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* McpInstructions.Service
|
||||
const generation = yield* service
|
||||
.load([{ action: "execute", resource: "*", effect: "deny" }])
|
||||
.load(
|
||||
[{ action: "execute", resource: "*", effect: "deny" }],
|
||||
view([instructions("alpha", "Alpha instructions")], [tool("alpha")]),
|
||||
)
|
||||
.pipe(Effect.flatMap(readInitial))
|
||||
|
||||
expect(generation.text).toBe("")
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
layer(
|
||||
() => [instructions("alpha", "Alpha instructions")],
|
||||
() => [tool("alpha")],
|
||||
),
|
||||
),
|
||||
),
|
||||
}).pipe(Effect.provide(McpInstructions.layer)),
|
||||
)
|
||||
|
||||
it.effect("keeps MCP instructions when Code Mode is disabled and execute is denied", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* McpInstructions.Service
|
||||
const generation = yield* service
|
||||
.load([{ action: "execute", resource: "*", effect: "deny" }])
|
||||
.load(
|
||||
[{ action: "execute", resource: "*", effect: "deny" }],
|
||||
view([instructions("alpha", "Alpha instructions")], [{ ...tool("alpha"), codemode: false }]),
|
||||
)
|
||||
.pipe(Effect.flatMap(readInitial))
|
||||
|
||||
expect(generation.text).toBe(
|
||||
@@ -99,31 +88,19 @@ describe("McpInstructions", () => {
|
||||
"</mcp_instructions>",
|
||||
].join("\n"),
|
||||
)
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
layer(
|
||||
() => [instructions("alpha", "Alpha instructions")],
|
||||
() => [
|
||||
{
|
||||
server: Mcp.ServerName.make("alpha"),
|
||||
name: "search",
|
||||
inputSchema: schema,
|
||||
codemode: false,
|
||||
} satisfies Mcp.Tool,
|
||||
],
|
||||
),
|
||||
),
|
||||
),
|
||||
}).pipe(Effect.provide(McpInstructions.layer)),
|
||||
)
|
||||
|
||||
it.effect("restates guidance when Code Mode is disabled for a server", () => {
|
||||
let tools: Mcp.Tool[] = [tool("alpha")]
|
||||
return Effect.gen(function* () {
|
||||
it.effect("restates guidance when Code Mode is disabled for a server", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* McpInstructions.Service
|
||||
const initialized = yield* service.load([]).pipe(Effect.flatMap(readInitial))
|
||||
const catalog = [instructions("alpha", "Alpha instructions")]
|
||||
const initialized = yield* service.load([], view(catalog, [tool("alpha")])).pipe(Effect.flatMap(readInitial))
|
||||
|
||||
tools = [{ ...tool("alpha"), codemode: false }]
|
||||
const changed = yield* readUpdate(yield* service.load([]), initialized)
|
||||
const changed = yield* readUpdate(
|
||||
yield* service.load([], view(catalog, [{ ...tool("alpha"), codemode: false }])),
|
||||
initialized,
|
||||
)
|
||||
expect(changed.text).toBe(
|
||||
[
|
||||
"The available MCP server instructions have changed. This list supersedes the previous one.",
|
||||
@@ -134,25 +111,20 @@ describe("McpInstructions", () => {
|
||||
"</mcp_instructions>",
|
||||
].join("\n"),
|
||||
)
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
layer(
|
||||
() => [instructions("alpha", "Alpha instructions")],
|
||||
() => tools,
|
||||
),
|
||||
),
|
||||
)
|
||||
})
|
||||
}).pipe(Effect.provide(McpInstructions.layer)),
|
||||
)
|
||||
|
||||
it.effect("renders additions, changes, and removal", () => {
|
||||
let catalog = [instructions("alpha", "Alpha instructions")]
|
||||
const tools = [tool("alpha"), tool("beta")]
|
||||
return Effect.gen(function* () {
|
||||
it.effect("renders additions, changes, and removal", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* McpInstructions.Service
|
||||
const initialized = yield* service.load([]).pipe(Effect.flatMap(readInitial))
|
||||
const tools = [tool("alpha"), tool("beta")]
|
||||
const load = (catalog: Mcp.ServerInstructions[]) => service.load([], view(catalog, tools))
|
||||
const initialized = yield* load([instructions("alpha", "Alpha instructions")]).pipe(Effect.flatMap(readInitial))
|
||||
|
||||
catalog = [instructions("alpha", "Alpha instructions"), instructions("beta", "Beta instructions")]
|
||||
const added = yield* readUpdate(yield* service.load([]), initialized)
|
||||
const added = yield* readUpdate(
|
||||
yield* load([instructions("alpha", "Alpha instructions"), instructions("beta", "Beta instructions")]),
|
||||
initialized,
|
||||
)
|
||||
expect(added.text).toBe(
|
||||
[
|
||||
"New MCP server instructions are available in addition to those previously listed:",
|
||||
@@ -163,8 +135,10 @@ describe("McpInstructions", () => {
|
||||
].join("\n"),
|
||||
)
|
||||
|
||||
catalog = [instructions("alpha", "Updated alpha"), instructions("beta", "Beta instructions")]
|
||||
const changed = yield* readUpdate(yield* service.load([]), added)
|
||||
const changed = yield* readUpdate(
|
||||
yield* load([instructions("alpha", "Updated alpha"), instructions("beta", "Beta instructions")]),
|
||||
added,
|
||||
)
|
||||
expect(changed.text).toBe(
|
||||
[
|
||||
"The available MCP server instructions have changed. This list supersedes the previous one.",
|
||||
@@ -181,21 +155,12 @@ describe("McpInstructions", () => {
|
||||
].join("\n"),
|
||||
)
|
||||
|
||||
catalog = [instructions("beta", "Beta instructions")]
|
||||
const removed = yield* readUpdate(yield* service.load([]), changed)
|
||||
const removed = yield* readUpdate(yield* load([instructions("beta", "Beta instructions")]), changed)
|
||||
expect(removed.text).toBe("Instructions for the following MCP servers are no longer available: alpha.")
|
||||
|
||||
catalog = []
|
||||
expect((yield* readUpdate(yield* service.load([]), removed)).text).toBe(
|
||||
expect((yield* readUpdate(yield* load([]), removed)).text).toBe(
|
||||
"MCP server instructions are no longer available.",
|
||||
)
|
||||
}).pipe(
|
||||
Effect.provide(
|
||||
layer(
|
||||
() => catalog,
|
||||
() => tools,
|
||||
),
|
||||
),
|
||||
)
|
||||
})
|
||||
}).pipe(Effect.provide(McpInstructions.layer)),
|
||||
)
|
||||
})
|
||||
|
||||
@@ -0,0 +1,243 @@
|
||||
import path from "node:path"
|
||||
import { mkdir, rm } from "node:fs/promises"
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Context, Effect, Layer, Schedule, Stream } from "effect"
|
||||
import { ConfigMCP } from "@opencode/schema/config/mcp"
|
||||
import { McpEvent } from "@opencode/schema/mcp-event"
|
||||
import { AppNodeBuilder } from "@opencode/core/effect/app-node-builder"
|
||||
import { LayerNode } from "@opencode/util/effect/layer-node"
|
||||
import { Database } from "@opencode/core/database/database"
|
||||
import { Bus } from "@opencode/core/bus"
|
||||
import { Instance } from "@opencode/core/instance/service"
|
||||
import { Location } from "@opencode/core/location"
|
||||
import { LocationServiceMap } from "@opencode/core/location-services"
|
||||
import { Mcp } from "@opencode/core/mcp/index"
|
||||
import { McpSession } from "@opencode/core/mcp/session"
|
||||
import { Project } from "@opencode/core/project"
|
||||
import { AbsolutePath } from "@opencode/core/schema"
|
||||
import { Session } from "@opencode/core/session"
|
||||
import { SessionEnvironment } from "@opencode/core/session/environment"
|
||||
import { SessionExecution } from "@opencode/core/session/execution"
|
||||
import { SessionModelTransport } from "@opencode/core/session/model-transport"
|
||||
import { SessionProjector } from "@opencode/core/session/projector"
|
||||
import { SessionStore } from "@opencode/core/session/store"
|
||||
import { Tool } from "@opencode/core/tool"
|
||||
import { McpTool } from "@opencode/core/tool/mcp"
|
||||
import { testEffect } from "./lib/effect"
|
||||
import { globalProjectNode } from "./lib/project"
|
||||
import { codeModeListings, waitForCodeModeTool } from "./lib/tool"
|
||||
import { offlineModels } from "./fixture/models"
|
||||
import { tmpdirScoped } from "./fixture/tmpdir"
|
||||
|
||||
const transport = Layer.succeed(
|
||||
SessionModelTransport.Service,
|
||||
SessionModelTransport.Service.of({
|
||||
bind: () => ({ execute: () => Effect.die("Unexpected WebSocket execution") }),
|
||||
close: () => Effect.void,
|
||||
closeAll: Effect.void,
|
||||
}),
|
||||
)
|
||||
|
||||
const it = testEffect(
|
||||
AppNodeBuilder.build(
|
||||
LayerNode.group([
|
||||
Database.node,
|
||||
Bus.node,
|
||||
SessionProjector.node,
|
||||
SessionStore.node,
|
||||
SessionEnvironment.node,
|
||||
Session.node,
|
||||
Instance.node,
|
||||
LocationServiceMap.node,
|
||||
]),
|
||||
[
|
||||
Project.node.replace(globalProjectNode),
|
||||
SessionExecution.node.replace(SessionExecution.noopLayer),
|
||||
SessionModelTransport.node.replace(transport),
|
||||
offlineModels,
|
||||
],
|
||||
),
|
||||
)
|
||||
|
||||
const fixture = (subdirectory?: string) =>
|
||||
Effect.gen(function* () {
|
||||
const temporary = yield* tmpdirScoped()
|
||||
const directory = subdirectory ? path.join(temporary.path, subdirectory) : temporary.path
|
||||
if (subdirectory) yield* Effect.promise(() => mkdir(directory))
|
||||
const location = Location.Ref.make({ directory: AbsolutePath.make(directory) })
|
||||
const locations = yield* LocationServiceMap.Service
|
||||
const context = yield* Effect.acquireRelease(locations.contextEffect(location), () =>
|
||||
locations.invalidate(location),
|
||||
)
|
||||
const pidFile = (identity: string) => path.join(temporary.path, `${identity}.pid`)
|
||||
return {
|
||||
root: AbsolutePath.make(temporary.path),
|
||||
directory,
|
||||
location,
|
||||
server: (identity: string) =>
|
||||
new ConfigMCP.Local({
|
||||
type: "local",
|
||||
command: [
|
||||
process.execPath,
|
||||
path.join(import.meta.dir, "fixture/mcp-identity.ts"),
|
||||
identity,
|
||||
pidFile(identity),
|
||||
],
|
||||
}),
|
||||
pid: (identity: string) => Effect.promise(() => Bun.file(pidFile(identity)).text()).pipe(Effect.map(Number)),
|
||||
mcp: Context.get(context, Mcp.Service),
|
||||
sessions: Context.get(context, McpSession.Service),
|
||||
mcpTools: Context.get(context, McpTool.Service),
|
||||
registry: Context.get(context, Tool.Service),
|
||||
}
|
||||
})
|
||||
|
||||
const text = (result: Mcp.ToolResult) =>
|
||||
result.content.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("")
|
||||
|
||||
// The harness still reads the user's global config, so assertions only look at the fixture's server.
|
||||
const names = (tools: ReadonlyArray<Mcp.Tool>) =>
|
||||
tools.filter((tool) => tool.server === "ctx").map((tool) => `${tool.server}.${tool.name}`)
|
||||
const owned = (view: McpSession.View) => names(view.owned.map((item) => item.tool))
|
||||
|
||||
const whoami = (view: McpSession.View, sessionID: Session.ID) =>
|
||||
Effect.gen(function* () {
|
||||
const tool = view.owned.find((item) => item.tool.server === "ctx" && item.tool.name === "whoami")
|
||||
if (!tool) return yield* Effect.die(new Error("whoami is not owned by this view"))
|
||||
return text(yield* tool.call({ sessionID }))
|
||||
})
|
||||
|
||||
const running = (pid: number) => {
|
||||
try {
|
||||
process.kill(pid, 0)
|
||||
return true
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
const exited = (pid: number) =>
|
||||
Effect.suspend(() => (running(pid) ? Effect.fail(`process ${pid} is still running`) : Effect.void)).pipe(
|
||||
Effect.retry({ times: 300, schedule: Schedule.spaced("10 millis") }),
|
||||
)
|
||||
|
||||
describe("Session-scoped MCP servers", () => {
|
||||
it.live("are visible only to their Session tree and shadow Location servers there", () =>
|
||||
Effect.gen(function* () {
|
||||
const test = yield* fixture()
|
||||
const sessions = yield* Session.Service
|
||||
const root = yield* sessions.create({ location: test.location })
|
||||
const child = yield* sessions.create({ parentID: root.id })
|
||||
const other = yield* sessions.create({ location: test.location })
|
||||
const bystander = yield* sessions.create({ location: test.location })
|
||||
const broken = yield* sessions.create({ location: test.location })
|
||||
|
||||
yield* test.sessions.add(root.id, "ctx", test.server("root"))
|
||||
|
||||
expect(owned(yield* test.sessions.view(root.id))).toEqual(["ctx.only_root", "ctx.whoami"])
|
||||
expect(owned(yield* test.sessions.view(child.id))).toEqual(["ctx.only_root", "ctx.whoami"])
|
||||
expect(owned(yield* test.sessions.view(other.id))).toEqual([])
|
||||
expect(names(yield* test.mcp.tools())).toEqual([])
|
||||
expect((yield* test.mcp.servers()).some((server) => server.name === "ctx")).toBe(false)
|
||||
expect((yield* test.sessions.view(child.id)).servers.some((server) => server.name === "ctx")).toBe(true)
|
||||
expect(yield* whoami(yield* test.sessions.view(child.id), child.id)).toBe("root")
|
||||
|
||||
yield* test.sessions.add(other.id, "ctx", test.server("other"))
|
||||
yield* test.sessions.add(broken.id, "ctx", new ConfigMCP.Local({ type: "local", command: ["false"] }))
|
||||
yield* test.mcp.add("ctx", test.server("location"))
|
||||
|
||||
expect(yield* whoami(yield* test.sessions.view(root.id), root.id)).toBe("root")
|
||||
expect(yield* whoami(yield* test.sessions.view(other.id), other.id)).toBe("other")
|
||||
expect(text(yield* test.mcp.callTool({ server: "ctx", name: "whoami", sessionID: bystander.id }))).toBe(
|
||||
"location",
|
||||
)
|
||||
const rootView = yield* test.sessions.view(root.id)
|
||||
expect(rootView.shadowed.has("ctx")).toBe(true)
|
||||
expect(names(rootView.tools)).toEqual([])
|
||||
const bystanderView = yield* test.sessions.view(bystander.id)
|
||||
expect(bystanderView.shadowed.size).toBe(0)
|
||||
expect(names(bystanderView.tools)).toEqual(["ctx.only_location", "ctx.whoami"])
|
||||
|
||||
yield* waitForCodeModeTool(test.registry, "ctx.only_location")
|
||||
const listed = (sessionID: Session.ID) =>
|
||||
test.sessions.view(sessionID).pipe(
|
||||
Effect.flatMap(test.mcpTools.overlay),
|
||||
Effect.flatMap((overlay) => test.registry.snapshot(undefined, overlay)),
|
||||
Effect.map((snapshot) =>
|
||||
(snapshot.codeModeCatalog ? codeModeListings(snapshot.codeModeCatalog) : [])
|
||||
.map((tool) => tool.path)
|
||||
.filter((path) => path.startsWith("ctx.")),
|
||||
),
|
||||
)
|
||||
expect(yield* listed(child.id)).toEqual(["ctx.only_root", "ctx.whoami"])
|
||||
expect(yield* listed(bystander.id)).toEqual(["ctx.only_location", "ctx.whoami"])
|
||||
expect(yield* listed(broken.id)).toEqual([])
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("releases servers on replacement, removal, and Session deletion without Location events", () =>
|
||||
Effect.gen(function* () {
|
||||
const test = yield* fixture()
|
||||
const sessions = yield* Session.Service
|
||||
const bus = yield* Bus.Service
|
||||
const events: string[] = []
|
||||
yield* bus
|
||||
.subscribe([McpEvent.StatusChanged, McpEvent.ToolsChanged, McpEvent.ResourcesChanged, Mcp.PromptsChanged])
|
||||
.pipe(
|
||||
Stream.runForEach((event) => Effect.sync(() => events.push(`${event.type}:${event.data.server}`))),
|
||||
Effect.forkScoped({ startImmediately: true }),
|
||||
)
|
||||
const first = yield* sessions.create({ location: test.location })
|
||||
const second = yield* sessions.create({ location: test.location })
|
||||
|
||||
yield* test.sessions.add(first.id, "ctx", test.server("first"))
|
||||
yield* test.sessions.add(second.id, "ctx", test.server("second"))
|
||||
const firstPid = yield* test.pid("first")
|
||||
const secondPid = yield* test.pid("second")
|
||||
|
||||
yield* test.sessions.add(first.id, "ctx", test.server("first"))
|
||||
expect(yield* test.pid("first")).toBe(firstPid)
|
||||
expect(running(firstPid)).toBe(true)
|
||||
|
||||
yield* test.sessions.add(first.id, "ctx", test.server("replaced"))
|
||||
yield* exited(firstPid)
|
||||
expect(yield* whoami(yield* test.sessions.view(first.id), first.id)).toBe("replaced")
|
||||
expect(yield* whoami(yield* test.sessions.view(second.id), second.id)).toBe("second")
|
||||
expect(running(secondPid)).toBe(true)
|
||||
|
||||
const replacedPid = yield* test.pid("replaced")
|
||||
yield* test.sessions.remove(first.id, "ctx")
|
||||
yield* exited(replacedPid)
|
||||
expect(owned(yield* test.sessions.view(first.id))).toEqual([])
|
||||
expect(yield* test.sessions.remove(first.id, "ctx").pipe(Effect.flip)).toBeInstanceOf(Mcp.NotFoundError)
|
||||
|
||||
yield* sessions.remove(second.id)
|
||||
yield* exited(secondPid)
|
||||
expect(owned(yield* test.sessions.view(second.id))).toEqual([])
|
||||
|
||||
yield* test.mcp.add("probe", new ConfigMCP.Local({ type: "local", command: ["unused"], disabled: true }))
|
||||
yield* Effect.suspend(() =>
|
||||
events.includes(`${McpEvent.StatusChanged.type}:probe`) ? Effect.void : Effect.fail("probe event missing"),
|
||||
).pipe(Effect.retry({ times: 100, schedule: Schedule.spaced("10 millis") }))
|
||||
expect(events.filter((event) => event.endsWith(":ctx"))).toEqual([])
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("releases servers in the former Location when their Session moves", () =>
|
||||
Effect.gen(function* () {
|
||||
const test = yield* fixture("source")
|
||||
const sessions = yield* Session.Service
|
||||
const session = yield* sessions.create({ location: test.location })
|
||||
|
||||
yield* test.sessions.add(session.id, "ctx", test.server("moved"))
|
||||
const pid = yield* test.pid("moved")
|
||||
// A missing source directory applies the move immediately instead of at a step boundary.
|
||||
yield* Effect.promise(() => rm(test.directory, { recursive: true }))
|
||||
yield* sessions.move({ sessionID: session.id, directory: test.root })
|
||||
|
||||
expect((yield* sessions.get(session.id)).location.directory).toBe(test.root)
|
||||
yield* exited(pid)
|
||||
expect(owned(yield* test.sessions.view(session.id))).toEqual([])
|
||||
}),
|
||||
)
|
||||
})
|
||||
@@ -25,11 +25,13 @@ import { Environment } from "@opencode/core/environment/index"
|
||||
import { EnvironmentUnavailable } from "@opencode/core/environment/unavailable"
|
||||
import { Location } from "@opencode/core/location"
|
||||
import { Mcp } from "@opencode/core/mcp/index"
|
||||
import { McpSession } from "@opencode/core/mcp/session"
|
||||
import { McpClient } from "@opencode/core/mcp/client"
|
||||
import { McpStdio } from "@opencode/core/mcp/stdio"
|
||||
import { Permission } from "@opencode/core/permission"
|
||||
import { AbsolutePath } from "@opencode/core/schema"
|
||||
import { Session } from "@opencode/core/session"
|
||||
import { SessionStore } from "@opencode/core/session/store"
|
||||
import { State } from "@opencode/core/state"
|
||||
import { McpTool } from "@opencode/core/tool/mcp"
|
||||
import { McpResourceTools } from "@opencode/core/tool/plugin/mcp-resource"
|
||||
@@ -283,6 +285,7 @@ function resourceMcpLayer(
|
||||
yield* ConfigMcpPlugin.register(bus.subscribe())
|
||||
}),
|
||||
).pipe(
|
||||
Layer.provideMerge(McpSession.layer(options)),
|
||||
Layer.provideMerge(Mcp.layer(options)),
|
||||
Layer.provideMerge(Form.layer),
|
||||
Layer.provide(
|
||||
@@ -348,6 +351,7 @@ function resourceMcpLayer(
|
||||
},
|
||||
}),
|
||||
Layer.mock(Credential.Service, {}),
|
||||
AppNodeBuilder.build(SessionStore.node),
|
||||
overrides?.environment ?? hostEnvironmentLayer,
|
||||
),
|
||||
),
|
||||
|
||||
@@ -4,6 +4,7 @@ import { PromptInput } from "@opencode/schema/prompt-input"
|
||||
import { Session } from "@opencode/schema/session"
|
||||
import { SessionStats } from "@opencode/schema/session-stats"
|
||||
import { InstructionEntry } from "@opencode/schema/instruction-entry"
|
||||
import { Mcp } from "@opencode/schema/mcp"
|
||||
import { Project } from "@opencode/schema/project"
|
||||
import {
|
||||
AbsolutePath,
|
||||
@@ -25,6 +26,7 @@ import {
|
||||
FormNotFoundError,
|
||||
InvalidCursorError,
|
||||
InvalidRequestError,
|
||||
McpServerNotFoundError,
|
||||
MessageNotFoundError,
|
||||
ServiceUnavailableError,
|
||||
SessionBusyError,
|
||||
@@ -702,6 +704,38 @@ export const makeSessionGroup = <
|
||||
}),
|
||||
),
|
||||
)
|
||||
.add(
|
||||
HttpApiEndpoint.put("session.mcp.add", "/api/experimental/session/:sessionID/mcp/:server", {
|
||||
params: { sessionID: Session.ID, server: Schema.String },
|
||||
payload: Schema.Struct({ config: Mcp.ServerConfig }),
|
||||
success: HttpApiSchema.NoContent,
|
||||
error: SessionNotFoundError,
|
||||
})
|
||||
.middleware(sessionLocationMiddleware)
|
||||
.annotateMerge(
|
||||
OpenApi.annotations({
|
||||
identifier: "experimental.session.mcp.add",
|
||||
summary: "Add session MCP server",
|
||||
description:
|
||||
"Add or replace an MCP server visible only to this session and its child sessions, connecting it immediately. It shadows a Location server of the same name for those sessions, is not persisted, and is released when removed or when the session is deleted or moved. Re-adding an identical config is a no-op.",
|
||||
}),
|
||||
),
|
||||
)
|
||||
.add(
|
||||
HttpApiEndpoint.delete("session.mcp.remove", "/api/experimental/session/:sessionID/mcp/:server", {
|
||||
params: { sessionID: Session.ID, server: Schema.String },
|
||||
success: HttpApiSchema.NoContent,
|
||||
error: [SessionNotFoundError, McpServerNotFoundError],
|
||||
})
|
||||
.middleware(sessionLocationMiddleware)
|
||||
.annotateMerge(
|
||||
OpenApi.annotations({
|
||||
identifier: "experimental.session.mcp.remove",
|
||||
summary: "Remove session MCP server",
|
||||
description: "Stop an MCP server registered for this session and remove it.",
|
||||
}),
|
||||
),
|
||||
)
|
||||
.add(
|
||||
HttpApiEndpoint.post("session.generate", "/api/session/:sessionID/generate", {
|
||||
params: { sessionID: Session.ID },
|
||||
|
||||
@@ -5,7 +5,7 @@ import { HttpApiBuilder, HttpApiSchema } from "effect/unstable/httpapi"
|
||||
import { Api } from "../api"
|
||||
import { response } from "../location"
|
||||
|
||||
const notFound = <A, R>(effect: Effect.Effect<A, Mcp.NotFoundError, R>) =>
|
||||
export const notFound = <A, R>(effect: Effect.Effect<A, Mcp.NotFoundError, R>) =>
|
||||
effect.pipe(Effect.mapError((error) => new McpServerNotFoundError({ server: error.server, message: error.message })))
|
||||
|
||||
export const McpHandler = HttpApiBuilder.group(Api, "server.mcp", (handlers) =>
|
||||
|
||||
@@ -4,6 +4,7 @@ import { SessionTitle } from "@opencode/core/session/title"
|
||||
import { SessionTransfer } from "@opencode/core/session/transfer"
|
||||
import { InstructionEntry } from "@opencode/core/session/instruction-entry"
|
||||
import { Form } from "@opencode/core/form"
|
||||
import { McpSession } from "@opencode/core/mcp/session"
|
||||
import { DateTime, Effect, Stream } from "effect"
|
||||
import { HttpApiBuilder, HttpApiSchema } from "effect/unstable/httpapi"
|
||||
import { Api } from "../api"
|
||||
@@ -24,6 +25,7 @@ import {
|
||||
} from "@opencode/protocol/errors"
|
||||
import { AbsolutePath } from "@opencode/core/schema"
|
||||
import { failedMessageDecode, failedSnapshot, missingMessage, missingSession } from "./session-error"
|
||||
import { notFound } from "./mcp"
|
||||
|
||||
const DefaultSessionsLimit = 50
|
||||
|
||||
@@ -583,6 +585,22 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
|
||||
return HttpApiSchema.NoContent.make()
|
||||
}),
|
||||
)
|
||||
.handle(
|
||||
"session.mcp.add",
|
||||
Effect.fn(function* (ctx) {
|
||||
const mcp = yield* McpSession.Service
|
||||
yield* mcp.add(ctx.params.sessionID, ctx.params.server, ctx.payload.config)
|
||||
return HttpApiSchema.NoContent.make()
|
||||
}),
|
||||
)
|
||||
.handle(
|
||||
"session.mcp.remove",
|
||||
Effect.fn(function* (ctx) {
|
||||
const mcp = yield* McpSession.Service
|
||||
yield* notFound(mcp.remove(ctx.params.sessionID, ctx.params.server))
|
||||
return HttpApiSchema.NoContent.make()
|
||||
}),
|
||||
)
|
||||
.handle(
|
||||
"session.generate",
|
||||
Effect.fn(function* (ctx) {
|
||||
|
||||
@@ -21,6 +21,7 @@ import { SessionTransfer } from "@opencode/core/session/transfer"
|
||||
import { ShellSelect } from "@opencode/core/shell/select"
|
||||
import { Job } from "@opencode/core/job"
|
||||
import { Mcp } from "@opencode/core/mcp/index"
|
||||
import { McpSession } from "@opencode/core/mcp/session"
|
||||
import { Global } from "@opencode/util/global"
|
||||
import { InstructionDiscovery } from "@opencode/core/instruction-discovery"
|
||||
import { LocationServiceMap } from "@opencode/core/location-service-map"
|
||||
@@ -112,6 +113,7 @@ function makeRoutes<AuthError, AuthServices>(
|
||||
overrides: LayerNode.Replacements,
|
||||
instances?: InstanceNode,
|
||||
) {
|
||||
const clientInfo = { name: options.app?.name ?? "opencode", version: options.app?.version ?? "unknown" }
|
||||
const standard: LayerNode.Replacements = [
|
||||
Database.node.replace(Database.configured(options.database)),
|
||||
PersistentPty.node.replace(PersistentPty.configured(options.pty)),
|
||||
@@ -130,14 +132,8 @@ function makeRoutes<AuthError, AuthServices>(
|
||||
),
|
||||
InstructionDiscovery.node.replace(InstructionDiscovery.configured({ project: options.config?.project })),
|
||||
ShellSelect.node.replace(ShellSelect.configured({ gitbash: options.windows?.gitbash })),
|
||||
Mcp.node.replace(
|
||||
Mcp.configured({
|
||||
clientInfo: {
|
||||
name: options.app?.name ?? "opencode",
|
||||
version: options.app?.version ?? "unknown",
|
||||
},
|
||||
}),
|
||||
),
|
||||
Mcp.node.replace(Mcp.configured({ clientInfo })),
|
||||
McpSession.node.replace(McpSession.configured({ clientInfo })),
|
||||
]
|
||||
const build = (overrides: LayerNode.Replacements) => {
|
||||
const replacements: LayerNode.Replacements = [
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
import { expect } from "bun:test"
|
||||
import { Session } from "@opencode/schema/session"
|
||||
import { Effect, Schema } from "effect"
|
||||
import { it } from "../../core/test/lib/effect"
|
||||
import { ServerFetch } from "../src/fetch"
|
||||
|
||||
const SessionResponse = Schema.Struct({ data: Schema.toEncoded(Session.Info) })
|
||||
|
||||
it.live("adds and removes session-scoped MCP servers", () =>
|
||||
Effect.gen(function* () {
|
||||
const handler = yield* ServerFetch.make({
|
||||
app: { version: "test" },
|
||||
database: { path: ":memory:" },
|
||||
fs: { filewatcher: false },
|
||||
models: { fetch: false },
|
||||
})
|
||||
const request = (path: string, method: string, body?: unknown) =>
|
||||
Effect.promise(async () => {
|
||||
const response = await handler(
|
||||
new Request(`http://opencode.local${path}`, {
|
||||
method,
|
||||
headers: { "content-type": "application/json" },
|
||||
body: body === undefined ? undefined : JSON.stringify(body),
|
||||
}),
|
||||
)
|
||||
return { status: response.status, body: response.status === 204 ? undefined : await response.json() }
|
||||
})
|
||||
const config = { config: { type: "local", command: ["unused"], disabled: true } }
|
||||
|
||||
const created = Schema.decodeUnknownSync(SessionResponse)((yield* request("/api/session", "POST", {})).body)
|
||||
const path = `/api/experimental/session/${created.data.id}/mcp/ctx`
|
||||
|
||||
expect((yield* request(path, "PUT", config)).status).toBe(204)
|
||||
expect((yield* request(path, "DELETE")).status).toBe(204)
|
||||
expect(yield* request(path, "DELETE")).toMatchObject({
|
||||
status: 404,
|
||||
body: { _tag: "McpServerNotFoundError", server: "ctx" },
|
||||
})
|
||||
expect(yield* request("/api/experimental/session/ses_missing/mcp/ctx", "PUT", config)).toMatchObject({
|
||||
status: 404,
|
||||
body: { _tag: "SessionNotFoundError" },
|
||||
})
|
||||
}).pipe(Effect.scoped),
|
||||
)
|
||||
Reference in New Issue
Block a user