返回源码地图

packages/mcp/mcp-client/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 12;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * MCP client bridge plugin: connects to an external MCP server and registers
3 * its tools on `ctx.tools` under server-qualified public names
4 * (`mcp__<serverName>__<rawName>`). Each plugin instance connects to one MCP
5 * server; load multiple instances in `cordis.yml` for multiple servers.
6 *
7 * Namespace plugin (named exports, no default export). Lifecycle is
8 * effect-scoped: disposal disconnects from the server, unregisters all tools,
9 * and releases the `serverName` namespace reservation. HMR hot-swaps by
10 * disposing the old instance and creating a new one; identical `serverName`
11 * reproduces identical public tool names.
12 *
13 * @module @deepseek-ai/dsh-mcp-client
14 */
15
16import type { Context } from '@deepseek-ai/cordis'
17import z from '@deepseek-ai/schemastery'
18import { scopeOf } from '@deepseek-ai/dsh-scope'
19import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
20import { DEFAULT_MAX_INSTRUCTION_BYTES, RECONNECT_DEFAULTS, resolveReconnectPolicy, startConnection } from './connection.ts'
21import type { ReconnectConfig } from './connection.ts'
22import { registerServerContext } from './server-context.ts'
23// Side-effect type import: declaration-merges `ctx.tools` onto Context.
24import type {} from '@deepseek-ai/dsh-tools'
25
26export { createMcpToolDefinition } from './tools.ts'
27export type { McpResult, McpToolDefinitionOptions } from './tools.ts'
28export type { ReconnectConfig, ResolvedReconnectPolicy } from './connection.ts'
29
30/** Cordis plugin name used by loader diagnostics. */
31export const name = 'mcp-client'
32
33/** Services required by this plugin. */
34export const inject = ['tools']
35
36/** Default timeout for individual MCP tool calls and resource requests (ms). */
37const DEFAULT_TOOL_CALL_TIMEOUT_MS = 60_000
38
39/** Valid `serverName`, kept below the public tool-name budget. */
40const SERVER_NAME_PATTERN = /^[A-Za-z0-9_-]{1,32}$/
41
42/**
43 * Live `serverName` reservations per registration scope. Agent-scoped MCP
44 * servers may reuse a namespace in another Agent, while global instances and
45 * duplicates inside one Agent remain mutually exclusive.
46 */
47const activeServerNames = new WeakMap<object, Set<string>>()
48
49// ---- Config ----
50
51/** Config for connecting to an MCP server via a spawned child process over stdio. */
52export interface StdioConfig {
53 /** Selects child-process stdio transport. */
54 transport: 'stdio'
55 /**
56 * Stable local namespace for this server's model-facing tool names
57 * (`mcp__<serverName>__<rawName>`). Must match `[A-Za-z0-9_-]{1,32}` and be
58 * unique across live mcp-client instances.
59 */
60 serverName: string
61 /** Executable used to start the server. */
62 command: string
63 /** Arguments passed directly, without shell interpolation. */
64 args: string[]
65 /** Extra env vars merged on top of scrubbed ambient env. */
66 env: Record<string, string>
67 /** Working directory for the child process. */
68 cwd: string
69 /** Timeout per tool call or resource request in milliseconds. */
70 toolCallTimeoutMs: number
71 /** Fail plugin activation when the initial connection or tool synchronization fails. */
72 failOnStartupError: boolean
73 /** Maximum UTF-8 bytes of attributed server instructions (default 32768). */
74 maxInstructionBytes?: number
75 /** Automatic reconnect policy after a lost connection; omission uses the defaults. */
76 reconnect?: ReconnectConfig
77}
78
79/** Config for connecting to an MCP server over Streamable HTTP (SSE). */
80export interface StreamableHttpConfig {
81 /** Selects Streamable HTTP transport. */
82 transport: 'streamable-http'
83 /**
84 * Stable local namespace for this server's model-facing tool names
85 * (`mcp__<serverName>__<rawName>`). Must match `[A-Za-z0-9_-]{1,32}` and be
86 * unique across live mcp-client instances.
87 */
88 serverName: string
89 /** MCP endpoint URL. */
90 url: string
91 /** Additional headers attached to MCP requests. */
92 headers: Record<string, string>
93 /** Timeout per tool call or resource request in milliseconds. */
94 toolCallTimeoutMs: number
95 /** Fail plugin activation when the initial connection or tool synchronization fails. */
96 failOnStartupError: boolean
97 /** Maximum UTF-8 bytes of attributed server instructions (default 32768). */
98 maxInstructionBytes?: number
99 /** Automatic reconnect policy after a lost connection; omission uses the defaults. */
100 reconnect?: ReconnectConfig
101}
102
103/** Configuration for one stdio or Streamable HTTP MCP server. */
104export type Config = StdioConfig | StreamableHttpConfig
105
106type StdioConfigInput = Omit<StdioConfig, 'args' | 'env' | 'cwd' | 'toolCallTimeoutMs' | 'failOnStartupError'>
107 & Partial<Pick<StdioConfig, 'args' | 'env' | 'cwd' | 'toolCallTimeoutMs' | 'failOnStartupError'>>
108type StreamableHttpConfigInput = Omit<StreamableHttpConfig, 'headers' | 'toolCallTimeoutMs' | 'failOnStartupError'>
109 & Partial<Pick<StreamableHttpConfig, 'headers' | 'toolCallTimeoutMs' | 'failOnStartupError'>>
110type ConfigInput = StdioConfigInput | StreamableHttpConfigInput
111
112const Reconnect: z<ReconnectConfig> = z.object({
113 enabled: z.boolean().default(RECONNECT_DEFAULTS.enabled),
114 initialDelayMs: z.number().min(1).max(MAX_TIMER_DELAY_MS).default(RECONNECT_DEFAULTS.initialDelayMs),
115 maxDelayMs: z.number().min(1).max(MAX_TIMER_DELAY_MS).default(RECONNECT_DEFAULTS.maxDelayMs),
116 maxAttempts: z.number().step(1).min(1).max(Number.MAX_SAFE_INTEGER).default(RECONNECT_DEFAULTS.maxAttempts),
117})
118
119export const Config = z.union([
120 z.object({
121 transport: z.const('stdio'),
122 serverName: z.string().required().pattern(SERVER_NAME_PATTERN),
123 command: z.string().required(),
124 args: z.array(String).default([]),
125 env: z.dict(String).default({}),
126 cwd: z.string().default(''),
127 toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
128 failOnStartupError: z.boolean().default(false),
129 maxInstructionBytes: z.number().step(1).min(1).default(DEFAULT_MAX_INSTRUCTION_BYTES),
130 reconnect: Reconnect,
131 }),
132 z.object({
133 transport: z.const('streamable-http'),
134 serverName: z.string().required().pattern(SERVER_NAME_PATTERN),
135 url: z.string().required(),
136 headers: z.dict(String).default({}),
137 toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
138 failOnStartupError: z.boolean().default(false),
139 maxInstructionBytes: z.number().step(1).min(1).default(DEFAULT_MAX_INSTRUCTION_BYTES),
140 reconnect: Reconnect,
141 }),
142]) as z<ConfigInput, Config>
143
144// ---- Plugin apply ----
145
146/**
147 * Connect one MCP server and publish its initial tool generation before activation.
148 * This entry remains explicitly `async`: Cordis treats a prototype-bearing
149 * ordinary function as a constructor, whose returned Promise is not startup work.
150 * @param ctx - plugin context carrying the tool registry.
151 * @param config - resolved transport and server namespace configuration.
152 * @returns startup readiness after connection and initial tool discovery settle.
153 */
154export async function apply(ctx: Context, config: Config): Promise<void> {
155 // Fail loud at load: reconnect misconfiguration (including programmatic
156 // construction that bypassed Schemastery) rejects THIS instance before any
157 // effect registers.
158 const reconnect = resolveReconnectPolicy(config.reconnect, `mcp-client(${config.serverName}): reconnect`)
159
160 // Reserve the namespace next: a duplicate `serverName` fails THIS instance
161 // at load with an actionable error and leaves the earlier instance intact.
162 ctx.effect(() => {
163 const owner = scopeOf(ctx) ?? ctx.root
164 let names = activeServerNames.get(owner)
165 if (!names) {
166 names = new Set()
167 activeServerNames.set(owner, names)
168 }
169 if (names.has(config.serverName)) {
170 throw new Error(
171 `mcp-client: serverName "${config.serverName}" is already in use by another mcp-client instance — pick a unique serverName in cordis.yml`,
172 )
173 }
174 names.add(config.serverName)
175 return () => void names.delete(config.serverName)
176 }, 'mcp-client.serverName')
177
178 // The supervisor owns the client/transport generations, the reconnect
179 // loop, and the live tool registrations; disposal stops reconnection,
180 // quiesces in-flight work, and unregisters the current generation.
181 const connection = startConnection(ctx, config, reconnect)
182 registerServerContext(ctx, config.serverName, connection)
183 let stopping: Promise<void> | undefined
184 const dispose = (): Promise<void> => stopping ??= connection.dispose()
185 // Cordis announces unload before awaiting an unfinished apply(). Closing
186 // the transport here releases startup requests that are still awaiting a reply.
187 // oxlint-disable-next-line typescript/no-misused-promises -- Cordis contains observer failures; the effect also awaits this promise.
188 ctx.on('internal/plugin', (fiber) => {
189 if (fiber !== ctx.fiber || fiber.uid !== null) return
190 return dispose()
191 }, { global: true })
192 ctx.effect(() => dispose, 'mcp-client.connection')
193
194 // Block plugin activation on the initial connection + tool discovery so
195 // Cordis consumers observe the tools immediately after the fiber activates.
196 // When failOnStartupError is true, a failed initial attempt rejects the
197 // fiber (Cordis rolls it back); otherwise the error is logged and the
198 // supervisor enters its reconnect loop.
199 const outcome = await connection.ready
200 if (outcome.error !== undefined && config.failOnStartupError) {
201 throw new Error(`mcp-client(${config.serverName}): initial connection or tool synchronization failed`, { cause: outcome.error })
202 }
203}