feat(mcp-client): auto-reconnect with bounded backoff after transport close

A per-instance connection supervisor restarts the original server config
with exponential backoff when the transport closes, re-runs tool discovery
on success, and atomically replaces the previous generation. Default policy
retries for ~2.5 minutes (10 attempts, 500ms→30s doubling) before giving up
and unregistering the server's tools.

New config block reconnect { enabled, initialDelayMs, maxDelayMs, maxAttempts }
on both transports; misconfiguration fails plugin load. A connection that
survives past the stability window (maxDelayMs) resets the attempt budget,
so occasional crashes recover indefinitely while a crash loop still exhausts
the cap.

Integrates with the upstream failOnStartupError: the initial sync uses
registrationFailure:'throw' when that flag is set so a squatted namespace
still rejects activation.

Fixes #1746
This commit is contained in:
lintianle
2026-08-10 18:50:03 +08:00
parent d2321d210a
commit 00f68d7e93
24 changed files with 1044 additions and 94 deletions

View File

@@ -0,0 +1,292 @@
/**
* Connection supervisor: owns the MCP client/transport generations for one
* plugin instance, keeps the harness tool registry in sync with the live
* generation, and — when the connection drops — restarts the configured
* server with bounded exponential backoff.
*
* One outage shares one attempt budget (`maxAttempts` consecutive failed
* attempts, delays doubling from `initialDelayMs` up to `maxDelayMs`). A
* connection that stays up past the stability window closes the outage, so
* the next disconnect starts a fresh budget while a crash-looping server —
* even one whose connects briefly succeed — still exhausts the cap instead of
* restarting forever. Exhaustion unregisters the server's tools and stops;
* disposal (including HMR) is the only way back from that state.
*
* @module
*/
import { Client } from '@modelcontextprotocol/sdk/client/index.js'
import { ToolListChangedNotificationSchema } from '@modelcontextprotocol/sdk/types.js'
import type { Context } from 'cordis'
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
import { createTransport } from './transport.ts'
import { syncTools } from './tools.ts'
import type { ToolBridgeOptions, ToolDisposers } from './tools.ts'
import type { Config } from './index.ts'
/** Automatic reconnect policy for one MCP server connection. */
export interface ReconnectConfig {
/** Reconnect automatically after a lost connection (default true). */
enabled?: boolean
/** First reconnect delay in milliseconds; doubles per consecutive failed attempt (default 500). */
initialDelayMs?: number
/** Backoff ceiling in milliseconds; also the uptime after which the attempt budget resets (default 30000). */
maxDelayMs?: number
/** Consecutive failed attempts per outage before giving up for good (default 10). */
maxAttempts?: number
}
/** Defaults shared by the Config schema and {@link resolveReconnectPolicy}. */
export const RECONNECT_DEFAULTS: Required<ReconnectConfig> = Object.freeze({
enabled: true,
initialDelayMs: 500,
maxDelayMs: 30_000,
maxAttempts: 10,
})
/** Fully resolved reconnect policy captured at plugin load. */
export type ResolvedReconnectPolicy = Readonly<Required<ReconnectConfig>>
/**
* The one explicit resolve step from raw reconnect config to the policy the
* supervisor runs. Programmatic construction may bypass Schemastery
* normalization, so every default and bound is re-judged here — misconfiguration
* fails the plugin instance at load.
*
* @param config - Raw `reconnect` config; omission uses the defaults.
* @param path - Diagnostic prefix naming the config location in thrown messages.
* @returns The frozen resolved policy.
*/
export function resolveReconnectPolicy(config: ReconnectConfig | undefined, path: string): ResolvedReconnectPolicy {
if (config !== undefined) {
for (const key of Object.keys(config)) {
if (!Object.hasOwn(RECONNECT_DEFAULTS, key)) throw new Error(`${path}.${key} is not a reconnect option`)
}
}
const enabled = config?.enabled ?? RECONNECT_DEFAULTS.enabled
const initialDelayMs = config?.initialDelayMs ?? RECONNECT_DEFAULTS.initialDelayMs
const maxDelayMs = config?.maxDelayMs ?? RECONNECT_DEFAULTS.maxDelayMs
const maxAttempts = config?.maxAttempts ?? RECONNECT_DEFAULTS.maxAttempts
/* jscpd:ignore-start — domain-specific delay validation parallels llm retry-policy; not extractable */
if (!Number.isFinite(initialDelayMs) || initialDelayMs <= 0 || initialDelayMs > MAX_TIMER_DELAY_MS) {
throw new Error(`${path}.initialDelayMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`)
}
if (!Number.isFinite(maxDelayMs) || maxDelayMs <= 0 || maxDelayMs > MAX_TIMER_DELAY_MS) {
throw new Error(`${path}.maxDelayMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`)
}
if (initialDelayMs > maxDelayMs) {
throw new Error(`${path}.initialDelayMs must be less than or equal to maxDelayMs`)
}
if (!Number.isInteger(maxAttempts) || maxAttempts < 1) {
throw new Error(`${path}.maxAttempts must be a positive integer`)
}
/* jscpd:ignore-end */
return Object.freeze({ enabled, initialDelayMs, maxDelayMs, maxAttempts })
}
/** Result from the initial connection attempt, for startup-await semantics. */
export interface ConnectionOutcome {
/** If the initial connection or tool sync failed, the error; otherwise absent. */
error?: unknown
}
/** Handle for one plugin instance's supervised connection. */
export interface ConnectionHandle {
/**
* Settles when the first connection attempt completes (success or failure).
* The supervisor enters its reconnect loop regardless; the caller decides
* whether a failed startup is fatal via `failOnStartupError`.
*/
ready: Promise<ConnectionOutcome>
/**
* Stop reconnection, close the live client, wait for the in-flight attempt
* and queued tool syncs to quiesce, then unregister every tool this server
* still owns.
*/
dispose(): Promise<void>
}
/**
* Start the supervised connection for one MCP server and keep it alive per
* the reconnect policy.
*
* @param ctx - Cordis context providing the `tools` registry and logger.
* @param config - Resolved plugin config selecting the transport and server identity.
* @param policy - Resolved reconnect policy from {@link resolveReconnectPolicy}.
* @returns Handle with a `ready` promise for startup-await and a `dispose` for teardown.
*/
export function startConnection(ctx: Context, config: Config, policy: ResolvedReconnectPolicy): ConnectionHandle {
const label = `mcp-client(${config.serverName})`
const opts: ToolBridgeOptions = {
registrationFailure: 'contain',
serverName: config.serverName,
toolCallTimeoutMs: config.toolCallTimeoutMs,
}
// The initial sync uses 'throw' when failOnStartupError is configured, so
// a registration conflict propagates to the startup-await path. Re-syncs
// and reconnect syncs always contain conflicts.
const startupOpts: ToolBridgeOptions = config.failOnStartupError
? { ...opts, registrationFailure: 'throw' }
: opts
let isFirstSync = true
let disposed = false
/** Current generation: the connecting or connected client; undefined during backoff waits and after final failure. */
let client: Client | undefined
/** Live tool registrations owned by this server; only {@link enqueueSync} and dispose swap it. */
let disposers: ToolDisposers = new Map()
let reconnectTimer: NodeJS.Timeout | undefined
/** Consecutive failed connection attempts within the current outage. */
let failedAttempts = 0
/** When the current generation finished connect + initial sync; undefined while down. */
let connectedAt: number | undefined
/** The real error from the first connection attempt, for startup-await diagnostics. */
let firstAttemptError: unknown
/** A generation may act only while it is the current one on a live plugin. */
const isCurrent = (generation: Client): boolean => !disposed && client === generation
/**
* Serializes every syncTools call — initial syncs and notification re-syncs
* across all generations — so two syncs can never interleave their
* dispose-previous/register-next swap (which would double-dispose one
* generation and leak another).
*/
let syncChain: Promise<void> = Promise.resolve()
function enqueueSync(generation: Client): Promise<void> {
const syncOpts = isFirstSync ? startupOpts : opts
isFirstSync = false
const run = syncChain.then(async () => {
if (!isCurrent(generation)) return
disposers = await syncTools(generation, ctx, syncOpts, disposers)
})
// The chain tail must survive a failed sync; the enqueuing caller owns reporting.
syncChain = run.catch(() => {})
return run
}
/** One disconnect decision per generation: the isCurrent guard makes racing close/error signals idempotent. */
function generationDown(generation: Client): void {
if (!isCurrent(generation)) return
client = undefined
scheduleReconnect()
}
function scheduleReconnect(): void {
if (!policy.enabled) {
const detail = connectedAt !== undefined
? 'registered tools will fail until an HMR reload or Host restart'
: 'no tools were registered; reload the plugin or restart the Host to connect'
ctx.logger.error(`${label}: connection lost and reconnect is disabled — ${detail}`)
return
}
// A connection that stayed up past the stability window (= maxDelayMs, the
// longest backoff spacing) ended the previous outage: start a fresh budget.
if (connectedAt !== undefined && Date.now() - connectedAt >= policy.maxDelayMs) failedAttempts = 0
connectedAt = undefined
failedAttempts += 1
if (failedAttempts > policy.maxAttempts) {
// Enqueue the give-up disposal so it cannot race an in-flight sync's
// phase-2 swap (which checks isCurrent inside the queue).
syncChain = syncChain.then(() => {
for (const dispose of disposers.values()) dispose()
disposers = new Map()
})
ctx.logger.error(`${label}: giving up after ${policy.maxAttempts} consecutive failed reconnect attempts — tools unregistered; reload the plugin or restart the Host to reconnect`)
return
}
const delayMs = Math.min(policy.maxDelayMs, policy.initialDelayMs * 2 ** (failedAttempts - 1))
ctx.logger.warn(`${label}: connection lost; reconnecting in ${delayMs}ms (attempt ${failedAttempts}/${policy.maxAttempts})`)
reconnectTimer = setTimeout(() => {
reconnectTimer = undefined
settling = connectGeneration()
}, delayMs)
// An armed reconnect timer must never hold the process open on its own.
reconnectTimer.unref()
}
/**
* One connection attempt: fresh transport + client (the MCP SDK binds a
* Protocol to one transport for life), connect, then queue the initial tool
* sync. Every failure funnels through {@link generationDown}; success arms
* the onclose-driven disconnect path. Never rejects.
*/
async function connectGeneration(): Promise<void> {
const generation = new Client(
{ name: 'dsh-mcp-client', version: '0.0.1' },
{ capabilities: {} },
)
client = generation
generation.onclose = () => { generationDown(generation) }
// Registered before connect so a list change during the initial sync is
// queued behind it rather than dropped.
generation.setNotificationHandler(
ToolListChangedNotificationSchema,
async () => {
if (!isCurrent(generation)) return
ctx.logger.info(`${label}: tool list changed, re-syncing`)
try {
await enqueueSync(generation)
} catch (error) {
// Fetch-phase failure: the previous generation is still registered
// and `disposers` still owns it — keep serving the last good list.
if (!disposed) ctx.logger.error(`${label}: tool re-sync failed: ${String(error)}`)
}
},
)
try {
await generation.connect(createTransport(config))
await enqueueSync(generation)
} catch (error) {
if (firstAttemptError === undefined) firstAttemptError = error
// When the transport closed first, onclose already logged and scheduled.
if (isCurrent(generation)) ctx.logger.warn(`${label}: connection attempt failed: ${String(error)}`)
try { await generation.close() } catch { /* transport already gone */ }
generationDown(generation)
return
}
if (!isCurrent(generation)) return
connectedAt = Date.now()
if (failedAttempts > 0) ctx.logger.info(`${label}: reconnected and re-synced tools (attempt ${failedAttempts}/${policy.maxAttempts})`)
}
/** The in-flight (or last settled) connection attempt; dispose awaits it for quiescence. */
let settling = connectGeneration()
// The ready promise settles when the first attempt finishes (regardless of
// success). If the first attempt fails and reconnect is enabled, the
// supervisor is already scheduling a retry — ready just reports the outcome.
const ready: Promise<ConnectionOutcome> = settling.then(() => {
// After settling: if client is set the initial connect+sync succeeded.
// If not, the supervisor either scheduled a retry (error logged) or gave
// up (error logged). Either way the outcome is reported with the real error.
// Note: settling.then() is a microtask; stdio onclose is a macrotask — so
// a server that crashes AFTER a successful initial sync cannot flip client
// to undefined before this continuation runs.
if (client !== undefined) return {}
/* v8 ignore next -- defensive: firstAttemptError is always set when connect/sync fails */
return { error: firstAttemptError ?? new Error(`${label}: initial connection failed`) }
})
return {
ready,
async dispose(): Promise<void> {
disposed = true
if (reconnectTimer !== undefined) {
clearTimeout(reconnectTimer)
reconnectTimer = undefined
}
const current = client
client = undefined
if (current !== undefined) {
try { await current.close() } catch { /* transport already gone */ }
}
// Quiesce, don't just request it: the in-flight attempt enqueues its
// sync before settling, so awaiting both leaves `disposers` final.
await settling
await syncChain
for (const dispose of disposers.values()) dispose()
disposers = new Map()
},
}
}

