feat: add MCP client plugin (dsh-mcp-client)
Connects to an external MCP server and registers its tools on ctx.tools. Supports stdio (child process) and Streamable HTTP transports. Credential-shaped env vars are scrubbed before forwarding to child processes. - Plugin lifecycle: connect, sync tools, re-sync on ToolListChanged, dispose unregisters and closes - Full JSDoc on all exports (@param/@returns on functions) - 100% per-file coverage (apply lifecycle, args coercion, env scrubbing) - Config catalog regenerated
This commit is contained in:
128
packages/mcp/mcp-client/src/index.ts
Normal file
128
packages/mcp/mcp-client/src/index.ts
Normal file
@@ -0,0 +1,128 @@
|
||||
/**
|
||||
* MCP client bridge plugin: connects to an external MCP server and registers
|
||||
* its tools on `ctx.tools`. Each plugin instance connects to one MCP server;
|
||||
* load multiple instances in `cordis.yml` for multiple servers.
|
||||
*
|
||||
* Namespace plugin (named exports, no default export). Lifecycle is
|
||||
* effect-scoped: disposal disconnects from the server and unregisters all
|
||||
* tools. HMR hot-swaps by disposing the old instance and creating a new one.
|
||||
*
|
||||
* @module @deepseek-ai/dsh-mcp-client
|
||||
*/
|
||||
|
||||
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'
|
||||
// Side-effect type import: declaration-merges `ctx.tools` onto Context.
|
||||
import type {} from '@deepseek-ai/dsh-tools'
|
||||
|
||||
/** Cordis plugin name used by loader diagnostics. */
|
||||
export const name = 'mcp-client'
|
||||
|
||||
/** Services required by this plugin. */
|
||||
export const inject = ['tools']
|
||||
|
||||
/** Default timeout for individual MCP tool calls (ms). */
|
||||
const DEFAULT_TOOL_CALL_TIMEOUT_MS = 60_000
|
||||
|
||||
// ---- Config ----
|
||||
|
||||
/** Config for connecting to an MCP server via a spawned child process over stdio. */
|
||||
export interface StdioConfig {
|
||||
/** Transport type: spawn a child process and communicate over stdio. */
|
||||
transport: 'stdio'
|
||||
/** Executable to spawn. */
|
||||
command: string
|
||||
/** Arguments passed to the command. */
|
||||
args: string[]
|
||||
/** Extra env vars merged on top of scrubbed ambient env. */
|
||||
env: Record<string, string>
|
||||
/** Working directory for the child process. */
|
||||
cwd: string
|
||||
/** Prefix prepended to each tool name before registration. */
|
||||
toolPrefix: string
|
||||
/** Timeout per callTool invocation (ms). */
|
||||
toolCallTimeoutMs: number
|
||||
}
|
||||
|
||||
/** Config for connecting to an MCP server over Streamable HTTP (SSE). */
|
||||
export interface StreamableHttpConfig {
|
||||
/** Transport type: connect to an MCP server over Streamable HTTP (SSE). */
|
||||
transport: 'streamable-http'
|
||||
/** MCP server URL. */
|
||||
url: string
|
||||
/** Extra headers (e.g. auth tokens). */
|
||||
headers: Record<string, string>
|
||||
/** Prefix prepended to each tool name before registration. */
|
||||
toolPrefix: string
|
||||
/** Timeout per callTool invocation (ms). */
|
||||
toolCallTimeoutMs: number
|
||||
}
|
||||
|
||||
/** Discriminated union of all supported MCP transport configurations. */
|
||||
export type Config = StdioConfig | StreamableHttpConfig
|
||||
|
||||
export const Config = z.union([
|
||||
z.object({
|
||||
transport: z.const('stdio'),
|
||||
command: z.string().required(),
|
||||
args: z.array(String).default([]),
|
||||
env: z.dict(String).default({}),
|
||||
cwd: z.string().default(''),
|
||||
toolPrefix: z.string().default(''),
|
||||
toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
|
||||
}),
|
||||
z.object({
|
||||
transport: z.const('streamable-http'),
|
||||
url: z.string().required(),
|
||||
headers: z.dict(String).default({}),
|
||||
toolPrefix: z.string().default(''),
|
||||
toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
|
||||
}),
|
||||
]) as unknown as z<Config>
|
||||
|
||||
// ---- Plugin apply ----
|
||||
|
||||
export function apply(ctx: Context, config: Config): void {
|
||||
const transport = createTransport(config)
|
||||
const client = new Client(
|
||||
{ name: 'dsh-mcp-client', version: '0.0.1' },
|
||||
{ capabilities: {} },
|
||||
)
|
||||
|
||||
// Connect and set up tools. Errors during connect are logged, not thrown
|
||||
// (the plugin simply has no tools registered).
|
||||
const ready = (async () => {
|
||||
await client.connect(transport)
|
||||
|
||||
let disposers = await syncTools(client, ctx, {
|
||||
toolPrefix: config.toolPrefix,
|
||||
toolCallTimeoutMs: config.toolCallTimeoutMs,
|
||||
}, new Map())
|
||||
|
||||
client.setNotificationHandler(
|
||||
ToolListChangedNotificationSchema,
|
||||
async () => {
|
||||
ctx.logger.info('mcp-client: tool list changed, re-syncing')
|
||||
disposers = await syncTools(client, ctx, {
|
||||
toolPrefix: config.toolPrefix,
|
||||
toolCallTimeoutMs: config.toolCallTimeoutMs,
|
||||
}, disposers)
|
||||
},
|
||||
)
|
||||
|
||||
return disposers
|
||||
})().catch((error: unknown) => {
|
||||
ctx.logger.error(`mcp-client: failed to connect: ${String(error)}`)
|
||||
return new Map<string, () => void>()
|
||||
})
|
||||
|
||||
ctx.effect(() => async () => {
|
||||
const disposers = await ready
|
||||
for (const dispose of disposers.values()) dispose()
|
||||
try { await client.close() } catch { /* transport already gone */ }
|
||||
}, 'mcp-client.connection')
|
||||
}
|
||||
168
packages/mcp/mcp-client/src/tools.ts
Normal file
168
packages/mcp/mcp-client/src/tools.ts
Normal file
@@ -0,0 +1,168 @@
|
||||
/**
|
||||
* Tool bridge: discovers MCP tools, registers them on the harness ToolRegistry,
|
||||
* and handles re-sync when the server's tool list changes.
|
||||
*
|
||||
* @module
|
||||
*/
|
||||
|
||||
import type { Client } from '@modelcontextprotocol/sdk/client/index.js'
|
||||
import type { Context } from 'cordis'
|
||||
import type { ToolDefinition, ToolExecution } from '@deepseek-ai/dsh-tools'
|
||||
|
||||
/** Resolved options relevant to tool bridging. */
|
||||
export interface ToolBridgeOptions {
|
||||
toolPrefix: string
|
||||
toolCallTimeoutMs: number
|
||||
}
|
||||
|
||||
/** State for one sync generation: the current set of disposers keyed by tool name. */
|
||||
type ToolDisposers = Map<string, () => void>
|
||||
|
||||
/**
|
||||
* Sync the MCP server's tool list into the harness ToolRegistry.
|
||||
*
|
||||
* - Calls `client.listTools()` (paginated: drains all pages).
|
||||
* - Registers each tool as a raw `ToolDefinition`.
|
||||
* - On name conflict: logs a warning and skips that tool.
|
||||
* - Returns a disposer map; call each value to unregister.
|
||||
*
|
||||
* @param client - Connected MCP Client instance used to list and call tools.
|
||||
* @param ctx - Cordis context providing the `tools` service for registration.
|
||||
* @param opts - Bridge options: tool name prefix and per-call timeout.
|
||||
* @param previous - Disposer map from a prior sync generation; all entries are
|
||||
* disposed before re-registering.
|
||||
* @returns A map of registered tool names to their unregister disposers.
|
||||
*/
|
||||
export async function syncTools(
|
||||
client: Client,
|
||||
ctx: Context,
|
||||
opts: ToolBridgeOptions,
|
||||
previous: ToolDisposers,
|
||||
): Promise<ToolDisposers> {
|
||||
for (const dispose of previous.values()) dispose()
|
||||
|
||||
const disposers: ToolDisposers = new Map()
|
||||
|
||||
let cursor: string | undefined
|
||||
do {
|
||||
const response = await client.listTools(cursor ? { cursor } : undefined)
|
||||
for (const tool of response.tools) {
|
||||
const registeredName = opts.toolPrefix + tool.name
|
||||
const definition: ToolDefinition = {
|
||||
name: registeredName,
|
||||
description: tool.description ?? '',
|
||||
parameters: tool.inputSchema,
|
||||
execute: createExecutor(client, tool.name, opts),
|
||||
}
|
||||
try {
|
||||
const dispose = ctx.tools.register(definition)
|
||||
disposers.set(registeredName, dispose)
|
||||
} catch {
|
||||
// Name conflict — another tool with this name is already registered.
|
||||
ctx.logger.warn(`mcp-client: skipping tool "${registeredName}" (name conflict)`)
|
||||
}
|
||||
}
|
||||
cursor = response.nextCursor
|
||||
} while (cursor)
|
||||
|
||||
return disposers
|
||||
}
|
||||
|
||||
/**
|
||||
* The shape we read from each MCP content block. Intentionally looser than the
|
||||
* SDK's `ContentBlock` type: we're at a network trust boundary (data arrives
|
||||
* from an external MCP server process via JSON-RPC), so fields that the SDK
|
||||
* declares required may be absent at runtime if the server is buggy.
|
||||
*/
|
||||
interface McpContentBlock {
|
||||
type: string
|
||||
text?: string
|
||||
mimeType?: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Create an execute function for one MCP tool. The executor calls
|
||||
* `client.callTool` with abort signal and timeout, then maps the result
|
||||
* to harness ContentBlocks.
|
||||
*
|
||||
* When the MCP server returns `isError: true`, the executor throws so that
|
||||
* the ToolRegistry's catch path produces an `isError` result for the model.
|
||||
*/
|
||||
function createExecutor(
|
||||
client: Client,
|
||||
mcpToolName: string,
|
||||
opts: ToolBridgeOptions,
|
||||
): ToolDefinition['execute'] {
|
||||
return async (args: unknown, exec: ToolExecution) => {
|
||||
// The agent loop passes `JSON.parse(model_arguments)` which is usually an
|
||||
// object, but can be any JSON value if the model misbehaves (outputs a bare
|
||||
// string/number/null). Fallback to {} lets the MCP server produce a
|
||||
// specific "missing required param" error the model can learn from.
|
||||
const argsObj = (typeof args === 'object' && args !== null ? args : {}) as Record<string, unknown>
|
||||
const result = await client.callTool(
|
||||
{ name: mcpToolName, arguments: argsObj },
|
||||
undefined,
|
||||
{
|
||||
...exec.signal ? { signal: exec.signal } : {},
|
||||
timeout: opts.toolCallTimeoutMs,
|
||||
},
|
||||
)
|
||||
|
||||
// The SDK may return a legacy `toolResult` shape; normalize to content array.
|
||||
if (!('content' in result) || !Array.isArray(result.content)) {
|
||||
const text = 'toolResult' in result
|
||||
? JSON.stringify(result.toolResult)
|
||||
: '(no output)'
|
||||
return [{ type: 'text' as const, text }]
|
||||
}
|
||||
|
||||
// Trust boundary: the SDK's return type erases to `any[]` due to the
|
||||
// union of CallToolResult | CompatibilityCallToolResult. We process each
|
||||
// element defensively in extractText (reading only .type/.text/.mimeType
|
||||
// with optional fallbacks).
|
||||
// eslint-disable-next-line @typescript-eslint/no-unsafe-assignment
|
||||
const content: McpContentBlock[] = result.content
|
||||
const text = extractText(content, mcpToolName)
|
||||
|
||||
// MCP isError → throw so ToolRegistry produces an isError result for the model.
|
||||
if ('isError' in result && result.isError === true) {
|
||||
throw new Error(text)
|
||||
}
|
||||
|
||||
return [{ type: 'text', text }]
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract text from an MCP content array into a single string.
|
||||
* - text blocks: join with '\n'
|
||||
* - image/audio/resource blocks: replaced with a placeholder
|
||||
*
|
||||
* Defensive: fields that the MCP spec declares required (mimeType, text) are
|
||||
* guarded with fallbacks because this is a network trust boundary.
|
||||
*/
|
||||
function extractText(mcpContent: McpContentBlock[], toolName: string): string {
|
||||
const parts: string[] = []
|
||||
|
||||
for (const block of mcpContent) {
|
||||
switch (block.type) {
|
||||
case 'text':
|
||||
if (block.text !== undefined) parts.push(block.text)
|
||||
break
|
||||
case 'image':
|
||||
parts.push(`[image: ${block.mimeType ?? 'unknown'}, content discarded]`)
|
||||
break
|
||||
case 'audio':
|
||||
parts.push(`[audio: ${block.mimeType ?? 'unknown'}, content discarded]`)
|
||||
break
|
||||
case 'resource':
|
||||
case 'resource_link':
|
||||
parts.push('[resource: content discarded]')
|
||||
break
|
||||
default:
|
||||
parts.push(`[unsupported content type: ${block.type}]`)
|
||||
}
|
||||
}
|
||||
|
||||
return parts.join('\n') || `(${toolName} returned no text content)`
|
||||
}
|
||||
56
packages/mcp/mcp-client/src/transport.ts
Normal file
56
packages/mcp/mcp-client/src/transport.ts
Normal file
@@ -0,0 +1,56 @@
|
||||
/**
|
||||
* Transport factory: creates the appropriate MCP transport based on the
|
||||
* plugin's resolved config. Stdio spawns a child process (with credential
|
||||
* scrubbing); Streamable HTTP connects to a URL.
|
||||
*
|
||||
* @module
|
||||
*/
|
||||
|
||||
import type { Transport } from '@modelcontextprotocol/sdk/shared/transport.js'
|
||||
import { StdioClientTransport } from '@modelcontextprotocol/sdk/client/stdio.js'
|
||||
import { StreamableHTTPClientTransport } from '@modelcontextprotocol/sdk/client/streamableHttp.js'
|
||||
import type { Config } from './index.ts'
|
||||
|
||||
/**
|
||||
* Credential-shaped ambient env vars are NOT forwarded to the child by default
|
||||
* (the parent harness's own secrets must not leak into a spawned process
|
||||
* implicitly). Same pattern as `dsh-subagent-acp`.
|
||||
*/
|
||||
const SENSITIVE_ENV_PATTERN = /KEY|SECRET|TOKEN/i
|
||||
|
||||
/** The ambient env minus credential-shaped vars, plus the spec's explicit env. */
|
||||
function buildChildEnv(extra: Record<string, string>): Record<string, string> {
|
||||
const env: Record<string, string> = {}
|
||||
for (const [key, value] of Object.entries(process.env)) {
|
||||
if (value !== undefined && !SENSITIVE_ENV_PATTERN.test(key)) env[key] = value
|
||||
}
|
||||
return { ...env, ...extra }
|
||||
}
|
||||
|
||||
/**
|
||||
* Create an MCP transport from the resolved plugin config.
|
||||
*
|
||||
* @param config - Resolved plugin config discriminated on `transport`.
|
||||
* @returns A connected-ready MCP Transport (stdio or Streamable HTTP).
|
||||
*/
|
||||
export function createTransport(config: Config): Transport {
|
||||
switch (config.transport) {
|
||||
case 'stdio':
|
||||
return new StdioClientTransport({
|
||||
command: config.command,
|
||||
args: config.args,
|
||||
env: buildChildEnv(config.env),
|
||||
cwd: config.cwd,
|
||||
})
|
||||
case 'streamable-http':
|
||||
// The MCP SDK's StreamableHTTPClientTransport has optional callback
|
||||
// properties typed without `| undefined` (exactOptionalPropertyTypes
|
||||
// mismatch with the Transport interface). The cast is safe — the SDK
|
||||
// constructed the object, it simply doesn't declare the optionals
|
||||
// strictly enough for our tsconfig.
|
||||
return new StreamableHTTPClientTransport(
|
||||
new URL(config.url),
|
||||
{ requestInit: { headers: config.headers } },
|
||||
) as Transport
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user