diff --git a/extensions/parallel/src/parallel-mcp-search.runtime.test.ts b/extensions/parallel/src/parallel-mcp-search.runtime.test.ts index f84e2661f1b6..e76e37ba8c78 100644 --- a/extensions/parallel/src/parallel-mcp-search.runtime.test.ts +++ b/extensions/parallel/src/parallel-mcp-search.runtime.test.ts @@ -44,6 +44,28 @@ function jsonResponse(body: unknown, headers?: Record): Response }); } +function cancelTrackedResponse( + text: string, + init: ResponseInit, +): { + response: Response; + wasCanceled: () => boolean; +} { + let canceled = false; + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(text)); + }, + cancel() { + canceled = true; + }, + }); + return { + response: new Response(stream, init), + wasCanceled: () => canceled, + }; +} + function readBody(call: EndpointCall): Record { if (typeof call.init.body !== "string") { throw new Error("Expected a JSON string body."); @@ -270,4 +292,23 @@ describe("runParallelMcpSearch", () => { /initialize failed \(500\)/, ); }); + + it("bounds initialize error bodies without using response.text()", async () => { + const tracked = cancelTrackedResponse(`${"parallel mcp unavailable ".repeat(1024)}tail`, { + status: 503, + headers: { "Content-Type": "text/plain" }, + }); + const textSpy = vi.spyOn(tracked.response, "text").mockRejectedValue(new Error("unbounded")); + endpointMockState.responses.push(tracked.response); + + const error = await runParallelMcpSearch({ searchQueries: ["x"], maxResults: 5 }).catch( + (cause: unknown) => cause, + ); + + expect(error).toBeInstanceOf(Error); + expect((error as Error).message).toMatch(/initialize failed \(503\): parallel mcp unavailable/); + expect((error as Error).message).not.toContain("tail"); + expect(tracked.wasCanceled()).toBe(true); + expect(textSpy).not.toHaveBeenCalled(); + }); }); diff --git a/extensions/parallel/src/parallel-mcp-search.runtime.ts b/extensions/parallel/src/parallel-mcp-search.runtime.ts index 989ee7f24120..0b0031f3a09a 100644 --- a/extensions/parallel/src/parallel-mcp-search.runtime.ts +++ b/extensions/parallel/src/parallel-mcp-search.runtime.ts @@ -1,6 +1,7 @@ import { randomUUID } from "node:crypto"; import { createRequire } from "node:module"; import { readPluginPackageVersion } from "openclaw/plugin-sdk/extension-shared"; +import { readResponseTextLimited } from "openclaw/plugin-sdk/provider-http"; import { withTrustedWebSearchEndpoint } from "openclaw/plugin-sdk/provider-web-search"; // Free hosted Search MCP. This keyless transport is used only after the user @@ -11,6 +12,7 @@ export const PARALLEL_MCP_SEARCH_URL = "https://search.parallel.ai/mcp"; // the server negotiates back on every follow-up request. const MCP_PROTOCOL_VERSION = "2025-06-18"; const MCP_TIMEOUT_SECONDS = 30; +const PARALLEL_MCP_ERROR_BODY_LIMIT_BYTES = 8 * 1024; const require = createRequire(import.meta.url); const PLUGIN_VERSION = readPluginPackageVersion({ require }); @@ -215,7 +217,9 @@ async function postMcp(params: { ok: response.ok, status: response.status, statusText: response.statusText, - text: await response.text(), + text: response.ok + ? await response.text() + : await readResponseTextLimited(response, PARALLEL_MCP_ERROR_BODY_LIMIT_BYTES), sessionIdHeader: response.headers.get("mcp-session-id"), }), ); diff --git a/extensions/parallel/src/parallel-web-search-provider.runtime.ts b/extensions/parallel/src/parallel-web-search-provider.runtime.ts index 0152bd816953..b55f2bc334ca 100644 --- a/extensions/parallel/src/parallel-web-search-provider.runtime.ts +++ b/extensions/parallel/src/parallel-web-search-provider.runtime.ts @@ -1,5 +1,6 @@ import { createRequire } from "node:module"; import { readPluginPackageVersion } from "openclaw/plugin-sdk/extension-shared"; +import { readResponseTextLimited } from "openclaw/plugin-sdk/provider-http"; import { DEFAULT_SEARCH_COUNT, mergeScopedSearchConfig, @@ -34,6 +35,7 @@ import { const PARALLEL_BASE_URL = "https://api.parallel.ai"; const PARALLEL_SEARCH_PATHNAME = "/v1/search"; +const PARALLEL_ERROR_BODY_LIMIT_BYTES = 8 * 1024; const require = createRequire(import.meta.url); const PLUGIN_VERSION = readPluginPackageVersion({ require }); @@ -144,7 +146,9 @@ async function runParallelSearch(params: { }, async (res) => { if (!res.ok) { - const detail = await res.text().catch(() => ""); + const detail = await readResponseTextLimited(res, PARALLEL_ERROR_BODY_LIMIT_BYTES).catch( + () => "", + ); throw new Error(`Parallel API error (${res.status}): ${detail || res.statusText}`); } try { @@ -277,6 +281,7 @@ export const testing = { resolveParallelConfig, resolveParallelSearchCount, resolveParallelSearchEndpoint, + PARALLEL_ERROR_BODY_LIMIT_BYTES, USER_AGENT, } as const; diff --git a/extensions/parallel/src/parallel-web-search-provider.test.ts b/extensions/parallel/src/parallel-web-search-provider.test.ts index 3670eb5bef02..c9d7fbe00453 100644 --- a/extensions/parallel/src/parallel-web-search-provider.test.ts +++ b/extensions/parallel/src/parallel-web-search-provider.test.ts @@ -37,6 +37,28 @@ function readMockedBody(call: EndpointCall | undefined): unknown { return JSON.parse(call.init.body); } +function cancelTrackedResponse( + text: string, + init: ResponseInit, +): { + response: Response; + wasCanceled: () => boolean; +} { + let canceled = false; + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(text)); + }, + cancel() { + canceled = true; + }, + }); + return { + response: new Response(stream, init), + wasCanceled: () => canceled, + }; +} + import { testing } from "../test-api.js"; import { createParallelWebSearchProvider as createContractParallelWebSearchProvider } from "../web-search-contract-api.js"; import { createParallelWebSearchProvider } from "./parallel-web-search-provider.js"; @@ -529,6 +551,38 @@ describe("parallel web search provider", () => { expect(body.advanced_settings?.max_results).toBe(5); }); + it("bounds Parallel API error bodies without using response.text()", async () => { + const tracked = cancelTrackedResponse(`${"parallel upstream unavailable ".repeat(1024)}tail`, { + status: 503, + headers: { "Content-Type": "text/plain" }, + }); + const textSpy = vi.spyOn(tracked.response, "text").mockRejectedValue(new Error("unbounded")); + endpointMockState.responses.push(tracked.response); + const provider = createParallelWebSearchProvider(); + const tool = provider.createTool({ + config: {}, + searchConfig: { parallel: { apiKey: "par-secret" } }, + }); + if (!tool) { + throw new Error("Expected tool definition"); + } + + const error = await tool + .execute({ + objective: `parallel-error-body-${Date.now()}`, + search_queries: ["openclaw"], + }) + .catch((cause: unknown) => cause); + + expect(error).toBeInstanceOf(Error); + expect((error as Error).message).toMatch( + /Parallel API error \(503\): parallel upstream unavailable/, + ); + expect((error as Error).message).not.toContain("tail"); + expect(tracked.wasCanceled()).toBe(true); + expect(textSpy).not.toHaveBeenCalled(); + }); + it("does not surface a Parallel-generated sessionId on a cache hit", async () => { // Unique objective so this test does not collide with the SDK's // module-level web-search cache across other cases.