返回源码地图

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

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

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

1/**
2 * Connection supervisor: owns the MCP client/transport generations for one
3 * plugin instance, keeps the harness tool registry in sync with the live
4 * generation, and — when the connection drops — restarts the configured
5 * server with bounded exponential backoff.
6 *
7 * One outage shares one attempt budget (`maxAttempts` consecutive failed
8 * attempts, delays doubling from `initialDelayMs` up to `maxDelayMs`). A
9 * connection that stays up past the stability window closes the outage, so
10 * the next disconnect starts a fresh budget while a crash-looping server —
11 * even one whose connects briefly succeed — still exhausts the cap instead of
12 * restarting forever. Exhaustion unregisters the server's tools and stops;
13 * disposal (including HMR) is the only way back from that state.
14 *
15 * @module
16 */
17
18import { Client, type Transport } from '@modelcontextprotocol/client'
19import type { Context } from '@deepseek-ai/cordis'
20import { assertNever, type JsonValue } from '@deepseek-ai/dsh-util-values'
21import type { ServerContext } from './server-context.ts'
22import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
23import { createTransport } from './transport.ts'
24import { syncTools } from './tools.ts'
25import type { ToolBridgeOptions, ToolDisposers } from './tools.ts'
26import type { Config } from './index.ts'
27
28/** Automatic reconnect policy for one MCP server connection. */
29export interface ReconnectConfig {
30 /** Reconnect automatically after a lost connection (default true). */
31 enabled?: boolean
32 /** First reconnect delay in milliseconds; doubles per consecutive failed attempt (default 500). */
33 initialDelayMs?: number
34 /** Backoff ceiling in milliseconds; also the uptime after which the attempt budget resets (default 30000). */
35 maxDelayMs?: number
36 /** Consecutive failed attempts per outage before giving up for good (default 10). */
37 maxAttempts?: number
38}
39
40/** Defaults shared by the Config schema and {@link resolveReconnectPolicy}. */
41export const RECONNECT_DEFAULTS: Required<ReconnectConfig> = Object.freeze({
42 enabled: true,
43 initialDelayMs: 500,
44 maxDelayMs: 30_000,
45 maxAttempts: 10,
46})
47
48/** Default UTF-8 byte limit for attributed server instructions. */
49export const DEFAULT_MAX_INSTRUCTION_BYTES = 32_768
50
51// The SDK's stdio transport owns two two-second termination grace periods.
52// Keep one additional second for the process-close event that proves the old
53// generation is gone; timing out fails closed instead of overlapping children.
54const GENERATION_CLOSE_TIMEOUT_MS = 5_000
55
56/** Fully resolved reconnect policy captured at plugin load. */
57export type ResolvedReconnectPolicy = Readonly<Required<ReconnectConfig>>
58
59/**
60 * The one explicit resolve step from raw reconnect config to the policy the
61 * supervisor runs. Programmatic construction may bypass Schemastery
62 * normalization, so every default and bound is re-judged here — misconfiguration
63 * fails the plugin instance at load.
64 *
65 * @param config - Raw `reconnect` config; omission uses the defaults.
66 * @param path - Diagnostic prefix naming the config location in thrown messages.
67 * @returns The frozen resolved policy.
68 */
69export function resolveReconnectPolicy(config: ReconnectConfig | undefined, path: string): ResolvedReconnectPolicy {
70 if (config !== undefined) {
71 for (const key of Object.keys(config)) {
72 if (!Object.hasOwn(RECONNECT_DEFAULTS, key)) throw new Error(`${path}.${key} is not a reconnect option`)
73 }
74 }
75 const enabled = config?.enabled ?? RECONNECT_DEFAULTS.enabled
76 const initialDelayMs = config?.initialDelayMs ?? RECONNECT_DEFAULTS.initialDelayMs
77 const maxDelayMs = config?.maxDelayMs ?? RECONNECT_DEFAULTS.maxDelayMs
78 const maxAttempts = config?.maxAttempts ?? RECONNECT_DEFAULTS.maxAttempts
79 /* jscpd:ignore-start — domain-specific delay validation parallels llm retry-policy; not extractable */
80 if (!Number.isFinite(initialDelayMs) || initialDelayMs <= 0 || initialDelayMs > MAX_TIMER_DELAY_MS) {
81 throw new Error(`${path}.initialDelayMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`)
82 }
83 if (!Number.isFinite(maxDelayMs) || maxDelayMs <= 0 || maxDelayMs > MAX_TIMER_DELAY_MS) {
84 throw new Error(`${path}.maxDelayMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`)
85 }
86 if (initialDelayMs > maxDelayMs) {
87 throw new Error(`${path}.initialDelayMs must be less than or equal to maxDelayMs`)
88 }
89 if (!Number.isInteger(maxAttempts) || maxAttempts < 1) {
90 throw new Error(`${path}.maxAttempts must be a positive integer`)
91 }
92 /* jscpd:ignore-end */
93 return Object.freeze({ enabled, initialDelayMs, maxDelayMs, maxAttempts })
94}
95
96/** Result from the initial connection attempt, for startup-await semantics. */
97export interface ConnectionOutcome {
98 /** If the initial connection or tool sync failed, the error; otherwise absent. */
99 error?: unknown
100}
101
102/** Handle for one plugin instance's supervised connection. */
103export interface ConnectionHandle extends ServerContext {
104 /**
105 * Settles when the first connection attempt completes (success or failure).
106 * The supervisor enters its reconnect loop regardless; the caller decides
107 * whether a failed startup is fatal via `failOnStartupError`.
108 */
109 ready: Promise<ConnectionOutcome>
110 /**
111 * Stop reconnection, close the negotiating transport or live client, wait
112 * for the in-flight attempt and queued tool syncs to quiesce, then
113 * unregister every tool this server still owns.
114 */
115 dispose(): Promise<void>
116}
117
118/**
119 * Start the supervised connection for one MCP server and keep it alive per
120 * the reconnect policy.
121 *
122 * @param ctx - Cordis context providing the `tools` registry and logger.
123 * @param config - Resolved plugin config selecting the transport and server identity.
124 * @param policy - Resolved reconnect policy from {@link resolveReconnectPolicy}.
125 * @returns Handle with a `ready` promise for startup-await and a `dispose` for teardown.
126 */
127export function startConnection(ctx: Context, config: Config, policy: ResolvedReconnectPolicy): ConnectionHandle {
128 const label = `mcp-client(${config.serverName})`
129 const incompleteDisposalMessage = `${label}: transport closure could not be confirmed during disposal — server shutdown may be incomplete`
130 const opts: ToolBridgeOptions = {
131 registrationFailure: 'contain',
132 serverName: config.serverName,
133 toolCallTimeoutMs: config.toolCallTimeoutMs,
134 }
135 // The initial sync uses 'throw' when failOnStartupError is configured, so
136 // a registration conflict propagates to the startup-await path. Re-syncs
137 // and reconnect syncs always contain conflicts.
138 const startupOpts: ToolBridgeOptions = config.failOnStartupError
139 ? { ...opts, registrationFailure: 'throw' }
140 : opts
141
142 let disposed = false
143 const maxInstructionBytes = config.maxInstructionBytes ?? DEFAULT_MAX_INSTRUCTION_BYTES
144 let serverInstructions = ''
145 /** Current generation: the connecting or connected client; undefined during backoff waits and after final failure. */
146 let client: Client | undefined
147 /** Transport-aware close operation paired with {@link client}. */
148 let closeClient: (() => Promise<boolean>) | undefined
149 /** Live tool registrations owned by this server; only {@link enqueueSync} and dispose swap it. */
150 let disposers: ToolDisposers = new Map()
151 let reconnectTimer: NodeJS.Timeout | undefined
152 /** Consecutive failed connection attempts within the current outage. */
153 let failedAttempts = 0
154 /** When the current generation finished connect + initial sync; undefined while down. */
155 let connectedAt: number | undefined
156 /** The real error from the first connection attempt, for startup-await diagnostics. */
157 let firstAttemptError: unknown
158
159 /** A generation may act only while it is the current one on a live plugin. */
160 const isCurrent = (generation: Client): boolean => !disposed && client === generation
161
162 /**
163 * Serializes every syncTools call — initial syncs and notification re-syncs
164 * across all generations — so two syncs can never interleave their
165 * dispose-previous/register-next swap (which would double-dispose one
166 * generation and leak another).
167 */
168 let syncChain: Promise<void> = Promise.resolve()
169 function enqueueSync(generation: Client, syncOpts: ToolBridgeOptions = opts): Promise<void> {
170 const run = syncChain.then(async () => {
171 if (!isCurrent(generation)) return
172 disposers = await syncTools(generation, ctx, syncOpts, disposers)
173 })
174 // The chain tail must survive a failed sync; the enqueuing caller owns reporting.
175 syncChain = run.catch(() => {})
176 return run
177 }
178
179 /** One disconnect decision per generation: the isCurrent guard makes racing close/error signals idempotent. */
180 function generationDown(generation: Client): void {
181 if (!isCurrent(generation)) return
182 client = undefined
183 closeClient = undefined
184 scheduleReconnect()
185 }
186
187 /** Decide retry ownership after a failed connection's close barrier settles. */
188 function settleFailedGeneration(generation: Client, quiesced: boolean): void {
189 if (!isCurrent(generation)) return
190 if (!quiesced) {
191 client = undefined
192 closeClient = undefined
193 ctx.logger.error(`${label}: failed generation could not confirm transport closure — reconnect stopped to avoid overlapping server processes; reload the plugin or restart the Host to retry`)
194 return
195 }
196 generationDown(generation)
197 }
198
199 /** Wait for the transport-owned close signal without letting a broken transport wedge teardown forever. */
200 function waitForClose(closed: Promise<void>): Promise<boolean> {
201 return new Promise((resolve) => {
202 const timeout = setTimeout(() => { resolve(false) }, GENERATION_CLOSE_TIMEOUT_MS)
203 timeout.unref()
204 void closed.then(() => {
205 clearTimeout(timeout)
206 resolve(true)
207 })
208 })
209 }
210
211 function scheduleReconnect(): void {
212 const lostEstablishedConnection = connectedAt !== undefined
213 if (!policy.enabled) {
214 const message = lostEstablishedConnection
215 ? 'connection lost and reconnect is disabled — registered tools will fail until an HMR reload or Host restart'
216 : 'connection failed and reconnect is disabled — no tools were registered; reload the plugin or restart the Host to connect'
217 ctx.logger.error(`${label}: ${message}`)
218 return
219 }
220 // A connection that stayed up past the stability window (= maxDelayMs, the
221 // longest backoff spacing) ended the previous outage: start a fresh budget.
222 if (connectedAt !== undefined && Date.now() - connectedAt >= policy.maxDelayMs) failedAttempts = 0
223 connectedAt = undefined
224 failedAttempts += 1
225 if (failedAttempts > policy.maxAttempts) {
226 // Enqueue the give-up disposal so it cannot race an in-flight sync's
227 // phase-2 swap (which checks isCurrent inside the queue).
228 syncChain = syncChain.then(() => {
229 for (const dispose of disposers.values()) dispose()
230 disposers = new Map()
231 serverInstructions = ''
232 })
233 ctx.logger.error(`${label}: giving up after ${policy.maxAttempts} consecutive failed reconnect attempts — tools unregistered; reload the plugin or restart the Host to reconnect`)
234 return
235 }
236 const delayMs = Math.min(policy.maxDelayMs, policy.initialDelayMs * 2 ** (failedAttempts - 1))
237 const action = lostEstablishedConnection ? 'connection lost; reconnecting' : 'connection failed; retrying'
238 ctx.logger.warn(`${label}: ${action} in ${delayMs}ms (attempt ${failedAttempts}/${policy.maxAttempts})`)
239 reconnectTimer = setTimeout(() => {
240 reconnectTimer = undefined
241 settling = connectGeneration(false)
242 }, delayMs)
243 // An armed reconnect timer must never hold the process open on its own.
244 reconnectTimer.unref()
245 }
246
247 /**
248 * One connection attempt: fresh transport + client (the MCP SDK binds a
249 * Protocol to one transport for life), connect, then queue the initial tool
250 * sync. The startup flag belongs to the attempt rather than the shared sync
251 * queue, so an early notification cannot consume strict startup semantics.
252 * Every failure funnels through {@link generationDown}; success arms the
253 * onclose-driven disconnect path. Never rejects.
254 *
255 * @param startup - Whether this is the plugin's activation attempt.
256 */
257 async function connectGeneration(startup: boolean): Promise<void> {
258 const generation = new Client(
259 { name: 'dsh-mcp-client', version: '0.0.1' },
260 {
261 capabilities: {},
262 versionNegotiation: { mode: 'auto' },
263 listChanged: {
264 tools: {
265 autoRefresh: false,
266 debounceMs: 0,
267 onChanged: () => { void refreshTools() },
268 },
269 },
270 },
271 )
272 const closed: PromiseWithResolvers<void> = Promise.withResolvers()
273 let attemptSettled = false
274 let closeObserved = false
275 let transport: Transport | undefined
276 const hasClosed = (): boolean => closeObserved
277 client = generation
278 closeClient = closeGeneration
279 generation.onclose = () => {
280 closeObserved = true
281 closed.resolve()
282 // A failed connect owns its close barrier in the catch path below. An
283 // established generation can transition down directly from this signal.
284 if (attemptSettled) generationDown(generation)
285 }
286 /** Unattached probes close through their transport; attached clients must also report transport closure. */
287 async function closeGeneration(): Promise<boolean> {
288 const attached = generation.transport !== undefined
289 try {
290 await (attached ? generation.close() : transport?.close())
291 } catch (_error) {
292 if (!attached) return hasClosed()
293 }
294 return !attached || hasClosed() || await waitForClose(closed.promise)
295 }
296 async function refreshTools(): Promise<void> {
297 if (!isCurrent(generation)) return
298 ctx.logger.info(`${label}: tool list changed, re-syncing`)
299 try {
300 await enqueueSync(generation)
301 } catch (error) {
302 if (!disposed) ctx.logger.error(`${label}: tool re-sync failed: ${String(error)}`)
303 }
304 }
305 let instructions: string
306 try {
307 transport = createTransport(config)
308 await generation.connect(transport)
309 if (hasClosed()) {
310 attemptSettled = true
311 generationDown(generation)
312 return
313 }
314 if (!isCurrent(generation)) {
315 if (!await closeGeneration()) ctx.logger.error(incompleteDisposalMessage)
316 return
317 }
318 const serverText = generation.getInstructions()?.trimEnd() ?? ''
319 instructions = serverText ? `### MCP server: ${config.serverName}\n\n${serverText}` : ''
320 if (Buffer.byteLength(instructions) > maxInstructionBytes) {
321 throw new Error(`${label}: server instructions exceed maxInstructionBytes (${maxInstructionBytes})`)
322 }
323 await enqueueSync(generation, startup ? startupOpts : opts)
324 } catch (error) {
325 if (firstAttemptError === undefined) firstAttemptError = error
326 // Disposal clears current ownership before it closes the generation, so
327 // only a live supervisor reports an attempt failure.
328 if (isCurrent(generation)) ctx.logger.warn(`${label}: connection attempt failed: ${String(error)}`)
329 const quiesced = await closeGeneration()
330 attemptSettled = true
331 settleFailedGeneration(generation, quiesced)
332 return
333 }
334 attemptSettled = true
335 if (hasClosed()) {
336 generationDown(generation)
337 return
338 }
339 if (!isCurrent(generation)) return
340 serverInstructions = instructions
341 connectedAt = Date.now()
342 if (failedAttempts > 0) ctx.logger.info(`${label}: reconnected and re-synced tools (attempt ${failedAttempts}/${policy.maxAttempts})`)
343 }
344
345 /** The in-flight (or last settled) connection attempt; dispose awaits it for quiescence. */
346 let settling = connectGeneration(true)
347
348 // The ready promise settles when the first attempt finishes (regardless of
349 // success). If the first attempt fails and reconnect is enabled, the
350 // supervisor is already scheduling a retry — ready just reports the outcome.
351 const ready: Promise<ConnectionOutcome> = settling.then(() => {
352 // After settling: if client is set the initial connect+sync succeeded.
353 // If not, the supervisor either scheduled a retry (error logged) or gave
354 // up (error logged). Either way the outcome is reported with the real error.
355 // Note: settling.then() is a microtask; stdio onclose is a macrotask — so
356 // a server that crashes AFTER a successful initial sync cannot flip client
357 // to undefined before this continuation runs.
358 if (client !== undefined) return {}
359 /* v8 ignore next -- defensive: firstAttemptError is always set when connect/sync fails */
360 return { error: firstAttemptError ?? new Error(`${label}: initial connection failed`) }
361 })
362
363 return {
364 ready,
365 instructions: () => serverInstructions,
366 resources: {
367 async request(request, exec): Promise<JsonValue> {
368 const generation = client
369 if (!generation || connectedAt === undefined) throw new Error(`${label}: server is disconnected`)
370 const options = { signal: exec.signal, timeout: config.toolCallTimeoutMs }
371 switch (request.method) {
372 case 'resources/list':
373 return await generation.listResources(
374 request.cursor === undefined ? undefined : { cursor: request.cursor }, options,
375 ) as JsonValue
376 case 'resources/templates/list':
377 return await generation.listResourceTemplates(
378 request.cursor === undefined ? undefined : { cursor: request.cursor }, options,
379 ) as JsonValue
380 case 'resources/read':
381 return await generation.readResource({ uri: request.uri }, options) as JsonValue
382 /* v8 ignore next 2 -- resource requests are the closed, typed tool operation union */
383 default:
384 return assertNever(request)
385 }
386 },
387 },
388 async dispose(): Promise<void> {
389 disposed = true
390 serverInstructions = ''
391 if (reconnectTimer !== undefined) {
392 clearTimeout(reconnectTimer)
393 reconnectTimer = undefined
394 }
395 const close = closeClient
396 client = undefined
397 closeClient = undefined
398 if (close !== undefined && !await close()) {
399 ctx.logger.error(incompleteDisposalMessage)
400 }
401 // Quiesce, don't just request it: the in-flight attempt enqueues its
402 // sync before settling, so awaiting both leaves `disposers` final.
403 await settling
404 await syncChain
405 for (const dispose of disposers.values()) dispose()
406 disposers = new Map()
407 },
408 }
409}