Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

62 changes: 45 additions & 17 deletions packages/browser/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,23 +51,51 @@ cache, and starts the worker loop; it resolves once the cache index is built. It

Only `harper` is required; everything else has a default.

| Option | Default | Purpose |
| ---------------------------- | ------------------------------------------------- | ------------------------------------------------------------------------------------------- |
| `harper` | _(required)_ | `{ mqttOrigin, user, pass, workerId }` — connection + identity |
| `queuePort` | `9926` | Port of the plugin's render-queue HTTP API |
| `bypass` | `{ header: x-harper-renderer-bypass, token: '' }` | Shared origin-bypass header/token (match the plugin) |
| `config` | built-in defaults | Rendering config (deep-partial object _or_ JSON file path) |
| `concurrency` | ~half the CPUs | Max concurrent page renders |
| `rps` | `8` | Max job starts per second (a job renders every device of one URL — see the queue protocol) |
| `jobClaimLimit` | `concurrency * 2` | Jobs claimed per batch |
| `browserExpirationThreshold` | `200` | Pages a browser renders before being retired |
| `incognitoPages` | `true` | Render each page in a fresh incognito context |
| `contentEncoding` | `gzip` | Encoding used when posting rendered HTML back |
| `chromeArgs` | hardened headless set | Chrome launch flags |
| `browserLaunchOptions` | built from `chromeArgs` | Full Puppeteer launch options (overrides `chromeArgs`) |
| `resourceCache` | enabled, ~8 GB in tmp | On-disk shared sub-resource cache (`enabled`/`dir`/limits) |
| `renderer` | the default renderer | Custom renderer (see below) |
| `installSignalHandlers` | `true` | Own SIGTERM/SIGINT (drain in-flight renders, then close Chrome); `false` to own the process |
| Option | Default | Purpose |
| ---------------------------- | ------------------------------------------------- | -------------------------------------------------------------------------------------------- |
| `harper` | _(required)_ | `{ mqttOrigin, user, pass, workerId }` — connection + identity |
| `queuePort` | `9926` | Port of the plugin's render-queue HTTP API |
| `bypass` | `{ header: x-harper-renderer-bypass, token: '' }` | Shared origin-bypass header/token (match the plugin) |
| `config` | built-in defaults | Rendering config (deep-partial object _or_ JSON file path) |
| `concurrency` | ~half the CPUs | Max concurrent page renders |
| `rps` | `8` | Max job starts per second (a job renders every device of one URL — see the queue protocol) |
| `jobClaimLimit` | `concurrency` (`admission.max` under `pressure`) | Jobs claimed per batch (under `pressure`, the ceiling; each claim is sized to free capacity) |
| `admission` | `{ mode: 'fixed' }` | How many renders run at once — `fixed` or CPU-`pressure`-stepped (see below) |
| `browserExpirationThreshold` | `200` | Pages a browser renders before being retired |
| `incognitoPages` | `true` | Render each page in a fresh incognito context |
| `contentEncoding` | `gzip` | Encoding used when posting rendered HTML back |
| `chromeArgs` | hardened headless set | Chrome launch flags |
| `browserLaunchOptions` | built from `chromeArgs` | Full Puppeteer launch options (overrides `chromeArgs`) |
| `resourceCache` | enabled, ~8 GB in tmp | On-disk shared sub-resource cache (`enabled`/`dir`/limits) |
| `renderer` | the default renderer | Custom renderer (see below) |
| `installSignalHandlers` | `true` | Own SIGTERM/SIGINT (drain in-flight renders, then close Chrome); `false` to own the process |

### `admission` — renders at once from CPU pressure

`fixed` (the default) always runs `concurrency` renders. `pressure` starts at `concurrency` (clamped to
[`min`, `max`]) and, every `intervalMs`, steps the limit from the container's CPU pressure — cgroup v2
PSI `cpu.pressure`, the share of that interval in which some task waited for a CPU: up one when
pressure is below `lowPressure` and the limit held a job back, down one above `highPressure`, down a
quarter above twice `highPressure`. Render cost varies ~20x by page, so a fixed slot count overloads a
pod that draws heavy pages; pressure is the waiting that overload causes.

Pressure mode also sizes each queue claim to what can start soon — the free slots, or the free
prefetch-pool room — up to `jobClaimLimit`, so a busy worker does not hold jobs an idle one could
start.

