From 44f3632cff9f3f78bcd6d18a3558677c2ee083fe Mon Sep 17 00:00:00 2001 From: roboomp Date: Wed, 24 Jun 2026 14:11:39 +0000 Subject: [PATCH] fix(tools): bounded tool asset downloads Stream fetched tool assets to disk under the existing download abort signal instead of passing the Response object to Bun.write. Remove partial files when a stalled body is aborted and cover completed plus stalled downloads with regression tests. Fixes #3369 --- packages/coding-agent/CHANGELOG.md | 1 + .../coding-agent/src/utils/tools-manager.ts | 77 ++++++++++++++++--- .../test/tools-manager-download.test.ts | 49 ++++++++++++ 3 files changed, 117 insertions(+), 10 deletions(-) create mode 100644 packages/coding-agent/test/tools-manager-download.test.ts diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 03ec72a40..6d2ec8f06 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -4,6 +4,7 @@ ### Fixed +- Fixed lazy tool auto-downloads hanging when `Bun.write(dest, response)` receives a streaming `fetch()` `Response`; tool assets now stream the response body to disk with the existing download abort signal and remove partial files on abort. ([#3369](https://github.com/can1357/oh-my-pi/issues/3369)) - Fixed all extension loading silently failing on the cross-compiled `omp-darwin-arm64` release binary (downloaded directly or via a Homebrew tap wrapper) because `__computeBunfsPackageRoot` mis-handled `import.meta.dir = "//root/omp-darwin-arm64"`. Bun 1.3.14 reports `/` for the compiled entry's `import.meta.dir`, but the pre-fix function joined `metaDir + "packages"` and produced `/root/omp-darwin-arm64/packages` — the binary basename was baked into every bunfs path, so the TypeBox/legacy-pi shims and every `@oh-my-pi/pi-*` package-root override failed `existsSync` validation and `resolveCanonicalPiSpecifier` fell through to a bunfs `Bun.resolveSync` that also could not find the module. The function now detects the bunfs-root + binary-basename shape (`path.basename(path.dirname(metaDir)) === "root"`) and strips the trailing binary segment by slicing the original `metaDir`; the production bunfs shim join path also preserves Bun's bunfs-native `//root` / `B:\~BUN\root` prefix that `path.join` would otherwise collapse. ([#3329](https://github.com/can1357/oh-my-pi/issues/3329)) - Fixed llama.cpp discovery to prefer per-model `/v1/models` `meta.n_ctx`/`meta.n_ctx_train` values, refresh selected models after lazy load, and bypass fresh-cache reuse so server restarts update context windows. ([#3310](https://github.com/can1357/oh-my-pi/issues/3310)) - Fixed `task.maxConcurrency: 0` serializing subagent spawns instead of running them unbounded. The settings UI labels `0` as "Unlimited", but the session-scoped spawn `Semaphore` clamped `max` via `Math.max(1, max)`, so the second subagent body in a batch always waited for the first to release the seat. The constructor now treats `max <= 0` (and any non-finite input) as unbounded via `Number.POSITIVE_INFINITY`, matching the eval `parallel()`/`pipeline()` worker-pool semantics ([#3305](https://github.com/can1357/oh-my-pi/issues/3305)). diff --git a/packages/coding-agent/src/utils/tools-manager.ts b/packages/coding-agent/src/utils/tools-manager.ts index 7d6921ec6..5afa45a92 100644 --- a/packages/coding-agent/src/utils/tools-manager.ts +++ b/packages/coding-agent/src/utils/tools-manager.ts @@ -8,6 +8,62 @@ const TOOLS_DIR = getToolsDir(); const TOOL_DOWNLOAD_TIMEOUT_MS = 120_000; const TOOL_METADATA_TIMEOUT_MS = 5000; +type BodyReadResult = Bun.ReadableStreamDefaultReadResult; +type BodyReader = { + read(): Promise; + cancel(reason?: unknown): Promise; +}; + +function isAbortLikeError(error: unknown): boolean { + return error instanceof Error && (error.name === "AbortError" || error.name === "TimeoutError"); +} + +function abortReason(signal: AbortSignal): unknown { + return signal.reason ?? new DOMException("The operation was aborted.", "AbortError"); +} + +async function readBodyChunk(reader: BodyReader, signal: AbortSignal | undefined): Promise { + if (!signal) return await reader.read(); + if (signal.aborted) throw abortReason(signal); + + const abort = Promise.withResolvers(); + const onAbort = () => abort.reject(abortReason(signal)); + signal.addEventListener("abort", onAbort, { once: true }); + try { + return await Promise.race([reader.read(), abort.promise]); + } finally { + signal.removeEventListener("abort", onAbort); + } +} + +async function writeResponseBody( + dest: string, + body: NonNullable, + signal?: AbortSignal, +): Promise { + const reader = body.getReader(); + const sink = Bun.file(dest).writer(); + let completed = false; + + try { + while (true) { + const { done, value } = await readBodyChunk(reader, signal); + if (done) break; + if (value) { + await sink.write(value); + } + } + await sink.end(); + completed = true; + } finally { + if (!completed) { + await reader.cancel().catch(() => {}); + await Promise.resolve(sink.end()).catch(() => {}); + await fs.promises.rm(dest, { force: true }).catch(() => {}); + } + } +} + interface ToolConfig { name: string; repo: string; // GitHub repo (e.g., "sharkdp/fd") @@ -154,25 +210,26 @@ async function getLatestVersion(repo: string, signal?: AbortSignal): Promise { +/** Download a tool asset without handing the streaming Response to Bun.write. */ +export async function downloadFile(url: string, dest: string, signal?: AbortSignal): Promise { + const downloadSignal = ptree.combineSignals(signal, TOOL_DOWNLOAD_TIMEOUT_MS); let response: Response; try { response = await fetch(url, { - signal: ptree.combineSignals(signal, TOOL_DOWNLOAD_TIMEOUT_MS), + signal: downloadSignal, }); + if (!response.ok) { + throw new Error(`Failed to download: ${response.status}`); + } else if (!response.body) { + throw new Error("No response body"); + } + await writeResponseBody(dest, response.body, downloadSignal); } catch (err) { - if (err instanceof Error && err.name === "AbortError") { + if (isAbortLikeError(err)) { throw new Error(`Download timed out: ${url}`); } throw err; } - if (!response.ok) { - throw new Error(`Failed to download: ${response.status}`); - } else if (!response.body) { - throw new Error("No response body"); - } - await Bun.write(dest, response); } // Download and install a tool diff --git a/packages/coding-agent/test/tools-manager-download.test.ts b/packages/coding-agent/test/tools-manager-download.test.ts new file mode 100644 index 000000000..ecc53af66 --- /dev/null +++ b/packages/coding-agent/test/tools-manager-download.test.ts @@ -0,0 +1,49 @@ +import { afterEach, describe, expect, it, vi } from "bun:test"; +import { downloadFile } from "@oh-my-pi/pi-coding-agent/utils/tools-manager"; +import { TempDir } from "@oh-my-pi/pi-utils"; + +function mockDownloadResponse(response: Response): void { + const fetchMock: typeof globalThis.fetch = Object.assign(async () => response, { + preconnect: globalThis.fetch.preconnect, + }); + vi.spyOn(globalThis, "fetch").mockImplementation(fetchMock); +} + +describe("tool asset downloads", () => { + afterEach(() => { + vi.restoreAllMocks(); + }); + + it("writes a completed response body to disk", async () => { + using tempDir = TempDir.createSync("@omp-tool-download-"); + const dest = tempDir.join("tool.bin"); + mockDownloadResponse(new Response("tool-bytes")); + + await downloadFile("https://example.test/tool.bin", dest); + + expect(await Bun.file(dest).text()).toBe("tool-bytes"); + }); + + it("aborts a stalled response body and removes the partial file", async () => { + using tempDir = TempDir.createSync("@omp-tool-download-stall-"); + const dest = tempDir.join("tool.bin"); + const stalled = Promise.withResolvers(); + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("partial")); + }, + pull() { + stalled.resolve(); + }, + }); + mockDownloadResponse(new Response(body)); + const controller = new AbortController(); + + const download = downloadFile("https://example.test/tool.bin", dest, controller.signal); + await stalled.promise; + controller.abort(new DOMException("The operation timed out.", "TimeoutError")); + + await expect(download).rejects.toThrow("Download timed out: https://example.test/tool.bin"); + expect(await Bun.file(dest).exists()).toBe(false); + }); +});