import { spawn } from "node:child_process"; import { ENV_DENY, type Step } from "./trim"; // --------------------------------------------------------------------------- // Running a chain of steps, with the guardrails that make it safe to put behind // a button. // // WHAT WAS DELIBERATELY NEVER HERE: the whole-song build scripts. // // mk-rebuild.sh, pkmn-rebuild.sh, yoshi-rebuild.sh, rpg-remake-v5.sh and the // ms2-*.sh family were one-off shell run logs that hardcoded their paths: each // began with a literal `cd` into the job temp dir the song was mined in and // wrote to a literal reports path, ignoring SONG_DIR and SONG_REPORTS_DIR, so // no e2e fixture could contain them -- a spec that clicked a button would have // rendered into the real deliverables tree. They also ran 20+ minutes, and // arrange-poly is non-deterministic across a changed palette, so one click // could invalidate a set of judgements already made by ear. They are not in // the tree any more (release 12), and nothing here can run them. // // What runs here is the .mjs chain, which honours SONG_DATA through // song/paths.mjs and finishes in seconds to minutes. // --------------------------------------------------------------------------- export type JobState = "running" | "done" | "failed"; export type Job = { id: string; kind: string; /** * The project this job works on, as an id, when it works on one. * * Stamped by the caller rather than parsed back out of `kind`: the dashboard * links a running build to its project page, and a label is not an id. */ project: string | null; state: JobState; startedAt: number; endedAt: number | null; steps: Step[]; stepIndex: number; log: string[]; /** NDJSON progress from steps that emit it. See push(). */ events: Record[]; /** Every child pid started, so a stray one can be named after the fact. */ pids: number[]; /** Set by cancel(); the running step notices and kills its group. */ cancelling: boolean; error: string | null; }; /** Kept in memory, like lib/mix.ts: a job that does not survive a restart is a * job worth running again. * * ON globalThis, not at module scope. Next compiles each route entry -- a * page, an API route -- into its own module graph in dev, so a module-scope * Map is one Map PER ENTRY: the dashboard page read an empty registry while * /api/jobs, in another entry, held the finished build. One process, one * registry, whichever entry asks. */ type Registry = { jobs: Map; current: Job | null; seq: number }; const g = globalThis as typeof globalThis & { __umtoolJobs?: Registry }; const reg: Registry = (g.__umtoolJobs ??= { jobs: new Map(), current: null, seq: 0 }); const jobs = reg.jobs; /** Two concurrent `pick-take --apply` runs would interleave cp calls into the * same TAKE_DIR, so this is process-wide rather than per-recipe. */ export const runningJob = () => (reg.current && reg.current.state === "running" ? reg.current : null); export const getJob = (id: string) => jobs.get(id) ?? null; export const recentJobs = (n = 10) => [...jobs.values()].sort((a, b) => b.startedAt - a.startedAt).slice(0, n); /** The accidental-hour-long-job guard. None of the old .sh builds fitted under * any cap worth setting, which was the other reason they were never here. * * A step may ask for more (Step.timeoutMs). A 19-clip crossfaded report build * runs 20 to 40 minutes and would otherwise be SIGKILLed at 15 -- but raising * this for everything would remove the guard from the jobs that need it. */ const STEP_TIMEOUT_MS = 15 * 60 * 1000; /** How long a killed process group gets to go quietly before SIGKILL. */ const KILL_GRACE_MS = 5000; /** Mirrors mix's stderr clamp: enough to diagnose, not enough to blow up RAM. */ const LOG_LIMIT = 400; /** One build emits a handful of events per clip; this is generous for any of them. */ const EVENT_LIMIT = 2000; function push(job: Job, line: string, ndjson = false) { for (const l of line.split("\n")) { if (!l.trim()) continue; // A step declared as NDJSON emits one JSON object per line. They are kept // separately so the UI can render per-clip state from them, and kept OUT of // the rolling log so twenty clips of events cannot push the command that // started the job off the top of it. if (ndjson && l.startsWith("{")) { try { job.events.push(JSON.parse(l)); if (job.events.length > EVENT_LIMIT) job.events.splice(0, job.events.length - EVENT_LIMIT); continue; } catch { /* not an event after all; fall through and log it */ } } job.log.push(l); } if (job.log.length > LOG_LIMIT) job.log.splice(0, job.log.length - LOG_LIMIT); } function runStep(job: Job, step: Step): Promise { return new Promise((resolve, reject) => { push(job, `$ ${step.argv.join(" ")}`); const denied = Object.keys(step.env).filter((k) => ENV_DENY.includes(k)); if (denied.length) push(job, ` (recipe sets ${denied.join(", ")} — by design)`); const child = spawn(step.argv[0], step.argv.slice(1), { cwd: step.cwd, env: { ...process.env, ...step.env }, stdio: ["ignore", "pipe", "pipe"], // Its own process GROUP, so it can be killed as one. // // build-video.mjs shells out to yt-dlp and ffmpeg through execFile, so the // thing actually burning CPU (or holding a download open) is a GRANDCHILD. // child.kill() reaps the node process and leaves those running -- the same // failure the diarize backfill had, where killing the CLI left // diarize-sherpa.py burning four threads. detached: true, }); job.pids.push(child.pid ?? 0); const stop = (why: string) => { push(job, ` ** ${why}`); killGroup(child.pid); }; const limit = step.timeoutMs ?? STEP_TIMEOUT_MS; const timer = setTimeout( () => stop(`killed after ${Math.round(limit / 60000)} minutes`), limit, ); // A cancel that arrives mid-step is what the abort flag is for; the step // itself has no other way to hear about it. const cancelTimer = setInterval(() => { if (job.cancelling) { clearInterval(cancelTimer); stop("cancelled"); } }, 250); const done = () => { clearTimeout(timer); clearInterval(cancelTimer); }; child.stdout.on("data", (b: Buffer) => push(job, b.toString(), step.ndjson)); child.stderr.on("data", (b: Buffer) => push(job, b.toString())); child.on("error", (e) => { done(); reject(e); }); child.on("close", (code) => { done(); if (job.cancelling) reject(new Error("cancelled")); else if (code === 0) resolve(); else reject(new Error(`${step.argv[0]} exited ${code}`)); }); }); } /** * SIGTERM the process group, then SIGKILL what is left. * * The negative pid is the whole point: it addresses the GROUP, which is what * `detached: true` created and what contains the yt-dlp and ffmpeg grandchildren. */ function killGroup(pid: number | undefined) { if (!pid) return; try { process.kill(-pid, "SIGTERM"); } catch { /* already gone */ } setTimeout(() => { try { process.kill(-pid, "SIGKILL"); } catch { /* already gone */ } }, KILL_GRACE_MS); } export function startJob(kind: string, steps: Step[], { project = null }: { project?: string | null } = {}): Job { const running = runningJob(); if (running) throw new Error(`a job is already running (${running.id})`); reg.seq += 1; const job: Job = { id: `job-${reg.seq}-${process.pid}`, kind, project, state: "running", startedAt: Date.now(), endedAt: null, steps, stepIndex: 0, log: [], events: [], pids: [], cancelling: false, error: null, }; jobs.set(job.id, job); reg.current = job; void (async () => { try { for (let i = 0; i < steps.length; i += 1) { job.stepIndex = i; push(job, `--- step ${i + 1}/${steps.length}: ${steps[i].label}`); // Sequential on purpose: step 2 reads what step 1 wrote, and step 3 // reads what step 2 wrote. A non-zero exit aborts the chain rather than // running the next step against a half-made directory. await runStep(job, steps[i]); } job.state = "done"; } catch (e) { job.state = "failed"; job.error = e instanceof Error ? e.message : String(e); push(job, `** ${job.error}`); } finally { job.endedAt = Date.now(); if (reg.current === job) reg.current = null; } })(); return job; } /** * Cancel the running job. * * Safe to do at any point, and worth saying why: every artefact a report build * makes is content-addressed -- a fetched window by its window, a segment by its * clip id -- so re-running skips whatever finished. A cancelled build is a * paused one. */ export function cancelJob(id: string): boolean { const job = jobs.get(id); if (!job || job.state !== "running") return false; job.cancelling = true; push(job, "** cancel requested"); return true; } /** What the client sees. The steps are included so the command is inspectable. */ export function jobView(job: Job, since = 0, sinceEvent = 0) { return { id: job.id, kind: job.kind, project: job.project, state: job.state, startedAt: job.startedAt, endedAt: job.endedAt, stepIndex: job.stepIndex, steps: job.steps.map((s) => ({ label: s.label, argv: s.argv, cwd: s.cwd, env: s.env })), error: job.error, log: job.log.slice(Math.max(0, since)), next: job.log.length, events: job.events.slice(Math.max(0, sinceEvent)), nextEvent: job.events.length, }; }