diff --git a/webdriver/src/prj/package.json b/webdriver/src/prj/package.json index 52706881..f3b12978 100644 --- a/webdriver/src/prj/package.json +++ b/webdriver/src/prj/package.json @@ -7,7 +7,9 @@ "scripts": { "start": "node --loader ts-node/esm src/index.ts", "dev": "ts-node src/index.ts", - "build": "tsc" + "build": "tsc", + "test": "node --loader ts-node/esm --test src/**/*.test.ts", + "typecheck": "tsc --noEmit" }, "dependencies": { "puppeteer-core": "24.16.0", diff --git a/webdriver/src/prj/src/duration.ts b/webdriver/src/prj/src/duration.ts index ec7ddaad..36128511 100644 --- a/webdriver/src/prj/src/duration.ts +++ b/webdriver/src/prj/src/duration.ts @@ -67,6 +67,21 @@ export function envInt(envName: string, defaultValue: number): number { return result; } +/** + * Like `envInt`, but for limits where a non-positive value is nonsense rather + * than merely unusual — a zero response cap, for instance, would reject every + * page. Fails at startup instead of silently disabling the feature. + */ +export function envPositiveInt(envName: string, defaultValue: number): number { + const value = envInt(envName, defaultValue); + if (!Number.isInteger(value) || value < 1) { + throw new Error( + `env ${envName} must be a positive integer, got "${process.env[envName]}"`, + ); + } + return value; +} + export function envStr(envName: string, defaultValue: string): string { const raw = process.env[envName]; const value = raw !== undefined && raw !== '' ? raw : defaultValue; diff --git a/webdriver/src/prj/src/index.ts b/webdriver/src/prj/src/index.ts index d948c6ac..a6ae7c25 100644 --- a/webdriver/src/prj/src/index.ts +++ b/webdriver/src/prj/src/index.ts @@ -5,7 +5,18 @@ import { Command } from 'commander'; import * as logger from './logging.js'; import * as chromeBrowser from './browser/chrome.js'; import * as ssrf from './ssrf.js'; -import { envDurationMs, envInt, formatDurationMs } from './duration.js'; +import { + envDurationMs, + envInt, + envPositiveInt, + formatDurationMs, +} from './duration.js'; +import { + Semaphore, + SemaphoreAborted, + SemaphoreQueueFull, + SemaphoreTimeout, +} from './semaphore.js'; interface NavigationOptions { waitUntil?: pup.PuppeteerLifeCycleEvent; @@ -17,6 +28,7 @@ interface RenderOptions { waitAfterLoaded?: number; waitUntil?: pup.PuppeteerLifeCycleEvent; maxPageHeapMB?: number; + maxResponseChars?: number; } const program = new Command(); @@ -30,9 +42,48 @@ program const options = program.opts(); const STATUS_I_AM_A_TEAPOT = 418; +const STATUS_INSUFFICIENT_STORAGE = 507; +const STATUS_SERVICE_UNAVAILABLE = 503; const DEFAULT_MAX_PAGE_HEAP_MB = envInt('GVM_WEBDRIVER_MAX_PAGE_HEAP_MB', 1024); +// Peak memory of this service is bounded by the product of the two limits +// below, so both have to exist for either to mean anything: +// +// peak ~= MAX_CONCURRENT_RENDERS * (page heap + extracted response) + browser baseline +// +// Every in-flight render is guaranteed its full per-request allocation. A +// caller that has a slot is never squeezed by how many other callers are +// active; callers beyond the limit queue instead. +const MAX_CONCURRENT_RENDERS = envPositiveInt( + 'GVM_WEBDRIVER_MAX_CONCURRENT_RENDERS', + 4, +); + +// Counted in characters, which is what bounds the string materialized in this +// process. One character encodes to at most 4 UTF-8 bytes on the wire, and +// occupies 2 bytes in V8 for the common (BMP) case. +const MAX_RESPONSE_MB = envPositiveInt('GVM_WEBDRIVER_MAX_RESPONSE_MB', 32); +const MAX_RESPONSE_CHARS = MAX_RESPONSE_MB * 1024 * 1024; + +// Waiting callers are cheap but not free: each holds a live connection. Past +// this depth the service is not going to work through the backlog before +// callers give up anyway, so refusing immediately beats refusing slowly. +const MAX_RENDER_QUEUE = envInt('GVM_WEBDRIVER_MAX_RENDER_QUEUE', 64); + +// Kept well below the caller's own transport timeout: a queue wait that +// outlives it converts a condition we can report cleanly into a transport +// failure, which callers treat far more harshly. +const RENDER_QUEUE_TIMEOUT_MS = envDurationMs( + 'GVM_WEBDRIVER_RENDER_QUEUE_TIMEOUT', + '120s', +); + +const renderSlots = new Semaphore({ + permits: MAX_CONCURRENT_RENDERS, + maxQueued: MAX_RENDER_QUEUE, +}); + const MAX_WAIT_AFTER_LOADED_MS = envDurationMs( 'GVM_WEBDRIVER_MAX_WAIT_AFTER_LOADED', '60s', @@ -119,22 +170,53 @@ async function navigateToPage( } } -async function asText(page: pup.Page) { - const bodyText = await page.evaluate(() => { - return document.body.innerText; - }); +/** + * Read `innerText` or `innerHTML`, rejecting oversized documents *in the page*. + * + * The size test has to happen on the browser side. Pulling the string across + * first and measuring it here would already have paid the memory cost we are + * trying to avoid — this process would hold the whole document no matter what + * any downstream limit says. + */ +async function extractBounded( + page: pup.Page, + kind: 'innerText' | 'innerHTML', + maxChars: number, +): Promise { + const extracted = await page.evaluate( + (k, limit) => { + const raw = + k === 'innerText' ? document.body.innerText : document.body.innerHTML; + return raw.length > limit + ? { tooLarge: true as const, length: raw.length } + : { tooLarge: false as const, content: raw }; + }, + kind, + maxChars, + ); + + if (extracted.tooLarge) { + throw new ResponseLimitExceeded(extracted.length, maxChars); + } + return extracted.content; +} - return normalizeWhitespace(bodyText); +async function asText(page: pup.Page, maxChars: number) { + return normalizeWhitespace(await extractBounded(page, 'innerText', maxChars)); } -async function asHTML(page: pup.Page) { - return await page.evaluate(() => { - return document.body.innerHTML; - }); +async function asHTML(page: pup.Page, maxChars: number) { + return await extractBounded(page, 'innerHTML', maxChars); } -async function asScreenshot(page: pup.Page) { - return await page.screenshot(); +async function asScreenshot(page: pup.Page, maxChars: number) { + // Screenshots are viewport-sized rather than full-page, so this is a + // backstop rather than a limit anyone should reach. + const image = await page.screenshot(); + if (image.length > maxChars) { + throw new ResponseLimitExceeded(image.length, maxChars); + } + return image; } class HeapLimitExceeded extends Error { @@ -143,6 +225,12 @@ class HeapLimitExceeded extends Error { } } +class ResponseLimitExceeded extends Error { + constructor(size: number, maxSize: number) { + super(`Rendered response ${size} exceeds limit ${maxSize}`); + } +} + const HEAP_CHECK_INTERVAL_MS = envInt( 'GVM_WEBDRIVER_HEAP_CHECK_INTERVAL_MS', 200, @@ -212,18 +300,49 @@ async function renderPage( targetUrl: string, mode: 'text' | 'html' | 'screenshot', options: RenderOptions = {}, + signal?: AbortSignal, ): Promise<{ status: number; body: any }> { - const browserManager = await chromeBrowser.INSTANCE; - const browserInstance = browserManager.getBrowser(); + // The slot is taken before the browser is touched, so a queued request costs + // nothing but the open connection. try { - return renderPageWithBrowser( - browserInstance.get(), - targetUrl, - mode, - options, - ); + await renderSlots.acquire(RENDER_QUEUE_TIMEOUT_MS, signal); + } catch (e) { + // Saturation escapes as a transport failure rather than a render result, + // and the distinction is deliberate. How busy this host happens to be is + // not a property of the page, so it must not reach the contract: two + // hosts under different load would otherwise return different answers + // for the same request and each would look legitimate. Failing at the + // transport level keeps the outcome an honest "this host could not run + // it" instead of a fabricated observation about the web. + if (e instanceof SemaphoreTimeout || e instanceof SemaphoreQueueFull) { + logger.log('warn', 'render rejected: saturated', { + url: targetUrl, + reason: e.name, + inFlight: renderSlots.inFlight, + queued: renderSlots.queued, + timeout: formatDurationMs(RENDER_QUEUE_TIMEOUT_MS), + }); + } else if (e instanceof SemaphoreAborted) { + logger.log('debug', 'render abandoned while queued', { url: targetUrl }); + } + throw e; + } + + try { + const browserManager = await chromeBrowser.INSTANCE; + const browserInstance = browserManager.getBrowser(); + try { + return await renderPageWithBrowser( + browserInstance.get(), + targetUrl, + mode, + options, + ); + } finally { + browserInstance.close(); + } } finally { - browserInstance.close(); + renderSlots.release(); } } @@ -238,13 +357,25 @@ async function renderPageWithBrowser( waitAfterLoaded = 0, waitUntil = 'domcontentloaded', maxPageHeapMB = DEFAULT_MAX_PAGE_HEAP_MB, + maxResponseChars = MAX_RESPONSE_CHARS, } = options; // Each render runs in its own browser context so cookies, localStorage, // IndexedDB, service workers and HSTS state never leak between tenants // sharing this long-lived browser. const context = await browserInstance.createBrowserContext(); - const page = await context.newPage(); + + // `newPage` fails when the shared browser is under pressure, and the cleanup + // below only runs once execution is inside the try. Without this the context + // survives for the browser's whole lifetime — and since pressure is exactly + // what makes `newPage` fail, the response to it would be to leak more. + let page: pup.Page; + try { + page = await context.newPage(); + } catch (e) { + await context.close().catch(() => {}); + throw e; + } try { await ssrf.installSsrfGuard(page); @@ -287,13 +418,13 @@ async function renderPageWithBrowser( let data; switch (mode) { case 'text': - data = await asText(page); + data = await asText(page, maxResponseChars); break; case 'html': - data = await asHTML(page); + data = await asHTML(page, maxResponseChars); break; case 'screenshot': - data = await asScreenshot(page); + data = await asScreenshot(page, maxResponseChars); break; default: data = 'Invalid mode'; @@ -302,12 +433,28 @@ async function renderPageWithBrowser( return { status: statusCode, body: data }; }); } catch (e) { + // These escape as transport failures rather than render results. The + // tempting reading is that they describe the page — but the test is not + // "the page is big", it is "the page is bigger than *this host's* + // limit", and that threshold is local configuration. Letting it reach + // the contract would make hosts with different settings answer the same + // request differently, each answer looking legitimate, and would hand + // anyone who knows the spread a way to craft a page that lands in the + // gap. + // + // They could become contract-visible, but only once the values are + // agreed protocol constants rather than per-host settings. if (e instanceof HeapLimitExceeded) { logger.log('warn', 'page heap limit exceeded', { url: targetUrl, error: e.message, }); - return { status: 507, body: e.message }; + } else if (e instanceof ResponseLimitExceeded) { + logger.log('warn', 'rendered response limit exceeded', { + url: targetUrl, + mode, + error: e.message, + }); } throw e; } finally { @@ -317,6 +464,23 @@ async function renderPageWithBrowser( } } +/** + * Fires when the caller goes away before it has been answered, so a request + * nobody is waiting for stops occupying a queue slot. + */ +function disconnectSignal( + req: http.IncomingMessage, + res: http.ServerResponse, +): AbortSignal { + const controller = new AbortController(); + req.on('close', () => { + if (!res.writableEnded) { + controller.abort(); + } + }); + return controller.signal; +} + async function handleRenderRequest( parsedUrl: URL, req: http.IncomingMessage, @@ -355,7 +519,12 @@ async function handleRenderRequest( : {}), }; - const result = await renderPage(targetUrl, mode, options); + const result = await renderPage( + targetUrl, + mode, + options, + disconnectSignal(req, res), + ); if (statusIsGood(result.status)) { updateLastSuccessfulRenderTime(); @@ -370,6 +539,40 @@ async function handleRenderRequest( } res.end(result.body); } catch (error) { + if (error instanceof SemaphoreAborted) { + // The caller hung up while queued; there is nobody left to answer. + return; + } + if ( + error instanceof SemaphoreTimeout || + error instanceof SemaphoreQueueFull + ) { + res.writeHead(STATUS_SERVICE_UNAVAILABLE, { + 'Content-Type': 'application/json', + }); + res.end( + JSON.stringify({ + error: 'Service unavailable', + message: (error as Error).message, + }), + ); + return; + } + if ( + error instanceof HeapLimitExceeded || + error instanceof ResponseLimitExceeded + ) { + res.writeHead(STATUS_INSUFFICIENT_STORAGE, { + 'Content-Type': 'application/json', + }); + res.end( + JSON.stringify({ + error: 'Resource limit exceeded', + message: (error as Error).message, + }), + ); + return; + } res.writeHead(500, { 'Content-Type': 'application/json' }); res.end( JSON.stringify({ @@ -380,7 +583,11 @@ async function handleRenderRequest( } } -async function handleHealthcheck(parsedUrl: URL, res: http.ServerResponse) { +async function handleHealthcheck( + parsedUrl: URL, + req: http.IncomingMessage, + res: http.ServerResponse, +) { const now = Date.now(); const sinceLastSuccessMs = now - lastSuccessfulRenderTime; if (sinceLastSuccessMs < HEALTHCHECK_CACHE_DURATION_MS) { @@ -415,7 +622,12 @@ async function handleHealthcheck(parsedUrl: URL, res: http.ServerResponse) { cacheDuration: formatDurationMs(HEALTHCHECK_CACHE_DURATION_MS), }); try { - const result = await renderPage(targetUrl, mode, {}); + const result = await renderPage( + targetUrl, + mode, + {}, + disconnectSignal(req, res), + ); if (statusIsGood(result.status)) { updateLastSuccessfulRenderTime(); logger.log('info', 'healthcheck ok', { status: result.status }); @@ -427,6 +639,13 @@ async function handleHealthcheck(parsedUrl: URL, res: http.ServerResponse) { res.end('unhealthy'); } } catch (error) { + if (error instanceof SemaphoreAborted) { + return; + } + // Saturation lands here too, and reporting unhealthy is right: the probe + // only renders when nothing has succeeded for the cache duration, so a + // host that is both saturated and has completed nothing in that window + // genuinely is not serving. logger.log('error', 'healthcheck error', { error: (error as Error).message, }); @@ -442,7 +661,7 @@ const server = http.createServer(async (req, res) => { if (pathname === '/render') { await handleRenderRequest(parsedUrl, req, res); } else if (pathname === '/healthcheck') { - await handleHealthcheck(parsedUrl, res); + await handleHealthcheck(parsedUrl, req, res); } else if (pathname === '/log-level') { if (req.method === 'GET') { res.writeHead(200, { 'Content-Type': 'text/plain' }); diff --git a/webdriver/src/prj/src/semaphore.test.ts b/webdriver/src/prj/src/semaphore.test.ts new file mode 100644 index 00000000..afc6d3fa --- /dev/null +++ b/webdriver/src/prj/src/semaphore.test.ts @@ -0,0 +1,270 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; + +import { + Semaphore, + SemaphoreAborted, + SemaphoreQueueFull, + SemaphoreTimeout, +} from './semaphore.js'; + +const NEVER = 60_000; + +/** Queue depth is irrelevant to most cases, so default it out of the way. */ +function sem(permits: number, maxQueued = 1024) { + return new Semaphore({ permits, maxQueued }); +} + +function deferred() { + let resolve!: (v: T) => void; + const promise = new Promise((r) => { + resolve = r; + }); + return { promise, resolve }; +} + +describe('Semaphore', () => { + it('rejects invalid construction', () => { + assert.throws(() => sem(0)); + assert.throws(() => sem(-1)); + assert.throws(() => sem(1.5)); + assert.throws(() => new Semaphore({ permits: 1, maxQueued: -1 })); + assert.throws(() => new Semaphore({ permits: 1, maxQueued: 1.5 })); + }); + + it('grants up to the permit count without waiting', async () => { + const s = sem(3); + await s.acquire(NEVER); + await s.acquire(NEVER); + await s.acquire(NEVER); + assert.equal(s.inFlight, 3); + assert.equal(s.queued, 0); + }); + + it('queues the caller past the limit until a permit is released', async () => { + const s = sem(1); + await s.acquire(NEVER); + + let granted = false; + const waiting = s.acquire(NEVER).then(() => { + granted = true; + }); + + await new Promise((r) => setTimeout(r, 10)); + assert.equal(granted, false, 'should still be queued'); + assert.equal(s.queued, 1); + + s.release(); + await waiting; + assert.equal(granted, true); + assert.equal(s.inFlight, 1, 'permit passed to the waiter, not returned'); + }); + + it('hands permits to waiters in FIFO order', async () => { + const s = sem(1); + await s.acquire(NEVER); + + const order: number[] = []; + const waiters = [1, 2, 3].map((n) => + s.acquire(NEVER).then(() => { + order.push(n); + }), + ); + + await new Promise((r) => setTimeout(r, 10)); + s.release(); + s.release(); + s.release(); + await Promise.all(waiters); + + assert.deepEqual(order, [1, 2, 3]); + }); + + it('times out a caller that waits too long, and stops queueing it', async () => { + const s = sem(1); + await s.acquire(NEVER); + + await assert.rejects(() => s.acquire(20), SemaphoreTimeout); + assert.equal(s.queued, 0, 'timed-out waiter must leave the queue'); + + s.release(); + await s.acquire(NEVER); + assert.equal(s.inFlight, 1); + }); + + it('does not grant a permit to a waiter that already timed out', async () => { + const s = sem(1); + await s.acquire(NEVER); + + const timedOut = assert.rejects(() => s.acquire(20), SemaphoreTimeout); + await new Promise((r) => setTimeout(r, 40)); + await timedOut; + + s.release(); + assert.equal(s.inFlight, 0, 'released permit must not be held by a ghost'); + }); + + it('never exceeds the permit count under concurrent load', async () => { + const permits = 4; + const s = sem(permits); + let active = 0; + let peak = 0; + + await Promise.all( + Array.from({ length: 50 }, () => + s.withPermit(NEVER, async () => { + active += 1; + peak = Math.max(peak, active); + await new Promise((r) => setTimeout(r, 1)); + active -= 1; + }), + ), + ); + + assert.equal(peak, permits, 'concurrency must saturate but not exceed'); + assert.equal(s.inFlight, 0); + assert.equal(s.queued, 0); + }); + + it('releases the permit when the guarded task throws', async () => { + const s = sem(1); + await assert.rejects( + () => + s.withPermit(NEVER, async () => { + throw new Error('boom'); + }), + /boom/, + ); + assert.equal(s.inFlight, 0, 'a failed task must not leak its permit'); + }); + + it('ignores a release with nothing outstanding', async () => { + const s = sem(2); + s.release(); + s.release(); + s.release(); + + // A leaked permit would let a third caller through here. + await s.acquire(NEVER); + await s.acquire(NEVER); + await assert.rejects(() => s.acquire(20), SemaphoreTimeout); + }); + + it('keeps a slow task from starving a queued one indefinitely', async () => { + const s = sem(1); + const slow = deferred(); + + const running = s.withPermit(NEVER, () => slow.promise); + const queued = s.acquire(NEVER); + + slow.resolve(); + await running; + await queued; + assert.equal(s.inFlight, 1); + }); + + describe('queue depth', () => { + it('refuses callers once the queue is full', async () => { + const s = new Semaphore({ permits: 1, maxQueued: 2 }); + await s.acquire(NEVER); + + const queued = [s.acquire(NEVER), s.acquire(NEVER)]; + assert.equal(s.queued, 2); + + await assert.rejects(() => s.acquire(NEVER), SemaphoreQueueFull); + assert.equal(s.queued, 2, 'a refused caller must not be queued'); + + s.release(); + s.release(); + await Promise.all(queued); + }); + + it('accepts again once the queue drains', async () => { + const s = new Semaphore({ permits: 1, maxQueued: 1 }); + await s.acquire(NEVER); + const first = s.acquire(NEVER); + await assert.rejects(() => s.acquire(NEVER), SemaphoreQueueFull); + + s.release(); + await first; + + // Slot freed: a new caller queues rather than being refused. + const second = s.acquire(NEVER); + assert.equal(s.queued, 1); + s.release(); + await second; + }); + + it('refuses immediately rather than waiting out the timeout', async () => { + const s = new Semaphore({ permits: 1, maxQueued: 0 }); + await s.acquire(NEVER); + + const started = Date.now(); + await assert.rejects(() => s.acquire(NEVER), SemaphoreQueueFull); + assert.ok( + Date.now() - started < 50, + 'a full queue must fail fast, not after the wait', + ); + }); + }); + + describe('abort', () => { + it('drops a waiter whose caller went away', async () => { + const s = sem(1); + await s.acquire(NEVER); + + const controller = new AbortController(); + const waiting = assert.rejects( + () => s.acquire(NEVER, controller.signal), + SemaphoreAborted, + ); + assert.equal(s.queued, 1); + + controller.abort(); + await waiting; + assert.equal(s.queued, 0, 'aborted waiter must leave the queue'); + + // The permit must go to a live caller, not the abandoned one. + s.release(); + assert.equal(s.inFlight, 0); + }); + + it('rejects at once when the signal is already aborted', async () => { + const s = sem(1); + await s.acquire(NEVER); + await assert.rejects( + () => s.acquire(NEVER, AbortSignal.abort()), + SemaphoreAborted, + ); + assert.equal(s.queued, 0); + }); + + it('does not reject a caller that already holds a permit', async () => { + const s = sem(1); + const controller = new AbortController(); + await s.acquire(NEVER, controller.signal); + controller.abort(); + // Aborting after the permit was granted is a no-op; release still works. + s.release(); + assert.equal(s.inFlight, 0); + }); + + it('frees the queue slot so a later caller is not refused', async () => { + const s = new Semaphore({ permits: 1, maxQueued: 1 }); + await s.acquire(NEVER); + + const controller = new AbortController(); + const abandoned = assert.rejects( + () => s.acquire(NEVER, controller.signal), + SemaphoreAborted, + ); + controller.abort(); + await abandoned; + + const queued = s.acquire(NEVER); + assert.equal(s.queued, 1, 'slot must have been reclaimed'); + s.release(); + await queued; + }); + }); +}); diff --git a/webdriver/src/prj/src/semaphore.ts b/webdriver/src/prj/src/semaphore.ts new file mode 100644 index 00000000..84197977 --- /dev/null +++ b/webdriver/src/prj/src/semaphore.ts @@ -0,0 +1,156 @@ +export interface SemaphoreOptions { + /** Concurrent holders. */ + permits: number; + /** Callers allowed to wait for a permit before further ones are refused. */ + maxQueued: number; +} + +/** + * Counting semaphore with FIFO waiters, a bounded queue and a bounded wait. + * + * Callers past the permit count queue rather than being refused outright: + * refusing on any contention would make the outcome depend on whatever else the + * host happened to be running at the time. + * + * The queue is still finite, and the wait is still bounded. Both exist because + * an unbounded queue only defers the problem — a caller parked indefinitely + * eventually trips its own transport timeout, which fails far more harshly than + * a refusal we issue ourselves, and until then it costs a live connection. + */ +export class Semaphore { + private readonly permits: number; + private readonly maxQueued: number; + private available: number; + private readonly waiters: Array<{ + resolve: () => void; + reject: (e: Error) => void; + timer: NodeJS.Timeout; + cleanup: () => void; + }> = []; + + constructor(options: SemaphoreOptions) { + const { permits, maxQueued } = options; + if (!Number.isInteger(permits) || permits < 1) { + throw new Error( + `semaphore permits must be a positive integer, got ${permits}`, + ); + } + if (!Number.isInteger(maxQueued) || maxQueued < 0) { + throw new Error( + `semaphore maxQueued must be a non-negative integer, got ${maxQueued}`, + ); + } + this.permits = permits; + this.maxQueued = maxQueued; + this.available = permits; + } + + get inFlight(): number { + return this.permits - this.available; + } + + get queued(): number { + return this.waiters.length; + } + + /** + * Acquire a permit. + * + * @param timeoutMs how long to wait before giving up. + * @param signal aborts the wait — pass the caller's disconnect signal so a + * vanished caller stops occupying a queue slot it will never use. + * @throws {SemaphoreQueueFull} if the queue is already at capacity. + * @throws {SemaphoreTimeout} if no permit became available in time. + * @throws {SemaphoreAborted} if `signal` fired while waiting. + */ + acquire(timeoutMs: number, signal?: AbortSignal): Promise { + if (this.available > 0) { + this.available -= 1; + return Promise.resolve(); + } + if (this.waiters.length >= this.maxQueued) { + return Promise.reject(new SemaphoreQueueFull(this.maxQueued)); + } + if (signal?.aborted) { + return Promise.reject(new SemaphoreAborted()); + } + + return new Promise((resolve, reject) => { + const remove = () => { + const idx = this.waiters.indexOf(entry); + if (idx !== -1) { + this.waiters.splice(idx, 1); + } + }; + const onAbort = () => { + remove(); + clearTimeout(entry.timer); + reject(new SemaphoreAborted()); + }; + const entry = { + resolve, + reject, + timer: setTimeout(() => { + remove(); + entry.cleanup(); + reject(new SemaphoreTimeout(timeoutMs)); + }, timeoutMs), + cleanup: () => signal?.removeEventListener('abort', onAbort), + }; + signal?.addEventListener('abort', onAbort, { once: true }); + this.waiters.push(entry); + }); + } + + release(): void { + const next = this.waiters.shift(); + if (next === undefined) { + // Guard against a double release leaking permits, which would silently + // raise the concurrency ceiling above the configured one. + if (this.available < this.permits) { + this.available += 1; + } + return; + } + clearTimeout(next.timer); + next.cleanup(); + // The permit passes straight to the waiter; `available` deliberately stays + // where it is. + next.resolve(); + } + + /** Run `fn` holding a permit, releasing it however `fn` settles. */ + async withPermit( + timeoutMs: number, + fn: () => Promise, + signal?: AbortSignal, + ): Promise { + await this.acquire(timeoutMs, signal); + try { + return await fn(); + } finally { + this.release(); + } + } +} + +export class SemaphoreTimeout extends Error { + constructor(timeoutMs: number) { + super(`timed out after ${timeoutMs}ms waiting for a render slot`); + this.name = 'SemaphoreTimeout'; + } +} + +export class SemaphoreQueueFull extends Error { + constructor(maxQueued: number) { + super(`render queue is full (${maxQueued} waiting)`); + this.name = 'SemaphoreQueueFull'; + } +} + +export class SemaphoreAborted extends Error { + constructor() { + super('caller went away while waiting for a render slot'); + this.name = 'SemaphoreAborted'; + } +}