```js
admission: { mode: 'pressure', min: 2, max: 5, lowPressure: 15, highPressure: 30, intervalMs: 5000 }
```

Defaults: `min` = `concurrency / 2`, `max` = `concurrency`, pressure 15 / 30, interval 5000 ms. With
the default `max` the limit only brakes. A `max` above `concurrency` lets it climb while CPU is idle,
but idle CPU with work waiting is also what a slow origin looks like, so a higher `max` means more
concurrent origin requests exactly then, and more Chrome pages in memory: set it against what the
origin and the pod's memory tolerate. Measured under a CFS quota, pressure 15 and 30 fell at about 80%
and 95% of the pod's cores; throughput peaked at 90–95%. Without PSI (cgroup v1, macOS) the limit
holds and one warning is logged at startup. Every worker in a container steps on the same signal,
from a random phase. The stats line reports the current limit as `saturation.concurrency` and the
last reading as `saturation.admission.pressure`.

## Queue protocol

Expand Down
2 changes: 1 addition & 1 deletion packages/browser/package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@harperfast/prerender-browser",
"version": "1.38.0",
"version": "1.39.0",
"type": "module",
"description": "Headless-browser render library for Harper Prerender: claims render jobs from the @harperfast/prerender queue, renders pages in headless Chrome (Puppeteer), and posts the HTML back. Embedded by a render service and configured entirely via startWorker() options.",
"keywords": [
Expand Down
8 changes: 6 additions & 2 deletions packages/browser/src/RenderQueueConsumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,11 @@ const toEpochMs = (value: unknown): number | undefined => {
return undefined;
};

export async function* RenderQueueConsumer(signal?: AbortSignal) {
/** @param claimLimit jobs to ask for in the next claim, read before each one (default `settings.jobClaimLimit`). */
export async function* RenderQueueConsumer(
signal?: AbortSignal,
claimLimit: () => number = () => settings.jobClaimLimit
) {
const mqttClient = await connectMqtt();
const health = getHostHealth();

Expand Down Expand Up @@ -96,7 +100,7 @@ export async function* RenderQueueConsumer(signal?: AbortSignal) {
continue;
}

const { jobs, outcome, retryAfterMs } = await claimJobs(host, settings.jobClaimLimit);
const { jobs, outcome, retryAfterMs } = await claimJobs(host, claimLimit());
switch (outcome) {
case 'jobs':
health.recordJobs(host);
Expand Down
134 changes: 125 additions & 9 deletions packages/browser/src/Worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,8 @@ import { setTimeout } from 'timers/promises';
import { noop } from './util/noop.js';
import { getResourceCache } from './ResourceCache.js';
import { settings } from './settings.js';
import { CpuSampler } from './util/cpu.js';
import { CpuSampler, pressureBetween, readCpuStallUs, type StallSample } from './util/cpu.js';
import { nextAdmissionLimit, type AdmissionSettings } from './admission.js';
import { renderPhaseOf } from './util/renderPhase.js';
import { JobDocumentCache } from './documentReuse.js';
import { closePrefetchAgent, prefetchDocument, type PrefetchOutcome } from './documentPrefetch.js';
Expand Down Expand Up @@ -76,6 +77,18 @@ export default class RenderWorker {

inflight: Set<Promise<void>> = new Set();

// Renders allowed at once right now: CONCURRENCY under fixed admission, stepped between the
// configured bounds under pressure admission (admission.ts). `admission` is null under fixed.
private readonly admission: AdmissionSettings | null;
private admissionLimit: number;
private maxActivePages: number;
private admissionTimer: NodeJS.Timeout | null = null;
private lastStall: StallSample | null = null;
private lastPressure: number | null = null;
private demandSinceStep = false;
private waitingForSlot = false;
private limitWaiters: Set<() => void> = new Set();

lastRenderStartTime = Date.now();

// Per-interval counters, snapshotted-and-reset by logStats() so each log line is a delta
Expand Down Expand Up @@ -237,6 +250,25 @@ export default class RenderWorker {
this.BROWSER_MAX_TOTAL_PAGES = config.browserExpirationThreshold ?? 5000;
this.renderFn = config.renderer;

const admission = settings.admission;
if (admission.mode === 'pressure') {
this.admission = { ...admission };
this.admissionLimit = Math.min(admission.max, Math.max(admission.min, this.CONCURRENCY));
this.maxActivePages = admission.max;
this.samplePressure();
if (this.lastStall === null) {
logger.warn(
{ limit: this.admissionLimit },
'pressure admission: no CPU pressure reading (needs cgroup v2 PSI) — the render limit holds'
);
}
this.startAdmission(this.admission);
Comment on lines +259 to +265

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

When CPU pressure monitoring (PSI) is unavailable (e.g., on macOS or cgroup v1 hosts), this.lastStall remains null. In this case, the admission limit will never change. To avoid the unnecessary overhead of starting the admission timer and repeatedly attempting to read the non-existent pressure file every few seconds, we can conditionally start the admission timer only when this.lastStall is not null.

Suggested change
if (this.lastStall === null) {
logger.warn(
{ limit: this.admissionLimit },
'pressure admission: no CPU pressure reading (needs cgroup v2 PSI) — the render limit holds'
);
}
this.startAdmission(this.admission);
if (this.lastStall === null) {
logger.warn(
{ limit: this.admissionLimit },
'pressure admission: no CPU pressure reading (needs cgroup v2 PSI) — the render limit holds'
);
} else {
this.startAdmission(this.admission);
}

} else {
this.admission = null;
this.admissionLimit = this.CONCURRENCY;
this.maxActivePages = this.CONCURRENCY;
}

this.browserCleanupInterval = setInterval(() => {
this.closeRetiredBrowsers();
}, 10000);
Expand Down Expand Up @@ -283,7 +315,7 @@ export default class RenderWorker {
* dropped: no variant of it was attempted, so there is nothing to post, and its lease expires and
* the queue re-grants it — the same fate a job sitting unclaimed in the consumer's batch always had.
*/
async run(jobs: AsyncIterable<RenderJob> = RenderQueueConsumer(this.consumerAbort.signal)) {
async run(jobs: AsyncIterable<RenderJob> = RenderQueueConsumer(this.consumerAbort.signal, () => this.claimSize())) {
const prefetch = settings.config.documentReuse.prefetch;
if (!prefetch.enabled) {
for await (const job of jobs) {
Expand Down Expand Up @@ -355,15 +387,92 @@ export default class RenderWorker {
await filling;
}

/** Hold until a render slot is free. */
/** Hold until a render slot is free: a render finishing, or the admission limit rising. */
private async awaitSlot() {
// wait for slot to open up
if (this.inflight.size >= this.CONCURRENCY) {
this.stats.concurrencyBlocked++;
await Promise.race(this.inflight);
if (this.inflight.size < this.admissionLimit) return;
this.stats.concurrencyBlocked++;
this.waitingForSlot = true;
try {
while (this.inflight.size >= this.admissionLimit) {
if (this.jobHeldBack()) this.demandSinceStep = true;
let wake = noop;
const raised = new Promise<void>((resolve) => (wake = resolve));
this.limitWaiters.add(wake);
try {
await Promise.race([...this.inflight, raised]);
} finally {
this.limitWaiters.delete(wake);
}
}
// A job that reached the pool during the wait waited too.
if (this.jobHeldBack()) this.demandSinceStep = true;
} finally {
this.waitingForSlot = false;
}
}

/**
* Whether the limit is holding a job back right now. The prefetch loop waits for a slot BEFORE
* taking a job, so there a full set of slots holds a job back only while the pool has one — and a
* pooled job with a slot free is prefetching ahead, not held back.
*/
private jobHeldBack(): boolean {
return this.waitingForSlot && (this.pool === null || this.pool.size > 0);
}

/**
* Jobs to ask for in the next claim. Under pressure admission, what can start soon — free slots, or
* free pool room when prefetching — so a busy worker does not hold jobs an idle one could start,
* and a worker with room fills it in one claim. Capped by `jobClaimLimit`, at least 1.
*/
private claimSize(): number {
if (this.admission === null) return settings.jobClaimLimit;
const free = this.pool ? this.pool.capacity - this.pool.size : this.admissionLimit - this.inflight.size;
return Math.min(settings.jobClaimLimit, Math.max(1, free));
}

/**
* Step the admission limit every `intervalMs`, from a random phase so the workers sharing a
* container step at different moments and each reads pressure that already reflects the others'
* last change.
*/
private startAdmission(admission: AdmissionSettings) {
const step = () => this.stepAdmission(admission, this.samplePressure());
this.admissionTimer = globalThis.setTimeout(
() => {
// Re-baseline rather than step: the window since construction is an arbitrary slice of an
// interval, and each step should read one whole interval.
this.samplePressure();
this.admissionTimer = setInterval(step, admission.intervalMs);
this.admissionTimer.unref();
},
Math.floor(Math.random() * admission.intervalMs)
);
Comment on lines +441 to +450

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

When starting the admission timer, if destroy() or shutdown() is called while the initial setTimeout is pending, there is a risk of starting an un-cleared setInterval if the callback executes. Although destroy() attempts to clear the timer, adding a defensive check for this.shuttingDown (or this.admissionTimer === null) inside the setTimeout callback ensures that we do not schedule a new setInterval during or after shutdown, preventing potential resource leaks.

Suggested change
this.admissionTimer = globalThis.setTimeout(
() => {
// Re-baseline rather than step: the window since construction is an arbitrary slice of an
// interval, and each step should read one whole interval.
this.samplePressure();
this.admissionTimer = setInterval(step, admission.intervalMs);
this.admissionTimer.unref();
},
Math.floor(Math.random() * admission.intervalMs)
);
this.admissionTimer = globalThis.setTimeout(
() => {
if (this.shuttingDown) return;
// Re-baseline rather than step: the window since construction is an arbitrary slice of an
// interval, and each step should read one whole interval.
this.samplePressure();
this.admissionTimer = setInterval(step, admission.intervalMs);
this.admissionTimer.unref();
},
Math.floor(Math.random() * admission.intervalMs)
);

this.admissionTimer.unref();
}

/** CPU pressure since the previous sample (util/cpu.ts `pressureBetween`), or null without one. */
private samplePressure(): number | null {
const stallUs = readCpuStallUs();
if (stallUs === null) return null;
const prev = this.lastStall;
this.lastStall = { stallUs, atMs: Date.now() };
return prev ? pressureBetween(prev, this.lastStall) : null;
}

private stepAdmission(admission: AdmissionSettings, pressure: number | null) {
this.lastPressure = pressure === null ? null : Number(pressure.toFixed(1));
const hasDemand = this.demandSinceStep || this.jobHeldBack();
this.demandSinceStep = false;
this.applyAdmissionLimit(nextAdmissionLimit(this.admissionLimit, pressure, hasDemand, admission));
}

private applyAdmissionLimit(next: number) {
const raised = next > this.admissionLimit;
this.admissionLimit = next;
if (raised) for (const wake of this.limitWaiters) wake();
}

/** Whether a claimed job still has enough lease to be worth rendering; counted and logged when not. */
private admit(job: RenderJob): boolean {
// Do not run expired jobs to prevent double rendering
Expand Down Expand Up @@ -622,7 +731,10 @@ export default class RenderWorker {
},
saturation: {
inflight: this.inflight.size,
concurrency: this.CONCURRENCY,
concurrency: this.admissionLimit,
admission: this.admission
? { min: this.admission.min, max: this.admission.max, pressure: this.lastPressure }
: undefined,
concurrencyBlocked: s.concurrencyBlocked,
rpsDelayed: s.rpsDelayed,
expiredSkipped: s.expiredSkipped,
Expand Down Expand Up @@ -695,6 +807,10 @@ export default class RenderWorker {
// the browser .close() promises run and Chrome is orphaned (the whole point of closing here).
async destroy() {
clearInterval(this.logStatsInterval);
if (this.admissionTimer !== null) {
clearInterval(this.admissionTimer);
this.admissionTimer = null;
}
if (this.browserCleanupInterval !== null) {
clearInterval(this.browserCleanupInterval);
this.browserCleanupInterval = null;
Expand Down Expand Up @@ -1063,7 +1179,7 @@ export default class RenderWorker {
}
logger.info({ event: 'launching browser', retired: this.retiredBrowsers.size });
this.browserPromise = ManagedBrowser.launch({
maxActivePages: this.CONCURRENCY,
maxActivePages: this.maxActivePages,
puppeteerLaunchOptions: this.browserLaunchOptions,
}).finally(() => (this.browserPromise = null));
this.browser = await this.browserPromise;
Expand Down
Loading