#!/usr/bin/env node // 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 // both a resource decision (the serial suite is ~24 min and e2e:sharded fans // out to 4 containers) and a correctness one: playwright.config's // `reuseExistingServer: !CI` means a second run that finds port 3011 already // bound silently drives the *other* worktree's test server, and resetData() // then wipes that session's fixtures with no error at all. // // node scripts/queue-lock.mjs [--name e2e] [--ports PORT:3011,...] -- // node scripts/queue-lock.mjs --heavy -- (`pnpm heavy -- `) // // 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=. 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 // -o/--close flag exists to undo — so the fd is inherited by every descendant, // transitively. Running the suite *under* flock would therefore hold the lock // until the last grandchild dies, not until playwright exits. This suite leaks // grandchildren for real (see editor/e2e/fixtureProcs.ts, which exists because // 27 strays once accumulated, and e2e/fixtures/bin/fake-parakeet-stitch.mjs, // which ignores SIGTERM on purpose), and one leaked stray would deadlock every // worktree on the machine forever — strictly worse than the bug being fixed. // // So the lock is held by a dedicated child that has no children of its own, // and the real command runs as its SIBLING: // // node queue-lock.mjs (this process) // |-- flock -F node queue-lock.mjs --hold (owns the fd; no kids) // `-- playwright test -> next dev -> fixtures (no lock fd anywhere) // // Acquisition is signalled by the holder writing ACQUIRED on a pipe, so the // waiting banner is cancelled deterministically rather than by polling. // Release is by closing the holder's stdin: it exits, and the kernel drops the // lock when the fd closes. That makes the queue crash-proof in both // directions — a SIGKILLed run can never strand it, so there is no stale-lock // recovery code here on purpose. import { execFileSync, spawn } from "node:child_process"; import fs from "node:fs"; import net from "node:net"; import os from "node:os"; import path from "node:path"; import { fileURLToPath } from "node:url"; const SELF = fileURLToPath(import.meta.url); // Set in the environment of the command we run, so a nested invocation (root // `pnpm e2e` -> `pnpm --filter editor run e2e`, both of which go through this // 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` const EXIT_PORTS = 1; // preflight found a bound port const EXIT_TIMEOUT = 3; // gave up waiting for the queue const TICK_MS = 30_000; const DEFAULT_PORT_GRACE_MS = 3_000; // ---------------------------------------------------------------- arguments function parseArgv(argv) { const sep = argv.indexOf("--"); const flags = sep === -1 ? argv : argv.slice(0, sep); const cmd = sep === -1 ? [] : argv.slice(sep + 1); const opts = { name: "e2e", portSpec: null, hold: false, 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; else if (f === "--port-grace") opts.portGraceMs = Number(flags[++i]); else { process.stderr.write(`queue-lock: unknown flag '${f}'\n`); process.exit(2); } } return { opts, cmd }; } // 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(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; } // "PORT:3011,EXPORT_PORT:3010" -> [{envVar, port}]. The literal is the same // default the matching playwright config uses, so the check and the run always // agree; an env var (e.g. from `wt run`) wins, exactly as it does in the config. function parsePorts(spec) { if (!spec || process.env.E2E_PORT_CHECK === "0") return []; return spec .split(",") .map((s) => s.trim()) .filter(Boolean) .map((entry) => { const [envVar, fallback] = entry.split(":"); const port = Number(process.env[envVar] ?? fallback); if (!Number.isInteger(port) || port <= 0) { process.stderr.write(`queue-lock: bad --ports entry '${entry}'\n`); process.exit(2); } return { envVar, port }; }); } // ------------------------------------------------------------- lock location function gitOut(args) { try { return execFileSync("git", args, { encoding: "utf8", stdio: ["ignore", "pipe", "ignore"], }).trim(); } catch { return null; } } // The same .git for every checkout: `--git-common-dir` resolves to the main // repository's .git from the main worktree, from ../feat-* siblings, and from // .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) { // 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(); return path.join(dir, `${name}-queue.lock`); } // ------------------------------------------------------- holder bookkeeping function holderFileFor(lock) { return `${lock}.holder.json`; } function writeHolderJson(lock, name, cmd) { const info = { name, pid: process.pid, worktree: gitOut(["rev-parse", "--show-toplevel"]) ?? process.cwd(), branch: gitOut(["rev-parse", "--abbrev-ref", "HEAD"]) ?? "?", cmd: cmd.join(" "), startedAt: new Date().toISOString(), }; const file = holderFileFor(lock); try { // tmp+rename so a reader never sees a half-written file. const tmp = `${file}.${process.pid}.tmp`; fs.writeFileSync(tmp, `${JSON.stringify(info, null, 2)}\n`); fs.renameSync(tmp, file); } catch { /* informational only — never fail a run over it */ } return file; } // Purely informational ("who am I waiting behind"): the authoritative // acquire/release signal is the lock itself, so a missing, stale, or corrupt // file must never break the queue. function readHolderJson(lock) { try { const info = JSON.parse(fs.readFileSync(holderFileFor(lock), "utf8")); let alive = true; try { process.kill(info.pid, 0); } catch { alive = false; } return { ...info, alive }; } catch { return null; } } 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"; // 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) { if (!Number.isFinite(startedMs)) return "?"; return humanDuration(Date.now() - startedMs); } function humanDuration(ms) { const s = Math.max(0, Math.round(ms / 1000)); return s < 60 ? `${s}s` : `${Math.floor(s / 60)}m${String(s % 60).padStart(2, "0")}s`; } // -------------------------------------------------------------- the holder // Runs with the flock fd inherited from `flock -F`, and deliberately has no // children of its own. function runHolder() { // Ctrl-C reaches the whole foreground process group at once. If this process // died here the lock would drop while the run it is guarding is still // shutting down and still bound to :3011/:3010 — the next waiter would wake // up and abort on ports that are a second away from being free. The wrapper // decides when to let go, by closing our stdin. for (const sig of ["SIGINT", "SIGTERM", "SIGHUP"]) process.on(sig, () => {}); process.stdout.write("ACQUIRED\n"); process.stdin.resume(); const bye = () => process.exit(0); process.stdin.on("end", bye); process.stdin.on("close", bye); process.stdin.on("error", bye); } // ------------------------------------------------------------ acquire/release function flockArgs(lock, timeoutMs) { // -F/--no-fork execs directly, so the spawned process IS the holder and we // can talk to it over its own stdio. const args = ["-F"]; if (timeoutMs > 0) { args.push("-w", String(Math.ceil(timeoutMs / 1000)), "-E", String(WAIT_TIMEOUT_EXIT)); } args.push(lock, process.execPath, SELF, "--hold"); return args; } // Non-blocking probe, purely so the banner can name who is ahead of us. A // false "free" here is harmless: the blocking acquire below is the real gate. function isHeld(lock) { try { execFileSync("flock", ["-n", "-E", String(PROBE_HELD_EXIT), lock, "true"], { stdio: "ignore", }); return false; } catch (err) { if (err.status === PROBE_HELD_EXIT) return true; if (err.code === "ENOENT") throw missingFlock(); return false; } } function missingFlock() { return new Error( "queue-lock: `flock` not found — install util-linux, or set E2E_QUEUE=0 to run unqueued", ); } function acquire(lock, timeoutMs) { const holder = spawn("flock", flockArgs(lock, timeoutMs), { stdio: ["pipe", "pipe", "inherit"], }); return new Promise((resolve, reject) => { let buf = ""; let settled = false; holder.stdout.setEncoding("utf8"); holder.stdout.on("data", (chunk) => { buf += chunk; if (!settled && buf.includes("ACQUIRED")) { settled = true; resolve(holder); } }); holder.on("error", (err) => reject(err.code === "ENOENT" ? missingFlock() : err), ); holder.on("exit", (code) => { if (settled) { // Expected: release() closed its stdin because our run finished. if (holder.releasing) return; // Unexpected: the lock is gone but our run is still going — say so // loudly, because another run can now start on top of this one. process.stderr.write( `\nqueue-lock: WARNING the lock holder died (exit ${code}) mid-run\n`, ); return; } if (code === WAIT_TIMEOUT_EXIT) { reject(Object.assign(new Error("queue-lock: timed out waiting"), { timeout: true })); } else { reject(new Error(`queue-lock: flock exited ${code} before acquiring`)); } }); }); } function release(holder) { return new Promise((resolve) => { holder.releasing = true; // distinguishes this from the holder crashing if (holder.exitCode != null || holder.signalCode != null) return resolve(); const hard = setTimeout(() => { try { holder.kill("SIGKILL"); } catch { /* already gone */ } }, 2000); holder.on("exit", () => { clearTimeout(hard); resolve(); }); try { holder.stdin.end(); } catch { try { holder.kill("SIGKILL"); } catch { /* already gone */ } } }); } // ------------------------------------------------------------ port preflight function tryBind(port, host) { return new Promise((resolve) => { const server = net.createServer(); server.once("error", (err) => resolve(err.code ?? "EUNKNOWN")); server.listen({ port, host, exclusive: true }, () => server.close(() => resolve(null)), ); }); } // Only EADDRINUSE counts as occupied. A 0.0.0.0 bind fails even when the // listener is bound to 127.0.0.1 only (SO_REUSEADDR does not permit // overlapping listens), verified against the live ollama stub on // 127.0.0.1:11435 and against next-server on *:3011. The `::` probe closes the // IPv6-loopback-only hole; EAFNOSUPPORT there on a v6-less host must not abort. async function isOccupied(port) { for (const host of ["0.0.0.0", "::"]) { if ((await tryBind(port, host)) === "EADDRINUSE") return true; } return false; } // Best-effort owner lookup for the abort message. `ss` may be absent. function culpritOf(port) { let out; try { out = execFileSync("ss", ["-ltnpH", `sport = :${port}`], { encoding: "utf8", stdio: ["ignore", "pipe", "ignore"], }); } catch { return null; } const pid = /pid=(\d+)/.exec(out)?.[1]; if (!pid) return null; let cwd = "?"; let cmd = "?"; try { cwd = fs.readlinkSync(`/proc/${pid}/cwd`); } catch { /* gone or not ours */ } try { cmd = fs .readFileSync(`/proc/${pid}/cmdline`, "utf8") .split("\0") .filter(Boolean) .join(" ") .slice(0, 140); } catch { /* gone or not ours */ } return { pid, cwd, cmd }; } // Runs only AFTER the lock is held. Before that, a busy port is most likely a // legitimately running queued suite, and aborting on it would be nonsense. async function preflight(ports, graceMs) { if (!ports.length) return; const deadline = Date.now() + Math.max(0, graceMs); for (;;) { const busy = []; for (const p of ports) if (await isOccupied(p.port)) busy.push(p); if (!busy.length) return; // The previous run's `next dev` may still be releasing the port in the // moment we take the lock. Poll briefly so a clean handoff is a one-second // wait rather than a spurious failure. if (Date.now() >= deadline) { process.stderr.write( "\nqueue-lock: ABORTING — this suite's ports are already in use.\n" + "We hold the queue lock, so no queued run is using them: these are\n" + "leftovers from a killed run, or servers started by hand.\n\n", ); for (const p of busy) { const c = culpritOf(p.port); process.stderr.write( c ? ` ${p.envVar}=${p.port} pid ${c.pid} cwd ${c.cwd}\n ${c.cmd}\n kill ${c.pid}\n` : ` ${p.envVar}=${p.port} (owner unknown — try: ss -ltnp | grep ${p.port})\n`, ); } process.stderr.write( "\nSet E2E_PORT_CHECK=0 to run anyway (the suite will attach to those servers).\n", ); process.exit(EXIT_PORTS); } await new Promise((r) => setTimeout(r, 250)); } } // ---------------------------------------------------------------- run + wait function runCommand(cmd, name, { forwardTerm = false } = {}) { return new Promise((resolve) => { const child = spawn(cmd[0], cmd.slice(1), { stdio: "inherit", // 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, forwardTerm); child.on("error", (err) => { process.stderr.write(`queue-lock: ${err.message}\n`); resolve(1); }); child.on("exit", (code, signal) => { // Never re-raise onto ourselves the way worktree.mjs's cmdRun does: we // have handlers installed, so re-raising would loop. resolve(signal ? 128 + (os.constants.signals[signal] ?? 0) : (code ?? 1)); }); }); } let signalHits = 0; // `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 // how this repo gets orphaned fixtures. We also stay alive, so the // lock is not released while that tree is still shutting down. process.stderr.write( `\nqueue-lock: ${sig} — waiting for the run to stop (again to force)\n`, ); return; } process.stderr.write("queue-lock: forcing shutdown\n"); try { child?.kill("SIGKILL"); } catch { /* already gone */ } }); } } function startWaitBanner(lock, what = E2E_LOCK) { const t0 = Date.now(); process.stderr.write( `queue-lock: waiting for ${what.label} — held by ${describeHolder(lock, what.showCmd)}\n` + ` ${what.hint}\n`, ); const timer = setInterval(() => { process.stderr.write( `queue-lock: still waiting (${humanDuration(Date.now() - t0)})\n`, ); }, TICK_MS); timer.unref?.(); return () => { clearInterval(timer); process.stderr.write( `queue-lock: acquired after ${humanDuration(Date.now() - t0)}\n`, ); }; } // ------------------------------------------------------------ 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, what); let holder; try { holder = await acquire(lock, timeoutMs); } catch (err) { if (err.timeout) { process.stderr.write( `\nqueue-lock: gave up after ${humanDuration(timeoutMs)} waiting for ${lock}\n` + ` ${what.timeoutHint}\n`, ); process.exit(EXIT_TIMEOUT); } throw err; } stopBanner?.(); const holderFile = writeHolderJson(lock, name, cmd); const cleanup = () => { try { fs.rmSync(holderFile, { force: true }); } catch { /* best effort */ } try { holder.kill("SIGKILL"); } catch { /* already gone */ } }; process.on("exit", cleanup); 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= 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: rawCmd } = parseArgv(process.argv.slice(2)); if (opts.hold) return runHolder(); // `pnpm heavy -- ` 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,...] -- \n" + " queue-lock.mjs --heavy [--timeout ] -- \n", ); process.exit(2); } 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); } // Guarded so `import { withQueue }` has no side effects. if (process.argv[1] && fs.realpathSync(process.argv[1]) === fs.realpathSync(SELF)) { main().catch((err) => { process.stderr.write(`${err?.message ?? err}\n`); process.exit(1); }); }