Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 2c13c9da898c003820e64b604594a206fd4ad286
parent a5b31f1df255a4e5cb122d3d6c61b35838066720
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  9 Oct 2026 11:50:15 -0400

scripts: the heavy slot — one heavy job machine-wide, above a 6 GB memory floor (`pnpm heavy`)

Two OOMs took the desktop session down: two `next build`s, or a build beside
an e2e suite, or a render beside either. queue-lock.mjs grows a second
machine-global lock, the heavy slot (`<git-common-dir>/heavy-queue.lock`),
and a memory floor: once the slot is held, wait until /proc/meminfo's
MemAvailable is at least HEAVY_MIN_FREE_MB (6000). Slot first, then the
floor, so nobody slips in while the holder waits for memory.

- `pnpm heavy -- <cmd>` (`queue-lock.mjs --heavy`) runs any command in it;
  a render is `pnpm heavy -- node umtool/report-to-video/build-video.mjs …`.
- Every e2e entry point (withQueue: the CLI wrapper and run-sharded-e2e)
  takes the heavy slot FIRST, then the e2e queue — one order, so nothing
  deadlocks; HEAVY_HELD lets a heavy command inside one pass through, as
  QUEUE_LOCK_HELD always has. E2E_QUEUE=0 skips both locks; the floor
  still applies. The e2e queue's semantics and bypasses are otherwise
  unchanged.
- The publish stages' real `next build` (a site's, the hub's, the
  homepage's) runs through it (build.ts heavyGated); the e2e fake
  EXPORT_NEXT_BIN does not. The gate forwards SIGTERM, so a stage's Cancel
  still stops the build. In a docker-runner container the slot is the
  container's own and the floor reads the host's memory.
- A waiter is told whom it waits behind and what they run; a MemTotal under
  the floor runs with a note instead of waiting forever; no usable flock
  warns and runs on the floor alone. HEAVY=0, HEAVY_MIN_FREE_MB=0 and
  HEAVY_TIMEOUT are the bypasses (ENVIRONMENT.md regenerated).

The queue's tests: 11 → 23 (two heavy contenders, an e2e run behind a heavy
job, nesting, the floor with a fake meminfo, the stale (SIGKILLed) holder,
SIGTERM forwarding, waitForMemory with an injected reader and clock). The
FIFO and banner tests now wait for each contender to queue instead of a
fixed stagger, which raced a loaded machine.

