From 12f577ec98e29aa17b823293c802fa079f8c0451 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Sun, 20 Sep 2026 13:09:47 +0900 Subject: [PATCH 1/2] fix(client): bound total hub catalog response lifetime A catalog response that keeps trickling bytes previously had no total deadline: only the inactivity window applied, so a slow-drip body could hold the request open indefinitely. Give the catalog body read an overall AbortSignal budget (24x the inactivity window, capped at 120s) and stop awaiting body cancel on error paths so rejection is not delayed by the stream. --- src/client/hub-client.ts | 21 ++++++++++------ tests/clients/remote-catalog.test.ts | 36 ++++++++++++++++++++++++++++ 2 files changed, 50 insertions(+), 7 deletions(-) diff --git a/src/client/hub-client.ts b/src/client/hub-client.ts index 47604e33982..be6c79431bc 100644 --- a/src/client/hub-client.ts +++ b/src/client/hub-client.ts @@ -103,7 +103,7 @@ async function fetchBounded( }); headerDeadline?.clear(); if (response.status >= 300 && response.status < 400 && response.status !== 304) { - try { await response.body?.cancel(); } catch { /* best effort */ } + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } throw new HubClientError("redirect_refused", "Hub request redirect was refused", response.status); } return response; @@ -118,15 +118,16 @@ async function fetchBounded( async function boundedText( response: Response, maxBytes: number, - options: { inactivityTimeoutMs?: number } = {}, + options: { signal?: AbortSignal; inactivityTimeoutMs?: number } = {}, ): Promise { const declared = Number(response.headers.get("content-length") ?? "0"); if (Number.isFinite(declared) && declared > maxBytes) { - try { await response.body?.cancel(); } catch { /* best effort */ } + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } throw new HubClientError("body_too_large", "Hub response exceeded the allowed size", response.status); } const result = await readBoundedResponseBytes(response, { maxBytes, + ...(options.signal === undefined ? {} : { signal: options.signal }), ...(options.inactivityTimeoutMs === undefined ? {} : { inactivityTimeoutMs: options.inactivityTimeoutMs }), }); if (result.oversized) { @@ -446,20 +447,26 @@ export async function downloadClientCatalog( headers, }, options.timeoutMs, "headers"); if (response.status === 304) { + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } throw new HubClientError("catalog_unexpected_304", "Hub answered 304 to an unconditional catalog request", 304); } if (!response.ok) { + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } const code = response.status === 401 ? "catalog_unauthorized" : `catalog_http_${response.status}`; throw new HubClientError(code, `Hub catalog request failed (${response.status})`, response.status); } if (!jsonCompatibleContentType(response)) { - try { await response.body?.cancel(); } catch { /* best effort */ } + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } throw new HubClientError("catalog_content_type_invalid", "Hub catalog response was not JSON", response.status); } let body: string; try { + const inactivityTimeoutMs = safeTimeout(options.timeoutMs); body = await boundedText(response, options.maxBytes ?? MAX_REMOTE_CATALOG_BYTES, { - inactivityTimeoutMs: safeTimeout(options.timeoutMs), + // Permit active catalog transfers to span multiple inactivity windows, + // while retaining the client's established maximum request lifetime. + signal: AbortSignal.timeout(Math.min(inactivityTimeoutMs * 24, 120_000)), + inactivityTimeoutMs, }); } catch (error) { if (error instanceof DOMException && error.name === "TimeoutError") { @@ -608,11 +615,11 @@ export async function downloadDesktop3pModels( }), }, options.timeoutMs); if (!response.ok || response.status === 304) { - try { await response.body?.cancel(); } catch { /* best effort */ } + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } throw new HubClientError(`desktop_snapshot_http_${response.status}`, "Hub Desktop model snapshot request failed", response.status); } if (!jsonCompatibleContentType(response)) { - try { await response.body?.cancel(); } catch { /* best effort */ } + try { void response.body?.cancel().catch(() => {}); } catch { /* best effort */ } throw new HubClientError("desktop_snapshot_invalid", "Hub Desktop model snapshot was invalid"); } const body = await boundedText(response, DESKTOP_SNAPSHOT_MAX_BYTES, { diff --git a/tests/clients/remote-catalog.test.ts b/tests/clients/remote-catalog.test.ts index 78ade3d974a..b8c2bf32002 100644 --- a/tests/clients/remote-catalog.test.ts +++ b/tests/clients/remote-catalog.test.ts @@ -250,6 +250,42 @@ describe("remote Desktop snapshot consumer", () => { }); describe("remote catalog adversarial consumer", () => { + test("enforces a total deadline even while catalog bytes keep arriving", async () => { + let timer: ReturnType | undefined; + let cancelled = false; + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode('{"models":[')); + timer = setInterval(() => controller.enqueue(new Uint8Array([0x20])), 5); + }, + cancel() { + cancelled = true; + clearInterval(timer); + }, + }); + + const startedAt = performance.now(); + await expect(downloadClientCatalog("https://hub.example.test", "ocx_data_test", { + timeoutMs: 20, + fetchImpl: async () => new Response(body, { headers: JSON_HEADERS }), + })).rejects.toMatchObject({ code: "unreachable" }); + expect(performance.now() - startedAt).toBeLessThan(1_500); + expect(cancelled).toBe(true); + }); + + test("cancels an HTTP error body before rejecting", async () => { + let cancelled = false; + const body = new ReadableStream({ + start(controller) { controller.enqueue(new Uint8Array([0x20])); }, + cancel() { cancelled = true; }, + }); + + await expect(downloadClientCatalog("https://hub.example.test", "ocx_data_test", { + fetchImpl: async () => new Response(body, { status: 500 }), + })).rejects.toMatchObject({ code: "catalog_http_500" }); + expect(cancelled).toBe(true); + }); + test("allows a catalog download to exceed five seconds while bytes keep arriving", async () => { const chunks = ['{"models":[', '{"slug":"provider/model"}', ']}']; const server = Bun.serve({ From 779ef91996b7a17052d1e5a72caef844095493ac Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Tue, 22 Sep 2026 08:45:09 +0900 Subject: [PATCH 2/2] test(client): bound error-body cancellation in the catalog test The error-path test now cancels through a promise that never settles, so a regression that awaits body cancellation fails the bounded race instead of passing silently. --- tests/clients/remote-catalog.test.ts | 21 +++++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/tests/clients/remote-catalog.test.ts b/tests/clients/remote-catalog.test.ts index b8c2bf32002..5d963881e1a 100644 --- a/tests/clients/remote-catalog.test.ts +++ b/tests/clients/remote-catalog.test.ts @@ -277,12 +277,25 @@ describe("remote catalog adversarial consumer", () => { let cancelled = false; const body = new ReadableStream({ start(controller) { controller.enqueue(new Uint8Array([0x20])); }, - cancel() { cancelled = true; }, + // A cancel that never settles must not hold the error path: the download + // still has to reject with the HTTP status error inside the bound below. + cancel() { cancelled = true; return new Promise(() => {}); }, }); - await expect(downloadClientCatalog("https://hub.example.test", "ocx_data_test", { - fetchImpl: async () => new Response(body, { status: 500 }), - })).rejects.toMatchObject({ code: "catalog_http_500" }); + let timer: ReturnType | undefined; + const bounded = Promise.race([ + downloadClientCatalog("https://hub.example.test", "ocx_data_test", { + fetchImpl: async () => new Response(body, { status: 500 }), + }), + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error("catalog download outlived a never-resolving body cancel")), 1_000); + }), + ]); + try { + await expect(bounded).rejects.toMatchObject({ code: "catalog_http_500" }); + } finally { + clearTimeout(timer); + } expect(cancelled).toBe(true); });