mirror of
https://github.com/modelstudioai/cli.git
synced 2026-09-14 19:49:23 +08:00
fix: fixed sse error
This commit is contained in:
@@ -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;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
@@ -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<Uint8Array>({
|
||||
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<Uint8Array>({
|
||||
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;
|
||||
}
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user