View File

@@ -15,14 +15,14 @@
import type { Context } from 'cordis'
import z from 'schemastery'
import { Client } from '@modelcontextprotocol/sdk/client/index.js'
import { ToolListChangedNotificationSchema } from '@modelcontextprotocol/sdk/types.js'
import { createTransport } from './transport.ts'
import { syncTools } from './tools.ts'
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
import { RECONNECT_DEFAULTS, resolveReconnectPolicy, startConnection } from './connection.ts'
import type { ReconnectConfig } from './connection.ts'
// Side-effect type import: declaration-merges `ctx.tools` onto Context.
import type {} from '@deepseek-ai/dsh-tools'
export type { McpResult } from './tools.ts'
export type { ReconnectConfig, ResolvedReconnectPolicy } from './connection.ts'
/** Cordis plugin name used by loader diagnostics. */
export const name = 'mcp-client'
@@ -74,6 +74,8 @@ export interface StdioConfig {
toolCallTimeoutMs: number
/** Fail plugin activation when the initial connection or tool synchronization fails. */
failOnStartupError: boolean
/** Automatic reconnect policy after a lost connection; omission uses the defaults. */
reconnect?: ReconnectConfig
}
/** Config for connecting to an MCP server over Streamable HTTP (SSE). */
@@ -94,11 +96,20 @@ export interface StreamableHttpConfig {
toolCallTimeoutMs: number
/** Fail plugin activation when the initial connection or tool synchronization fails. */
failOnStartupError: boolean
/** Automatic reconnect policy after a lost connection; omission uses the defaults. */
reconnect?: ReconnectConfig
}
/** Configuration for one stdio or Streamable HTTP MCP server. */
export type Config = StdioConfig | StreamableHttpConfig
const Reconnect: z<ReconnectConfig> = z.object({
enabled: z.boolean().default(RECONNECT_DEFAULTS.enabled),
initialDelayMs: z.number().min(1).max(MAX_TIMER_DELAY_MS).default(RECONNECT_DEFAULTS.initialDelayMs),
maxDelayMs: z.number().min(1).max(MAX_TIMER_DELAY_MS).default(RECONNECT_DEFAULTS.maxDelayMs),
maxAttempts: z.number().step(1).min(1).max(Number.MAX_SAFE_INTEGER).default(RECONNECT_DEFAULTS.maxAttempts),
})
export const Config = z.union([
z.object({
transport: z.const('stdio'),
@@ -109,6 +120,7 @@ export const Config = z.union([
cwd: z.string().default(''),
toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
failOnStartupError: z.boolean().default(false),
reconnect: Reconnect,
}),
z.object({
transport: z.const('streamable-http'),
@@ -117,6 +129,7 @@ export const Config = z.union([
headers: z.dict(String).default({}),
toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
failOnStartupError: z.boolean().default(false),
reconnect: Reconnect,
}),
]) as unknown as z<Config>
@@ -131,7 +144,12 @@ export const Config = z.union([
* @returns startup readiness after connection and initial tool discovery settle.
*/
export async function apply(ctx: Context, config: Config): Promise<void> {
// Reserve the namespace first: a duplicate `serverName` fails THIS instance
// Fail loud at load: reconnect misconfiguration (including programmatic
// construction that bypassed Schemastery) rejects THIS instance before any
// effect registers.
const reconnect = resolveReconnectPolicy(config.reconnect, `mcp-client(${config.serverName}): reconnect`)
// Reserve the namespace next: a duplicate `serverName` fails THIS instance
// at load with an actionable error and leaves the earlier instance intact.
ctx.effect(() => {
let names = activeServerNames.get(ctx.root)
@@ -148,58 +166,22 @@ export async function apply(ctx: Context, config: Config): Promise<void> {
return () => void names.delete(config.serverName)
}, 'mcp-client.serverName')
const transport = createTransport(config)
const client = new Client(
{ name: 'dsh-mcp-client', version: '0.0.1' },
{ capabilities: {} },
)
// The supervisor owns the client/transport generations, the reconnect
// loop, and the live tool registrations; disposal stops reconnection,
// quiesces in-flight work, and unregisters the current generation.
const connection = startConnection(ctx, config, reconnect)
const opts = {
registrationFailure: 'contain' as const,
serverName: config.serverName,
toolCallTimeoutMs: config.toolCallTimeoutMs,
}
// Connect and set up tools. `ready` always settles to an outcome so rollback
// can close a partially opened client even when strict startup later rejects.
// Its accessor returns the CURRENT disposer generation, so disposal always
// unregisters the live set, not the first one.
const ready = (async () => {
await client.connect(transport)
let disposers = await syncTools(client, ctx, {
...opts,
registrationFailure: config.failOnStartupError ? 'throw' : 'contain',
}, new Map())
client.setNotificationHandler(
ToolListChangedNotificationSchema,
async () => {
ctx.logger.info(`mcp-client(${config.serverName}): tool list changed, re-syncing`)
try {
disposers = await syncTools(client, ctx, opts, disposers)
} catch (error) {
// Fetch-phase failure: the previous generation is still registered
// and `disposers` still owns it — keep serving the last good list.
ctx.logger.error(`mcp-client(${config.serverName}): tool re-sync failed: ${String(error)}`)
}
},
)
return { getDisposers: () => disposers }
})().catch((error: unknown) => {
ctx.logger.error(`mcp-client(${config.serverName}): startup failed: ${String(error)}`)
return { getDisposers: () => new Map<string, () => void>(), error }
})
ctx.effect(() => async () => {
const outcome = await ready
for (const dispose of outcome.getDisposers().values()) dispose()
try { await client.close() } catch { /* transport already gone */ }
ctx.effect(() => {
return () => connection.dispose()
}, 'mcp-client.connection')
const outcome = await ready
if ('error' in outcome && config.failOnStartupError) {
// Block plugin activation on the initial connection + tool discovery so
// Cordis consumers observe the tools immediately after the fiber activates.
// When failOnStartupError is true, a failed initial attempt rejects the
// fiber (Cordis rolls it back); otherwise the error is logged and the
// supervisor enters its reconnect loop.
const outcome = await connection.ready
if (outcome.error !== undefined && config.failOnStartupError) {
throw new Error(`mcp-client(${config.serverName}): initial connection or tool synchronization failed`, { cause: outcome.error })
}
}