Compare commits

...
28 changed files with 1534 additions and 709 deletions
+17 -20
View File
@@ -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,
+23 -14
View File
@@ -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) {
+2 -2
View File
@@ -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" } },
},
+16 -1
View File
@@ -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 }>(
{
+54 -8
View File
@@ -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"]
+2
View File
@@ -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,
+145
View File
@@ -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
View File
@@ -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]
+41 -45
View File
@@ -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: [] })
+229
View File
@@ -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()
+3
View File
@@ -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,
+6 -2
View File
@@ -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,
+26 -7
View File
@@ -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: [
+3 -1
View File
@@ -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
View File
@@ -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())
+55 -90
View File
@@ -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)),
)
})
+243
View File
@@ -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([])
}),
)
})
+4
View File
@@ -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,
),
),
+34
View File
@@ -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 },
+1 -1
View File
@@ -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) =>
+18
View File
@@ -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) {
+4 -8
View File
@@ -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 = [
+44
View File
@@ -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),
)