From 28393d4bbd82327aefe7db8e2d5f41721bb77eab Mon Sep 17 00:00:00 2001 From: Dirk <0668000837@xydigit.com> Date: Sun, 23 Aug 2026 18:47:26 +0800 Subject: [PATCH] fix(link-understanding): cancel unread error response bodies (#118683) Co-authored-by: Dirk <279172199+xzh-icenter@users.noreply.github.com> Punchcard-Session: quiet-river-timber-wh Co-authored-by: Vincent Koc --- .../runner.transport.test.ts | 130 ++++++++++++++++++ src/link-understanding/runner.ts | 4 +- 2 files changed, 133 insertions(+), 1 deletion(-) create mode 100644 src/link-understanding/runner.transport.test.ts diff --git a/src/link-understanding/runner.transport.test.ts b/src/link-understanding/runner.transport.test.ts new file mode 100644 index 000000000000..7362505b5af4 --- /dev/null +++ b/src/link-understanding/runner.transport.test.ts @@ -0,0 +1,130 @@ +import { createServer } from "node:http"; +import type { AddressInfo, Socket } from "node:net"; +import { describe, expect, it, vi } from "vitest"; +import type { MsgContext } from "../auto-reply/templating.js"; +import type { OpenClawConfig } from "../config/types.openclaw.js"; +import { runLinkUnderstanding } from "./runner.js"; + +const mocks = vi.hoisted(() => ({ + bodyCancel: vi.fn(), + releaseAfterCancel: vi.fn(), + runCommandWithTimeout: vi.fn(), +})); + +vi.mock("../infra/net/fetch-guard.js", async () => { + const actual = await vi.importActual( + "../infra/net/fetch-guard.js", + ); + return { + ...actual, + fetchWithSsrFGuard: async (params: Parameters[0]) => { + // Keep the real guarded transport while allowing only this test's loopback server. + const result = await actual.fetchWithSsrFGuard({ + ...params, + lookupFn: async () => [{ address: "127.0.0.1", family: 4 }], + policy: { ...params.policy, allowPrivateNetwork: true }, + }); + const body = result.response.body; + if (body) { + const cancel = body.cancel.bind(body); + vi.spyOn(body, "cancel").mockImplementation(async (reason?: unknown) => { + mocks.bodyCancel(reason); + await cancel(reason); + }); + } + const release = result.release; + result.release = async () => { + mocks.releaseAfterCancel(mocks.bodyCancel.mock.calls.length > 0); + await release(); + }; + return result; + }, + }; +}); + +vi.mock("../process/exec.js", async () => { + const actual = await vi.importActual("../process/exec.js"); + return { + ...actual, + runCommandWithTimeout: mocks.runCommandWithTimeout, + }; +}); + +function deferred() { + let resolve!: () => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +async function within(promise: Promise, timeoutMs: number, message: string): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(message)), timeoutMs); + }), + ]); + } finally { + if (timer) { + clearTimeout(timer); + } + } +} + +describe("runLinkUnderstanding transport cleanup", () => { + it("cancels a non-OK response body before releasing its guarded transport", async () => { + const sockets = new Set(); + const requestSocketClosed = deferred(); + const server = createServer((request, response) => { + request.socket.once("close", requestSocketClosed.resolve); + response.writeHead(500, { + "content-length": "1000000", + "content-type": "text/plain", + }); + // Leave the declared body unfinished so cleanup must actively cancel it. + response.write("error"); + }); + server.on("connection", (socket) => { + sockets.add(socket); + socket.once("close", () => sockets.delete(socket)); + }); + + try { + await new Promise((resolve) => { + server.listen(0, "127.0.0.1", resolve); + }); + const port = (server.address() as AddressInfo).port; + const url = `http://loopback.test:${port}/error`; + + const resultPromise = runLinkUnderstanding({ + cfg: { + tools: { + links: { + enabled: true, + models: [{ type: "cli", command: "summarize" }], + }, + }, + } as OpenClawConfig, + ctx: { Body: `see ${url}` } as MsgContext, + }); + + const result = await within(resultPromise, 1000, "link understanding did not finish"); + expect(result).toEqual({ urls: [url], outputs: [] }); + await within(requestSocketClosed.promise, 1000, "loopback socket stayed open"); + + expect(mocks.bodyCancel).toHaveBeenCalledOnce(); + expect(mocks.releaseAfterCancel).toHaveBeenCalledWith(true); + expect(mocks.runCommandWithTimeout).not.toHaveBeenCalled(); + } finally { + for (const socket of sockets) { + socket.destroy(); + } + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + } + }); +}); diff --git a/src/link-understanding/runner.ts b/src/link-understanding/runner.ts index 7aba0ec2f7fd..fec777029dc6 100644 --- a/src/link-understanding/runner.ts +++ b/src/link-understanding/runner.ts @@ -4,7 +4,7 @@ import type { OpenClawConfig } from "../config/types.openclaw.js"; import type { LinkModelConfig, LinkToolsConfig } from "../config/types.tools.js"; import { logVerbose, shouldLogVerbose } from "../globals.js"; // Link-understanding runner fetches allowed URLs and invokes configured commands with bounded content. -import { readResponseWithLimit } from "../infra/http-body.js"; +import { cancelUnreadResponseBody, readResponseWithLimit } from "../infra/http-body.js"; import { fetchWithSsrFGuard, GUARDED_FETCH_MODE } from "../infra/net/fetch-guard.js"; import { CLI_OUTPUT_MAX_BUFFER } from "../media-understanding/defaults.js"; import { resolveTimeoutMs } from "../media-understanding/resolve.js"; @@ -102,6 +102,8 @@ async function fetchLinkContent(params: { }); try { if (!response.ok) { + // Do not await: a debug-capture tee settles only after its sibling branch cancels. + void cancelUnreadResponseBody(response); throw new Error(`Link fetch failed with HTTP ${response.status}`); } const buffer = await readResponseWithLimit(response, CLI_OUTPUT_MAX_BUFFER);