Merge origin/master into xjt/incremental-proofreading-275-apply

This commit is contained in:
xjt
2026-08-09 12:17:59 +08:00
47 changed files with 1433 additions and 202 deletions

View File

@@ -2,5 +2,5 @@
# side as of the last confirmed-consistent state. Both languages carry equal authority;
# after editing either side, bring the other along and re-record with:
# pnpm run verify-translation-pairing --write packages/mcp/mcp-client/README.md
README.md: d7966595c68ff1ec4a288caf5d9fe4b0bf580cc5
README.zh.md: ced1ada5e4a317e37d8601ffe7beaed0e4056aff
README.md: 76d1271f6f7a3e9c959bdcf5e969906f25563c56
README.zh.md: b2da1119af3a8d52761a5040059e7a9d922aa567

View File

@@ -44,6 +44,7 @@ The model sees `mcp__github__create_issue`, `mcp__web__search`, … — the same
| `url` | http | yes | MCP server URL |
| `headers` | http | no | Extra headers (e.g. auth tokens) |
| `toolCallTimeoutMs` | both | no | Timeout per `callTool` invocation (default 60000) |
| `failOnStartupError` | both | no | Reject plugin activation when initial connection or tool synchronization fails (default `false`) |
## Tool naming
@@ -56,8 +57,8 @@ Every MCP tool has two names: the raw MCP name (sent on the wire in `tools/call`
## Behavior
- On connect: `listTools()` registers each tool via `ctx.tools.register()` under its public name.
- Listens for `notifications/tools/list_changed` → re-syncs; a failed re-sync keeps the previous generation registered.
- On connect: plugin activation awaits `listTools()` and registers each tool via `ctx.tools.register()` under its public name before the composition starts its first turn. Initial connection, discovery, or registration failure is always logged; it rejects activation when `failOnStartupError` is true and otherwise activates with no tools.
- Listens for `notifications/tools/list_changed` → re-syncs; a fetch-phase failure keeps the previous generation registered, while a registration conflict rolls back the attempted generation and leaves no tools from that server.
- Tool execute: `client.callTool({ name: rawName, arguments }, { signal })` with timeout + abort support—the public name is never sent to the server.
- Canonical success is `{ content: JsonValue[], structuredContent? }`; complete JSON MCP blocks survive for programmatic callers. A supported advertised `outputSchema` validates `structuredContent`; unsupported schema vocabulary falls back to unconstrained `JsonValue`.
- Native/model rendering keeps the existing text projection: text blocks join with newlines while image, audio, resource, and unsupported blocks become placeholders.
@@ -101,8 +102,8 @@ Append-only; newly visible content follows the reusable request prefix and does
## Known Limitations and Deferred Work
- **Initial discovery is asynchronous** — plugin load does not wait for connection and `listTools()`, so a turn started immediately after boot or HMR can assemble before the MCP tools are registered.
- **Tools are the only bridged MCP capability** — Resources and Prompts have no harness consumption surface and are deferred.
- **Startup timeout is inherited from the MCP SDK** — DSH does not yet expose a connection/discovery timeout. Each initialize or paginated `tools/list` request uses the SDK's 60-second default, so an unresponsive server or cursor chain can delay both activation and teardown while the initial synchronization settles.
- **Crash recovery is manual** — transport closure does not auto-reconnect; registered tools can remain visible but fail against the closed transport until an HMR reload or Host restart.
- **Native non-text rendering is lossy** — image, audio, and resource payloads become placeholders in model context even though the execution-local canonical value preserves their JSON blocks. Richer Native multimedia projection is deferred.
- **Unsupported MCP output schemas are not enforced** — `structuredContent` falls back to `JsonValue` when the advertised schema uses vocabulary outside the harness subset.

View File

@@ -44,6 +44,7 @@ MCP 客户端桥接插件:连接外部 [Model Context Protocol](https://modelc
| `url` | http | 是 | MCP 服务器 URL |
| `headers` | http | 否 | 额外标头(例如认证 token |
| `toolCallTimeoutMs` | 两者 | 否 | 每次 `callTool` 调用的超时(默认 60000 |
| `failOnStartupError` | 两者 | 否 | 初始连接或工具同步失败时拒绝插件激活(默认 `false` |
## 工具命名
@@ -56,8 +57,8 @@ MCP 客户端桥接插件:连接外部 [Model Context Protocol](https://modelc
## 行为
- 连接时:`listTools()`通过 `ctx.tools.register()` 使用各自公开名称注册每个工具。
- 监听 `notifications/tools/list_changed` → 重新同步;同步失败时保留上一世代的注册。
- 连接时:插件激活会等待 `listTools()`,并在组合开始首个轮次前通过 `ctx.tools.register()` 公开名称注册每个工具。初始连接、发现或注册失败始终会记录日志;`failOnStartupError` 为 true 时拒绝激活,否则插件仍会激活但不注册工具。
- 监听 `notifications/tools/list_changed` → 重新同步;获取阶段失败时保留上一世代的注册,注册冲突则会回滚本次尝试的世代,并且不保留该服务器的任何工具
- 工具执行:`client.callTool({ name: rawName, arguments }, { signal })`,支持超时 + 中止;公开名称绝不会发给服务器。
- 规范成功值是 `{ content: JsonValue[], structuredContent? }`;完整的 JSON MCP 块会保留给编程调用方。受支持且已声明的 `outputSchema` 会验证 `structuredContent`;不受支持的 schema 词汇会回退为不受约束的 `JsonValue`
- Native模型渲染保留现有文本投影文本块以换行连接图片、音频、资源和不受支持的块会变成占位符。
@@ -101,8 +102,8 @@ MCP 客户端桥接插件:连接外部 [Model Context Protocol](https://modelc
## 已知限制与暂缓事项
- **初始发现是异步的**:插件加载不会等待连接和 `listTools()`,因此在启动或 HMR 后立即开始的轮次可能在 MCP 工具注册前完成组装。
- **只桥接 MCP 的工具能力**:资源和提示词没有 harness 消费接口,暂缓实现。
- **启动超时继承自 MCP SDK**DSH 尚未公开连接/发现超时。每次 initialize 请求或分页 `tools/list` 请求都使用 SDK 默认的 60 秒,因此在初始同步完成期间,无响应的 server 或 cursor chain 可能同时延迟激活与 teardown。
- **崩溃恢复需要手动触发**:传输关闭后不会自动重新连接;已注册工具可能仍然可见,但会因传输已关闭而调用失败,直到 HMR 重载或重启 Host。
- **Native 非文本渲染有损**:图片、音频与资源载荷在模型上下文中会变成占位符,即使执行局部的规范值保留了其 JSON 块。更丰富的 Native 多媒体投影暂缓实现。
- **不强制执行不受支持的 MCP 输出 schema**:已声明 schema 使用 harness 子集之外的词汇时,`structuredContent` 会回退到 `JsonValue`

View File

@@ -72,6 +72,8 @@ export interface StdioConfig {
cwd: string
/** Per-tool-call timeout in milliseconds. */
toolCallTimeoutMs: number
/** Fail plugin activation when the initial connection or tool synchronization fails. */
failOnStartupError: boolean
}
/** Config for connecting to an MCP server over Streamable HTTP (SSE). */
@@ -90,6 +92,8 @@ export interface StreamableHttpConfig {
headers: Record<string, string>
/** Per-tool-call timeout in milliseconds. */
toolCallTimeoutMs: number
/** Fail plugin activation when the initial connection or tool synchronization fails. */
failOnStartupError: boolean
}
/** Configuration for one stdio or Streamable HTTP MCP server. */
@@ -104,6 +108,7 @@ export const Config = z.union([
env: z.dict(String).default({}),
cwd: z.string().default(''),
toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
failOnStartupError: z.boolean().default(false),
}),
z.object({
transport: z.const('streamable-http'),
@@ -111,12 +116,21 @@ export const Config = z.union([
url: z.string().required(),
headers: z.dict(String).default({}),
toolCallTimeoutMs: z.number().default(DEFAULT_TOOL_CALL_TIMEOUT_MS),
failOnStartupError: z.boolean().default(false),
}),
]) as unknown as z<Config>
// ---- Plugin apply ----
export function apply(ctx: Context, config: Config): void {
/**
* Connect one MCP server and publish its initial tool generation before activation.
* This entry remains explicitly `async`: Cordis treats a prototype-bearing
* ordinary function as a constructor, whose returned Promise is not startup work.
* @param ctx - plugin context carrying the tool registry.
* @param config - resolved transport and server namespace configuration.
* @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
// at load with an actionable error and leaves the earlier instance intact.
ctx.effect(() => {
@@ -141,18 +155,22 @@ export function apply(ctx: Context, config: Config): void {
)
const opts = {
registrationFailure: 'contain' as const,
serverName: config.serverName,
toolCallTimeoutMs: config.toolCallTimeoutMs,
}
// Connect and set up tools. Errors during connect/first sync are logged,
// not thrown (the plugin simply has no tools registered). `ready` resolves
// to an accessor for the CURRENT disposer generation, so the effect
// disposer below always unregisters the live set, not the first one.
// 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, new Map())
let disposers = await syncTools(client, ctx, {
...opts,
registrationFailure: config.failOnStartupError ? 'throw' : 'contain',
}, new Map())
client.setNotificationHandler(
ToolListChangedNotificationSchema,
@@ -168,15 +186,20 @@ export function apply(ctx: Context, config: Config): void {
},
)
return () => disposers
return { getDisposers: () => disposers }
})().catch((error: unknown) => {
ctx.logger.error(`mcp-client(${config.serverName}): failed to connect: ${String(error)}`)
return () => new Map<string, () => void>()
ctx.logger.error(`mcp-client(${config.serverName}): startup failed: ${String(error)}`)
return { getDisposers: () => new Map<string, () => void>(), error }
})
ctx.effect(() => async () => {
const live = await ready
for (const dispose of live().values()) dispose()
const outcome = await ready
for (const dispose of outcome.getDisposers().values()) dispose()
try { await client.close() } catch { /* transport already gone */ }
}, 'mcp-client.connection')
const outcome = await ready
if ('error' in outcome && config.failOnStartupError) {
throw new Error(`mcp-client(${config.serverName}): initial connection or tool synchronization failed`, { cause: outcome.error })
}
}

View File

@@ -23,6 +23,8 @@ import type { JsonSchemaNode, JsonValue } from '@deepseek-ai/dsh-tools'
/** Resolved options relevant to tool bridging. */
export interface ToolBridgeOptions {
/** Whether a registry conflict is contained or rejects this synchronization. */
registrationFailure: 'contain' | 'throw'
serverName: string
toolCallTimeoutMs: number
}
@@ -111,8 +113,9 @@ export function publicToolName(serverName: string, rawName: string): string {
* 2. Swap: dispose the previous generation, register the new one. A registry
* conflict here can only mean a foreign registration squats on this
* server's `mcp__<serverName>__` namespace — the partial generation is
* rolled back (zero tools from this server), the error is logged, and an
* empty map is returned.
* rolled back (zero tools from this server) and logged. Initial strict
* synchronization may propagate the conflict so its parent transaction
* rejects; ordinary clients and later re-syncs return an empty map.
*
* @param client - Connected MCP Client instance used to list and call tools.
* @param ctx - Cordis context providing the `tools` service for registration.
@@ -164,6 +167,7 @@ export async function syncTools(
// sees either the full generation or none of it — never a partial set.
for (const dispose of disposers.values()) dispose()
ctx.logger.error(`mcp-client(${opts.serverName}): tool registration failed, no tools registered: ${String(error)}`)
if (opts.registrationFailure === 'throw') throw error
return new Map()
}
return disposers

View File

@@ -82,6 +82,7 @@ const stdioConfig: Config = {
env: {},
cwd: '',
toolCallTimeoutMs: 60_000,
failOnStartupError: false,
}
// ---- Tests ----
@@ -141,8 +142,7 @@ describe('apply (plugin lifecycle)', () => {
})
it('connects, syncs tools under the namespace, and registers a notification handler', async () => {
apply(ctx, stdioConfig)
await sleep(50)
await apply(ctx, stdioConfig)
expect(mockConnect).toHaveBeenCalled()
expect(mockListTools).toHaveBeenCalled()
@@ -151,12 +151,30 @@ describe('apply (plugin lifecycle)', () => {
expect(ctx.tools.get('remote')).toBeUndefined()
})
it('keeps the Cordis plugin loading until initial discovery publishes its tools', async () => {
const connection: PromiseWithResolvers<void> = Promise.withResolvers()
mockConnect.mockImplementation(async () => {
await connection.promise
})
const fiber = ctx.plugin({ name: 'mcp-client-lifecycle', inject, apply }, stdioConfig)
let activated = false
const activation = Promise.resolve(fiber).then(() => { activated = true })
await vi.waitFor(() => { expect(mockConnect).toHaveBeenCalled() })
expect(activated).toBe(false)
expect(ctx.tools.get('mcp__srv__remote')).toBeUndefined()
connection.resolve()
await activation
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
await fiber.dispose()
})
it('rejects a duplicate serverName at load and leaves the first instance intact', async () => {
apply(ctx, stdioConfig)
await sleep(50)
await apply(ctx, stdioConfig)
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
expect(() => { apply(ctx, stdioConfig) }).toThrow(/serverName "srv" is already in use/)
await expect(apply(ctx, stdioConfig)).rejects.toThrow(/serverName "srv" is already in use/)
// First instance unaffected.
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
})
@@ -165,8 +183,7 @@ describe('apply (plugin lifecycle)', () => {
const first = new Context()
await first.plugin(SystemPrompt)
await first.plugin(ToolRegistry)
apply(first, stdioConfig)
await sleep(50)
await apply(first, stdioConfig)
await first.fiber.dispose()
await sleep(50)
@@ -176,26 +193,26 @@ describe('apply (plugin lifecycle)', () => {
const second = new Context()
await second.plugin(SystemPrompt)
await second.plugin(ToolRegistry)
expect(() => { apply(second, stdioConfig) }).not.toThrow()
await expect(apply(second, stdioConfig)).resolves.toBeUndefined()
await second.fiber.dispose()
})
it('scopes serverName reservations per app root', async () => {
const other = await mountRegistry()
apply(ctx, stdioConfig)
const first = apply(ctx, stdioConfig)
// Same serverName on a DIFFERENT root is fine.
expect(() => { apply(other, stdioConfig) }).not.toThrow()
await sleep(50)
const second = apply(other, stdioConfig)
await Promise.all([first, second])
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
expect(other.tools.get('mcp__srv__remote')).toBeDefined()
})
it('logs error and registers no tools when connect fails; dispose is a no-op', async () => {
it('logs error and registers no tools when connect fails; dispose closes the client', async () => {
mockConnect.mockRejectedValue(new Error('connection refused'))
apply(ctx, stdioConfig)
await sleep(50)
await apply(ctx, stdioConfig)
expect(mockListTools).not.toHaveBeenCalled()
expect(ctx.tools.get('mcp__srv__remote')).toBeUndefined()
@@ -207,9 +224,43 @@ describe('apply (plugin lifecycle)', () => {
expect(mockClose).toHaveBeenCalled()
})
it('rejects activation and still closes the client when startup failure is configured as fatal', async () => {
mockConnect.mockRejectedValue(new Error('connection refused'))
await expect(apply(ctx, {
...stdioConfig,
failOnStartupError: true,
})).rejects.toThrow('initial connection or tool synchronization failed')
expect(mockListTools).not.toHaveBeenCalled()
expect(ctx.tools.get('mcp__srv__remote')).toBeUndefined()
await ctx.fiber.dispose()
expect(mockClose).toHaveBeenCalled()
})
it('rejects strict startup when the initial tool generation cannot be registered', async () => {
ctx.tools.register({
name: 'mcp__srv__remote',
description: 'Foreign squatter',
parameters: { type: 'object' },
output: {
schema: { type: 'string' },
render: (_args, value) => [{ type: 'text', text: value as string }],
},
execute: async () => 'foreign',
})
await expect(apply(ctx, {
...stdioConfig,
failOnStartupError: true,
})).rejects.toThrow('initial connection or tool synchronization failed')
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
await ctx.fiber.dispose()
expect(mockClose).toHaveBeenCalled()
})
it('re-syncs tools on ToolListChanged notification', async () => {
apply(ctx, stdioConfig)
await sleep(50)
await apply(ctx, stdioConfig)
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
@@ -226,8 +277,7 @@ describe('apply (plugin lifecycle)', () => {
})
it('keeps the previous generation when a re-sync fails', async () => {
apply(ctx, stdioConfig)
await sleep(50)
await apply(ctx, stdioConfig)
expect(ctx.tools.get('mcp__srv__remote')).toBeDefined()
mockListTools.mockRejectedValue(new Error('flaky server'))
@@ -242,7 +292,7 @@ describe('apply (plugin lifecycle)', () => {
// Load through ctx.plugin so ONLY the plugin's fiber is disposed — the
// registry must survive to observe the unregistration.
const fiber = ctx.plugin({ name: 'mcp-client', inject: ['tools'], apply }, stdioConfig)
await sleep(50)
await fiber
// Advance to a second generation first.
mockListTools.mockResolvedValue({
@@ -264,8 +314,7 @@ describe('apply (plugin lifecycle)', () => {
it('effect disposer handles client.close failure gracefully', async () => {
mockClose.mockRejectedValue(new Error('already closed'))
apply(ctx, stdioConfig)
await sleep(50)
await apply(ctx, stdioConfig)
// Should not throw when dispose is triggered.
await ctx.fiber.dispose()
@@ -281,10 +330,10 @@ describe('apply (plugin lifecycle)', () => {
url: 'http://localhost:3000/mcp',
headers: { Authorization: 'Bearer x' },
toolCallTimeoutMs: 30_000,
failOnStartupError: false,
}
apply(ctx, httpConfig)
await sleep(50)
await apply(ctx, httpConfig)
expect(mockConnect).toHaveBeenCalled()
expect(ctx.tools.get('mcp__web__remote')).toBeDefined()

View File

@@ -43,21 +43,6 @@ async function mountRegistry(): Promise<Context> {
return ctx
}
/** Apply the MCP client plugin and wait for tools to be registered. */
async function applyAndWait(ctx: Context, config: Config, timeoutMs = 20_000): Promise<void> {
// Annotated bindings (not withResolvers<void>()): the tests lint layer runs
// no-invalid-void-type with default options, which rejects the explicit
// type argument in call position but accepts the inferred form.
const gate: PromiseWithResolvers<void> = Promise.withResolvers()
const timer = setTimeout(
() => { gate.reject(new Error(`applyAndWait timed out after ${timeoutMs}ms — no tools/change event`)) },
timeoutMs,
)
ctx.on('tools/change', () => { clearTimeout(timer); gate.resolve() })
apply(ctx, config)
await gate.promise
}
function sleep(ms: number): Promise<void> {
const gate: PromiseWithResolvers<void> = Promise.withResolvers()
setTimeout(gate.resolve, ms)
@@ -90,11 +75,12 @@ describe('fixture server — controlled scenarios', () => {
env: {},
cwd: packageDir,
toolCallTimeoutMs: 15_000,
failOnStartupError: false,
}
beforeAll(async () => {
ctx = await mountRegistry()
await applyAndWait(ctx, fixtureConfig)
await apply(ctx, fixtureConfig)
}, 30_000)
afterAll(async () => {
@@ -179,10 +165,11 @@ describe('fixture server — duplicate serverName', () => {
env: {},
cwd: packageDir,
toolCallTimeoutMs: 15_000,
failOnStartupError: false,
}
await applyAndWait(ctx, config)
await apply(ctx, config)
expect(() => { apply(ctx, config) }).toThrow(/serverName "dup" is already in use/)
await expect(apply(ctx, config)).rejects.toThrow(/serverName "dup" is already in use/)
await ctx.fiber.dispose()
await sleep(200)
@@ -192,7 +179,7 @@ describe('fixture server — duplicate serverName', () => {
describe('fixture server — disposal', () => {
it('disposes cleanly without error', async () => {
const ctx = await mountRegistry()
await applyAndWait(ctx, {
await apply(ctx, {
transport: 'stdio',
serverName: 'fixture',
command: process.execPath,
@@ -200,6 +187,7 @@ describe('fixture server — disposal', () => {
env: {},
cwd: packageDir,
toolCallTimeoutMs: 15_000,
failOnStartupError: false,
})
// Tools are registered before dispose.
@@ -225,11 +213,12 @@ describe('server-everything — official test server', () => {
env: {},
cwd: '',
toolCallTimeoutMs: 30_000,
failOnStartupError: false,
}
beforeAll(async () => {
ctx = await mountRegistry()
await applyAndWait(ctx, config)
await apply(ctx, config)
}, 60_000)
afterAll(async () => {
@@ -292,8 +281,9 @@ describe('server-filesystem — real filesystem operations', () => {
env: {},
cwd: '',
toolCallTimeoutMs: 30_000,
failOnStartupError: false,
}
await applyAndWait(ctx, config)
await apply(ctx, config)
}, 60_000)
afterAll(async () => {
@@ -408,8 +398,9 @@ describe('streamable-http — in-process MCP server', () => {
url: baseUrl,
headers: { Authorization: 'Bearer e2e-test-token' },
toolCallTimeoutMs: 15_000,
failOnStartupError: false,
}
await applyAndWait(ctx, config)
await apply(ctx, config)
}, 30_000)
afterAll(async () => {

View File

@@ -64,6 +64,7 @@ async function mountRegistry(): Promise<Context> {
}
const defaultOpts: ToolBridgeOptions = {
registrationFailure: 'contain',
serverName: 'srv',
toolCallTimeoutMs: 60_000,
}
@@ -713,6 +714,7 @@ describe('createTransport', () => {
env: {},
cwd: '/tmp',
toolCallTimeoutMs: 60_000,
failOnStartupError: false,
}
const transport = createTransport(config)
expect(transport).toBeDefined()
@@ -727,6 +729,7 @@ describe('createTransport', () => {
url: 'http://localhost:3000/mcp',
headers: {},
toolCallTimeoutMs: 60_000,
failOnStartupError: false,
}
const transport = createTransport(config)
expect(transport).toBeDefined()
@@ -741,6 +744,7 @@ describe('createTransport', () => {
url: 'http://localhost:3000/mcp',
headers: { Authorization: 'Bearer token' },
toolCallTimeoutMs: 60_000,
failOnStartupError: false,
}
const transport = createTransport(config)
expect(transport).toBeDefined()
@@ -764,6 +768,7 @@ describe('createTransport', () => {
env: { EXTRA: 'injected' },
cwd: '',
toolCallTimeoutMs: 60_000,
failOnStartupError: false,
}
// createTransport internally calls buildChildEnv; we verify by inspecting
// the constructed StdioClientTransport. Since we can't inspect private fields
@@ -791,6 +796,7 @@ describe('createTransport', () => {
env: { CUSTOM: 'value' },
cwd: '',
toolCallTimeoutMs: 60_000,
failOnStartupError: false,
}
const transport = createTransport(config)
expect(transport).toBeDefined()