From 9a2e58d61806eb1c539c40e3dec9911ecabec8a1 Mon Sep 17 00:00:00 2001 From: Mathieu Tarral Date: Wed, 9 Sep 2026 03:48:29 +0200 Subject: [PATCH 1/6] feat: serve Windows PE blobs from Winbindex + symbol server Add a fast path in front of GET /blob/:hash: when the request carries ?filename= for a Windows PE file (.exe/.dll/.sys), resolve the file on Winbindex by content SHA-1 and stream verified bytes straight from Microsoft's public symbol server instead of proxying MinIO, so the public corpus never has to re-host Windows binaries. - src/winbindex.ts: isWindowsPEFilename, symbolServerUrl, resolveEntry (24h bounded in-memory cache of the per-filename index, negative-caches 404s), streamFromSymbolServer (on-the-fly SHA-1 verification, destroys the response on a post-send mismatch), tryServeFromWinbindex orchestrator. - rest-routes: fifth createRestRouter arg WinbindexConfig; the handler tries Winbindex first and falls through to the unchanged MinIO logic on every miss/failure except once bytes have been streamed. - index.ts: WINBINDEX_ENABLED / _DATA_URL / _SYMBOL_SERVER_URL / _FETCH_TIMEOUT_MS env vars, all optional with defaults. MinIO stays a pure GetObject proxy; nothing is ever written back. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01NmATMvFbdupfLRq7Ldwabv --- src/index.ts | 26 ++- src/rest-routes.ts | 33 ++++ src/winbindex.ts | 323 +++++++++++++++++++++++++++++++++++ tests/winbindex.test.ts | 363 ++++++++++++++++++++++++++++++++++++++++ 4 files changed, 744 insertions(+), 1 deletion(-) create mode 100644 src/winbindex.ts create mode 100644 tests/winbindex.test.ts diff --git a/src/index.ts b/src/index.ts index 1aa7ba3..6358f17 100644 --- a/src/index.ts +++ b/src/index.ts @@ -6,7 +6,7 @@ import { ApolloServerPluginDrainHttpServer } from "@apollo/server/plugin/drainHt import { readFileSync } from "fs"; import neo4j from "neo4j-driver"; import * as dotenv from "dotenv"; -import { cleanEnv, str, url } from "envalid"; +import { cleanEnv, str, url, bool, num } from "envalid"; import { createConstraintsIfNotExists } from "./constraints.js"; import { resolvers } from "./resolvers.js"; import { createRestRouter } from "./rest-routes.js"; @@ -55,6 +55,24 @@ const env = cleanEnv(process.env, { default: "", desc: "Comma-separated list of sensitive registry value names to redact", }), + // Winbindex fast path for Windows PE blob downloads -- see + // docs/reference/winbindex-source.md + WINBINDEX_ENABLED: bool({ + default: true, + desc: "Enable the Winbindex fast path for Windows PE blob downloads", + }), + WINBINDEX_DATA_URL: url({ + default: "https://winbindex.m417z.com/data/by_filename_compressed", + desc: "Winbindex per-filename JSON index host", + }), + WINBINDEX_SYMBOL_SERVER_URL: url({ + default: "https://msdl.microsoft.com/download/symbols", + desc: "Microsoft public symbol server (serves PE binaries by timestamp+size)", + }), + WINBINDEX_FETCH_TIMEOUT_MS: num({ + default: 15000, + desc: "Timeout for each Winbindex index / symbol-server request", + }), }); // Neo4j driver instance @@ -322,6 +340,12 @@ async function main() { env.MINIO_ACCESS_KEY, env.MINIO_SECRET_KEY, env.MINIO_OBJECTS_BUCKET_NAME, + { + enabled: env.WINBINDEX_ENABLED, + dataUrl: env.WINBINDEX_DATA_URL, + symbolServerUrl: env.WINBINDEX_SYMBOL_SERVER_URL, + timeoutMs: env.WINBINDEX_FETCH_TIMEOUT_MS, + }, ), ); diff --git a/src/rest-routes.ts b/src/rest-routes.ts index 9a0e214..94f1828 100644 --- a/src/rest-routes.ts +++ b/src/rest-routes.ts @@ -2,12 +2,14 @@ import { Router, Request, Response } from "express"; import { BlobHashParamSchema } from "./validation.js"; import { ZodError } from "zod"; import { S3Client, GetObjectCommand } from "@aws-sdk/client-s3"; +import { WinbindexConfig, tryServeFromWinbindex } from "./winbindex.js"; export const createRestRouter = ( objectStorageUri: string, minioAccessKey: string, minioSecretKey: string, minioObjectsBucketName: string, + winbindexConfig: WinbindexConfig, ) => { const router = Router(); @@ -30,6 +32,37 @@ export const createRestRouter = ( console.log(`Blob download requested: ${hash}`); + // Winbindex fast path: for Windows PE files the frontend passes + // `?filename=`, letting us resolve the file on Winbindex and stream + // verified bytes from Microsoft's symbol server instead of MinIO. + // Any miss or failure falls through to MinIO unchanged; a failure + // after bytes were already streamed ends the response here. + const filename = + typeof req.query.filename === "string" + ? req.query.filename + : undefined; + if (filename) { + try { + const outcome = await tryServeFromWinbindex( + winbindexConfig, + hash, + filename, + res, + ); + if ( + outcome === "served" || + outcome === "failed_after_send" + ) { + return; + } + } catch (winbindexError) { + console.warn( + "Winbindex fast path threw, falling back to MinIO:", + winbindexError, + ); + } + } + // Fetch blob from MinIO using S3 SDK const command = new GetObjectCommand({ Bucket: minioObjectsBucketName, diff --git a/src/winbindex.ts b/src/winbindex.ts new file mode 100644 index 0000000..0965845 --- /dev/null +++ b/src/winbindex.ts @@ -0,0 +1,323 @@ +import { createHash } from "node:crypto"; +import { gunzipSync } from "node:zlib"; + +/** + * Winbindex fast path for Windows PE blob downloads. + * + * For a `GET /blob/:hash?filename=` request the API resolves the file on + * Winbindex (per-filename index) and streams the verified bytes straight from + * Microsoft's public symbol server, so the public corpus never re-hosts Windows + * binaries. Every miss or failure falls back to MinIO, except once bytes have + * already been streamed to the client. + * + * See docs/reference/winbindex-source.md. + */ + +export const PE_EXTENSIONS = new Set([".exe", ".dll", ".sys"]); + +/** Parsed per-filename Winbindex JSON is cached for 24h. */ +export const WINBINDEX_JSON_TTL_MS = 86_400_000; + +/** At most this many parsed per-filename indexes are kept in memory. */ +const WINBINDEX_JSON_CACHE_MAX = 200; + +const SYMBOL_SERVER_USER_AGENT = "Microsoft-Symbol-Server/10.0.0.0"; + +export interface WinbindexConfig { + enabled: boolean; + dataUrl: string; + symbolServerUrl: string; + timeoutMs: number; +} + +export interface WinbindexEntry { + timestamp: number; + virtualSize: number; +} + +/** + * Outcome of an attempt to serve a blob from Winbindex. + * - `not_available`: nothing was written to the response, the caller must fall + * back to MinIO. + * - `served`: the response has been fully sent, the caller must not touch it. + * - `failed_after_send`: bytes were already streamed then something failed + * (post-send hash mismatch, upstream stream error); the response has been + * destroyed and the caller must not retry MinIO. + */ +export type WinbindexOutcome = "not_available" | "served" | "failed_after_send"; + +/** Minimal view of the HTTP response the streaming step needs. */ +export interface BlobResponse { + setHeader(name: string, value: string | number): void; + write(chunk: Uint8Array): boolean; + end(): void; + destroy(): void; +} + +interface WinbindexFileInfo { + timestamp?: number; + virtualSize?: number; + sha1?: string; +} + +type WinbindexJson = Record< + string, + { fileInfo?: WinbindexFileInfo } | undefined +>; + +interface CacheRecord { + /** `null` is the sentinel for "filename is 404 on Winbindex". */ + value: WinbindexJson | null; + expiresAt: number; +} + +const jsonCache = new Map(); + +/** Test hook: drop the module-level per-filename JSON cache. */ +export function __clearWinbindexCache(): void { + jsonCache.clear(); +} + +/** Last path segment of `filename`, handling both `/` and `\` separators. */ +function basename(filename: string): string { + const parts = filename.split(/[/\\]/); + return parts[parts.length - 1] ?? ""; +} + +/** True iff the lowercased basename of `filename` ends with a PE extension. */ +export function isWindowsPEFilename(filename: string): boolean { + const name = basename(filename).toLowerCase(); + for (const ext of PE_EXTENSIONS) { + if (name.length > ext.length && name.endsWith(ext)) { + return true; + } + } + return false; +} + +/** + * Symbol-server URL for a PE file: `///` where + * `TS = timestamp` as uppercase hex zero-padded to 8 chars and + * `VS = virtualSize` as lowercase hex with no padding. + */ +export function symbolServerUrl( + baseUrl: string, + name: string, + timestamp: number, + virtualSize: number, +): string { + const ts = timestamp.toString(16).toUpperCase().padStart(8, "0"); + const vs = virtualSize.toString(16); + return `${baseUrl.replace(/\/+$/, "")}/${name}/${ts}${vs}/${name}`; +} + +function cacheGet(name: string): { hit: boolean; value: WinbindexJson | null } { + const record = jsonCache.get(name); + if (!record) { + return { hit: false, value: null }; + } + if (Date.now() >= record.expiresAt) { + jsonCache.delete(name); + return { hit: false, value: null }; + } + return { hit: true, value: record.value }; +} + +function cacheSet(name: string, value: WinbindexJson | null): void { + if (!jsonCache.has(name) && jsonCache.size >= WINBINDEX_JSON_CACHE_MAX) { + const oldest = jsonCache.keys().next().value; + if (oldest !== undefined) { + jsonCache.delete(oldest); + } + } + jsonCache.set(name, { + value, + expiresAt: Date.now() + WINBINDEX_JSON_TTL_MS, + }); +} + +/** + * Fetch, gunzip and parse the per-filename Winbindex index. A 404 is cached as a + * `null` sentinel so a missing filename is not re-fetched on every request; + * transient failures (network error, timeout, non-404 status, parse error) are + * not cached. Never throws. + */ +async function fetchIndex( + cfg: WinbindexConfig, + name: string, +): Promise { + const url = `${cfg.dataUrl.replace(/\/+$/, "")}/${name}.json.gz`; + let response: Response; + try { + response = await fetch(url, { + signal: AbortSignal.timeout(cfg.timeoutMs), + }); + } catch { + return null; + } + + if (response.status === 404) { + cacheSet(name, null); + return null; + } + if (!response.ok) { + return null; + } + + try { + const compressed = new Uint8Array(await response.arrayBuffer()); + const parsed = JSON.parse( + gunzipSync(compressed).toString("utf-8"), + ) as WinbindexJson; + cacheSet(name, parsed); + return parsed; + } catch { + return null; + } +} + +/** + * Resolve `sha1` (raw SHA-1 of the file content) to a `{ timestamp, virtualSize }` + * entry via the per-filename Winbindex index. Returns `null` on any miss or + * failure and never throws. + */ +export async function resolveEntry( + cfg: WinbindexConfig, + name: string, + sha1: string, +): Promise { + const key = name.toLowerCase(); + const cached = cacheGet(key); + const index = cached.hit ? cached.value : await fetchIndex(cfg, key); + if (!index) { + return null; + } + + const target = sha1.toLowerCase(); + for (const entry of Object.values(index)) { + const fileInfo = entry?.fileInfo; + if ( + fileInfo && + typeof fileInfo.sha1 === "string" && + fileInfo.sha1.toLowerCase() === target + ) { + if ( + typeof fileInfo.timestamp === "number" && + typeof fileInfo.virtualSize === "number" + ) { + return { + timestamp: fileInfo.timestamp, + virtualSize: fileInfo.virtualSize, + }; + } + return null; + } + } + return null; +} + +/** + * Fetch the PE file from the symbol server and stream it to `res` while + * verifying its SHA-1. Returns `not_available` (nothing written) on a non-2xx + * upstream response or a fetch error, `served` on a verified transfer, and + * `failed_after_send` if the stream errored or the digest did not match after + * bytes were already sent (the response is destroyed in that case). + */ +export async function streamFromSymbolServer( + cfg: WinbindexConfig, + entry: WinbindexEntry, + name: string, + expectedSha1: string, + res: BlobResponse, +): Promise { + const url = symbolServerUrl( + cfg.symbolServerUrl, + name, + entry.timestamp, + entry.virtualSize, + ); + + let upstream: Response; + try { + upstream = await fetch(url, { + headers: { "User-Agent": SYMBOL_SERVER_USER_AGENT }, + redirect: "follow", + signal: AbortSignal.timeout(cfg.timeoutMs), + }); + } catch { + return "not_available"; + } + + if (!upstream.ok || !upstream.body) { + return "not_available"; + } + + res.setHeader("Content-Type", "application/octet-stream"); + res.setHeader("Content-Disposition", `attachment; filename="${name}"`); + const contentLength = upstream.headers.get("content-length"); + if (contentLength) { + res.setHeader("Content-Length", contentLength); + } + + const hash = createHash("sha1"); + const reader = upstream.body.getReader(); + try { + let chunk = await reader.read(); + while (!chunk.done) { + hash.update(chunk.value); + res.write(chunk.value); + chunk = await reader.read(); + } + } catch (error) { + console.warn( + `Winbindex: symbol-server stream for ${name} errored mid-transfer, destroying response:`, + error, + ); + res.destroy(); + return "failed_after_send"; + } + + const digest = hash.digest("hex"); + if (digest !== expectedSha1.toLowerCase()) { + console.warn( + `Winbindex: SHA-1 mismatch for ${name} (expected ${expectedSha1.toLowerCase()}, got ${digest}), destroying response`, + ); + res.destroy(); + return "failed_after_send"; + } + + res.end(); + return "served"; +} + +/** + * Orchestrator for the Winbindex fast path. Returns `not_available` (caller + * falls back to MinIO) unless the feature is enabled, the filename is a Windows + * PE file, and Winbindex resolves the hash; otherwise delegates to + * {@link streamFromSymbolServer}. Never throws. + */ +export async function tryServeFromWinbindex( + cfg: WinbindexConfig, + hash: string, + filename: string, + res: BlobResponse, +): Promise { + if (!cfg.enabled || !isWindowsPEFilename(filename)) { + return "not_available"; + } + + const name = basename(filename).toLowerCase(); + + let entry: WinbindexEntry | null; + try { + entry = await resolveEntry(cfg, name, hash); + } catch (error) { + console.warn(`Winbindex: unexpected error resolving ${name}:`, error); + return "not_available"; + } + if (!entry) { + return "not_available"; + } + + return streamFromSymbolServer(cfg, entry, name, hash, res); +} diff --git a/tests/winbindex.test.ts b/tests/winbindex.test.ts new file mode 100644 index 0000000..2ece27e --- /dev/null +++ b/tests/winbindex.test.ts @@ -0,0 +1,363 @@ +import { describe, it, expect, jest, beforeEach, afterEach } from "@jest/globals"; +import { createHash } from "node:crypto"; +import { gzipSync } from "node:zlib"; +import { + isWindowsPEFilename, + symbolServerUrl, + resolveEntry, + tryServeFromWinbindex, + __clearWinbindexCache, + WinbindexConfig, + BlobResponse, +} from "../src/winbindex.js"; + +const CONFIG: WinbindexConfig = { + enabled: true, + dataUrl: "https://winbindex.example/data/by_filename_compressed", + symbolServerUrl: "https://symbols.example/download/symbols", + timeoutMs: 5000, +}; + +const KERNEL32_SHA1 = "20735ae6dd1fe416d7c8c08389df321d7c917db5"; + +const bytes = (text: string): Uint8Array => new TextEncoder().encode(text); + +const sha1Hex = (data: Uint8Array): string => + createHash("sha1").update(data).digest("hex"); + +const gzip = (text: string): Uint8Array => Uint8Array.from(gzipSync(bytes(text))); + +const toArrayBuffer = (data: Uint8Array): ArrayBuffer => { + const ab = new ArrayBuffer(data.byteLength); + new Uint8Array(ab).set(data); + return ab; +}; + +/** A fake of the parts of a `fetch` Response the index path reads. */ +const indexResponse = (gunzipped: string, status = 200) => ({ + status, + ok: status >= 200 && status < 300, + arrayBuffer: async (): Promise => toArrayBuffer(gzip(gunzipped)), +}); + +const jsonIndexResponse = (value: unknown): ReturnType => + indexResponse(JSON.stringify(value)); + +const notFoundResponse = () => ({ status: 404, ok: false }); + +const streamOf = (data: Uint8Array): ReadableStream => + new ReadableStream({ + start(controller) { + controller.enqueue(data); + controller.close(); + }, + }); + +/** A fake of the parts of a `fetch` Response the symbol-server path reads. */ +const symbolResponse = (body: Uint8Array | null, status = 200) => ({ + ok: status >= 200 && status < 300, + body: body ? streamOf(body) : null, + headers: { + get: (name: string): string | null => + body && name.toLowerCase() === "content-length" + ? String(body.byteLength) + : null, + }, +}); + +interface MockRes extends BlobResponse { + headers: Record; + body: Uint8Array[]; + ended: boolean; + destroyed: boolean; +} + +const makeMockRes = (): MockRes => { + const headers: Record = {}; + const body: Uint8Array[] = []; + const res: MockRes = { + headers, + body, + ended: false, + destroyed: false, + setHeader(name, value) { + headers[name] = value; + }, + write(chunk) { + body.push(chunk); + return true; + }, + end() { + res.ended = true; + }, + destroy() { + res.destroyed = true; + }, + }; + return res; +}; + +const receivedText = (res: MockRes): string => { + const total = res.body.reduce((n, chunk) => n + chunk.byteLength, 0); + const merged = new Uint8Array(total); + let offset = 0; + for (const chunk of res.body) { + merged.set(chunk, offset); + offset += chunk.byteLength; + } + return new TextDecoder().decode(merged); +}; + +type FetchMock = jest.MockedFunction<(...args: any[]) => Promise>; +let fetchMock: FetchMock; + +beforeEach(() => { + __clearWinbindexCache(); + fetchMock = jest.fn<(...args: any[]) => Promise>(); + globalThis.fetch = fetchMock as unknown as typeof fetch; + jest.spyOn(console, "warn").mockImplementation(() => {}); +}); + +afterEach(() => { + jest.restoreAllMocks(); +}); + +describe("isWindowsPEFilename", () => { + it.each(["KERNEL32.DLL", "ntoskrnl.exe", "a/b/c/win32k.sys", "foo.EXE"])( + "is true for %s", + (name) => { + expect(isWindowsPEFilename(name)).toBe(true); + }, + ); + + it.each(["foo.txt", "libc.so.6", "", "notes", ".dll"])( + "is false for %s", + (name) => { + expect(isWindowsPEFilename(name)).toBe(false); + }, + ); +}); + +describe("symbolServerUrl", () => { + it("builds the timestamp+virtualSize path segment", () => { + expect( + symbolServerUrl( + "https://symbols.example/download/symbols", + "kernel32.dll", + 1584069829, + 118784, + ), + ).toBe( + "https://symbols.example/download/symbols/kernel32.dll/5E6AFCC51d000/kernel32.dll", + ); + }); +}); + +describe("resolveEntry", () => { + it("returns { timestamp, virtualSize } for a matching sha1 (case-insensitive)", async () => { + const index = { + somesha256: { + fileInfo: { + timestamp: 1584069829, + virtualSize: 118784, + sha1: KERNEL32_SHA1.toUpperCase(), + }, + }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + await expect( + resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1), + ).resolves.toEqual({ timestamp: 1584069829, virtualSize: 118784 }); + }); + + it("returns null when no entry carries the requested sha1", async () => { + const index = { + a: { fileInfo: { timestamp: 1, virtualSize: 2, sha1: "deadbeef" } }, + b: { fileInfo: { timestamp: 3, virtualSize: 4 } }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + await expect( + resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1), + ).resolves.toBeNull(); + }); + + it("returns null on a 404 from Winbindex", async () => { + fetchMock.mockResolvedValue(notFoundResponse()); + + await expect( + resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1), + ).resolves.toBeNull(); + }); + + it("returns null when fetch rejects", async () => { + fetchMock.mockRejectedValue(new Error("network down")); + + await expect( + resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1), + ).resolves.toBeNull(); + }); + + it("returns null when the gzipped payload is not valid JSON", async () => { + fetchMock.mockResolvedValue(indexResponse("{ not json")); + + await expect( + resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1), + ).resolves.toBeNull(); + }); + + it("returns null when the matching entry lacks timestamp or virtualSize", async () => { + const index = { + a: { fileInfo: { virtualSize: 118784, sha1: KERNEL32_SHA1 } }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + await expect( + resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1), + ).resolves.toBeNull(); + }); + + it("caches the parsed index per filename and re-fetches after __clearWinbindexCache()", async () => { + const index = { + a: { + fileInfo: { + timestamp: 1584069829, + virtualSize: 118784, + sha1: KERNEL32_SHA1, + }, + }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + await resolveEntry(CONFIG, "KERNEL32.DLL", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(1); + + __clearWinbindexCache(); + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(2); + }); + + it("caches a 404 as a sentinel so the filename is not re-fetched", async () => { + fetchMock.mockResolvedValue(notFoundResponse()); + + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + + expect(fetchMock).toHaveBeenCalledTimes(1); + }); +}); + +describe("tryServeFromWinbindex", () => { + const entryIndex = (sha1: string) => ({ + k: { + fileInfo: { timestamp: 1584069829, virtualSize: 118784, sha1 }, + }, + }); + + it("returns not_available for a non-PE filename without calling fetch", async () => { + const res = makeMockRes(); + + await expect( + tryServeFromWinbindex(CONFIG, "a".repeat(40), "readme.txt", res), + ).resolves.toBe("not_available"); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it("returns not_available when the feature is disabled without calling fetch", async () => { + const res = makeMockRes(); + + await expect( + tryServeFromWinbindex( + { ...CONFIG, enabled: false }, + "a".repeat(40), + "kernel32.dll", + res, + ), + ).resolves.toBe("not_available"); + expect(fetchMock).not.toHaveBeenCalled(); + }); + + it("streams verified bytes and returns served on the happy path", async () => { + const bodyText = "MZ...fake portable executable bytes..."; + const body = bytes(bodyText); + const hash = sha1Hex(body); + fetchMock.mockImplementation(async (input: unknown) => { + return String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(body); + }); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("served"); + expect(receivedText(res)).toBe(bodyText); + expect(res.headers["Content-Type"]).toBe("application/octet-stream"); + expect(res.headers["Content-Disposition"]).toBe( + 'attachment; filename="kernel32.dll"', + ); + expect(res.headers["Content-Length"]).toBe(String(body.byteLength)); + expect(res.ended).toBe(true); + expect(res.destroyed).toBe(false); + expect(fetchMock).toHaveBeenLastCalledWith( + "https://symbols.example/download/symbols/kernel32.dll/5E6AFCC51d000/kernel32.dll", + expect.objectContaining({ + headers: { "User-Agent": "Microsoft-Symbol-Server/10.0.0.0" }, + redirect: "follow", + }), + ); + }); + + it("returns not_available and writes nothing when the symbol server 500s", async () => { + const hash = sha1Hex(bytes("kernel32 body")); + fetchMock.mockImplementation(async (input: unknown) => { + return String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(null, 500); + }); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("not_available"); + expect(res.body).toHaveLength(0); + expect(res.headers).toEqual({}); + expect(res.ended).toBe(false); + expect(res.destroyed).toBe(false); + }); + + it("destroys the response and returns failed_after_send on a post-send SHA-1 mismatch", async () => { + const requestedHash = sha1Hex(bytes("what the caller asked for")); + const servedBody = bytes("something else entirely"); + fetchMock.mockImplementation(async (input: unknown) => { + return String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(requestedHash)) + : symbolResponse(servedBody); + }); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + requestedHash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("failed_after_send"); + expect(res.destroyed).toBe(true); + expect(res.ended).toBe(false); + expect(console.warn).toHaveBeenCalled(); + }); +}); From 2a2799802be728eb49e01fc9f139c15fe70f946a Mon Sep 17 00:00:00 2001 From: Mathieu Tarral Date: Wed, 9 Sep 2026 03:48:29 +0200 Subject: [PATCH 2/6] docs: document the Winbindex blob source fast path Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01NmATMvFbdupfLRq7Ldwabv --- CLAUDE.md | 6 +++ README.md | 7 +++ docs/reference/winbindex-source.md | 83 ++++++++++++++++++++++++++++++ 3 files changed, 96 insertions(+) create mode 100644 docs/reference/winbindex-source.md diff --git a/CLAUDE.md b/CLAUDE.md index 5284dc8..8b83b1e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -87,6 +87,10 @@ Custom resolvers for operations not auto-generated by Neo4j GraphQL: - Express router for blob downloads at `/blob/:hash` - Uses S3 SDK to fetch blobs from MinIO with authenticated requests - No per-blob access control: every stored blob is downloadable by hash +- Winbindex fast path (`src/winbindex.ts`): when `?filename=` names a Windows PE + file (`.exe`/`.dll`/`.sys`), the blob is resolved on Winbindex and streamed + verified from Microsoft's symbol server instead of MinIO; every miss or failure + falls back to MinIO unchanged. See `docs/reference/winbindex-source.md`. #### Data Processing @@ -119,6 +123,7 @@ See README.md for complete `.env` setup. Key variables: - **Neo4j**: `NEO4J_URI`, `NEO4J_USER`, `NEO4J_PASSWORD` - **Storage**: `OBJECT_STORAGE_URI`, `MINIO_ACCESS_KEY`, `MINIO_SECRET_KEY`, `MINIO_OBJECTS_BUCKET_NAME` - **Optional**: `NODE_ENV`, `SENSITIVE_REGISTRY_VALUES` +- **Winbindex fast path** (optional, all defaulted): `WINBINDEX_ENABLED`, `WINBINDEX_DATA_URL`, `WINBINDEX_SYMBOL_SERVER_URL`, `WINBINDEX_FETCH_TIMEOUT_MS` — see `docs/reference/winbindex-source.md` ## Testing Strategy - Jest 30 with ts-jest for TypeScript support @@ -135,6 +140,7 @@ The project includes comprehensive documentation in the `docs/` directory follow - **Reference** (`docs/reference/`) - API specifications - [Data Model](docs/reference/data-model.md) - **Core graph structure, merkle tree architecture, CRITICAL for Cypher queries** - [Blob Download API](docs/reference/blob-api.md) - Complete REST API spec with error codes, rate limits, examples + - [Winbindex Blob Source](docs/reference/winbindex-source.md) - Winbindex + symbol-server fast path for Windows PE blob downloads - [Access Restrictions](docs/reference/access-restrictions.md) - Security controls and filtering mechanisms - **How-To Guides** (`docs/how-to/`) - Step-by-step task instructions - [Integrate Blob API](docs/how-to/integrate-blob-api.md) - Frontend integration guide diff --git a/README.md b/README.md index d2b4733..93dd596 100644 --- a/README.md +++ b/README.md @@ -22,6 +22,13 @@ MINIO_OBJECTS_BUCKET_NAME=objects # Optional - Environment NODE_ENV=development +# Optional - Winbindex fast path for Windows PE blob downloads +# (all defaulted; see docs/reference/winbindex-source.md) +WINBINDEX_ENABLED=true +WINBINDEX_DATA_URL=https://winbindex.m417z.com/data/by_filename_compressed +WINBINDEX_SYMBOL_SERVER_URL=https://msdl.microsoft.com/download/symbols +WINBINDEX_FETCH_TIMEOUT_MS=15000 + # Optional - Registry Filtering # Unused unless the plugin commented out in src/index.ts is re-enabled # (see docs/reference/access-restrictions.md) diff --git a/docs/reference/winbindex-source.md b/docs/reference/winbindex-source.md new file mode 100644 index 0000000..2b50a63 --- /dev/null +++ b/docs/reference/winbindex-source.md @@ -0,0 +1,83 @@ +# Winbindex Blob Source Reference + +## What it is + +A fast path in front of the `GET /blob/:hash` handler (`src/rest-routes.ts`). +For Windows PE files the API resolves the file on +[Winbindex](https://winbindex.m417z.com/) and streams the verified bytes +straight from Microsoft's public symbol server, instead of proxying them from +MinIO. The public corpus therefore never has to re-host Windows binaries. + +Implementation: `src/winbindex.ts`. + +## When it engages + +All of the following must hold, otherwise the request is served from MinIO +exactly as before: + +| Condition | Detail | +|-----------|--------| +| `WINBINDEX_ENABLED` is true | Global kill switch. | +| The request carries `?filename=` | A non-string or absent `filename` skips the fast path with no external call. | +| `filename` is a Windows PE file | The lowercased basename ends with `.exe`, `.dll` or `.sys`. Detection is purely by extension. | +| Winbindex resolves the hash | The per-filename index contains an entry whose `fileInfo.sha1` equals the requested `hash`, with numeric `timestamp` and `virtualSize`. | +| The symbol server returns 2xx | A non-2xx response leaves the HTTP response untouched. | + +`filename` drives Winbindex only. The MinIO path always keys on `hash` and +ignores `filename`. + +## Request flow + +1. `GET /blob/?filename=kernel32.dll` +2. `GET /kernel32.dll.json.gz` — gzipped JSON, top-level keys + are SHA-256, each value carries `fileInfo.timestamp`, `fileInfo.virtualSize` + and (about 70% of the time) `fileInfo.sha1`. The parsed index is cached in + memory for 24h, keyed by lowercased filename, bounded to 200 entries. A 404 is + cached as a negative result so a missing filename is not re-fetched on every + request. +3. The entry whose `fileInfo.sha1` matches `` yields `timestamp` and + `virtualSize`. +4. `GET ///` where + `TS = timestamp.toString(16).toUpperCase().padStart(8, "0")` and + `VS = virtualSize.toString(16)` (lowercase, unpadded). Sent with + `User-Agent: Microsoft-Symbol-Server/10.0.0.0`, redirects followed. +5. The body is streamed to the client as + `Content-Type: application/octet-stream`, + `Content-Disposition: attachment; filename=""` (and `Content-Length` + when the upstream provides it), while the SHA-1 is recomputed on the fly. If + the final digest does not match ``, the response is destroyed + mid-transfer so the client sees a failed download. + +The `symbols` plugin in the OSWatcher collector already downloads PDBs from this +same server, so no new external trust boundary is introduced. + +## Configuration + +All four variables are optional and have defaults (`src/index.ts`). + +| Variable | Default | Purpose | +|----------|---------|---------| +| `WINBINDEX_ENABLED` | `true` | Enable the Winbindex fast path for Windows PE blob downloads. | +| `WINBINDEX_DATA_URL` | `https://winbindex.m417z.com/data/by_filename_compressed` | Winbindex per-filename JSON index host. | +| `WINBINDEX_SYMBOL_SERVER_URL` | `https://msdl.microsoft.com/download/symbols` | Microsoft public symbol server (serves PE binaries by timestamp+size). | +| `WINBINDEX_FETCH_TIMEOUT_MS` | `15000` | Timeout for each Winbindex index / symbol-server request. | + +## Failure modes + +| Situation | Behaviour | +|-----------|-----------| +| Feature disabled, no `filename`, or non-PE extension | MinIO, no external call. | +| Winbindex index 404 / network error / timeout / bad JSON | MinIO. | +| No matching `fileInfo.sha1`, or match missing `timestamp` / `virtualSize` | MinIO. | +| Symbol server returns non-2xx, or the fetch fails before any byte is sent | MinIO; nothing was written to the response. | +| Symbol-server stream errors after bytes were sent | Response destroyed; MinIO is **not** retried. | +| Downloaded bytes hash to something other than `` | Response destroyed after the fact; a warning is logged; MinIO is **not** retried. | + +MinIO is never written to by this feature — it is a pure proxy, `GetObject` +only. + +## Source Code + +- **Winbindex path**: `src/winbindex.ts` (`tryServeFromWinbindex`) +- **Integration point**: `src/rest-routes.ts` (`GET /:hash` handler) +- **Configuration**: `src/index.ts` From 378d3a34ded1a7dec2c2cb09f58f9ab63f32cb3c Mon Sep 17 00:00:00 2001 From: Mathieu Tarral Date: Wed, 9 Sep 2026 04:13:50 +0200 Subject: [PATCH 3/6] fix(winbindex): validate filename, honour stream backpressure, guard headers Final whole-branch review fixes for the Winbindex blob source fast path: - Reject non-plain `filename` values (SAFE_PE_NAME) at the top of tryServeFromWinbindex, before any URL/cache/header work. Closes a blind single-host SSRF: `?`, `#` and percent-encoded segments survive basename() and previously flowed unencoded into the outbound Winbindex/symbol-server URLs. One gate makes the cache key, both URLs and Content-Disposition safe by construction. - Honour res.write() backpressure in streamFromSymbolServer: await `drain` (bailing on client `close`/`error`) instead of letting Node buffer the whole PE for a slow client. BlobResponse gains once()/off(). - Stream error before the first byte now returns not_available (MinIO fallback) instead of destroying the response; only a failure after bytes were sent yields failed_after_send + res.destroy(). Headers are set lazily on the first chunk so a pre-first-byte failure leaves the response pristine. - Forward upstream Content-Length only when the body is not content-encoded (undici may have transparently decompressed). - Set ETag: "" on the winbindex-served response, for parity with MinIO. - Widen tryServeFromWinbindex's try/catch to cover streamFromSymbolServer, matching its never-throws docblock. - tests/winbindex.test.ts: 23 -> 35 cases covering all the above, plus JSON-cache TTL expiry and the 200-entry eviction bound. console.warn/error stubbed; output pristine. - docs: blob-api.md query-parameter section, integrate-blob-api.md sample now sends ?filename=, winbindex-source.md negative-cache caveat + updated flow. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01NmATMvFbdupfLRq7Ldwabv --- docs/how-to/integrate-blob-api.md | 10 +- docs/reference/blob-api.md | 6 + docs/reference/winbindex-source.md | 20 ++- src/winbindex.ts | 120 +++++++++++-- tests/winbindex.test.ts | 270 +++++++++++++++++++++++++++-- 5 files changed, 389 insertions(+), 37 deletions(-) diff --git a/docs/how-to/integrate-blob-api.md b/docs/how-to/integrate-blob-api.md index 4329d4e..cf453b3 100644 --- a/docs/how-to/integrate-blob-api.md +++ b/docs/how-to/integrate-blob-api.md @@ -45,7 +45,11 @@ export class BlobService { hash: string, filename?: string ): Promise { - const url = `${this.blobBaseUrl}/${hash}`; + // Pass `filename` as a query param so the API can serve Windows PE files + // from the Winbindex fast path. It is ignored for non-PE names. + const url = filename + ? `${this.blobBaseUrl}/${hash}?filename=${encodeURIComponent(filename)}` + : `${this.blobBaseUrl}/${hash}`; try { const response = await fetch(url); @@ -115,7 +119,9 @@ export class BlobService { onProgress: (percent: number) => void, filename?: string ): Promise { - const url = `${this.blobBaseUrl}/${hash}`; + const url = filename + ? `${this.blobBaseUrl}/${hash}?filename=${encodeURIComponent(filename)}` + : `${this.blobBaseUrl}/${hash}`; const response = await fetch(url); diff --git a/docs/reference/blob-api.md b/docs/reference/blob-api.md index b7c1a24..20e8e3e 100644 --- a/docs/reference/blob-api.md +++ b/docs/reference/blob-api.md @@ -16,6 +16,12 @@ Downloads a blob by its SHA-1 hash. |------|------|----------|-------------| | `hash` | string | Yes | SHA-1 hash of the blob (40 hexadecimal characters) | +### Query Parameters + +| Name | Type | Required | Description | +|------|------|----------|-------------| +| `filename` | string | No | The file's basename (e.g. `kernel32.dll`). Used **only** to look the blob up on Winbindex for the Windows-PE fast path (see [Winbindex Blob Source](./winbindex-source.md)). Ignored for non-PE names, and by the MinIO path, which always keys on `hash`. A non-plain name (anything outside `[a-z0-9._+-]`) is ignored. | + ## Request Headers The endpoint is unauthenticated: no headers are required. diff --git a/docs/reference/winbindex-source.md b/docs/reference/winbindex-source.md index 2b50a63..6da4853 100644 --- a/docs/reference/winbindex-source.md +++ b/docs/reference/winbindex-source.md @@ -20,6 +20,7 @@ exactly as before: | `WINBINDEX_ENABLED` is true | Global kill switch. | | The request carries `?filename=` | A non-string or absent `filename` skips the fast path with no external call. | | `filename` is a Windows PE file | The lowercased basename ends with `.exe`, `.dll` or `.sys`. Detection is purely by extension. | +| `filename` is a plain safe name | The lowercased basename matches `^[a-z0-9._+-]{1,255}$`. Names with `?`, `#`, whitespace or percent-encoding are rejected (they would otherwise steer the outbound request), and the request goes to MinIO. | | Winbindex resolves the hash | The per-filename index contains an entry whose `fileInfo.sha1` equals the requested `hash`, with numeric `timestamp` and `virtualSize`. | | The symbol server returns 2xx | A non-2xx response leaves the HTTP response untouched. | @@ -34,7 +35,10 @@ ignores `filename`. and (about 70% of the time) `fileInfo.sha1`. The parsed index is cached in memory for 24h, keyed by lowercased filename, bounded to 200 entries. A 404 is cached as a negative result so a missing filename is not re-fetched on every - request. + request. This in-memory negative cache is best-effort only: it is per-process, + unbounded in eviction pressure, and can be flushed by requests for 200 other + filenames, so a cache hit is a latency optimisation, not a reliability or + rate-limit guarantee. 3. The entry whose `fileInfo.sha1` matches `` yields `timestamp` and `virtualSize`. 4. `GET ///` where @@ -43,10 +47,14 @@ ignores `filename`. `User-Agent: Microsoft-Symbol-Server/10.0.0.0`, redirects followed. 5. The body is streamed to the client as `Content-Type: application/octet-stream`, - `Content-Disposition: attachment; filename=""` (and `Content-Length` - when the upstream provides it), while the SHA-1 is recomputed on the fly. If - the final digest does not match ``, the response is destroyed - mid-transfer so the client sees a failed download. + `Content-Disposition: attachment; filename=""`, `ETag: ""`, and + `Content-Length` when the upstream provides it and did not content-encode the + body (undici may have transparently decompressed it). `res.write()` + backpressure is honoured so a slow client cannot force the whole PE to buffer + in memory. The SHA-1 is recomputed on the fly; if the final digest does not + match ``, the response is destroyed mid-transfer so the client sees a + failed download. A stream error *before* the first byte falls back to MinIO; + an error after bytes were sent ends the response as a failure. The `symbols` plugin in the OSWatcher collector already downloads PDBs from this same server, so no new external trust boundary is introduced. @@ -69,7 +77,7 @@ All four variables are optional and have defaults (`src/index.ts`). | Feature disabled, no `filename`, or non-PE extension | MinIO, no external call. | | Winbindex index 404 / network error / timeout / bad JSON | MinIO. | | No matching `fileInfo.sha1`, or match missing `timestamp` / `virtualSize` | MinIO. | -| Symbol server returns non-2xx, or the fetch fails before any byte is sent | MinIO; nothing was written to the response. | +| Symbol server returns non-2xx, the fetch fails, or the stream errors before any byte is sent | MinIO; nothing was written to the response. | | Symbol-server stream errors after bytes were sent | Response destroyed; MinIO is **not** retried. | | Downloaded bytes hash to something other than `` | Response destroyed after the fact; a warning is logged; MinIO is **not** retried. | diff --git a/src/winbindex.ts b/src/winbindex.ts index 0965845..aa23af1 100644 --- a/src/winbindex.ts +++ b/src/winbindex.ts @@ -15,6 +15,15 @@ import { gunzipSync } from "node:zlib"; export const PE_EXTENSIONS = new Set([".exe", ".dll", ".sys"]); +/** + * A plain, safe PE filename: lowercase letters, digits and `. _ + -` only. + * Every real Winbindex filename fits this. Names outside it (containing `?`, + * `#`, whitespace, percent-encoding, ...) are attacker attempts to steer the + * outbound Winbindex / symbol-server request to another path or query and are + * rejected before any URL, cache key or header is built from the name. + */ +export const SAFE_PE_NAME = /^[a-z0-9._+-]{1,255}$/; + /** Parsed per-filename Winbindex JSON is cached for 24h. */ export const WINBINDEX_JSON_TTL_MS = 86_400_000; @@ -50,6 +59,8 @@ export type WinbindexOutcome = "not_available" | "served" | "failed_after_send"; export interface BlobResponse { setHeader(name: string, value: string | number): void; write(chunk: Uint8Array): boolean; + once(event: string, listener: (...args: unknown[]) => void): void; + off(event: string, listener: (...args: unknown[]) => void): void; end(): void; destroy(): void; } @@ -216,12 +227,47 @@ export async function resolveEntry( return null; } +/** + * Resolve once the response has drained and can take more data, or reject if the + * client went away first (so the caller stops pulling from the upstream instead + * of buffering forever). Listeners are always removed before settling. + */ +function waitForDrain(res: BlobResponse): Promise { + return new Promise((resolve, reject) => { + const cleanup = (): void => { + res.off("drain", onDrain); + res.off("close", onClose); + res.off("error", onError); + }; + const onDrain = (): void => { + cleanup(); + resolve(); + }; + const onClose = (): void => { + cleanup(); + reject(new Error("response closed before drain")); + }; + const onError = (err: unknown): void => { + cleanup(); + reject(err instanceof Error ? err : new Error(String(err))); + }; + res.once("drain", onDrain); + res.once("close", onClose); + res.once("error", onError); + }); +} + /** * Fetch the PE file from the symbol server and stream it to `res` while * verifying its SHA-1. Returns `not_available` (nothing written) on a non-2xx - * upstream response or a fetch error, `served` on a verified transfer, and - * `failed_after_send` if the stream errored or the digest did not match after - * bytes were already sent (the response is destroyed in that case). + * upstream response, a fetch error, or a stream error before the first byte + * reached the client; `served` on a verified transfer; and `failed_after_send` + * if the stream errored or the digest did not match *after* bytes were already + * sent (the response is destroyed in that case). + * + * `res.write()` backpressure is honoured: on a `false` return the loop awaits + * `drain` before reading more, so a slow client cannot make Node buffer the + * whole (potentially tens-of-MB) PE in memory. */ export async function streamFromSymbolServer( cfg: WinbindexConfig, @@ -252,23 +298,52 @@ export async function streamFromSymbolServer( return "not_available"; } - res.setHeader("Content-Type", "application/octet-stream"); - res.setHeader("Content-Disposition", `attachment; filename="${name}"`); + // Only forward the upstream Content-Length when the bytes are not + // content-encoded: undici may have transparently decompressed the body, in + // which case the upstream length no longer matches what the client receives. const contentLength = upstream.headers.get("content-length"); - if (contentLength) { - res.setHeader("Content-Length", contentLength); - } + const contentEncoding = ( + upstream.headers.get("content-encoding") ?? "" + ).toLowerCase(); + const forwardContentLength = + contentLength !== null && + (contentEncoding === "" || contentEncoding === "identity"); + + // Headers are set only once bytes are actually in hand, so a failure before + // the first byte leaves the response pristine for the MinIO fallback. + const setStreamHeaders = (): void => { + res.setHeader("Content-Type", "application/octet-stream"); + res.setHeader("Content-Disposition", `attachment; filename="${name}"`); + // The requested hash is exactly what the streamed bytes are verified + // against below; this gives the winbindex path ETag parity with MinIO. + res.setHeader("ETag", `"${expectedSha1.toLowerCase()}"`); + if (forwardContentLength && contentLength !== null) { + res.setHeader("Content-Length", contentLength); + } + }; const hash = createHash("sha1"); const reader = upstream.body.getReader(); + let wrote = false; try { let chunk = await reader.read(); while (!chunk.done) { hash.update(chunk.value); - res.write(chunk.value); + if (!wrote) { + setStreamHeaders(); + } + const flushed = res.write(chunk.value); + wrote = true; + if (!flushed) { + await waitForDrain(res); + } chunk = await reader.read(); } } catch (error) { + if (!wrote) { + // Nothing reached the client yet: fall back to MinIO. + return "not_available"; + } console.warn( `Winbindex: symbol-server stream for ${name} errored mid-transfer, destroying response:`, error, @@ -286,6 +361,10 @@ export async function streamFromSymbolServer( return "failed_after_send"; } + if (!wrote) { + // Zero-byte body that still verified: emit the headers before ending. + setStreamHeaders(); + } res.end(); return "served"; } @@ -308,16 +387,23 @@ export async function tryServeFromWinbindex( const name = basename(filename).toLowerCase(); - let entry: WinbindexEntry | null; - try { - entry = await resolveEntry(cfg, name, hash); - } catch (error) { - console.warn(`Winbindex: unexpected error resolving ${name}:`, error); + // Single validation gate. `basename()` strips `/` and `\` but not `?`, `#` + // or percent-encoded segments, which would otherwise flow unencoded into the + // outbound Winbindex / symbol-server URLs (blind single-host SSRF). Rejecting + // here makes the cache key, both URLs and the Content-Disposition value all + // safe by construction, with no per-site encoding needed. + if (!SAFE_PE_NAME.test(name)) { return "not_available"; } - if (!entry) { + + try { + const entry = await resolveEntry(cfg, name, hash); + if (!entry) { + return "not_available"; + } + return await streamFromSymbolServer(cfg, entry, name, hash, res); + } catch (error) { + console.warn(`Winbindex: unexpected error serving ${name}:`, error); return "not_available"; } - - return streamFromSymbolServer(cfg, entry, name, hash, res); } diff --git a/tests/winbindex.test.ts b/tests/winbindex.test.ts index 2ece27e..21f9592 100644 --- a/tests/winbindex.test.ts +++ b/tests/winbindex.test.ts @@ -7,6 +7,7 @@ import { resolveEntry, tryServeFromWinbindex, __clearWinbindexCache, + WINBINDEX_JSON_TTL_MS, WinbindexConfig, BlobResponse, } from "../src/winbindex.js"; @@ -53,28 +54,72 @@ const streamOf = (data: Uint8Array): ReadableStream => }, }); +/** A stream that yields `data` then errors, i.e. fails *after* the first byte. */ +const streamThenError = (data: Uint8Array): ReadableStream => { + let sent = false; + return new ReadableStream({ + pull(controller) { + if (!sent) { + sent = true; + controller.enqueue(data); + } else { + controller.error(new Error("upstream stream reset")); + } + }, + }); +}; + +/** A stream that errors before yielding anything. */ +const erroringStream = (): ReadableStream => + new ReadableStream({ + start(controller) { + controller.error(new Error("upstream stream reset")); + }, + }); + /** A fake of the parts of a `fetch` Response the symbol-server path reads. */ -const symbolResponse = (body: Uint8Array | null, status = 200) => ({ - ok: status >= 200 && status < 300, - body: body ? streamOf(body) : null, - headers: { - get: (name: string): string | null => - body && name.toLowerCase() === "content-length" - ? String(body.byteLength) - : null, - }, -}); +const symbolResponse = ( + body: Uint8Array | ReadableStream | null, + status = 200, + extraHeaders: Record = {}, +) => { + const stream = + body instanceof ReadableStream ? body : body ? streamOf(body) : null; + const headers: Record = {}; + if (body instanceof Uint8Array) { + headers["content-length"] = String(body.byteLength); + } + for (const [k, v] of Object.entries(extraHeaders)) { + headers[k.toLowerCase()] = v; + } + return { + ok: status >= 200 && status < 300, + body: stream, + headers: { + get: (name: string): string | null => + headers[name.toLowerCase()] ?? null, + }, + }; +}; interface MockRes extends BlobResponse { headers: Record; body: Uint8Array[]; ended: boolean; destroyed: boolean; + emit(event: string, ...args: unknown[]): void; } -const makeMockRes = (): MockRes => { +/** + * `writeReturns` feeds the boolean each `write()` call returns (default `true`). + * A `false` schedules a `drain` on the next microtask so the code under test + * resumes, exercising the backpressure path. + */ +const makeMockRes = (writeReturns: boolean[] = []): MockRes => { const headers: Record = {}; const body: Uint8Array[] = []; + const pending = [...writeReturns]; + const listeners: Record void>> = {}; const res: MockRes = { headers, body, @@ -85,7 +130,26 @@ const makeMockRes = (): MockRes => { }, write(chunk) { body.push(chunk); - return true; + const ok = pending.length ? (pending.shift() as boolean) : true; + if (!ok) { + queueMicrotask(() => res.emit("drain")); + } + return ok; + }, + once(event, listener) { + (listeners[event] ??= []).push(listener); + }, + off(event, listener) { + listeners[event] = (listeners[event] ?? []).filter( + (l) => l !== listener, + ); + }, + emit(event, ...args) { + const ls = listeners[event] ?? []; + listeners[event] = []; + for (const l of ls) { + l(...args); + } }, end() { res.ended = true; @@ -116,6 +180,7 @@ beforeEach(() => { fetchMock = jest.fn<(...args: any[]) => Promise>(); globalThis.fetch = fetchMock as unknown as typeof fetch; jest.spyOn(console, "warn").mockImplementation(() => {}); + jest.spyOn(console, "error").mockImplementation(() => {}); }); afterEach(() => { @@ -247,6 +312,58 @@ describe("resolveEntry", () => { expect(fetchMock).toHaveBeenCalledTimes(1); }); + + it("re-fetches the index once the 24h TTL has expired", async () => { + jest.useFakeTimers(); + try { + const index = { + a: { + fileInfo: { + timestamp: 1584069829, + virtualSize: 118784, + sha1: KERNEL32_SHA1, + }, + }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(1); + + jest.advanceTimersByTime(WINBINDEX_JSON_TTL_MS + 1); + await resolveEntry(CONFIG, "kernel32.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(2); + } finally { + jest.useRealTimers(); + } + }); + + it("bounds the cache to 200 entries, evicting the oldest first", async () => { + const index = { + a: { + fileInfo: { + timestamp: 1584069829, + virtualSize: 118784, + sha1: KERNEL32_SHA1, + }, + }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + for (let i = 0; i < 201; i++) { + await resolveEntry(CONFIG, `f${i}.dll`, KERNEL32_SHA1); + } + expect(fetchMock).toHaveBeenCalledTimes(201); + + // f0 was evicted when f200 was inserted -> re-fetch. + await resolveEntry(CONFIG, "f0.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(202); + + // f200 is still cached -> no fetch. + await resolveEntry(CONFIG, "f200.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(202); + }); }); describe("tryServeFromWinbindex", () => { @@ -360,4 +477,133 @@ describe("tryServeFromWinbindex", () => { expect(res.ended).toBe(false); expect(console.warn).toHaveBeenCalled(); }); + + it.each(["a?b.dll", "x#y.dll", "mal ware.dll", "%2e%2e.dll", "café.dll"])( + "returns not_available without fetching for an unsafe filename %s", + async (unsafe) => { + const res = makeMockRes(); + + await expect( + tryServeFromWinbindex(CONFIG, "a".repeat(40), unsafe, res), + ).resolves.toBe("not_available"); + expect(fetchMock).not.toHaveBeenCalled(); + expect(res.body).toHaveLength(0); + }, + ); + + it("sets a quoted ETag of the requested hash on a successful serve", async () => { + const body = bytes("MZ pe bytes"); + const hash = sha1Hex(body); + fetchMock.mockImplementation(async (input: unknown) => + String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(body), + ); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("served"); + expect(res.headers["ETag"]).toBe(`"${hash}"`); + }); + + it("omits Content-Length when the symbol server sent Content-Encoding: gzip", async () => { + const body = bytes("MZ decompressed pe bytes"); + const hash = sha1Hex(body); + fetchMock.mockImplementation(async (input: unknown) => + String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(body, 200, { + "content-encoding": "gzip", + "content-length": "11", + }), + ); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("served"); + expect(res.headers).not.toHaveProperty("Content-Length"); + }); + + it("returns not_available (no destroy) when the stream errors before the first byte", async () => { + const hash = sha1Hex(bytes("kernel32 body")); + fetchMock.mockImplementation(async (input: unknown) => + String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(erroringStream()), + ); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("not_available"); + expect(res.body).toHaveLength(0); + expect(res.destroyed).toBe(false); + expect(res.ended).toBe(false); + }); + + it("returns failed_after_send and destroys the response when the stream errors after bytes were sent", async () => { + const first = bytes("first chunk of the pe"); + // Requested hash need not match; the stream error fires first. + const hash = sha1Hex(bytes("whole file")); + fetchMock.mockImplementation(async (input: unknown) => + String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(streamThenError(first)), + ); + + const res = makeMockRes(); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("failed_after_send"); + expect(res.destroyed).toBe(true); + expect(res.body).toHaveLength(1); + expect(console.warn).toHaveBeenCalled(); + }); + + it("honours write() backpressure and still serves the whole body", async () => { + const body = bytes("MZ...a portable executable that needs draining..."); + const hash = sha1Hex(body); + fetchMock.mockImplementation(async (input: unknown) => + String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(body), + ); + + // First write() reports the buffer is full -> code must await "drain". + const res = makeMockRes([false]); + const outcome = await tryServeFromWinbindex( + CONFIG, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("served"); + expect(receivedText(res)).toBe( + "MZ...a portable executable that needs draining...", + ); + expect(res.ended).toBe(true); + }); }); From 38e91bd44b2cf2a151874866bbfe612fafb898fe Mon Sep 17 00:00:00 2001 From: Mathieu Tarral Date: Wed, 9 Sep 2026 13:38:07 +0200 Subject: [PATCH 4/6] fix(winbindex): true LRU cache, async gunzip, drain timeout Addresses Copilot review feedback on #72: - The per-filename index cache now refreshes recency on read, so eviction is genuinely least-recently-used rather than least-recently-inserted; a hot filename survives churn from 200 other lookups. - Index decompression moves off the event loop (async gunzip instead of gunzipSync), so a multi-MB index cannot stall unrelated traffic. - waitForDrain now rejects after WINBINDEX_FETCH_TIMEOUT_MS, so a socket that never emits drain/close/error tears the transfer down instead of hanging the request handler and holding the upstream reader open. +2 tests (LRU recency, drain timeout -> failed_after_send); 37/37 in tests/winbindex.test.ts, full suite 89 passing. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01NmATMvFbdupfLRq7Ldwabv --- docs/reference/winbindex-source.md | 6 ++-- src/winbindex.ts | 41 ++++++++++++++++------ tests/winbindex.test.ts | 56 +++++++++++++++++++++++++++++- 3 files changed, 89 insertions(+), 14 deletions(-) diff --git a/docs/reference/winbindex-source.md b/docs/reference/winbindex-source.md index 6da4853..4eb2b4c 100644 --- a/docs/reference/winbindex-source.md +++ b/docs/reference/winbindex-source.md @@ -32,8 +32,10 @@ ignores `filename`. 1. `GET /blob/?filename=kernel32.dll` 2. `GET /kernel32.dll.json.gz` — gzipped JSON, top-level keys are SHA-256, each value carries `fileInfo.timestamp`, `fileInfo.virtualSize` - and (about 70% of the time) `fileInfo.sha1`. The parsed index is cached in - memory for 24h, keyed by lowercased filename, bounded to 200 entries. A 404 is + and (about 70% of the time) `fileInfo.sha1`. The index is decompressed off the + event loop (async gunzip). The parsed result is cached in memory for 24h, + keyed by lowercased filename, bounded to 200 entries evicted least-recently-used. + A 404 is cached as a negative result so a missing filename is not re-fetched on every request. This in-memory negative cache is best-effort only: it is per-process, unbounded in eviction pressure, and can be flushed by requests for 200 other diff --git a/src/winbindex.ts b/src/winbindex.ts index aa23af1..c57ebdc 100644 --- a/src/winbindex.ts +++ b/src/winbindex.ts @@ -1,5 +1,9 @@ import { createHash } from "node:crypto"; -import { gunzipSync } from "node:zlib"; +import { promisify } from "node:util"; +import { gunzip } from "node:zlib"; + +/** Off-loads decompression to the threadpool so the event loop is not blocked. */ +const gunzipAsync = promisify(gunzip); /** * Winbindex fast path for Windows PE blob downloads. @@ -131,14 +135,20 @@ function cacheGet(name: string): { hit: boolean; value: WinbindexJson | null } { jsonCache.delete(name); return { hit: false, value: null }; } + // Move to the most-recently-used end so eviction in `cacheSet` is LRU: a + // `Map` keeps insertion order, so re-inserting is the cheapest bump. + jsonCache.delete(name); + jsonCache.set(name, record); return { hit: true, value: record.value }; } function cacheSet(name: string, value: WinbindexJson | null): void { - if (!jsonCache.has(name) && jsonCache.size >= WINBINDEX_JSON_CACHE_MAX) { - const oldest = jsonCache.keys().next().value; - if (oldest !== undefined) { - jsonCache.delete(oldest); + // Refresh recency even when the key already exists (re-insert at the end). + jsonCache.delete(name); + if (jsonCache.size >= WINBINDEX_JSON_CACHE_MAX) { + const lru = jsonCache.keys().next().value; + if (lru !== undefined) { + jsonCache.delete(lru); } } jsonCache.set(name, { @@ -177,9 +187,8 @@ async function fetchIndex( try { const compressed = new Uint8Array(await response.arrayBuffer()); - const parsed = JSON.parse( - gunzipSync(compressed).toString("utf-8"), - ) as WinbindexJson; + const json = (await gunzipAsync(compressed)).toString("utf-8"); + const parsed = JSON.parse(json) as WinbindexJson; cacheSet(name, parsed); return parsed; } catch { @@ -230,11 +239,17 @@ export async function resolveEntry( /** * Resolve once the response has drained and can take more data, or reject if the * client went away first (so the caller stops pulling from the upstream instead - * of buffering forever). Listeners are always removed before settling. + * of buffering forever). Rejects after `timeoutMs` as well, so a socket that + * never emits `drain`, `close` or `error` cannot hang the request handler. + * Listeners and the timer are always cleared before settling. */ -function waitForDrain(res: BlobResponse): Promise { +function waitForDrain(res: BlobResponse, timeoutMs: number): Promise { return new Promise((resolve, reject) => { + let timer: ReturnType | undefined; const cleanup = (): void => { + if (timer !== undefined) { + clearTimeout(timer); + } res.off("drain", onDrain); res.off("close", onClose); res.off("error", onError); @@ -251,6 +266,10 @@ function waitForDrain(res: BlobResponse): Promise { cleanup(); reject(err instanceof Error ? err : new Error(String(err))); }; + timer = setTimeout(() => { + cleanup(); + reject(new Error("timed out waiting for the response to drain")); + }, timeoutMs); res.once("drain", onDrain); res.once("close", onClose); res.once("error", onError); @@ -335,7 +354,7 @@ export async function streamFromSymbolServer( const flushed = res.write(chunk.value); wrote = true; if (!flushed) { - await waitForDrain(res); + await waitForDrain(res, cfg.timeoutMs); } chunk = await reader.read(); } diff --git a/tests/winbindex.test.ts b/tests/winbindex.test.ts index 21f9592..62fbc8a 100644 --- a/tests/winbindex.test.ts +++ b/tests/winbindex.test.ts @@ -339,7 +339,7 @@ describe("resolveEntry", () => { } }); - it("bounds the cache to 200 entries, evicting the oldest first", async () => { + it("bounds the cache to 200 entries, evicting the least-recently-used first", async () => { const index = { a: { fileInfo: { @@ -364,6 +364,36 @@ describe("resolveEntry", () => { await resolveEntry(CONFIG, "f200.dll", KERNEL32_SHA1); expect(fetchMock).toHaveBeenCalledTimes(202); }); + + it("a cache read refreshes recency, sparing a hot entry under churn", async () => { + const index = { + a: { + fileInfo: { + timestamp: 1584069829, + virtualSize: 118784, + sha1: KERNEL32_SHA1, + }, + }, + }; + fetchMock.mockResolvedValue(jsonIndexResponse(index)); + + // Fill the cache exactly: f0 .. f199. + for (let i = 0; i < 200; i++) { + await resolveEntry(CONFIG, `f${i}.dll`, KERNEL32_SHA1); + } + expect(fetchMock).toHaveBeenCalledTimes(200); + + // Re-read f0: it becomes most-recently-used, f1 is now the coldest. + await resolveEntry(CONFIG, "f0.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(200); + + // A new entry evicts the LRU (f1), not the just-read f0. + await resolveEntry(CONFIG, "f200.dll", KERNEL32_SHA1); + await resolveEntry(CONFIG, "f0.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(201); + await resolveEntry(CONFIG, "f1.dll", KERNEL32_SHA1); + expect(fetchMock).toHaveBeenCalledTimes(202); + }); }); describe("tryServeFromWinbindex", () => { @@ -606,4 +636,28 @@ describe("tryServeFromWinbindex", () => { ); expect(res.ended).toBe(true); }); + + it("tears the response down if the client never drains", async () => { + const body = bytes("MZ...a client that stops reading..."); + const hash = sha1Hex(body); + fetchMock.mockImplementation(async (input: unknown) => + String(input).includes(".json.gz") + ? jsonIndexResponse(entryIndex(hash)) + : symbolResponse(body), + ); + + // write() reports backpressure but "drain" is never emitted. + const res = makeMockRes(); + res.write = () => false; + + const outcome = await tryServeFromWinbindex( + { ...CONFIG, timeoutMs: 10 }, + hash, + "kernel32.dll", + res, + ); + + expect(outcome).toBe("failed_after_send"); + expect(res.destroyed).toBe(true); + }); }); From 94ad52c2adc427e524b4803c75f3d78527ab90e2 Mon Sep 17 00:00:00 2001 From: Mathieu Tarral Date: Wed, 9 Sep 2026 13:44:33 +0200 Subject: [PATCH 5/6] fix(winbindex): release unread fetch bodies, restore fetch in tests, doc Second round of Copilot review feedback on #72: - fetchIndex and streamFromSymbolServer now best-effort cancel the response body on a 404 / non-OK / non-2xx return, so undici does not keep the socket and unread data in its pool under repeated misses. - winbindex.test.ts restores the real globalThis.fetch in afterEach instead of leaving the mock installed for later test files. - blob-api.md: clarify that the filename match is case-insensitive (the basename is lowercased before the [a-z0-9._+-] check), so KERNEL32.DLL is accepted. winbindex.test.ts 37/37, full suite 89 passing, build + ccode clean. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01NmATMvFbdupfLRq7Ldwabv --- docs/reference/blob-api.md | 2 +- src/winbindex.ts | 11 +++++++++++ tests/winbindex.test.ts | 2 ++ 3 files changed, 14 insertions(+), 1 deletion(-) diff --git a/docs/reference/blob-api.md b/docs/reference/blob-api.md index 20e8e3e..f40ce71 100644 --- a/docs/reference/blob-api.md +++ b/docs/reference/blob-api.md @@ -20,7 +20,7 @@ Downloads a blob by its SHA-1 hash. | Name | Type | Required | Description | |------|------|----------|-------------| -| `filename` | string | No | The file's basename (e.g. `kernel32.dll`). Used **only** to look the blob up on Winbindex for the Windows-PE fast path (see [Winbindex Blob Source](./winbindex-source.md)). Ignored for non-PE names, and by the MinIO path, which always keys on `hash`. A non-plain name (anything outside `[a-z0-9._+-]`) is ignored. | +| `filename` | string | No | The file's basename (e.g. `kernel32.dll`). Used **only** to look the blob up on Winbindex for the Windows-PE fast path (see [Winbindex Blob Source](./winbindex-source.md)). Ignored for non-PE names, and by the MinIO path, which always keys on `hash`. Matching is case-insensitive (the name is lowercased first); a name that then contains anything outside `[a-z0-9._+-]` (path separators, `?`, `#`, whitespace, ...) skips the fast path. | ## Request Headers diff --git a/src/winbindex.ts b/src/winbindex.ts index c57ebdc..bba8b0c 100644 --- a/src/winbindex.ts +++ b/src/winbindex.ts @@ -5,6 +5,14 @@ import { gunzip } from "node:zlib"; /** Off-loads decompression to the threadpool so the event loop is not blocked. */ const gunzipAsync = promisify(gunzip); +/** + * Best-effort release of a fetch body we are not going to read, so undici does + * not keep the socket and unread data around in its connection pool. + */ +function discardBody(response: Response): void { + response.body?.cancel().catch(() => {}); +} + /** * Winbindex fast path for Windows PE blob downloads. * @@ -178,10 +186,12 @@ async function fetchIndex( } if (response.status === 404) { + discardBody(response); cacheSet(name, null); return null; } if (!response.ok) { + discardBody(response); return null; } @@ -314,6 +324,7 @@ export async function streamFromSymbolServer( } if (!upstream.ok || !upstream.body) { + discardBody(upstream); return "not_available"; } diff --git a/tests/winbindex.test.ts b/tests/winbindex.test.ts index 62fbc8a..169c248 100644 --- a/tests/winbindex.test.ts +++ b/tests/winbindex.test.ts @@ -174,6 +174,7 @@ const receivedText = (res: MockRes): string => { type FetchMock = jest.MockedFunction<(...args: any[]) => Promise>; let fetchMock: FetchMock; +const REAL_FETCH = globalThis.fetch; beforeEach(() => { __clearWinbindexCache(); @@ -185,6 +186,7 @@ beforeEach(() => { afterEach(() => { jest.restoreAllMocks(); + globalThis.fetch = REAL_FETCH; }); describe("isWindowsPEFilename", () => { From 9beea04beeec926767b468e5bffb9e73e5296210 Mon Sep 17 00:00:00 2001 From: Mathieu Tarral Date: Wed, 9 Sep 2026 13:51:18 +0200 Subject: [PATCH 6/6] fix(winbindex): abort upstream on stream teardown, fix filename doc Third round of Copilot review feedback on #72: - streamFromSymbolServer cancels the upstream body reader in a finally block, so a mid-stream error or drain timeout aborts the symbol-server download instead of letting it run to completion in the background holding a socket. - blob-api.md: the server takes basename(filename) first, so a path like windows/system32/KERNEL32.DLL resolves as kernel32.dll and still engages the fast path; separators do not skip it. Corrected. Test asserts the upstream stream is cancelled on a drain timeout. winbindex.test.ts 37/37, full suite 89 passing, build + ccode clean. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01NmATMvFbdupfLRq7Ldwabv --- docs/reference/blob-api.md | 2 +- src/winbindex.ts | 5 +++++ tests/winbindex.test.ts | 22 ++++++++++++++++++++-- 3 files changed, 26 insertions(+), 3 deletions(-) diff --git a/docs/reference/blob-api.md b/docs/reference/blob-api.md index f40ce71..e8dc757 100644 --- a/docs/reference/blob-api.md +++ b/docs/reference/blob-api.md @@ -20,7 +20,7 @@ Downloads a blob by its SHA-1 hash. | Name | Type | Required | Description | |------|------|----------|-------------| -| `filename` | string | No | The file's basename (e.g. `kernel32.dll`). Used **only** to look the blob up on Winbindex for the Windows-PE fast path (see [Winbindex Blob Source](./winbindex-source.md)). Ignored for non-PE names, and by the MinIO path, which always keys on `hash`. Matching is case-insensitive (the name is lowercased first); a name that then contains anything outside `[a-z0-9._+-]` (path separators, `?`, `#`, whitespace, ...) skips the fast path. | +| `filename` | string | No | The file's basename (e.g. `kernel32.dll`). Used **only** to look the blob up on Winbindex for the Windows-PE fast path (see [Winbindex Blob Source](./winbindex-source.md)). Ignored for non-PE names, and by the MinIO path, which always keys on `hash`. The server takes the basename and lowercases it (so `windows/system32/KERNEL32.DLL` resolves as `kernel32.dll`); if that basename then contains anything outside `[a-z0-9._+-]` (`?`, `#`, `%`, whitespace, ...) the fast path is skipped. | ## Request Headers diff --git a/src/winbindex.ts b/src/winbindex.ts index bba8b0c..7c55910 100644 --- a/src/winbindex.ts +++ b/src/winbindex.ts @@ -380,6 +380,11 @@ export async function streamFromSymbolServer( ); res.destroy(); return "failed_after_send"; + } finally { + // Abort the symbol-server download on any exit: a no-op once the body is + // fully read, but on a mid-stream error or drain timeout it stops the + // upstream transfer instead of leaving it running in the background. + void reader.cancel().catch(() => {}); } const digest = hash.digest("hex"); diff --git a/tests/winbindex.test.ts b/tests/winbindex.test.ts index 169c248..7dbb00a 100644 --- a/tests/winbindex.test.ts +++ b/tests/winbindex.test.ts @@ -69,6 +69,22 @@ const streamThenError = (data: Uint8Array): ReadableStream => { }); }; +/** A stream that yields `data` then stalls, recording whether it was cancelled. */ +const stallingStream = ( + data: Uint8Array, +): { stream: ReadableStream; cancelled: () => boolean } => { + let cancelled = false; + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(data); + }, + cancel() { + cancelled = true; + }, + }); + return { stream, cancelled: () => cancelled }; +}; + /** A stream that errors before yielding anything. */ const erroringStream = (): ReadableStream => new ReadableStream({ @@ -639,13 +655,14 @@ describe("tryServeFromWinbindex", () => { expect(res.ended).toBe(true); }); - it("tears the response down if the client never drains", async () => { + it("tears the response and the upstream down if the client never drains", async () => { const body = bytes("MZ...a client that stops reading..."); const hash = sha1Hex(body); + const upstream = stallingStream(body); fetchMock.mockImplementation(async (input: unknown) => String(input).includes(".json.gz") ? jsonIndexResponse(entryIndex(hash)) - : symbolResponse(body), + : symbolResponse(upstream.stream), ); // write() reports backpressure but "drain" is never emitted. @@ -661,5 +678,6 @@ describe("tryServeFromWinbindex", () => { expect(outcome).toBe("failed_after_send"); expect(res.destroyed).toBe(true); + expect(upstream.cancelled()).toBe(true); }); });