WORKTREES.md documents the slot and the render recipe that holds the
transcription lane (`pnpm ops lane`) for a render's duration; AGENTS.md
says it in one paragraph.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
MAGENTS.md | 10+++++++++-
MENVIRONMENT.md | 7+++++++
MWORKTREES.md | 61++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/lib/envVars.test.ts | 11++++++++++-
Mcommon/lib/envVars.ts | 7+++++++
Mcommon/publish/build.test.ts | 28++++++++++++++++++++++++++++
Mcommon/publish/build.ts | 35++++++++++++++++++++++++++++++++---
Meditor/CHANGELOG.md | 1+
Mpackage.json | 1+
Mscripts/queue-lock.mjs | 349++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mscripts/queue-lock.test.mjs | 296++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
11 files changed, 731 insertions(+), 75 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -39,7 +39,15 @@ specs that judge a clip read it through `e2e/capabilities.ts` and skip themselve capability is a directory that exists, decided once by the builder — not re-guessed per spec. -See [WORKTREES.md](WORKTREES.md) for the port scheme, the queue, and the shared-data caveat. +**Heavy work takes the heavy slot.** e2e, the publish stages' `next build` and video renders +share ONE machine-global slot and start only above a 6000 MB MemAvailable floor (two OOMs +took the desktop session down). e2e and the publish builds take it on their own; a render or +any other heavy command runs as `pnpm heavy -- <cmd>`. `queue-lock: waiting for the heavy slot +— held by …` or `heavy: waiting for memory …` is the gate working, not a hang. Bypasses: +`HEAVY=0`, `HEAVY_MIN_FREE_MB=<MB>`, `HEAVY_TIMEOUT=<seconds>`. + +See [WORKTREES.md](WORKTREES.md) for the port scheme, the queue, the heavy slot, and the +shared-data caveat. # Working this repo with no local corpus diff --git a/ENVIRONMENT.md b/ENVIRONMENT.md @@ -111,6 +111,9 @@ Tokens, credentials and knobs a running process reads. Most configuration is not | `PARAKEET_DECODER` | parakeet-cli's | `ctc` or `tdt`, passed through to parakeet-cli. | scripts/parakeet-stitch.mjs | | `PARAKEET_LANG` | parakeet-cli's | A locale, passed through to parakeet-cli. | scripts/parakeet-stitch.mjs | | `PARAKEET_DEVICE` | parakeet-cli's | Compute device (`cpu`, `CUDA0`, `Vulkan1`, …), exported to parakeet-cli. | scripts/parakeet-stitch.mjs | +| `HEAVY` | on | `0` skips the heavy slot AND the memory floor: the machine-wide one-at-a-time gate that `pnpm heavy -- <cmd>`, every e2e entry point and the publish stages' `next build` go through. | scripts/queue-lock.mjs | +| `HEAVY_MIN_FREE_MB` | `6000` | The memory floor: a heavy job, once it holds the slot, waits until /proc/meminfo's MemAvailable is at least this many MB. `0` turns the floor off; a machine whose MemTotal is under it runs without waiting. | scripts/queue-lock.mjs | +| `HEAVY_TIMEOUT` | wait forever | Seconds a `pnpm heavy` run waits for the slot, and then for the floor, before giving up (exit 3). An e2e run uses `E2E_QUEUE_TIMEOUT` for both. | scripts/queue-lock.mjs | ## Ports @@ -187,6 +190,10 @@ Read only by a test harness, a fake binary or a test-mode branch. Never set one | `E2E_PORT_GRACE_MS` | `3000` | How long the port check waits for a just-freed port. | scripts/queue-lock.mjs | | `E2E_QUEUE_LOCK_FILE` | one per machine | The queue's lock file; the queue's own tests point it elsewhere. | scripts/queue-lock.mjs | | `QUEUE_LOCK_HELD` | — | Set by the queue for the command it runs, so a nested wrapper passes through. | scripts/queue-lock.mjs | +| `HEAVY_HELD` | — | Set by the heavy slot for the command it runs, so a heavy command inside it (a `pnpm heavy -- pnpm e2e`, a build stage under an e2e suite's editor) passes through. | scripts/queue-lock.mjs | +| `HEAVY_LOCK_FILE` | one per machine | The heavy slot's lock file; the gate's own tests point it elsewhere. | scripts/queue-lock.mjs | +| `HEAVY_MEMINFO_FILE` | `/proc/meminfo` | Where the memory floor reads MemAvailable; the gate's tests hand it a fake. | scripts/queue-lock.mjs | +| `HEAVY_POLL_MS` | `5000` | How often a run waiting for the memory floor re-reads it. | scripts/queue-lock.mjs | | `PLAYWRIGHT_BASE_URL` | `http://localhost:<PORT>` | The editor test server's URL; the worktree injector sets it. | editor/playwright.config.ts, editor/e2e/baseUrl.ts | | `E2E_AUDIO_CHECK_INTERVAL_MS` | the real cadence | Shrinks the mid-download audio check so the e2e suite sees it fire. | common/ytdlp/audioCheckedDownload.ts | | `E2E_AUDIO_CHECK_SIZE_GATE` | the real gate | Likewise, the size gate. | common/ytdlp/audioCheckedDownload.ts | diff --git a/WORKTREES.md b/WORKTREES.md @@ -141,12 +141,71 @@ audio checks run at production pace, which reads as real failures. | Variable | Effect | |---|---| -| `E2E_QUEUE=0` | Skip the queue entirely (the port preflight still runs) | +| `E2E_QUEUE=0` | Skip the queue and the heavy slot (the port preflight and the memory floor still run) | | `E2E_PORT_CHECK=0` | Skip the port preflight, reusing whatever servers are up (a hand-started editor test server needs the config's `E2E_SERVER_ENV`, above) | | `E2E_QUEUE_TIMEOUT=<seconds>` | Give up waiting after N seconds (default: wait forever) | Verify the queue with `pnpm test:scripts`. +## The heavy slot (`pnpm heavy`) + +An e2e suite, a `next build` and a video render each want several GB, and two of them at +once is how this machine OOMed (taking the desktop session with it). So all three go +through ONE more machine-global lock, the **heavy slot**, and start only once +`/proc/meminfo`'s MemAvailable is at least a floor (6000 MB): + +```sh +pnpm heavy -- <cmd…> # any heavy command, by hand +pnpm heavy -- node umtool/report-to-video/build-video.mjs <manifest> … # a render +``` + +Who takes it, and in what order: + +- **Every e2e entry point** (the same ones the e2e queue covers): the heavy slot FIRST, + then the e2e queue. One order everywhere, so nothing can deadlock — and a heavy command + that starts another (`pnpm heavy -- pnpm e2e`, a build stage run by an e2e suite's + editor) passes straight through (`HEAVY_HELD`), as a nested e2e run always has. +- **The publish stages' `next build`** — a site's, the hub's, the homepage's + (`common/publish/build.ts`, `heavyGated`). The wait shows in the stage's log, and a + Cancel still stops the build (the gate forwards SIGTERM). In a docker-runner container + the slot is the container's own; the floor still reads the host's memory, which + throttles a fan-out when the host runs low. +- **A render**, by hand, as above. + +The slot is taken before the floor is waited for, so nobody slips in while the holder +waits for memory. A waiter is told what it waits behind: + +``` +queue-lock: waiting for the heavy slot — held by feature-x (feature-x, pid 31337) for 2m10s: playwright test +heavy: waiting for memory — 4210 MB available, the floor is 6000 MB +``` + +The lock is `<git-common-dir>/heavy-queue.lock`, released by the kernel like the e2e one. +A machine whose MemTotal is under the floor is told so and runs; with no usable `flock` the +gate warns and runs on the floor alone — it is a safety net, not a correctness lock. + +| Variable | Effect | +|---|---| +| `HEAVY=0` | Skip the heavy slot AND the memory floor | +| `HEAVY_MIN_FREE_MB=<MB>` | Move the floor (`0` turns it off) | +| `HEAVY_TIMEOUT=<seconds>` | Give up waiting (slot, then floor) after N seconds; an e2e run uses `E2E_QUEUE_TIMEOUT` | + +### A render and the transcription lane + +A render competes with the transcription lane for memory and CPU, and the lane is not a heavy +slot holder. +Hold the lane for the render's duration, and release it whatever the render's outcome: + +```sh +pnpm ops lane --json '{"lane":"transcription","held":true}' +pnpm heavy -- node umtool/report-to-video/build-video.mjs <manifest> … ; \ + pnpm ops lane --json '{"lane":"transcription","held":false}' +``` + +A hold stops new dispatches, not a transcription already running; and the release +resumes the lane even if someone else held it for another reason — check `/operations` +first. + ## Data directories By default each worktree is **fully isolated**: `common/lib/paths.ts` resolves diff --git a/common/lib/envVars.test.ts b/common/lib/envVars.test.ts @@ -153,7 +153,16 @@ test("the docker audience is the ARCHILYZER_ set", () => { // (one-core Phase 4 slice 3). The two exceptions keep names others depend on: // Playwright's own convention, and the machine-global queue's nesting marker // (its protocol is shared with checkouts on older code). -const UNPREFIXED_TEST_VARS = new Set(["PLAYWRIGHT_BASE_URL", "QUEUE_LOCK_HELD"]); +// The heavy slot's seams are named for the gate, not for e2e: the gate runs +// builds and renders too, and its tests are what set them. +const UNPREFIXED_TEST_VARS = new Set([ + "PLAYWRIGHT_BASE_URL", + "QUEUE_LOCK_HELD", + "HEAVY_HELD", + "HEAVY_LOCK_FILE", + "HEAVY_MEMINFO_FILE", + "HEAVY_POLL_MS", +]); test("every test-only variable carries the E2E_ prefix", () => { const bad = ENV_VARS.filter( diff --git a/common/lib/envVars.ts b/common/lib/envVars.ts @@ -147,6 +147,9 @@ const DECLARED: EnvVarDecl[] = [ { name: "PARAKEET_DECODER", audience: "runtime", default: "parakeet-cli's", readBy: "scripts/parakeet-stitch.mjs", doc: "`ctc` or `tdt`, passed through to parakeet-cli." }, { name: "PARAKEET_LANG", audience: "runtime", default: "parakeet-cli's", readBy: "scripts/parakeet-stitch.mjs", doc: "A locale, passed through to parakeet-cli." }, { name: "PARAKEET_DEVICE", audience: "runtime", default: "parakeet-cli's", readBy: "scripts/parakeet-stitch.mjs", doc: "Compute device (`cpu`, `CUDA0`, `Vulkan1`, …), exported to parakeet-cli." }, + { name: "HEAVY", audience: "runtime", default: "on", readBy: "scripts/queue-lock.mjs", doc: "`0` skips the heavy slot AND the memory floor: the machine-wide one-at-a-time gate that `pnpm heavy -- <cmd>`, every e2e entry point and the publish stages' `next build` go through." }, + { name: "HEAVY_MIN_FREE_MB", audience: "runtime", default: "`6000`", readBy: "scripts/queue-lock.mjs", doc: "The memory floor: a heavy job, once it holds the slot, waits until /proc/meminfo's MemAvailable is at least this many MB. `0` turns the floor off; a machine whose MemTotal is under it runs without waiting." }, + { name: "HEAVY_TIMEOUT", audience: "runtime", default: "wait forever", readBy: "scripts/queue-lock.mjs", doc: "Seconds a `pnpm heavy` run waits for the slot, and then for the floor, before giving up (exit 3). An e2e run uses `E2E_QUEUE_TIMEOUT` for both." }, // ── internal: the pipeline sets these for a process it spawns ────────── { name: "SITE_ID", audience: "internal", default: "—", readBy: "common/bin/compose-site.ts, export/app/lib/site.ts", doc: "Which site a compose or an export build is for. The build stage (`archilyzer publish build <id>`) sets it for its children; `compose site`, `build site` and `deploy site` fall back to it when no id is given." }, @@ -185,6 +188,10 @@ const DECLARED: EnvVarDecl[] = [ { name: "E2E_PORT_GRACE_MS", audience: "test", default: "`3000`", readBy: "scripts/queue-lock.mjs", doc: "How long the port check waits for a just-freed port." }, { name: "E2E_QUEUE_LOCK_FILE", audience: "test", default: "one per machine", readBy: "scripts/queue-lock.mjs", doc: "The queue's lock file; the queue's own tests point it elsewhere." }, { name: "QUEUE_LOCK_HELD", audience: "test", default: "—", readBy: "scripts/queue-lock.mjs", doc: "Set by the queue for the command it runs, so a nested wrapper passes through." }, + { name: "HEAVY_HELD", audience: "test", default: "—", readBy: "scripts/queue-lock.mjs", doc: "Set by the heavy slot for the command it runs, so a heavy command inside it (a `pnpm heavy -- pnpm e2e`, a build stage under an e2e suite's editor) passes through." }, + { name: "HEAVY_LOCK_FILE", audience: "test", default: "one per machine", readBy: "scripts/queue-lock.mjs", doc: "The heavy slot's lock file; the gate's own tests point it elsewhere." }, + { name: "HEAVY_MEMINFO_FILE", audience: "test", default: "`/proc/meminfo`", readBy: "scripts/queue-lock.mjs", doc: "Where the memory floor reads MemAvailable; the gate's tests hand it a fake." }, + { name: "HEAVY_POLL_MS", audience: "test", default: "`5000`", readBy: "scripts/queue-lock.mjs", doc: "How often a run waiting for the memory floor re-reads it." }, { name: "PLAYWRIGHT_BASE_URL", audience: "test", default: "`http://localhost:<PORT>`", readBy: "editor/playwright.config.ts, editor/e2e/baseUrl.ts", doc: "The editor test server's URL; the worktree injector sets it." }, { name: "E2E_AUDIO_CHECK_INTERVAL_MS", audience: "test", default: "the real cadence", readBy: "common/ytdlp/audioCheckedDownload.ts", doc: "Shrinks the mid-download audio check so the e2e suite sees it fire." }, { name: "E2E_AUDIO_CHECK_SIZE_GATE", audience: "test", default: "the real gate", readBy: "common/ytdlp/audioCheckedDownload.ts", doc: "Likewise, the size gate." }, diff --git a/common/publish/build.test.ts b/common/publish/build.test.ts @@ -20,6 +20,7 @@ import { homepageOutDir, dockerSiteOutDir, dockerSiteStagingDir, + heavyGated, resolveOutDir, } from "./build"; @@ -130,6 +131,33 @@ test("EXPORT_NEXT_BIN replaces `pnpm exec next build` in a site's and the hub's assert.deepEqual(plain[1].args, ["exec", "next", "build"]); }); +// A real `next build` runs through the heavy slot (scripts/queue-lock.mjs +// --heavy): one heavy job machine-wide, above the memory floor. The e2e fake is +// never gated, and a root without the gate script (every pure test above, whose +// /repo does not exist) runs the step unchanged. +test("a real next build goes through the heavy slot; the e2e fake and a root without the gate do not", () => { + const root = mkdtempSync(path.join(os.tmpdir(), "build-heavy-")); + try { + mkdirSync(path.join(root, "scripts")); + const gate = path.join(root, "scripts", "queue-lock.mjs"); + writeFileSync(gate, ""); + const p = { ...paths, monorepoRoot: root, exportDir: path.join(root, "export") } as Paths; + const [, next] = buildSiteSteps({ siteId: "jer", paths: p, skipData: true, baseEnv: {} }); + assert.equal(next.command, process.execPath); + assert.deepEqual(next.args, [gate, "--heavy", "--", "pnpm", "exec", "next", "build"]); + assert.equal(next.cwd, path.join(root, "export")); + const hub = buildHubSteps({ paths: p, baseEnv: {} }); + assert.deepEqual(hub[1].args.slice(0, 3), [gate, "--heavy", "--"]); + assert.equal(hub[1].env.INSTANCE_MODE, "hub"); + const fake = buildSiteSteps({ siteId: "jer", paths: p, skipData: true, baseEnv: { EXPORT_NEXT_BIN: "/bin/fake-next" } }); + assert.deepEqual([fake[1].command, ...fake[1].args], ["/bin/fake-next", "build"]); + const step = { command: "pnpm", args: ["x"], cwd: "/", env: {} }; + assert.equal(heavyGated({ monorepoRoot: path.join(root, "nope") }, step), step); + } finally { + rmSync(root, { recursive: true, force: true }); + } +}); + test("buildHubSteps: compose:hub, then next build with INSTANCE_MODE=hub, in export/", () => { const steps = buildHubSteps({ paths, baseEnv: { PATH: "/bin" } }); const env = { diff --git a/common/publish/build.ts b/common/publish/build.ts @@ -114,7 +114,36 @@ export function nextBuildStep(paths: Paths, env: NodeJS.ProcessEnv): BuildStep { const bin = env.EXPORT_NEXT_BIN?.trim(); return bin ? { command: bin, args: ["build"], cwd: paths.exportDir, env } - : { command: "pnpm", args: ["exec", "next", "build"], cwd: paths.exportDir, env }; + : heavyGated(paths, { + command: "pnpm", + args: ["exec", "next", "build"], + cwd: paths.exportDir, + env, + }); +} + +/** + * A real `next build` run through the HEAVY SLOT (scripts/queue-lock.mjs + * --heavy, the same gate as `pnpm heavy -- <cmd>`): one heavy job — an e2e + * run, a build, a render — at a time machine-wide, started only once + * MemAvailable is at least HEAVY_MIN_FREE_MB (6000). Two concurrent builds, or a + * build beside an e2e suite, is how this machine OOMed. The wait is logged + * ("waiting for the heavy slot — held by …") into the stage's own log, and a + * Cancel still stops it: the gate forwards SIGTERM to the build. + * + * Inside a docker runner container the slot is the container's own (no shared + * lock); the floor still reads the host's /proc/meminfo, which throttles a + * fan-out when the host runs low. HEAVY=0 bypasses both; a checkout without + * the gate script (a test's temp root) runs the step as it was. + */ +export function heavyGated(paths: Pick<Paths, "monorepoRoot">, step: BuildStep): BuildStep { + const gate = path.join(paths.monorepoRoot ?? "", "scripts", "queue-lock.mjs"); + if (!paths.monorepoRoot || !existsSync(gate)) return step; + return { + ...step, + command: process.execPath, + args: [gate, "--heavy", "--", step.command, ...step.args], + }; } // Run a list of child steps in order, streaming into `onLog`, stopping at the @@ -695,12 +724,12 @@ export async function buildHomepage( } } return runSteps(onLog, signal, [ - { + heavyGated(paths, { command: "pnpm", args: ["exec", "next", "build"], cwd: homepageDir(paths), env: homepageEnv(paths), - }, + }), ]); } diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **Heavy work takes turns, above a memory floor.** A publish stage's `next build` — a site's, the hub's, the homepage's — now waits for the machine's one heavy slot, which every e2e run and any `pnpm heavy -- <cmd>` (a video render) take too, and then until at least 6000 MB is available; the stage's log says whom it waits behind ("waiting for the heavy slot — held by …") or how much memory there is ("waiting for memory — 4210 MB available, the floor is 6000 MB"). Cancel still stops it. `HEAVY_MIN_FREE_MB` moves the floor (`0` turns it off) and `HEAVY=0` skips the gate. Needs a restart of the editor. - **A curated tag can exist on some sites only.** A tag's new **Sites** field on /tags (`sites` in `transcripts/tags.json`; `pnpm ops tags` takes it in a define) names the sites it exists on. Its rules then fire, and its pins apply, only to videos on those sites' channels, and every other site drops it from its records, its counts and its `/tags.json` — where **Hidden** only hid the chip. Empty is every site, as before. Setting it, or changing the channels of those sites, re-derives the corpus's tags once at the next index update. The Eva tags are what this is for: they belong on Anilyzer alone. - **The publish lane.** Publishing can run itself: turn it on at **/operations/publish** (the runner's Start, Drain and Stop, the hold, and the lane's settings; or `publish.enabled` in settings) and the lane checks every `checkEveryMinutes` (10) whether the index is stale; when it is — and its last update is at least `refreshEveryMinutes` (360) old — it updates it, then builds every site whose channels changed or whose data the new index moved, one stage at a time on the `publish` queue. What it may do with a site is the site's own — the **Publish policy** on the site's settings form, `site.json` `publish.auto` —: `off` (the default: left alone), `build`, `preview` (built and deployed to the preview branch `publish.previewBranch`) or `production`; the hub and the homepage have `publish.hub` and `publish.homepage`. A private site is only ever built, and a site needs its Cloudflare Pages project before it may deploy. Hold the lane and the stage running finishes and no next one starts; quiet hours (`publish.quietHours`) do the same; Drain finishes the stage and ends the runner. The lane never forces a stage: a stage that finds its target current does nothing. On /jobs every stage of one run reads `run <id> · <target>`, and a stage still queued when the editor restarts is cancelled, never re-queued — the lane works out again what is stale from what is on disk. `archilyzer publish now` runs the same plan from the command line, one stage after another in its own process. - **One index for every site.** The index is updated once and every site, the hub and the homepage are built from it; `archilyzer publish status` says, per site, whether its build is current — "stale: 3 channels changed (a, b, c)" as soon as a download, transcription or digest on one of its channels finishes, before any index runs; "stale: data changed" once the index has run and the site's data moved; "stale: config changed" after its site.json, tags or aliases changed — and whether what is deployed is that build, with a build made by older code marked "code newer" but not stale. diff --git a/package.json b/package.json @@ -20,6 +20,7 @@ "dev:umtool": "node scripts/worktree.mjs run -- pnpm --filter umtool run dev", "deploy:homepage": "pnpm --filter homepage run deploy", "e2e": "node scripts/worktree.mjs run -- pnpm --filter editor run e2e", + "heavy": "node scripts/queue-lock.mjs --heavy --", "wt": "node scripts/worktree.mjs", "e2e:sharded": "node scripts/run-sharded-e2e.mjs", "test:scripts": "node --test scripts/*.test.mjs umtool/report-to-video/*.test.mjs umtool/lib/report/*.test.mjs umtool/lib/annotations/*.test.mjs umtool/lib/articles/*.test.mjs", diff --git a/scripts/queue-lock.mjs b/scripts/queue-lock.mjs @@ -1,5 +1,7 @@ #!/usr/bin/env node -// Global e2e queue: exactly one e2e run at a time, machine-wide. +// Global e2e queue: exactly one e2e run at a time, machine-wide — and the +// HEAVY SLOT: one heavy job (an e2e run, a `next build`, a video render) at a +// time, machine-wide, started only above a free-memory floor. // // Every checkout of this repo shares one lock file, so a suite started in a // second worktree waits for the first to finish instead of racing it. That is @@ -10,6 +12,27 @@ // then wipes that session's fixtures with no error at all. // // node scripts/queue-lock.mjs [--name e2e] [--ports PORT:3011,...] -- <cmd...> +// node scripts/queue-lock.mjs --heavy -- <cmd...> (`pnpm heavy -- <cmd>`) +// +// THE HEAVY SLOT. Two OOMs on this machine (2026-10-08) took Xwayland and dbus +// with them: two `next build`s, or a build beside an e2e suite, or a render +// beside either. So every heavy entry point goes through `withHeavy`: one +// machine-wide lock (`heavy-queue.lock`, beside the e2e one), and once it is +// held, a wait until /proc/meminfo's MemAvailable is at least +// HEAVY_MIN_FREE_MB (6000 by default). The slot is taken FIRST, then the floor +// is waited for, so a later contender cannot slip in while the holder waits +// for memory. Callers: +// - every e2e entry point (this file's CLI and run-sharded-e2e.mjs, through +// `withQueue`): the heavy slot first, then the e2e queue. One order +// everywhere, so nesting cannot deadlock: an e2e run holds both, and a +// `pnpm heavy -- pnpm e2e` passes through its own slot (HEAVY_HELD). +// - the publish stages' `next build` (common/publish/build.ts, heavyGated). +// - a render: `pnpm heavy -- node umtool/report-to-video/build-video.mjs …`. +// Bypasses: HEAVY=0 (no slot, no floor), HEAVY_MIN_FREE_MB=0 (no floor), +// HEAVY_TIMEOUT=<seconds>. E2E_QUEUE=0 skips the slot as it skips the queue; +// the floor still applies. A machine whose MemTotal is under the floor is not +// made to wait forever: it is told, and runs. The slot is a safety net, not a +// correctness lock: with no usable `flock` it warns and runs on the floor alone. // // WHY THE LOCK IS HELD BY A SEPARATE CHILD. // flock(1) deliberately keeps its lock fd open across exec — that is what the @@ -48,6 +71,14 @@ const SELF = fileURLToPath(import.meta.url); // wrapper) passes straight through instead of deadlocking against the lock its // own parent is holding. Verified to survive nested `pnpm --filter` calls. export const HELD_ENV = "QUEUE_LOCK_HELD"; +// The same, for the heavy slot: set while a command runs inside it, so a heavy +// command that starts another (a `pnpm heavy -- pnpm e2e`, a build stage run +// from inside an e2e suite's editor) passes through instead of waiting on +// itself. +export const HEAVY_HELD_ENV = "HEAVY_HELD"; +const HEAVY_NAME = "heavy"; +export const DEFAULT_MIN_FREE_MB = 6000; +const DEFAULT_MEM_POLL_MS = 5_000; const PROBE_HELD_EXIT = 91; // `flock -n -E 91`: distinguishes held from failed const WAIT_TIMEOUT_EXIT = 92; // `flock -w N -E 92` @@ -66,12 +97,16 @@ function parseArgv(argv) { name: "e2e", portSpec: null, hold: false, - timeoutMs: defaultTimeoutMs(), + heavy: false, + // Unset unless --timeout: the e2e queue defaults it from E2E_QUEUE_TIMEOUT, + // the heavy slot from HEAVY_TIMEOUT. + timeoutMs: undefined, portGraceMs: Number(process.env.E2E_PORT_GRACE_MS ?? DEFAULT_PORT_GRACE_MS), }; for (let i = 0; i < flags.length; i++) { const f = flags[i]; if (f === "--hold") opts.hold = true; + else if (f === "--heavy") opts.heavy = true; else if (f === "--name") opts.name = flags[++i]; else if (f === "--ports") opts.portSpec = flags[++i]; else if (f === "--timeout") opts.timeoutMs = Number(flags[++i]) * 1000; @@ -86,8 +121,8 @@ function parseArgv(argv) { // Waiting forever is the point of the feature, so that is the default. A // bounded wait is available for anything that would rather fail than block. -function defaultTimeoutMs() { - const raw = process.env.E2E_QUEUE_TIMEOUT; +function defaultTimeoutMs(envVar = "E2E_QUEUE_TIMEOUT") { + const raw = process.env[envVar]; if (raw == null || raw === "") return 0; const n = Number(raw); return Number.isFinite(n) && n > 0 ? n * 1000 : 0; @@ -131,7 +166,13 @@ function gitOut(args) { // .claude/worktrees/* alike — so one file is genuinely machine-global. Living // inside .git/ it is also untracked by construction (no .gitignore entry). function lockFileFor(name) { - if (process.env.E2E_QUEUE_LOCK_FILE) return process.env.E2E_QUEUE_LOCK_FILE; + // Each lock has its own override: one file for both would make an e2e run + // (which takes the heavy slot, then the queue) wait on itself. + const override = + name === HEAVY_NAME + ? process.env.HEAVY_LOCK_FILE + : process.env.E2E_QUEUE_LOCK_FILE; + if (override) return override; const dir = gitOut(["rev-parse", "--path-format=absolute", "--git-common-dir"]) ?? os.tmpdir(); @@ -183,13 +224,15 @@ function readHolderJson(lock) { } } -function describeHolder(lock) { +function describeHolder(lock, showCmd = false) { const h = readHolderJson(lock); if (!h) return "another run (details unavailable)"; const where = h.worktree ? path.basename(h.worktree) : "?"; const age = h.startedAt ? ` for ${humanAge(Date.parse(h.startedAt))}` : ""; const dead = h.alive ? "" : " — pid gone, releasing"; - return `${where} (${h.branch}, pid ${h.pid})${age}${dead}`; + // The heavy slot is shared by kinds of work, so say which one is ahead. + const what = showCmd && h.cmd ? `: ${String(h.cmd).slice(0, 100)}` : ""; + return `${where} (${h.branch}, pid ${h.pid})${age}${what}${dead}`; } function humanAge(startedMs) { @@ -414,13 +457,15 @@ async function preflight(ports, graceMs) { // ---------------------------------------------------------------- run + wait -function runCommand(cmd, name) { +function runCommand(cmd, name, { forwardTerm = false } = {}) { return new Promise((resolve) => { const child = spawn(cmd[0], cmd.slice(1), { stdio: "inherit", - env: { ...process.env, [HELD_ENV]: name }, + // A heavy-only run (name null) leaves QUEUE_LOCK_HELD alone: it holds no + // e2e queue for a nested e2e run to pass through. + env: name ? { ...process.env, [HELD_ENV]: name } : { ...process.env }, }); - installSignalHandlers(() => child); + installSignalHandlers(() => child, forwardTerm); child.on("error", (err) => { process.stderr.write(`queue-lock: ${err.message}\n`); resolve(1); @@ -434,12 +479,24 @@ function runCommand(cmd, name) { } let signalHits = 0; -function installSignalHandlers(getChild) { +// `forwardTerm`: a heavy-slot run is usually started by a PROCESS, not a +// terminal — a publish stage whose Cancel SIGTERMs this wrapper alone, then +// SIGKILLs it, which would orphan the `next build` under it. So a heavy run +// passes SIGTERM/SIGHUP on to its command (a terminal's SIGINT already reached +// the whole group, and is still never forwarded). +function installSignalHandlers(getChild, forwardTerm = false) { for (const sig of ["SIGINT", "SIGTERM", "SIGHUP"]) { process.on(sig, () => { signalHits++; const child = getChild(); if (signalHits === 1) { + if (forwardTerm && sig !== "SIGINT") { + try { + child?.kill(sig); + } catch { + /* already gone */ + } + } // The terminal already delivered this to the whole foreground process // group, child included. We deliberately do not forward it: a second // SIGINT is precisely how playwright skips globalTeardown, which is @@ -460,11 +517,11 @@ function installSignalHandlers(getChild) { } } -function startWaitBanner(lock) { +function startWaitBanner(lock, what = E2E_LOCK) { const t0 = Date.now(); process.stderr.write( - `queue-lock: waiting for the e2e queue — held by ${describeHolder(lock)}\n` + - " (one e2e run at a time, machine-wide; E2E_QUEUE=0 to bypass)\n", + `queue-lock: waiting for ${what.label} — held by ${describeHolder(lock, what.showCmd)}\n` + + ` ${what.hint}\n`, ); const timer = setInterval(() => { process.stderr.write( @@ -480,28 +537,29 @@ function startWaitBanner(lock) { }; } -// ------------------------------------------------------------- the entry point - -/** - * Run `fn` with the global queue lock held, after checking `ports` are free. - * Used both by the CLI below and directly by scripts/run-sharded-e2e.mjs. - */ -export async function withQueue(opts, fn) { - const name = opts.name ?? "e2e"; - const ports = parsePorts(opts.portSpec ?? null); - const graceMs = opts.portGraceMs ?? Number(process.env.E2E_PORT_GRACE_MS ?? DEFAULT_PORT_GRACE_MS); - const timeoutMs = opts.timeoutMs ?? defaultTimeoutMs(); - - // "Don't queue" never means "don't check the ports": the preflight is what - // turns a silent cross-worktree data wipe into a loud abort. - if (process.env.E2E_QUEUE === "0" || process.env[HELD_ENV] === name) { - await preflight(ports, graceMs); - return fn(); - } - - const lock = lockFileFor(name); +// ------------------------------------------------------------ the two locks + +const E2E_LOCK = { + label: "the e2e queue", + hint: "(one e2e run at a time, machine-wide; E2E_QUEUE=0 to bypass)", + timeoutHint: + "Raise or unset E2E_QUEUE_TIMEOUT, or set E2E_QUEUE=0 to bypass the queue.", + showCmd: false, +}; + +const HEAVY_LOCK = { + label: "the heavy slot", + hint: + "(one heavy job — an e2e run, a next build, a render — at a time, machine-wide; HEAVY=0 to bypass)", + timeoutHint: "Raise or unset HEAVY_TIMEOUT, or set HEAVY=0 to bypass the heavy slot.", + showCmd: true, +}; + +// Take `lock`, announcing whom we wait behind; returns the release function. +// A timeout exits EXIT_TIMEOUT, as it always has for the e2e queue. +async function holdLock(lock, name, cmd, timeoutMs, what) { let stopBanner = null; - if (isHeld(lock)) stopBanner = startWaitBanner(lock); + if (isHeld(lock)) stopBanner = startWaitBanner(lock, what); let holder; try { @@ -510,7 +568,7 @@ export async function withQueue(opts, fn) { if (err.timeout) { process.stderr.write( `\nqueue-lock: gave up after ${humanDuration(timeoutMs)} waiting for ${lock}\n` + - " Raise or unset E2E_QUEUE_TIMEOUT, or set E2E_QUEUE=0 to bypass the queue.\n", + ` ${what.timeoutHint}\n`, ); process.exit(EXIT_TIMEOUT); } @@ -518,7 +576,7 @@ export async function withQueue(opts, fn) { } stopBanner?.(); - const holderFile = writeHolderJson(lock, name, opts.cmd ?? [name]); + const holderFile = writeHolderJson(lock, name, cmd); const cleanup = () => { try { fs.rmSync(holderFile, { force: true }); @@ -532,30 +590,231 @@ export async function withQueue(opts, fn) { } }; process.on("exit", cleanup); - - try { - await preflight(ports, graceMs); - return await fn(); - } finally { + return async () => { + process.off("exit", cleanup); try { fs.rmSync(holderFile, { force: true }); } catch { /* best effort */ } await release(holder); + }; +} + +// --------------------------------------------------------- the memory floor + +// /proc/meminfo's MemAvailable and MemTotal, in MB, or null when the text has +// neither (not Linux, or a reader handed something else). +export function parseMeminfo(text) { + const kb = (key) => { + const m = new RegExp(`^${key}:\\s+(\\d+)\\s*kB`, "m").exec(String(text)); + return m ? Number(m[1]) : null; + }; + const available = kb("MemAvailable"); + const total = kb("MemTotal"); + if (available == null || total == null) return null; + return { + availableMb: Math.floor(available / 1024), + totalMb: Math.floor(total / 1024), + }; +} + +// HEAVY_MEMINFO_FILE is the tests' seam: a file they rewrite to move the +// "available" figure under a waiting run. +export function readMeminfo(file = process.env.HEAVY_MEMINFO_FILE || "/proc/meminfo") { + try { + return parseMeminfo(fs.readFileSync(file, "utf8")); + } catch { + return null; } } +export function minFreeMb(env = process.env) { + const raw = env.HEAVY_MIN_FREE_MB; + if (raw == null || raw === "") return DEFAULT_MIN_FREE_MB; + const n = Number(raw); + return Number.isFinite(n) && n >= 0 ? n : DEFAULT_MIN_FREE_MB; +} + +function memPollMs(env = process.env) { + const n = Number(env.HEAVY_POLL_MS); + return Number.isFinite(n) && n > 0 ? n : DEFAULT_MEM_POLL_MS; +} + +/** + * Wait until MemAvailable >= `minMb`. Every input is injectable — `read` + * returns `{availableMb, totalMb}` or null — so the unit tests drive it with + * no real memory pressure. Resolves `{waitedMs}` or `{skipped}` (why the floor + * was not waited for); throws `{timeout: true}` past `timeoutMs` (0 = never). + */ +export async function waitForMemory({ + minMb = minFreeMb(), + read = readMeminfo, + pollMs = memPollMs(), + tickMs = TICK_MS, + timeoutMs = 0, + log = (line) => process.stderr.write(line), + now = Date.now, + sleep = (ms) => new Promise((r) => setTimeout(r, ms)), +} = {}) { + if (!(minMb > 0)) return { skipped: "off" }; + let m = read(); + if (!m) { + log("heavy: /proc/meminfo is not readable — the memory floor is not checked\n"); + return { skipped: "unreadable" }; + } + // A floor the machine cannot reach would be a wait forever. Say so and run. + if (m.totalMb < minMb) { + log( + `heavy: this machine has ${m.totalMb} MB in all, under the ${minMb} MB floor — not waiting for it\n`, + ); + return { skipped: "total" }; + } + if (m.availableMb >= minMb) return { waitedMs: 0 }; + const t0 = now(); + let lastTick = t0; + log( + `heavy: waiting for memory — ${m.availableMb} MB available, the floor is ${minMb} MB\n` + + " (HEAVY_MIN_FREE_MB=<MB> to change it; 0, or HEAVY=0, to skip it)\n", + ); + for (;;) { + await sleep(pollMs); + m = read() ?? m; + const waited = now() - t0; + if (m.availableMb >= minMb) { + log(`heavy: ${m.availableMb} MB available after ${humanDuration(waited)}\n`); + return { waitedMs: waited }; + } + if (timeoutMs > 0 && waited >= timeoutMs) { + throw Object.assign( + new Error( + `heavy: gave up after ${humanDuration(waited)} waiting for ${minMb} MB available (${m.availableMb} MB)`, + ), + { timeout: true }, + ); + } + if (now() - lastTick >= tickMs) { + lastTick = now(); + log( + `heavy: still waiting for memory (${m.availableMb} MB available, ${humanDuration(waited)})\n`, + ); + } + } +} + +async function memoryFloorOrExit(timeoutMs) { + try { + await waitForMemory({ timeoutMs }); + } catch (err) { + if (!err.timeout) throw err; + process.stderr.write( + `\n${err.message}\n Raise or unset the timeout, lower HEAVY_MIN_FREE_MB, or set HEAVY=0.\n`, + ); + process.exit(EXIT_TIMEOUT); + } +} + +// ------------------------------------------------------------ the entry points + +/** + * Run `fn` in the heavy slot: one heavy job machine-wide, started only once + * MemAvailable is at or above the floor. `opts.cmd` names the work in the + * holder file (what a waiter is told it waits behind). + */ +export async function withHeavy(opts, fn) { + if (process.env.HEAVY === "0" || process.env[HEAVY_HELD_ENV]) return fn(); + const timeoutMs = opts.timeoutMs ?? defaultTimeoutMs("HEAVY_TIMEOUT"); + const cmd = opts.cmd ?? [HEAVY_NAME]; + + let releaseSlot = null; + try { + releaseSlot = await holdLock( + lockFileFor(HEAVY_NAME), + HEAVY_NAME, + cmd, + timeoutMs, + HEAVY_LOCK, + ); + } catch (err) { + // No usable flock (a container image without util-linux): the slot is a + // safety net, so run on the floor alone rather than not at all. + process.stderr.write( + `heavy: the heavy slot is not held (${err?.message ?? err}) — the memory floor still applies\n`, + ); + } + try { + await memoryFloorOrExit(timeoutMs); + process.env[HEAVY_HELD_ENV] = "1"; + return await fn(); + } finally { + delete process.env[HEAVY_HELD_ENV]; + await releaseSlot?.(); + } +} + +/** + * Run `fn` with the global queue lock held, after checking `ports` are free. + * Used both by the CLI below and directly by scripts/run-sharded-e2e.mjs. + * Unless `opts.heavySlot === false`, the run takes the heavy slot first. + */ +export async function withQueue(opts, fn) { + const name = opts.name ?? "e2e"; + const ports = parsePorts(opts.portSpec ?? null); + const graceMs = opts.portGraceMs ?? Number(process.env.E2E_PORT_GRACE_MS ?? DEFAULT_PORT_GRACE_MS); + const timeoutMs = opts.timeoutMs ?? defaultTimeoutMs(); + + // A nested invocation: the parent holds the queue (and the slot) already. + if (process.env[HELD_ENV] === name) { + await preflight(ports, graceMs); + return fn(); + } + // "Don't queue" never means "don't check the ports": the preflight is what + // turns a silent cross-worktree data wipe into a loud abort. Nor does it + // mean "ignore the memory floor" — HEAVY=0 is that switch. + if (process.env.E2E_QUEUE === "0") { + await preflight(ports, graceMs); + if (process.env.HEAVY !== "0" && !process.env[HEAVY_HELD_ENV]) { + await memoryFloorOrExit(timeoutMs); + } + return fn(); + } + + const queued = async () => { + const releaseQueue = await holdLock( + lockFileFor(name), + name, + opts.cmd ?? [name], + timeoutMs, + E2E_LOCK, + ); + try { + await preflight(ports, graceMs); + return await fn(); + } finally { + await releaseQueue(); + } + }; + if (opts.heavySlot === false) return queued(); + return withHeavy({ timeoutMs, cmd: opts.cmd ?? [name] }, queued); +} + async function main() { - const { opts, cmd } = parseArgv(process.argv.slice(2)); + const { opts, cmd: rawCmd } = parseArgv(process.argv.slice(2)); if (opts.hold) return runHolder(); + // `pnpm heavy -- <cmd>` may hand the separator through as the first word. + const cmd = opts.heavy && rawCmd[0] === "--" ? rawCmd.slice(1) : rawCmd; if (cmd.length === 0) { process.stderr.write( - "usage: queue-lock.mjs [--name e2e] [--ports PORT:3011,...] -- <cmd...>\n", + "usage: queue-lock.mjs [--name e2e] [--ports PORT:3011,...] -- <cmd...>\n" + + " queue-lock.mjs --heavy [--timeout <s>] -- <cmd...>\n", ); process.exit(2); } - const code = await withQueue({ ...opts, cmd }, () => runCommand(cmd, opts.name)); + const code = opts.heavy + ? await withHeavy({ timeoutMs: opts.timeoutMs, cmd }, () => + runCommand(cmd, null, { forwardTerm: true }), + ) + : await withQueue({ ...opts, cmd }, () => runCommand(cmd, opts.name)); process.exit(code); } diff --git a/scripts/queue-lock.test.mjs b/scripts/queue-lock.test.mjs @@ -1,8 +1,10 @@ -// Tests for the global e2e queue (scripts/queue-lock.mjs). +// Tests for the global e2e queue and the heavy slot (scripts/queue-lock.mjs). // -// Every test drives the real CLI against a throwaway lock file via -// E2E_QUEUE_LOCK_FILE, so none of them can touch the actual .git/e2e-queue.lock -// or any real port. Run with: pnpm test:scripts +// Every test drives the real CLI against throwaway lock files via +// E2E_QUEUE_LOCK_FILE / HEAVY_LOCK_FILE, so none of them can touch the actual +// .git/e2e-queue.lock or .git/heavy-queue.lock or any real port. The e2e-queue +// tests run with HEAVY=0 (they are about the queue); the heavy tests turn it +// on with their own lock and a fake /proc/meminfo. Run with: pnpm test:scripts import assert from "node:assert/strict"; import { spawn } from "node:child_process"; import fs from "node:fs"; @@ -11,6 +13,7 @@ import os from "node:os"; import path from "node:path"; import test from "node:test"; import { fileURLToPath } from "node:url"; +import { parseMeminfo, waitForMemory } from "./queue-lock.mjs"; const SCRIPT = fileURLToPath(new URL("./queue-lock.mjs", import.meta.url)); @@ -18,21 +21,58 @@ function tmpDir() { return fs.mkdtempSync(path.join(os.tmpdir(), "queue-lock-test-")); } -// Run the wrapper to completion, capturing output. `env` is merged over the -// current environment; E2E_PORT_CHECK defaults off so tests that are not about -// the preflight never probe a port. -function runLock(args, env = {}) { - return new Promise((resolve) => { - const child = spawn(process.execPath, [SCRIPT, ...args], { - env: { E2E_PORT_CHECK: "0", ...process.env, ...env }, - stdio: ["ignore", "pipe", "pipe"], - }); - let stdout = ""; - let stderr = ""; - child.stdout.on("data", (d) => (stdout += d)); - child.stderr.on("data", (d) => (stderr += d)); - child.on("exit", (code) => resolve({ code, stdout, stderr })); +// The environment a run sees: `env` over the current one. E2E_PORT_CHECK +// defaults off so tests that are not about the preflight never probe a port; +// HEAVY defaults off so the e2e-queue tests never take a heavy slot. A +// pass-through marker inherited from whatever runs this suite (a +// `pnpm heavy -- pnpm test:scripts`) is dropped unless the test sets it. +function lockEnv(env) { + const out = { E2E_PORT_CHECK: "0", HEAVY: "0", ...process.env, ...env }; + for (const k of ["HEAVY_HELD", "QUEUE_LOCK_HELD", "E2E_QUEUE"]) { + if (!(k in env)) delete out[k]; + } + return out; +} + +// Start the wrapper; `done` resolves at exit with the captured output, and +// `waitFor(re)` resolves once stderr matches — how a test knows a contender is +// queued (its banner is out) rather than guessing with a delay. +function startLock(args, env = {}) { + const child = spawn(process.execPath, [SCRIPT, ...args], { + env: lockEnv(env), + stdio: ["ignore", "pipe", "pipe"], }); + let stdout = ""; + let stderr = ""; + const waiters = []; + child.stdout.on("data", (d) => (stdout += d)); + child.stderr.on("data", (d) => { + stderr += d; + for (const w of waiters) if (w.re.test(stderr)) w.resolve(); + }); + const done = new Promise((resolve) => + child.on("exit", (code) => resolve({ code, stdout, stderr })), + ); + const waitFor = (re) => + new Promise((resolve, reject) => { + if (re.test(stderr)) return resolve(); + waiters.push({ re, resolve }); + done.then(() => reject(new Error(`exited before stderr matched ${re}: ${stderr}`))); + }); + return { child, done, waitFor }; +} + +function runLock(args, env = {}) { + return startLock(args, env).done; +} + +// Resolves once `check()` is true (a holder.json written: a run has acquired). +async function until(check, ms = 10_000) { + const t0 = Date.now(); + while (!check()) { + if (Date.now() - t0 > ms) throw new Error("timed out waiting"); + await delay(20); + } } // A command that records "S<id>" when it starts and "E<id>" when it ends, so @@ -72,10 +112,17 @@ test("serves waiters in arrival order (FIFO)", async () => { const log = path.join(dir, "fifo.log"); const env = { E2E_QUEUE_LOCK_FILE: lock }; - const runs = []; - for (const id of ["1", "2", "3"]) { - runs.push(runLock(["--", ...markerCmd(log, id, 300)], env)); - await delay(150); // stagger arrivals so the intended order is unambiguous + // Each arrival waits until the one before it is in line: the first holds + // (its holder.json is written), the next two have printed their banner and + // had a moment to block in flock. A fixed stagger raced a loaded machine. + const first = startLock(["--", ...markerCmd(log, "1", 600)], env); + await until(() => fs.existsSync(`${lock}.holder.json`)); + const runs = [first.done]; + for (const id of ["2", "3"]) { + const r = startLock(["--", ...markerCmd(log, id, 300)], env); + await r.waitFor(/waiting for the e2e queue/); + await delay(150); + runs.push(r.done); } await Promise.all(runs); @@ -88,7 +135,7 @@ test("prints a banner naming the holder while waiting", async () => { const env = { E2E_QUEUE_LOCK_FILE: lock }; const first = runLock(["--", process.execPath, "-e", "setTimeout(()=>{},600)"], env); - await delay(200); + await until(() => fs.existsSync(`${lock}.holder.json`)); const second = await runLock(["--", process.execPath, "-e", "0"], env); await first; @@ -139,7 +186,7 @@ test("E2E_QUEUE=0 bypasses the queue entirely", async () => { test("a SIGKILLed run releases the lock immediately", async () => { const dir = tmpDir(); const lock = path.join(dir, "q.lock"); - const env = { E2E_PORT_CHECK: "0", ...process.env, E2E_QUEUE_LOCK_FILE: lock }; + const env = lockEnv({ E2E_QUEUE_LOCK_FILE: lock }); const victim = spawn( process.execPath, @@ -242,3 +289,204 @@ test("removes its holder.json when the run finishes", async () => { await runLock(["--", process.execPath, "-e", "0"], { E2E_QUEUE_LOCK_FILE: lock }); assert.equal(fs.existsSync(`${lock}.holder.json`), false); }); + +// ------------------------------------------------------------- the heavy slot + +// A fake /proc/meminfo: `availableMb` free of `totalMb`. +function meminfo(file, availableMb, totalMb = 32_000) { + fs.writeFileSync( + file, + `MemTotal: ${totalMb * 1024} kB\nMemFree: 1024 kB\nMemAvailable: ${availableMb * 1024} kB\n`, + ); +} + +// The heavy slot on, against its own lock and a roomy fake meminfo. +function heavyEnv(dir, extra = {}) { + const mem = path.join(dir, "meminfo"); + if (!fs.existsSync(mem)) meminfo(mem, 20_000); + return { + HEAVY: "1", + HEAVY_LOCK_FILE: path.join(dir, "heavy.lock"), + HEAVY_MEMINFO_FILE: mem, + HEAVY_POLL_MS: "50", + ...extra, + }; +} + +test("parseMeminfo reads MemAvailable and MemTotal in MB", () => { + assert.deepEqual( + parseMeminfo("MemTotal: 32768000 kB\nMemFree: 1 kB\nMemAvailable: 6144000 kB\n"), + { availableMb: 6000, totalMb: 32000 }, + ); + assert.equal(parseMeminfo("nothing here"), null); +}); + +test("waitForMemory waits for the floor, polling the injected reader", async () => { + const readings = [2000, 4000, 5999, 6000]; + const lines = []; + let clock = 0; + const res = await waitForMemory({ + minMb: 6000, + read: () => ({ availableMb: readings.shift() ?? 6000, totalMb: 32_000 }), + pollMs: 1000, + tickMs: 2000, + log: (l) => lines.push(l), + now: () => clock, + sleep: async (ms) => { + clock += ms; + }, + }); + assert.equal(res.waitedMs, 3000); + assert.match(lines[0], /waiting for memory — 2000 MB available, the floor is 6000 MB/); + assert.ok(lines.some((l) => /still waiting for memory \(5999 MB/.test(l)), lines.join("")); + assert.match(lines.at(-1), /6000 MB available after 3s/); +}); + +test("waitForMemory: no wait above the floor, at 0, or under a MemTotal that can never reach it", async () => { + const never = () => { + throw new Error("must not sleep"); + }; + const read = (a, t = 32_000) => () => ({ availableMb: a, totalMb: t }); + assert.deepEqual(await waitForMemory({ minMb: 6000, read: read(9000), sleep: never }), { waitedMs: 0 }); + assert.deepEqual(await waitForMemory({ minMb: 0, read: read(10), sleep: never }), { skipped: "off" }); + const lines = []; + assert.deepEqual( + await waitForMemory({ minMb: 6000, read: read(100, 4000), sleep: never, log: (l) => lines.push(l) }), + { skipped: "total" }, + ); + assert.match(lines.join(""), /4000 MB in all, under the 6000 MB floor/); + assert.deepEqual( + await waitForMemory({ minMb: 6000, read: () => null, sleep: never, log: () => {} }), + { skipped: "unreadable" }, + ); +}); + +test("waitForMemory gives up past its timeout", async () => { + let clock = 0; + await assert.rejects( + waitForMemory({ + minMb: 6000, + read: () => ({ availableMb: 100, totalMb: 32_000 }), + pollMs: 1000, + timeoutMs: 3000, + log: () => {}, + now: () => clock, + sleep: async (ms) => { + clock += ms; + }, + }), + (err) => err.timeout === true, + ); +}); + +test("two heavy contenders run one at a time; the second names what it waits behind", async () => { + const dir = tmpDir(); + const log = path.join(dir, "order.log"); + const env = heavyEnv(dir); + const first = startLock(["--heavy", "--", ...markerCmd(log, "1", 600)], env); + await until(() => fs.existsSync(`${env.HEAVY_LOCK_FILE}.holder.json`)); + const second = await runLock(["--heavy", "--", ...markerCmd(log, "2", 50)], env); + assert.equal((await first.done).code, 0); + assert.equal(second.code, 0); + assert.equal(fs.readFileSync(log, "utf8"), "S1E1S2E2"); + assert.match(second.stderr, /waiting for the heavy slot — held by .*pid \d+.*: .*appendFileSync/); + assert.equal(fs.existsSync(`${env.HEAVY_LOCK_FILE}.holder.json`), false); +}); + +test("an e2e run takes the heavy slot too: it waits for a heavy job, then runs", async () => { + const dir = tmpDir(); + const log = path.join(dir, "order.log"); + const env = heavyEnv(dir, { E2E_QUEUE_LOCK_FILE: path.join(dir, "q.lock") }); + const build = startLock(["--heavy", "--", ...markerCmd(log, "b", 600)], env); + await until(() => fs.existsSync(`${env.HEAVY_LOCK_FILE}.holder.json`)); + const e2e = await runLock(["--", ...markerCmd(log, "e", 50)], env); + await build.done; + assert.equal(e2e.code, 0); + assert.equal(fs.readFileSync(log, "utf8"), "SbEbSeEe"); + assert.match(e2e.stderr, /waiting for the heavy slot/); +}); + +test("a heavy run inside a heavy run passes through (no self-deadlock), and so does an e2e run inside one", async () => { + const dir = tmpDir(); + const env = heavyEnv(dir, { E2E_QUEUE_LOCK_FILE: path.join(dir, "q.lock") }); + const inner = `${JSON.stringify(process.execPath)} ${JSON.stringify(SCRIPT)}`; + const res = await runLock( + ["--heavy", "--", "sh", "-c", `${inner} --heavy -- true && ${inner} -- true && echo nested-ok`], + env, + ); + assert.equal(res.code, 0, res.stderr); + assert.match(res.stdout, /nested-ok/); + assert.doesNotMatch(res.stderr, /waiting for the heavy slot/); +}); + +test("a heavy run waits under the memory floor and starts once memory is back", async () => { + const dir = tmpDir(); + const env = heavyEnv(dir); + meminfo(env.HEAVY_MEMINFO_FILE, 1500); + const run = startLock(["--heavy", "--", process.execPath, "-e", 'console.log("ran")'], env); + await run.waitFor(/waiting for memory — 1500 MB available, the floor is 6000 MB/); + meminfo(env.HEAVY_MEMINFO_FILE, 7000); + const res = await run.done; + assert.equal(res.code, 0); + assert.match(res.stdout, /ran/); + assert.match(res.stderr, /7000 MB available after/); +}); + +test("HEAVY_MIN_FREE_MB moves the floor; HEAVY=0 skips slot and floor", async () => { + const dir = tmpDir(); + const env = heavyEnv(dir); + meminfo(env.HEAVY_MEMINFO_FILE, 1500); + const lowered = await runLock(["--heavy", "--", "true"], { ...env, HEAVY_MIN_FREE_MB: "1000" }); + assert.equal(lowered.code, 0); + assert.doesNotMatch(lowered.stderr, /waiting for memory/); + const off = await runLock(["--heavy", "--", "true"], { ...env, HEAVY: "0" }); + assert.equal(off.code, 0); + assert.equal(off.stderr, ""); +}); + +test("a SIGKILLed heavy holder hands the slot on at once (the stale holder)", async () => { + const dir = tmpDir(); + const env = heavyEnv(dir); + const victim = startLock(["--heavy", "--", process.execPath, "-e", "setTimeout(()=>{},30000)"], env); + await until(() => fs.existsSync(`${env.HEAVY_LOCK_FILE}.holder.json`)); + victim.child.kill("SIGKILL"); + await victim.done; + // Its holder.json is left behind (SIGKILL runs no cleanup); the lock is not. + const t0 = Date.now(); + const next = await runLock(["--heavy", "--", "true"], env); + assert.equal(next.code, 0); + assert.ok(Date.now() - t0 < 5000, "the heavy slot survived a SIGKILLed holder"); +}); + +test("pnpm's `--` separator is accepted before a heavy command", async () => { + const dir = tmpDir(); + const res = await runLock( + ["--heavy", "--", "--", process.execPath, "-e", 'console.log("ran")'], + heavyEnv(dir), + ); + assert.equal(res.code, 0, res.stderr); + assert.match(res.stdout, /ran/); +}); + +test("SIGTERM to a heavy wrapper alone stops its command (a stage's Cancel)", async () => { + const dir = tmpDir(); + const pidFile = path.join(dir, "cmd.pid"); + const run = startLock( + [ + "--heavy", + "--", + process.execPath, + "-e", + `require("fs").writeFileSync(${JSON.stringify(pidFile)}, String(process.pid)); setTimeout(()=>{},30000)`, + ], + heavyEnv(dir), + ); + await until(() => fs.existsSync(pidFile) && fs.readFileSync(pidFile, "utf8") !== ""); + const cmdPid = Number(fs.readFileSync(pidFile, "utf8")); + const t0 = Date.now(); + run.child.kill("SIGTERM"); + const res = await run.done; + assert.ok(Date.now() - t0 < 5000, "the wrapper outlived its SIGTERM"); + assert.equal(res.code, 128 + os.constants.signals.SIGTERM); + assert.throws(() => process.kill(cmdPid, 0), "the command survived the wrapper's SIGTERM"); +});