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..5d963881e1a 100644 --- a/tests/clients/remote-catalog.test.ts +++ b/tests/clients/remote-catalog.test.ts @@ -250,6 +250,55 @@ 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])); }, + // 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(() => {}); }, + }); + + 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); + }); + test("allows a catalog download to exceed five seconds while bytes keep arriving", async () => { const chunks = ['{"models":[', '{"slug":"provider/model"}', ']}']; const server = Bun.serve({