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
This commit is contained in:
roboomp
2026-06-24 14:11:39 +00:00
parent c053afa087
commit 44f3632cff
3 changed files with 117 additions and 10 deletions
+1
View File
@@ -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 `<bunfs-root>/<binary-name>` 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)).
@@ -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<Uint8Array>;
type BodyReader = {
read(): Promise<BodyReadResult>;
cancel(reason?: unknown): Promise<void>;
};
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<BodyReadResult> {
if (!signal) return await reader.read();
if (signal.aborted) throw abortReason(signal);
const abort = Promise.withResolvers<BodyReadResult>();
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<Response["body"]>,
signal?: AbortSignal,
): Promise<void> {
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<str
return data.tag_name.replace(/^v/, "");
}
// Download a file from URL
async function downloadFile(url: string, dest: string, signal?: AbortSignal): Promise<void> {
/** Download a tool asset without handing the streaming Response to Bun.write. */
export async function downloadFile(url: string, dest: string, signal?: AbortSignal): Promise<void> {
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
@@ -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<void>();
const body = new ReadableStream<Uint8Array>({
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);
});
});