diff --git a/packages/core/src/client/mcp-sse.ts b/packages/core/src/client/mcp-sse.ts index 97188b5..a156557 100644 --- a/packages/core/src/client/mcp-sse.ts +++ b/packages/core/src/client/mcp-sse.ts @@ -153,12 +153,8 @@ export class McpSseClient { if (headerTimedOut) { throw new BailianError("MCP SSE timed out waiting for response headers.", ExitCode.TIMEOUT); } - throw new BailianError( - `MCP SSE request failed: ${error instanceof Error ? error.message : String(error)}`, - ExitCode.NETWORK, - undefined, - { cause: error }, - ); + // Rethrow fetch failures so runtime can surface errno (e.g. ENOTFOUND) in JSON/text. + throw error; } if (this.deps.settings.verbose) { @@ -352,35 +348,47 @@ export class McpSseClient { const requestSignal = createLinkedAbortSignal(timeoutMs, this.abortController?.signal); let res: Response; try { - res = await fetch(this.messageUrl, { - method: "POST", - headers, - body: JSON.stringify(body), - signal: requestSignal.signal, - }); - } catch (error) { - if (this.closed) { - throw new BailianError("MCP SSE session closed.", ExitCode.GENERAL); + try { + res = await fetch(this.messageUrl, { + method: "POST", + headers, + body: JSON.stringify(body), + signal: requestSignal.signal, + }); + } catch (error) { + if (this.closed) { + throw new BailianError("MCP SSE session closed.", ExitCode.GENERAL); + } + throw error; + } + + if (this.deps.settings.verbose) { + console.error(`< ${res.status} ${res.statusText}`); + } + + if (!res.ok) { + // Keep signal until error body is read (same class of bug as GET openSse). + let errMsg = `MCP request failed: ${res.status} ${res.statusText}`; + try { + const errBody = await res.text(); + if (errBody) errMsg += ` - ${errBody.slice(0, 500)}`; + } catch (error) { + if (this.closed) { + throw new BailianError("MCP SSE session closed.", ExitCode.GENERAL); + } + if (requestSignal.timedOut) { + throw new BailianError( + "MCP SSE timed out reading error response body.", + ExitCode.TIMEOUT, + ); + } + throw new BailianError(errMsg, ExitCode.GENERAL, undefined, { cause: error }); + } + throw new BailianError(errMsg, ExitCode.GENERAL); } - throw error; } finally { requestSignal.cleanup(); } - - if (this.deps.settings.verbose) { - console.error(`< ${res.status} ${res.statusText}`); - } - - if (!res.ok) { - let errMsg = `MCP request failed: ${res.status} ${res.statusText}`; - try { - const errBody = await res.text(); - if (errBody) errMsg += ` - ${errBody.slice(0, 500)}`; - } catch { - /* ignore */ - } - throw new BailianError(errMsg, ExitCode.GENERAL); - } } } @@ -439,9 +447,13 @@ function cancellableTimeoutReject( function createLinkedAbortSignal( timeoutMs: number, parentSignal?: AbortSignal, -): { signal: AbortSignal; cleanup: () => void } { +): { signal: AbortSignal; cleanup: () => void; timedOut: boolean } { const controller = new AbortController(); - const timeout = setTimeout(() => controller.abort(), timeoutMs); + const state = { timedOut: false }; + const timeout = setTimeout(() => { + state.timedOut = true; + controller.abort(); + }, timeoutMs); const abortFromParent = () => controller.abort(parentSignal?.reason); const cleanup = () => { clearTimeout(timeout); @@ -452,5 +464,11 @@ function createLinkedAbortSignal( else parentSignal?.addEventListener("abort", abortFromParent, { once: true }); controller.signal.addEventListener("abort", cleanup, { once: true }); - return { signal: controller.signal, cleanup }; + return { + signal: controller.signal, + cleanup, + get timedOut() { + return state.timedOut; + }, + }; } diff --git a/packages/core/tests/mcp.test.ts b/packages/core/tests/mcp.test.ts index 55e1034..357f419 100644 --- a/packages/core/tests/mcp.test.ts +++ b/packages/core/tests/mcp.test.ts @@ -563,7 +563,7 @@ test("McpSseClient:非 2xx 读 body 仍受 --timeout 约束", async () => { } }); -test("McpSseClient:fetch 失败保留 cause", async () => { +test("McpSseClient:fetch 失败抛出原始 TypeError(保留 ENOTFOUND)", async () => { const originalFetch = globalThis.fetch; const root = Object.assign(new Error("getaddrinfo ENOTFOUND example.test"), { code: "ENOTFOUND", @@ -576,11 +576,9 @@ test("McpSseClient:fetch 失败保留 cause", async () => { try { const client = new McpSseClient(testDeps(), "https://example.test/sse", "sk-test"); - await expect(client.initialize()).rejects.toMatchObject({ - message: expect.stringMatching(/MCP SSE request failed:\s*fetch failed/i), - exitCode: 6, - cause: fetchFailed, - }); + const error = await client.initialize().catch((reason: unknown) => reason); + expect(error).toBe(fetchFailed); + expect((error as TypeError & { cause?: NodeJS.ErrnoException }).cause?.code).toBe("ENOTFOUND"); client.close(); } finally { globalThis.fetch = originalFetch; @@ -684,3 +682,55 @@ test("McpSseClient:close 可中止进行中的 POST", async () => { globalThis.fetch = originalFetch; } }); + +test("McpSseClient:POST 非 2xx 读 body 仍受 --timeout 约束", async () => { + const originalFetch = globalThis.fetch; + const encoder = new TextEncoder(); + + globalThis.fetch = async (input, init) => { + const url = requestUrl(input); + const method = init?.method ?? "GET"; + if (method === "GET" || url.endsWith("/sse")) { + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode("event: endpoint\ndata: /message\n\n")); + }, + }); + return new Response(stream, { + status: 200, + headers: { "Content-Type": "text/event-stream" }, + }); + } + + const signal = init?.signal; + const body = new ReadableStream({ + start(controller) { + if (!signal) return; + const onAbort = () => { + try { + controller.error(new DOMException("This operation was aborted.", "AbortError")); + } catch { + /* ignore */ + } + }; + if (signal.aborted) onAbort(); + else signal.addEventListener("abort", onAbort, { once: true }); + }, + }); + return new Response(body, { status: 500, statusText: "Internal Server Error" }); + }; + + try { + const client = new McpSseClient( + testDeps({ timeout: 1 }), + "https://example.test/sse", + "sk-test", + ); + const started = Date.now(); + await expect(client.initialize()).rejects.toThrow(/timed out reading error response body/i); + expect(Date.now() - started).toBeLessThan(2500); + client.close(); + } finally { + globalThis.fetch = originalFetch; + } +});