From bba027abab12ccc20affe3b1554c8c8bfb3e596f Mon Sep 17 00:00:00 2001 From: Dallin Romney Date: Thu, 6 Aug 2026 18:51:20 +0800 Subject: [PATCH] fix(providers): bound successful response reads Backport of 0a14444924e34e031133c997b00d30573353c5a0; retains the target's existing provider-reader SDK baseline. --- extensions/duckduckgo/src/ddg-client.ts | 8 +++- .../src/ddg-search-provider.test.ts | 19 +++++++++ extensions/firecrawl/src/firecrawl-client.ts | 21 +++++----- .../firecrawl/src/firecrawl-tools.test.ts | 22 ++++++++++ .../ollama/src/embedding-provider.test.ts | 33 ++++++++++++++- extensions/ollama/src/embedding-provider.ts | 14 +++---- .../ollama/src/web-search-provider.test.ts | 28 ++++++++++++- extensions/ollama/src/web-search-provider.ts | 7 +--- .../src/parallel-mcp-search.runtime.test.ts | 24 +++++++++++ .../src/parallel-mcp-search.runtime.ts | 7 +++- .../src/parallel-web-search-provider.test.ts | 42 ++++--------------- .../perplexity-web-search-provider.runtime.ts | 7 +--- .../perplexity-web-search-provider.test.ts | 21 +++++++++- .../qqbot/src/engine/api/api-client.test.ts | 32 ++++++++++++++ extensions/qqbot/src/engine/api/api-client.ts | 7 +++- .../src/engine/tools/channel-api.test.ts | 30 +++++++++++++ .../qqbot/src/engine/tools/channel-api.ts | 7 +++- extensions/tavily/src/tavily-client.test.ts | 24 +++++++++++ extensions/tavily/src/tavily-client.ts | 17 +++++--- .../test-support/streaming-error-response.ts | 41 ++++++++++++++---- 20 files changed, 325 insertions(+), 86 deletions(-) diff --git a/extensions/duckduckgo/src/ddg-client.ts b/extensions/duckduckgo/src/ddg-client.ts index bcecbc7b4b0a..d25ab7231807 100644 --- a/extensions/duckduckgo/src/ddg-client.ts +++ b/extensions/duckduckgo/src/ddg-client.ts @@ -1,5 +1,6 @@ // Duckduckgo plugin module implements ddg client behavior. import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import { readProviderTextResponse } from "openclaw/plugin-sdk/provider-http"; import { DEFAULT_CACHE_TTL_MINUTES, DEFAULT_SEARCH_COUNT, @@ -85,6 +86,10 @@ function isBotChallenge(html: string): boolean { return /g-recaptcha|are you a human|id="challenge-form"|name="challenge"/i.test(html); } +async function readDuckDuckGoHtmlResponse(response: Response): Promise { + return await readProviderTextResponse(response, "DuckDuckGo search"); +} + function parseDuckDuckGoHtml(html: string): DuckDuckGoResult[] { const results: DuckDuckGoResult[] = []; const resultRegex = /]*\bclass="[^"]*\bresult__a\b[^"]*")([^>]*)>([\s\S]*?)<\/a>/gi; @@ -174,7 +179,7 @@ export async function runDuckDuckGoSearch(params: { ); } - const html = await response.text(); + const html = await readDuckDuckGoHtmlResponse(response); if (isBotChallenge(html)) { throw new Error("DuckDuckGo returned a bot-detection challenge."); } @@ -210,5 +215,6 @@ export const testing = { decodeHtmlEntities, isBotChallenge, parseDuckDuckGoHtml, + readDuckDuckGoHtmlResponse, }; export { testing as __testing }; diff --git a/extensions/duckduckgo/src/ddg-search-provider.test.ts b/extensions/duckduckgo/src/ddg-search-provider.test.ts index 2728f3ca9a43..13b89dd80162 100644 --- a/extensions/duckduckgo/src/ddg-search-provider.test.ts +++ b/extensions/duckduckgo/src/ddg-search-provider.test.ts @@ -1,5 +1,6 @@ // Duckduckgo tests cover ddg search provider plugin behavior. import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; import { createDuckDuckGoWebSearchProvider as createDuckDuckGoWebSearchContractProvider } from "../web-search-contract-api.js"; import { DEFAULT_DDG_SAFE_SEARCH, resolveDdgRegion, resolveDdgSafeSearch } from "./config.js"; @@ -104,6 +105,24 @@ describe("duckduckgo web search provider", () => { expect(runDuckDuckGoSearch).not.toHaveBeenCalled(); }); + it("bounds successful DuckDuckGo HTML bodies without using response.text()", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "Content-Type": "text/html" }, + }); + const textSpy = vi.spyOn(streamed.response, "text").mockRejectedValue(new Error("unbounded")); + + await expect(ddgClientTesting.readDuckDuckGoHtmlResponse(streamed.response)).rejects.toThrow( + "DuckDuckGo search: text response exceeds 16777216 bytes", + ); + + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(textSpy).not.toHaveBeenCalled(); + }); + it("reads region from plugin config and normalizes empty values away", () => { expect( resolveDdgRegion({ diff --git a/extensions/firecrawl/src/firecrawl-client.ts b/extensions/firecrawl/src/firecrawl-client.ts index 162bf6042d27..fdbb4c14ff25 100644 --- a/extensions/firecrawl/src/firecrawl-client.ts +++ b/extensions/firecrawl/src/firecrawl-client.ts @@ -1,5 +1,6 @@ // Firecrawl plugin module implements firecrawl client behavior. import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; import { DEFAULT_CACHE_TTL_MINUTES, markdownToText, @@ -41,6 +42,7 @@ const SCRAPE_CACHE = new Map< >(); const DEFAULT_SEARCH_COUNT = 5; const DEFAULT_SCRAPE_MAX_CHARS = 50_000; +const FIRECRAWL_SCRAPE_RESPONSE_MAX_BYTES = 64 * 1024 * 1024; const ALLOWED_FIRECRAWL_HOSTS = new Set(["api.firecrawl.dev"]); const FIRECRAWL_SELF_HOSTED_PRIVATE_ERROR = "Firecrawl custom baseUrl must target a private or internal self-hosted endpoint."; @@ -65,12 +67,9 @@ type FirecrawlSearchItem = { async function readFirecrawlJsonResponse( response: Response, label: string, + opts?: { maxBytes?: number }, ): Promise> { - try { - return (await response.json()) as Record; - } catch (cause) { - throw new Error(`${label}: malformed JSON response`, { cause }); - } + return await readProviderJsonResponse>(response, label, opts); } export type FirecrawlSearchParams = { @@ -220,11 +219,9 @@ async function postFirecrawlJson( const readJsonPayload = async (): Promise | null> => { const candidate = response as Response & { clone?: () => Response }; const jsonResponse = typeof candidate.clone === "function" ? candidate.clone() : response; - if (typeof jsonResponse.json !== "function") { - return null; - } try { - const payload = await jsonResponse.json(); + const body = await readResponseText(jsonResponse, { maxBytes: 64_000 }); + const payload = JSON.parse(body.text) as unknown; return payload && typeof payload === "object" && !Array.isArray(payload) ? (payload as Record) : null; @@ -579,7 +576,10 @@ export async function runFirecrawlScrape( }, }, async (response) => { - const payloadLocal = await readFirecrawlJsonResponse(response, "Firecrawl fetch failed"); + const payloadLocal = await readFirecrawlJsonResponse(response, "Firecrawl fetch failed", { + // Scrape can legitimately return page bodies before maxChars truncates parsed output. + maxBytes: FIRECRAWL_SCRAPE_RESPONSE_MAX_BYTES, + }); if (payloadLocal.success === false) { const detail = typeof payloadLocal.error === "string" @@ -613,6 +613,7 @@ export const testing = { assertFirecrawlScrapeTargetAllowed, parseFirecrawlScrapePayload, postFirecrawlJson, + readFirecrawlJsonResponse, resolveEndpoint, validateFirecrawlBaseUrl, resolveSearchItems, diff --git a/extensions/firecrawl/src/firecrawl-tools.test.ts b/extensions/firecrawl/src/firecrawl-tools.test.ts index fa2168e7b829..8ea338dfabb5 100644 --- a/extensions/firecrawl/src/firecrawl-tools.test.ts +++ b/extensions/firecrawl/src/firecrawl-tools.test.ts @@ -2,6 +2,7 @@ import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { mockPinnedHostnameResolution } from "openclaw/plugin-sdk/test-env"; import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; import { DEFAULT_FIRECRAWL_BASE_URL, DEFAULT_FIRECRAWL_MAX_AGE_MS, @@ -966,6 +967,27 @@ describe("firecrawl tools", () => { ).rejects.toThrow("Firecrawl Search API error: malformed JSON response"); }); + it("bounds successful Firecrawl JSON bodies before parsing", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const jsonSpy = vi.spyOn(streamed.response, "json").mockRejectedValue(new Error("unbounded")); + + await expect( + firecrawlClientTesting.readFirecrawlJsonResponse( + streamed.response, + "Firecrawl Search API error", + ), + ).rejects.toThrow("Firecrawl Search API error: JSON response exceeds 16777216 bytes"); + + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(jsonSpy).not.toHaveBeenCalled(); + }); + it("reports malformed Firecrawl scrape JSON with a stable provider error", async () => { global.fetch = vi.fn( async () => diff --git a/extensions/ollama/src/embedding-provider.test.ts b/extensions/ollama/src/embedding-provider.test.ts index 95d00edd0efc..52957c2792bb 100644 --- a/extensions/ollama/src/embedding-provider.test.ts +++ b/extensions/ollama/src/embedding-provider.test.ts @@ -1,6 +1,7 @@ // Ollama tests cover embedding provider plugin behavior. import type { OpenClawConfig } from "openclaw/plugin-sdk/provider-auth"; import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; const { fetchConfiguredLocalOriginWithSsrFGuardMock } = vi.hoisted(() => ({ fetchConfiguredLocalOriginWithSsrFGuardMock: vi.fn( @@ -412,10 +413,40 @@ describe("ollama embedding provider", () => { }); await expect(provider.embedQuery("hello")).rejects.toThrow( - "Ollama embed response returned malformed JSON", + "Ollama embed response: malformed JSON response", ); }); + it("bounds successful embed JSON bodies before parsing", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const jsonSpy = vi.spyOn(streamed.response, "json").mockRejectedValue(new Error("unbounded")); + vi.stubGlobal( + "fetch", + vi.fn(async () => streamed.response), + ); + + const { provider } = await createOllamaEmbeddingProvider({ + config: {} as OpenClawConfig, + provider: "ollama", + model: "nomic-embed-text", + fallback: "none", + remote: { baseUrl: "http://127.0.0.1:11434" }, + }); + + await expect(provider.embedQuery("hello")).rejects.toThrow( + "Ollama embed response: JSON response exceeds 16777216 bytes", + ); + + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(jsonSpy).not.toHaveBeenCalled(); + }); + it("rejects non-number embedding values instead of zeroing them", async () => { vi.stubGlobal( "fetch", diff --git a/extensions/ollama/src/embedding-provider.ts b/extensions/ollama/src/embedding-provider.ts index 9471fae8bc17..1fd0521e6c60 100644 --- a/extensions/ollama/src/embedding-provider.ts +++ b/extensions/ollama/src/embedding-provider.ts @@ -6,7 +6,10 @@ import { normalizeOptionalSecretInput, } from "openclaw/plugin-sdk/provider-auth"; import { resolveEnvApiKey } from "openclaw/plugin-sdk/provider-auth-runtime"; -import { readResponseTextLimited } from "openclaw/plugin-sdk/provider-http"; +import { + readProviderJsonResponse, + readResponseTextLimited, +} from "openclaw/plugin-sdk/provider-http"; import { normalizeProviderId } from "openclaw/plugin-sdk/provider-model-shared"; import { hasConfiguredSecretInput, @@ -117,14 +120,9 @@ async function withRemoteHttpResponse(params: { } async function readOllamaEmbeddingJsonResponse( - response: Pick, + response: Response, ): Promise<{ embeddings?: unknown }> { - let payload: unknown; - try { - payload = await response.json(); - } catch (cause) { - throw new Error("Ollama embed response returned malformed JSON", { cause }); - } + const payload = await readProviderJsonResponse(response, "Ollama embed response"); if (typeof payload !== "object" || payload === null || Array.isArray(payload)) { throw new Error("Ollama embed response returned a non-object JSON payload"); } diff --git a/extensions/ollama/src/web-search-provider.test.ts b/extensions/ollama/src/web-search-provider.test.ts index 7b98109276dc..7e4dcc515780 100644 --- a/extensions/ollama/src/web-search-provider.test.ts +++ b/extensions/ollama/src/web-search-provider.test.ts @@ -1,6 +1,7 @@ // Ollama tests cover web search provider plugin behavior. import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; import { createOllamaWebSearchProvider as createContractOllamaWebSearchProvider } from "../web-search-contract-api.js"; import { testing, @@ -403,7 +404,32 @@ describe("ollama web search provider", () => { config: createOllamaConfig(), query: "openclaw", }), - ).rejects.toThrow("Ollama web search returned malformed JSON"); + ).rejects.toThrow("Ollama web search: malformed JSON response"); + }); + + it("bounds successful Ollama web search JSON bodies before parsing", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const jsonSpy = vi.spyOn(streamed.response, "json").mockRejectedValue(new Error("unbounded")); + fetchWithSsrFGuardMock.mockResolvedValueOnce({ + response: streamed.response, + release: vi.fn(async () => {}), + }); + + await expect( + runOllamaWebSearch({ + config: createOllamaConfig(), + query: "openclaw", + }), + ).rejects.toThrow("Ollama web search: JSON response exceeds 16777216 bytes"); + + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(jsonSpy).not.toHaveBeenCalled(); }); it("warns when Ollama is not reachable during setup without cancelling", async () => { diff --git a/extensions/ollama/src/web-search-provider.ts b/extensions/ollama/src/web-search-provider.ts index 1ef4ca8be3a8..55f13b037b68 100644 --- a/extensions/ollama/src/web-search-provider.ts +++ b/extensions/ollama/src/web-search-provider.ts @@ -5,6 +5,7 @@ import { normalizeOptionalSecretInput, } from "openclaw/plugin-sdk/provider-auth"; import { resolveEnvApiKey } from "openclaw/plugin-sdk/provider-auth-runtime"; +import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; import { enablePluginInConfig, readPositiveIntegerParam, @@ -67,11 +68,7 @@ type OllamaWebSearchAttempt = { }; async function readOllamaWebSearchResponse(response: Response): Promise { - try { - return (await response.json()) as OllamaWebSearchResponse; - } catch (cause) { - throw new Error("Ollama web search returned malformed JSON", { cause }); - } + return await readProviderJsonResponse(response, "Ollama web search"); } function isOllamaCloudBaseUrl(baseUrl: string): boolean { diff --git a/extensions/parallel/src/parallel-mcp-search.runtime.test.ts b/extensions/parallel/src/parallel-mcp-search.runtime.test.ts index e76e37ba8c78..597eb952f247 100644 --- a/extensions/parallel/src/parallel-mcp-search.runtime.test.ts +++ b/extensions/parallel/src/parallel-mcp-search.runtime.test.ts @@ -1,4 +1,5 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; type EndpointCall = { url: string; @@ -311,4 +312,27 @@ describe("runParallelMcpSearch", () => { expect(tracked.wasCanceled()).toBe(true); expect(textSpy).not.toHaveBeenCalled(); }); + + it("bounds successful MCP bodies without using response.text()", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "Content-Type": "application/json" }, + }); + const textSpy = vi.spyOn(streamed.response, "text").mockRejectedValue(new Error("unbounded")); + endpointMockState.responses.push(streamed.response); + + const error = await runParallelMcpSearch({ searchQueries: ["x"], maxResults: 5 }).catch( + (cause: unknown) => cause, + ); + + expect(error).toBeInstanceOf(Error); + expect((error as Error).message).toContain( + "Parallel MCP: text response exceeds 16777216 bytes", + ); + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.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 0b0031f3a09a..91153dd074a8 100644 --- a/extensions/parallel/src/parallel-mcp-search.runtime.ts +++ b/extensions/parallel/src/parallel-mcp-search.runtime.ts @@ -1,7 +1,10 @@ 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 { + readProviderTextResponse, + 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 @@ -218,7 +221,7 @@ async function postMcp(params: { status: response.status, statusText: response.statusText, text: response.ok - ? await response.text() + ? await readProviderTextResponse(response, "Parallel MCP") : 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.test.ts b/extensions/parallel/src/parallel-web-search-provider.test.ts index b83081548680..3f19400e1024 100644 --- a/extensions/parallel/src/parallel-web-search-provider.test.ts +++ b/extensions/parallel/src/parallel-web-search-provider.test.ts @@ -1,4 +1,5 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; type EndpointCall = { url: string; @@ -59,40 +60,6 @@ function cancelTrackedResponse( }; } -function streamedJsonResponse(params: { chunkCount: number; chunkSize: number }): { - response: Response; - getReadCount: () => number; - wasCanceled: () => boolean; -} { - // Multi-chunk fixture: proves the bounded read stops pulling chunks before - // the whole (here syntactically broken / unbounded) body is buffered, and - // that the stream is cancelled on overflow. - let reads = 0; - let canceled = false; - const encoder = new TextEncoder(); - const stream = new ReadableStream({ - pull(controller) { - if (reads >= params.chunkCount) { - controller.close(); - return; - } - reads += 1; - controller.enqueue(encoder.encode("a".repeat(params.chunkSize))); - }, - cancel() { - canceled = true; - }, - }); - return { - response: new Response(stream, { - status: 200, - headers: { "Content-Type": "application/json" }, - }), - getReadCount: () => reads, - 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"; @@ -621,7 +588,12 @@ describe("parallel web search provider", () => { // 200-chunk x 1 MiB body (~200 MiB) caps at 16 MiB: the bounded reader must // stop pulling chunks and cancel the stream well before draining it, then // surface a bounded error rather than buffering the whole payload. - const streamed = streamedJsonResponse({ chunkCount: 200, chunkSize: 1024 * 1024 }); + const streamed = createStreamingResponse({ + chunkCount: 200, + chunkSize: 1024 * 1024, + text: "a", + headers: { "Content-Type": "application/json" }, + }); endpointMockState.responses.push(streamed.response); const provider = createParallelWebSearchProvider(); const tool = provider.createTool({ diff --git a/extensions/perplexity/src/perplexity-web-search-provider.runtime.ts b/extensions/perplexity/src/perplexity-web-search-provider.runtime.ts index 2fe4de817fa0..0d7d056138b6 100644 --- a/extensions/perplexity/src/perplexity-web-search-provider.runtime.ts +++ b/extensions/perplexity/src/perplexity-web-search-provider.runtime.ts @@ -1,3 +1,4 @@ +import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; // Perplexity provider module implements model/runtime integration. import { readPositiveIntegerParam, @@ -142,11 +143,7 @@ function buildPerplexityRequestHeaders(apiKey: string, acceptJson = false): Reco } async function readPerplexityJsonResponse(response: Response, label: string): Promise { - try { - return (await response.json()) as T; - } catch (cause) { - throw new Error(`${label}: malformed JSON response`, { cause }); - } + return await readProviderJsonResponse(response, label); } function resolvePerplexityTransport(perplexity?: PerplexityConfig): { diff --git a/extensions/perplexity/src/perplexity-web-search-provider.test.ts b/extensions/perplexity/src/perplexity-web-search-provider.test.ts index 52ebd6bd85ca..8998f9f60b7f 100644 --- a/extensions/perplexity/src/perplexity-web-search-provider.test.ts +++ b/extensions/perplexity/src/perplexity-web-search-provider.test.ts @@ -1,6 +1,7 @@ // Perplexity tests cover perplexity web search provider plugin behavior. import { withEnv, withEnvAsync } from "openclaw/plugin-sdk/test-env"; -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; import { createPerplexityWebSearchProvider } from "./perplexity-web-search-provider.js"; import { testing } from "./perplexity-web-search-provider.runtime.js"; @@ -171,4 +172,22 @@ describe("perplexity web search provider", () => { testing.readPerplexityJsonResponse(new Response("{ nope"), "Perplexity"), ).rejects.toThrow("Perplexity: malformed JSON response"); }); + + it("bounds successful Perplexity JSON bodies before parsing", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const jsonSpy = vi.spyOn(streamed.response, "json").mockRejectedValue(new Error("unbounded")); + + await expect( + testing.readPerplexityJsonResponse(streamed.response, "Perplexity Search"), + ).rejects.toThrow("Perplexity Search: JSON response exceeds 16777216 bytes"); + + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(jsonSpy).not.toHaveBeenCalled(); + }); }); diff --git a/extensions/qqbot/src/engine/api/api-client.test.ts b/extensions/qqbot/src/engine/api/api-client.test.ts index 521e28936308..9beba470eb33 100644 --- a/extensions/qqbot/src/engine/api/api-client.test.ts +++ b/extensions/qqbot/src/engine/api/api-client.test.ts @@ -1,5 +1,6 @@ // Qqbot tests cover api-client plugin behavior. import { afterEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../../../test-support/streaming-error-response.js"; const fetchWithSsrFGuardMock = vi.hoisted(() => vi.fn()); @@ -88,4 +89,35 @@ describe("ApiClient", () => { }, }); }); + + it("bounds successful response bodies without using response.text()", async () => { + const release = vi.fn(async () => {}); + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const textSpy = vi.spyOn(streamed.response, "text").mockRejectedValue(new Error("unbounded")); + fetchWithSsrFGuardMock.mockResolvedValueOnce({ + response: streamed.response, + release, + }); + + const client = new ApiClient({ baseUrl: "https://qqbot.test" }); + + let error: unknown; + try { + await client.request("token-1", "GET", "/v2/users/@me"); + } catch (caught) { + error = caught; + } + + expect(error).toBeInstanceOf(ApiError); + expect(String(error)).toContain("QQBot API response: text response exceeds 16777216 bytes"); + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(textSpy).not.toHaveBeenCalled(); + expect(release).toHaveBeenCalledTimes(1); + }); }); diff --git a/extensions/qqbot/src/engine/api/api-client.ts b/extensions/qqbot/src/engine/api/api-client.ts index 1036f08289fb..a8ce31309be4 100644 --- a/extensions/qqbot/src/engine/api/api-client.ts +++ b/extensions/qqbot/src/engine/api/api-client.ts @@ -9,7 +9,10 @@ * - `redactBodyKeys` replaces the hardcoded `file_data` redaction. */ -import { readResponseTextLimited } from "openclaw/plugin-sdk/provider-http"; +import { + readProviderTextResponse, + readResponseTextLimited, +} from "openclaw/plugin-sdk/provider-http"; import { fetchWithSsrFGuard, type SsrFPolicy } from "openclaw/plugin-sdk/ssrf-runtime"; import { ApiError, type ApiClientConfig, type EngineLogger } from "../types.js"; import { formatErrorMessage } from "../utils/format.js"; @@ -162,7 +165,7 @@ export class ApiClient { const readBody = async (limitBytes?: number): Promise => { try { return limitBytes === undefined - ? await res.text() + ? await readProviderTextResponse(res, "QQBot API response") : await readResponseTextLimited(res, limitBytes); } catch (err) { throw new ApiError( diff --git a/extensions/qqbot/src/engine/tools/channel-api.test.ts b/extensions/qqbot/src/engine/tools/channel-api.test.ts index ea64e764e2c0..a8cabd040ccc 100644 --- a/extensions/qqbot/src/engine/tools/channel-api.test.ts +++ b/extensions/qqbot/src/engine/tools/channel-api.test.ts @@ -1,5 +1,6 @@ // Qqbot tests cover channel-api tool behavior. import { afterEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../../../test-support/streaming-error-response.js"; const fetchWithSsrFGuardMock = vi.hoisted(() => vi.fn()); @@ -109,4 +110,33 @@ describe("executeChannelApi", () => { expect(textSpy).not.toHaveBeenCalled(); expect(release).toHaveBeenCalledTimes(1); }); + + it("bounds successful response bodies without using response.text()", async () => { + const release = vi.fn(async () => {}); + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const textSpy = vi.spyOn(streamed.response, "text").mockRejectedValue(new Error("unbounded")); + fetchWithSsrFGuardMock.mockResolvedValueOnce({ + response: streamed.response, + release, + }); + + const result = await executeChannelApi( + { method: "GET", path: "/guilds/123/channels" }, + { accessToken: "token-1" }, + ); + + expect(result.details).toMatchObject({ + error: "QQ channel API response: text response exceeds 16777216 bytes", + path: "/guilds/123/channels", + }); + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(textSpy).not.toHaveBeenCalled(); + expect(release).toHaveBeenCalledTimes(1); + }); }); diff --git a/extensions/qqbot/src/engine/tools/channel-api.ts b/extensions/qqbot/src/engine/tools/channel-api.ts index 4d0df88ed909..b9b8173eb016 100644 --- a/extensions/qqbot/src/engine/tools/channel-api.ts +++ b/extensions/qqbot/src/engine/tools/channel-api.ts @@ -8,7 +8,10 @@ * validation, fetch, and structured response formatting. */ -import { readResponseTextLimited } from "openclaw/plugin-sdk/provider-http"; +import { + readProviderTextResponse, + readResponseTextLimited, +} from "openclaw/plugin-sdk/provider-http"; import { fetchWithSsrFGuard, type SsrFPolicy } from "openclaw/plugin-sdk/ssrf-runtime"; import { formatErrorMessage } from "../utils/format.js"; import { debugLog, debugError } from "../utils/log.js"; @@ -216,7 +219,7 @@ export async function executeChannelApi( debugLog(`[qqbot-channel-api] <<< Status: ${res.status} ${res.statusText}`); const rawBody = res.ok - ? await res.text() + ? await readProviderTextResponse(res, "QQ channel API response") : await readResponseTextLimited(res, CHANNEL_API_ERROR_BODY_LIMIT_BYTES); if (!rawBody || rawBody.trim() === "") { if (res.ok) { diff --git a/extensions/tavily/src/tavily-client.test.ts b/extensions/tavily/src/tavily-client.test.ts index d5e6c3cd26ad..d0dce289e708 100644 --- a/extensions/tavily/src/tavily-client.test.ts +++ b/extensions/tavily/src/tavily-client.test.ts @@ -1,5 +1,6 @@ // Tavily tests cover tavily client plugin behavior. import { beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { createStreamingResponse } from "../../test-support/streaming-error-response.js"; // Capture every call to postTrustedWebToolsJson so we can assert on extraHeaders. const postTrustedWebToolsJson = vi.fn(); @@ -61,6 +62,29 @@ describe("tavily client X-Client-Source header", () => { ); }); + it("bounds successful Tavily JSON bodies before parsing", async () => { + const streamed = createStreamingResponse({ + chunkCount: 32, + chunkSize: 1024 * 1024, + text: "x", + headers: { "content-type": "application/json" }, + }); + const jsonSpy = vi.spyOn(streamed.response, "json").mockRejectedValue(new Error("unbounded")); + + postTrustedWebToolsJson.mockImplementationOnce( + async (_params: unknown, parse: (r: Response) => Promise) => + parse(streamed.response), + ); + + await expect(runTavilySearch({ query: "test query" })).rejects.toThrow( + "Tavily Search: JSON response exceeds 16777216 bytes", + ); + + expect(streamed.getReadCount()).toBeLessThan(32); + expect(streamed.wasCanceled()).toBe(true); + expect(jsonSpy).not.toHaveBeenCalled(); + }); + it("runTavilyExtract sends X-Client-Source: openclaw", async () => { await runTavilyExtract({ urls: ["https://example.com"] }); diff --git a/extensions/tavily/src/tavily-client.ts b/extensions/tavily/src/tavily-client.ts index 63b1337cbcba..d3704b4fe0f6 100644 --- a/extensions/tavily/src/tavily-client.ts +++ b/extensions/tavily/src/tavily-client.ts @@ -1,5 +1,6 @@ // Tavily plugin module implements tavily client behavior. import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http"; import { DEFAULT_CACHE_TTL_MINUTES, normalizeCacheKey, @@ -26,6 +27,7 @@ const EXTRACT_CACHE = new Map< { value: Record; expiresAt: number; insertedAt: number } >(); const DEFAULT_SEARCH_COUNT = 5; +const TAVILY_EXTRACT_RESPONSE_MAX_BYTES = 64 * 1024 * 1024; export type TavilySearchParams = { cfg?: OpenClawConfig; @@ -73,6 +75,7 @@ async function postTavilyJson(params: { apiKey: string; body: Record; errorLabel: string; + responseMaxBytes?: number; }): Promise> { return postTrustedWebToolsJson( { @@ -83,19 +86,19 @@ async function postTavilyJson(params: { errorLabel: params.errorLabel, extraHeaders: { "X-Client-Source": "openclaw" }, }, - async (response) => readTavilyJsonResponse(response, params.errorLabel), + async (response) => + readTavilyJsonResponse(response, params.errorLabel, { + maxBytes: params.responseMaxBytes, + }), ); } async function readTavilyJsonResponse( response: Response, label: string, + opts?: { maxBytes?: number }, ): Promise> { - try { - return (await response.json()) as Record; - } catch (cause) { - throw new Error(`${label}: malformed JSON response`, { cause }); - } + return await readProviderJsonResponse>(response, label, opts); } export async function runTavilySearch( @@ -255,6 +258,8 @@ export async function runTavilyExtract( apiKey, body, errorLabel: "Tavily Extract", + // Extract can include raw page content and image lists, unlike search metadata. + responseMaxBytes: TAVILY_EXTRACT_RESPONSE_MAX_BYTES, }); const rawResults = Array.isArray(payload.results) ? payload.results : []; diff --git a/extensions/test-support/streaming-error-response.ts b/extensions/test-support/streaming-error-response.ts index ea2777c42d3b..e7f2ad155965 100644 --- a/extensions/test-support/streaming-error-response.ts +++ b/extensions/test-support/streaming-error-response.ts @@ -1,11 +1,21 @@ -// Test Support plugin module implements streaming error response behavior. -export function createStreamingErrorResponse(params: { - status: number; +// Test Support plugin module implements streaming response fixtures. +export type StreamingResponseFixture = { + response: Response; + getReadCount: () => number; + wasCanceled: () => boolean; +}; + +export function createStreamingResponse(params: { + status?: number; chunkCount: number; chunkSize: number; - byte: number; -}): { response: Response; getReadCount: () => number } { + byte?: number; + text?: string; + headers?: HeadersInit; +}): StreamingResponseFixture { let reads = 0; + let canceled = false; + const encoder = new TextEncoder(); const stream = new ReadableStream({ pull(controller) { if (reads >= params.chunkCount) { @@ -13,11 +23,28 @@ export function createStreamingErrorResponse(params: { return; } reads += 1; - controller.enqueue(new Uint8Array(params.chunkSize).fill(params.byte)); + const chunk = + params.text !== undefined + ? encoder.encode(params.text.repeat(params.chunkSize)) + : new Uint8Array(params.chunkSize).fill(params.byte ?? 120); + controller.enqueue(chunk); + }, + cancel() { + canceled = true; }, }); return { - response: new Response(stream, { status: params.status }), + response: new Response(stream, { status: params.status ?? 200, headers: params.headers }), getReadCount: () => reads, + wasCanceled: () => canceled, }; } + +export function createStreamingErrorResponse(params: { + status: number; + chunkCount: number; + chunkSize: number; + byte: number; +}): StreamingResponseFixture { + return createStreamingResponse(params); +}