#!/usr/bin/env tsx // Capture speaker diarization for videos whose audio is still on disk, headless. // // WHY THIS EXISTS. `runBackfillBatch` already accepts an `ids` list // (controller/operationBatch.ts), but nothing passes one — the only caller, // editor/app/channels/[slug]/backfillActions.ts, sends `kindIds` and nothing // else. So there is no way to say "diarize exactly these videos" without a // headless entrypoint, and this is the repo's established shape for corpus work // (digest-plan.ts, plans/tools/digest-bakeoff.ts, plans/tools/attribution-pilot.ts). // // WHY IT WRITES TO THE CORPUS, unlike attribution-bakeoff.ts which refuses to. // There the sidecar was a side effect of measuring a variant that had won // nothing. Here the sidecar IS the deliverable: audio is deleted once a video is // transcribed (cleanAudioFromTranscribed), and diarization is the one derived // artifact that cannot be reconstructed afterwards. Everything downstream — // attribution, naming, badges — can be redone from the JSON at any time. // // THE URGENT SET is `--scope transcribed` (the default): videos that are // transcribed AND still have audio. That is precisely the state the cleanup // sweep deletes, so those videos can lose their audio at any moment. Videos // still awaiting transcription are not in danger yet, which is why they are not // the default. // // SAFE TO KILL AND RESTART. diarizeOneVideo short-circuits on // isDiarizationFresh(existing, diarizationTarget(cfg)), and that identity is the // segmentation/embedding model BASENAMES plus the threshold — so a re-run skips // what is already done and resumes rather than redoing. Ctrl-C at hour 30 costs // one video. // // NO SETTINGS WRITE. The lane is enabled through an in-memory override, the way // attribution-pilot.ts does it, so this runs whether or not settings.json has // been armed and never depends on ambient state. The resolved config is printed. // // UNIT DISCIPLINE. Diarization is audio-bounded, so seconds-per-audio-hour IS // its honest unit — unlike the attribution/digest lanes, whose unit is the chunk // and for which plans/STATE.md records a retracted seconds-per-audio-hour // headline. Do not copy this script's rate reporting into a chunked lane. // // Examples: // tsx bin/diarize-backfill.ts --dry-run // tsx bin/diarize-backfill.ts --scope transcribed // tsx bin/diarize-backfill.ts --scope channel:destiny --limit 5 // tsx bin/diarize-backfill.ts --scope video:ObviousRises-rumble/v6z1o2g --force import os from "node:os"; import path from "node:path"; import { readdir, readFile, stat } from "node:fs/promises"; import pLimit from "p-limit"; import { getPaths, type Paths } from "../lib/paths"; import { parseFlags } from "./_parseFlags"; import { getSettings, type DiarizationSettings } from "../lib/settings"; import { diarizeOneVideo, resolveDiarizableMedia, type DiarizeOneOutcome, } from "../controller/diarizeOne"; import { isVideoTranscribed, readVideoFiles } from "../lib/videoStatus"; // Where `.diarize/README.md` says the runtime lives. Used only when settings // leave a field empty, so an armed settings.json always wins. const DEFAULT_PYTHON = ".diarize/env/bin/python"; const DEFAULT_SEG = ".diarize/models/sherpa-onnx-pyannote-segmentation-3-0/model.onnx"; const DEFAULT_EMB = ".diarize/models/titanet.onnx"; // The measured clean rate on this box: 124.5 s of wall for 766 s of audio, 4 // threads, loadavg 2.7. Used ONLY to price a run up front; the report prints the // rate this run actually achieved. const MEASURED_SECONDS_PER_AUDIO_HOUR = 585; type Candidate = { slug: string; channelSlug: string; videoDir: string; videoId: string; dir: string; transcribed: boolean; hasDiarization: boolean; media: string; // From metadata.info.json. 0 when unknown — reported as such rather than // guessed, since it is the denominator of every rate below. durationSeconds: number; }; function fmtHms(totalSeconds: number): string { const s = Math.max(0, Math.round(totalSeconds)); const h = Math.floor(s / 3600); const m = Math.floor((s % 3600) / 60); return `${h}h${String(m).padStart(2, "0")}m`; } async function readDurationSeconds(dir: string): Promise { try { const raw = await readFile(path.join(dir, "metadata.info.json"), "utf8"); const meta = JSON.parse(raw) as { duration?: unknown }; const d = Number(meta?.duration); return Number.isFinite(d) && d > 0 ? d : 0; } catch { return 0; } } // Every video dir that has something diarizable, annotated with what we need to // filter and price. Deliberately uses the SAME resolver diarizeOneVideo uses, so // "this script thinks there is audio" and "the runner found audio" cannot // disagree. async function collectCandidates( paths: Paths, onlyChannel: string | null, log: (m: string) => void, ): Promise { const out: Candidate[] = []; const channels = (await readdir(paths.channelsDir, { withFileTypes: true })) .filter((e) => e.isDirectory()) .map((e) => e.name) .filter((name) => !onlyChannel || name === onlyChannel) .sort(); let scanned = 0; for (const channelSlug of channels) { const dataDir = path.join(paths.channelsDir, channelSlug, "data"); let entries; try { entries = await readdir(dataDir, { withFileTypes: true }); } catch { continue; } for (const e of entries) { if (!e.isDirectory()) continue; scanned++; if (scanned % 20000 === 0) log(` …scanned ${scanned} video dir(s)`); const dir = path.join(dataDir, e.name); const media = await resolveDiarizableMedia(dir); if (!media) continue; const files = await readVideoFiles(dir); out.push({ slug: `${channelSlug}/${e.name}`, channelSlug, videoDir: e.name, videoId: e.name, dir, transcribed: isVideoTranscribed(files), hasDiarization: files.hasDiarization, media, durationSeconds: await readDurationSeconds(dir), }); } } return out; } function boxState(): Record { return { loadavg: os.loadavg().map((n) => Number(n.toFixed(2))), freeMemGb: Number((os.freemem() / 2 ** 30).toFixed(2)), totalMemGb: Number((os.totalmem() / 2 ** 30).toFixed(2)), }; } async function main(): Promise { const flags = parseFlags(process.argv.slice(2)); const paths = getPaths(); const scope = flags.scope ?? "transcribed"; const dryRun = flags["dry-run"] === "true"; const force = flags.force === "true"; const limitCount = flags.limit ? Number(flags.limit) : undefined; const settings = getSettings().diarization; const repo = paths.monorepoRoot; const resolve = (flagValue: string | undefined, fromSettings: string, fallback: string) => flagValue?.trim() || fromSettings.trim() || path.join(repo, fallback); // Enabled IN MEMORY. Nothing here writes settings.json — arming the lane is a // separate, explicit decision, and a backfill must not be the thing that makes // it. `python3` is treated as unset: it is the default, and sherpa ships no // wheel for this box's system 3.14. const cfg: DiarizationSettings = { ...settings, enabled: true, python: resolve( flags.python, settings.python === "python3" ? "" : settings.python, DEFAULT_PYTHON, ), segModel: resolve(flags.seg, settings.segModel, DEFAULT_SEG), embModel: resolve(flags.emb, settings.embModel, DEFAULT_EMB), threshold: flags.threshold ? Number(flags.threshold) : settings.threshold, threads: flags.threads ? Number(flags.threads) : settings.threads, }; const concurrency = Math.max( 1, flags.concurrency ? Number(flags.concurrency) : settings.concurrency, ); console.log("Diarize backfill — headless. Writes diarization.json; writes NO settings."); console.log(` engine ${paths.diarizeBin}`); console.log(` python ${cfg.python}`); console.log(` seg ${cfg.segModel}`); console.log(` emb ${cfg.embModel}`); console.log( ` threshold ${cfg.threshold} · threads ${cfg.threads} · concurrency ${concurrency}`, ); console.log(` scope ${scope}${force ? " · FORCE (redo fresh sidecars)" : ""}`); console.log(""); // FAIL FAST. A 40-hour job that discovers a bad model path on video 1 and then // returns "not-configured" 87 times is not a useful failure mode. for (const [label, p] of [ ["python", cfg.python], ["segModel", cfg.segModel], ["embModel", cfg.embModel], ["engine", paths.diarizeBin], ] as const) { try { await stat(p); } catch { console.error(`${label} does not exist: ${p}`); console.error("See .diarize/README.md for the rebuild recipe."); process.exitCode = 1; return; } } let onlyChannel: string | null = null; let onlyVideo: string | null = null; if (scope.startsWith("channel:")) onlyChannel = scope.slice("channel:".length); if (scope.startsWith("video:")) { onlyVideo = scope.slice("video:".length); onlyChannel = onlyVideo.slice(0, onlyVideo.indexOf("/")); } console.log("Scanning for diarizable media…"); const all = await collectCandidates(paths, onlyChannel, (m) => console.log(m)); // --max-audio-hours: skip files longer than this. // // MEASURED, NOT PRECAUTIONARY. A 6h12m Twitch VOD was OOM-killed by the kernel // at 11.1 GB anon-RSS (34.6 GB virtual) on this 16 GB box, after burning 35 // minutes. Memory is NOT linear in duration — a 7h42m file completed fine at // ~3.9 GB — because agglomerative clustering holds a pairwise distance matrix // over speech-segment embeddings, which is O(n^2) in SEGMENT COUNT. A dense, // fast-turnover stream produces far more segments per hour than a monologue, // so duration is only a proxy. It is the proxy we have. // // An OOM costs the whole video's engine time and puts a 16 GB box under real // pressure with other work on it, so capping is how a long run is made to // finish rather than thrash. The skipped videos are NOT lost: they keep their // audio (the cleanup guard holds it) and can be run later, on a bigger box or // with a segment-bounded clusterer. const maxAudioHours = flags["max-audio-hours"] ? Number(flags["max-audio-hours"]) : undefined; let selected = all; if (onlyVideo) selected = all.filter((c) => c.slug === onlyVideo); else if (scope === "transcribed") selected = all.filter((c) => c.transcribed); else if (scope !== "all" && !onlyChannel) { console.error(`unknown --scope "${scope}" (transcribed | all | channel: | video:/)`); process.exitCode = 1; return; } // Already-fresh videos still go through diarizeOneVideo (it is the authority on // freshness, and re-deriving that here would be a second definition), but they // are excluded from the ETA so the estimate reflects real work. const todo = force ? selected : selected.filter((c) => !c.hasDiarization); // Deferred, not dropped — and SAID so, loudly. A cap that silently shrinks the // work reads as "the corpus is done" when it is not. const deferred = maxAudioHours === undefined ? [] : todo.filter((c) => c.durationSeconds > maxAudioHours * 3600); const capped = maxAudioHours === undefined ? todo : todo.filter((c) => !deferred.includes(c)); const ordered = limitCount ? capped.slice(0, limitCount) : capped; const audioSeconds = ordered.reduce((a, c) => a + c.durationSeconds, 0); const unknownDuration = ordered.filter((c) => c.durationSeconds === 0).length; console.log(""); console.log( `${all.length} video(s) with diarizable media · ${all.filter((c) => c.transcribed).length} transcribed ` + `· ${all.filter((c) => c.hasDiarization).length} already have a sidecar`, ); console.log( `Selected by scope: ${selected.length} · to run: ${ordered.length}` + (limitCount ? ` (--limit ${limitCount})` : ""), ); if (deferred.length) { const hours = deferred.reduce((a, c) => a + c.durationSeconds, 0) / 3600; console.log( `DEFERRED by --max-audio-hours ${maxAudioHours}: ${deferred.length} video(s), ` + `${hours.toFixed(0)} audio-hour(s). They keep their audio and can be run later — ` + "the cap exists because the clusterer is O(n^2) in segment count and long dense " + "streams get OOM-killed. Longest deferred:", ); for (const c of [...deferred].sort((a, b) => b.durationSeconds - a.durationSeconds).slice(0, 5)) { console.log(` ${fmtHms(c.durationSeconds)} ${c.slug}`); } } console.log( `Audio to process: ${(audioSeconds / 3600).toFixed(0)} audio-hour(s)` + (unknownDuration ? ` (+${unknownDuration} of unknown duration)` : ""), ); console.log( `At the measured ${MEASURED_SECONDS_PER_AUDIO_HOUR} s/audio-hour (clean box, ${cfg.threads} threads): ` + `~${fmtHms((audioSeconds / 3600) * MEASURED_SECONDS_PER_AUDIO_HOUR)} single-lane.`, ); console.log(`Box: ${JSON.stringify(boxState())}`); console.log(""); if (dryRun) { for (const c of ordered.slice(0, 20)) { console.log( ` ${c.slug} ${fmtHms(c.durationSeconds)} ${c.media}${c.transcribed ? " [transcribed]" : ""}`, ); } if (ordered.length > 20) console.log(` … and ${ordered.length - 20} more`); console.log(""); console.log("--dry-run: nothing was run and nothing was written."); return; } if (ordered.length === 0) { console.log("Nothing to do."); return; } const started = Date.now(); const outcomes: Record = {}; let done = 0; let secondsDone = 0; const failures: string[] = []; const limit = pLimit(concurrency); await Promise.all( ordered.map((c) => limit(async () => { const t0 = Date.now(); let outcome: DiarizeOneOutcome; try { outcome = await diarizeOneVideo({ paths, videoDir: c.dir, videoId: c.videoId, settings: cfg, force, // The engine's own chatter goes nowhere: it is one progress line per // video here, not 87 interleaved sherpa logs. onLog: () => {}, }); } catch (err) { outcome = "failed"; failures.push(`${c.slug}: ${(err as Error)?.message ?? String(err)}`); } outcomes[outcome] = (outcomes[outcome] ?? 0) + 1; if (outcome === "failed" && !failures.some((f) => f.startsWith(c.slug))) { failures.push(`${c.slug}: see log`); } done++; secondsDone += c.durationSeconds; const elapsed = (Date.now() - started) / 1000; const rate = secondsDone > 0 ? elapsed / (secondsDone / 3600) : 0; const remaining = audioSeconds - secondsDone; console.log( `[${done}/${ordered.length}] ${c.slug} → ${outcome} ` + `(${((Date.now() - t0) / 1000).toFixed(0)}s for ${fmtHms(c.durationSeconds)} audio)` + (rate > 0 ? ` · ${rate.toFixed(0)} s/audio-hour so far · ETA ${fmtHms((remaining / 3600) * rate)}` : ""), ); }), ), ); const elapsed = (Date.now() - started) / 1000; console.log(""); console.log( `Done in ${fmtHms(elapsed)}. Outcomes: ${ Object.entries(outcomes).map(([k, v]) => `${k} ${v}`).join(", ") || "none" }`, ); if (secondsDone > 0) { console.log( `Rate: ${(elapsed / (secondsDone / 3600)).toFixed(0)} s/audio-hour ` + `over ${(secondsDone / 3600).toFixed(1)} audio-hour(s) / ${ordered.length} video(s). ` + "Audio-bounded work, so this unit is honest here — it is NOT for the chunked lanes.", ); } console.log(`Box: ${JSON.stringify(boxState())}`); if (failures.length) { console.log(""); console.log(`${failures.length} failure(s):`); for (const f of failures.slice(0, 20)) console.log(` ${f}`); } // A run in which every video was disabled/not-configured is a silent no-op // dressed as success — the exact failure mode the model-path fail-fast above // is meant to prevent, caught again here in case settings changed mid-run. if ((outcomes.diarized ?? 0) === 0 && (outcomes["already-exists"] ?? 0) === 0) { console.error("Nothing was diarized and nothing was already fresh."); process.exitCode = 1; } } main().catch((err) => { console.error(err); process.exit(1); });