commit 07705a4c73333d6c578fe002a9aa22c20ae2ce20
parent 74dff241593b9d4303d848e4b50841d3ca55d3c8
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 7 Aug 2026 16:32:09 -0400
Merge feat/diarization-capture: capture speaker diarization, and stop e2e fixtures orphaning
Two changes that arrived together.
Diarization capture. Audio is the one input in this pipeline that goes away —
cleanAudioFromTranscribed deletes it once a video is transcribed, and nothing
downstream can reconstruct it. So this ships only the perishable half:
diarization.json per video with full engine provenance, off by default, plus the
cleanup guard that stops the sweep deleting audio for a transcribed video that
has no sidecar yet. Attribution, LLM speaker naming, badges and quote filtering
are all deliberately deferred; they can be redone from the saved JSON at any
time. Measured on this box: GPU transcription runs 221 s/audio-hour over 3,602
real videos, CPU diarization 500-680, so the inline hook ships OFF and the
intended shape for a batch is capture-on, inline-off, backfill after.
E2E fixture orphaning. /api/test/invalidate-cache dropped the job registry
without cancelling first, and that registry holds the only reachable kill path
for an in-flight child — so fixture processes leaked on every run, not just
interrupted ones. A full suite once left 27 alive at ~29% of a core each, load
46, and 13 phantom failures that read like real regressions. Fixed in four
layers: cancel-before-wipe, a graceful-shutdown handler, a self-terminating
watchdog in all nine fakes, and globalSetup/globalTeardown sweeps.
Verification: 434/435 e2e in 22.6 min (from a 30.1 min baseline) with zero strays
afterwards; common 510/511; tsc clean in common/editor/export; editor build clean.
The two residual failures are pre-existing and were each confirmed by reverting:
the digest corpus census (the corpus grew to 194,053 chunks) and
jobs-batch-tasks-drain "hard Cancel during a drain" (flaky, and worse on a fast
idle box).
Diffstat:
44 files changed, 2486 insertions(+), 7 deletions(-)
diff --git a/PLAN.md b/PLAN.md
@@ -404,6 +404,24 @@ posture is set here.
### Phase 9 — Attribution (two lanes) + quote filtering
+> **The CAPTURE half landed 2026-08-07** — `diarization.json` per video, off by default,
+> plus the cleanup guard that stops the Clean-audio sweep deleting audio for a
+> transcribed-but-undiarized video. That is deliberately all of it: capture is the only
+> part with a deadline, because audio is the only input that cannot be regenerated.
+> Everything below — `attribution.json`, the text-only lane, `attributionStatus.ts`,
+> snapshot buckets, badges, upgrade jobs, quote filtering — is still to do and can be
+> built from the saved sidecars at any time. Measurements that changed the design (and
+> corrected two assumptions in the text below) are in
+> [`plans/diarization-spike-results.md`](plans/diarization-spike-results.md); the
+> decisions are in `plans/STATE.md`.
+>
+> Two corrections to what this section says: **audio deletion is not automatic** —
+> `cleanAudioFromTranscribed` has one call site, an operator-triggered action, so the
+> deadline is disk pressure, not the transcribe job. And **diarization does not run
+> "right after transcription" by default**: it is 2–3× slower than the transcription it
+> would follow (221 s/audio-hour on the GPU vs ~500–680 on the CPU), so the inline hook
+> ships off and the intended shape for a batch is capture-on, inline-off, backfill after.
+
Largest and least certain. Ship 1–8 first; keep it off by default.
**Blocking constraint:** audio is deleted once a video is transcribed, except for dirs
diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts
@@ -474,6 +474,11 @@ export async function generateChannelSnapshot(
// will actually remove.
const protectedFromCleanup = new Set<string>([...doNotCleanIds, ...keptIds]);
+ // The diarization capture lane holds audio back until its sidecar exists. Read
+ // once per snapshot rather than per video — this and the sweep's own guard key
+ // off the same setting, and they have to agree or the reclaim estimate lies.
+ const diarizationEnabled = getSettings().diarization.enabled;
+
// Videos the user has opted out of the truncated/incomplete-transcript check
// (e.g. legitimately have no speech for the back half). Suppressed from the
// incompleteTranscript + shortAudio detection buckets below.
@@ -617,7 +622,11 @@ export async function generateChannelSnapshot(
if (
files.hasWhisper &&
files.audioFiles.length > 0 &&
- !protectedFromCleanup.has(id)
+ !protectedFromCleanup.has(id) &&
+ // Held back by the diarization guard in cleanAudioFromTranscribed. Kept
+ // out of the bucket AND its byte total so "Est. reclaim" never promises
+ // space the sweep is going to refuse to take.
+ !(diarizationEnabled && !files.hasDiarization)
) {
transcribedWithAudio.push(id);
for (const name of files.audioFiles) {
diff --git a/common/controller/cleanAudioFromTranscribed.ts b/common/controller/cleanAudioFromTranscribed.ts
@@ -2,6 +2,7 @@ import path from "node:path";
import fs from "fs-extra";
import type { Paths } from "../lib/paths";
import { isDoNotClean } from "../lib/doNotClean-server";
+import { hasDiarization } from "../lib/diarization-server";
import { getSettings } from "../lib/settings";
import { audioFilesToRemove } from "../lib/videoStatus";
import { readChannelConfig } from "./channels";
@@ -34,6 +35,12 @@ export type CleanAudioResult = {
// Left alone because availability could not be resolved — no marker written,
// retried on the next sweep.
unverified: number;
+ // Left alone because diarization is enabled and this video has no
+ // diarization.json yet. Counted separately from `skipped` because it is a
+ // TRANSIENT hold, not a protection: it clears itself as the diarize lane
+ // catches up, and a number that only ever grows is the signal that the CPU
+ // lane has fallen behind the GPU.
+ awaitingDiarization: number;
// Saved-store containers evicted because they rolled out of the keep-latest
// window (the retention prune runs alongside the audio sweep).
prunedSavedVideos: number;
@@ -64,6 +71,7 @@ export async function cleanAudioFromTranscribed({
skipped: 0,
pinned: 0,
unverified: 0,
+ awaitingDiarization: 0,
prunedSavedVideos: 0,
prunedBytes: 0,
};
@@ -82,6 +90,16 @@ export async function cleanAudioFromTranscribed({
let cleanedDirs = 0;
let removedFiles = 0;
let protectedCount = 0;
+ let awaitingDiarization = 0;
+
+ // While diarization is enabled, a transcribed video that still has audio but
+ // no diarization.json is NOT cleanable yet — its diarize job hasn't run.
+ //
+ // This guard is the whole reason deferring diarization to an async job is
+ // safe. Without it the sweep would delete the audio out from under a queued
+ // diarize job, silently and permanently: audio is the one input in this
+ // pipeline that cannot be regenerated, and the failure would leave no trace.
+ const diarization = getSettings().diarization;
// --- Discover -------------------------------------------------------------
// Nothing is deleted in this pass. Protected videos are filtered out here, so
@@ -104,6 +122,11 @@ export async function cleanAudioFromTranscribed({
protectedCount++;
continue;
}
+ if (diarization.enabled && !(await hasDiarization(videoDir))) {
+ log(`Skipped ${id} (awaiting diarization)`);
+ awaitingDiarization++;
+ continue;
+ }
candidates.push({ id, videoDir, audioFiles });
}
@@ -139,15 +162,21 @@ export async function cleanAudioFromTranscribed({
cleanedDirs++;
}
- const skipped = protectedCount + pinned + unverified;
+ const skipped = protectedCount + pinned + unverified + awaitingDiarization;
// Keep the historical wording when the gate had nothing to report, and break
// the total down only when it would otherwise be misleading.
+ // The diarization clause is appended only when it is non-zero, so the message
+ // stays byte-identical for every install that doesn't use the capture lane —
+ // a "0 awaiting diarization" on a box with diarization off is noise, and this
+ // exact string is asserted in pre-clean-availability.spec.
+ const diarNote =
+ awaitingDiarization > 0 ? `, ${awaitingDiarization} awaiting diarization` : "";
const skippedNote =
skipped === 0
? ""
- : pinned === 0 && unverified === 0
+ : pinned === 0 && unverified === 0 && awaitingDiarization === 0
? ` Skipped ${skipped} (do not clean).`
- : ` Skipped ${skipped} (${protectedCount} protected, ${pinned} gone-from-source pinned, ${unverified} unverified).`;
+ : ` Skipped ${skipped} (${protectedCount} protected, ${pinned} gone-from-source pinned, ${unverified} unverified${diarNote}).`;
log(
`Cleaned ${removedFiles} audio file(s) from ${cleanedDirs} of ${dirs.length} video dir(s).${skippedNote}`,
);
@@ -163,6 +192,7 @@ export async function cleanAudioFromTranscribed({
skipped,
pinned,
unverified,
+ awaitingDiarization,
prunedSavedVideos: prune.pruned,
prunedBytes: prune.bytesFreed,
};
diff --git a/common/controller/diarizeAll.ts b/common/controller/diarizeAll.ts
@@ -0,0 +1,142 @@
+// Corpus-wide diarization backfill: walk every channel/video that still has
+// audio and no diarization.json, and capture its speaker turns.
+//
+// This is what catches everything the post-transcribe hook cannot — videos
+// transcribed before the feature landed, videos transcribed with the inline hook
+// off (the recommended way to run a large batch), and anything whose inline
+// attempt failed. Modelled on normalizeAllTranscripts: pLimit over the channel
+// walk, with diarizeOneVideo short-circuiting on an existing sidecar so a
+// re-run does no rework.
+//
+// Concurrency defaults to settings.diarization.concurrency (1), NOT the 8 that
+// normalizeAll uses: normalizing cues is I/O, diarization is CPU-bound and
+// competes with GPU feeding and the digest sweep for the same 8 threads.
+
+import path from "node:path";
+import { readdir } from "node:fs/promises";
+import pLimit from "p-limit";
+import { listChannelStatsFromDisk } from "./channels";
+import { diarizeOneVideo } from "./diarizeOne";
+import { getSettings } from "../lib/settings";
+import type { Paths } from "../lib/paths";
+
+export type DiarizeAllOptions = {
+ paths: Paths;
+ // Restrict to these channel slugs. Empty/omitted = the whole corpus.
+ channelSlugs?: string[];
+ // Only diarize videos that already have a transcript. ON by default: the
+ // capture lane exists to pair speaker turns with a transcript, and an
+ // untranscribed video's audio is not at risk from the cleanup sweep yet.
+ transcribedOnly?: boolean;
+ force?: boolean;
+ concurrency?: number;
+ onLog?: (msg: string) => void;
+ signal?: AbortSignal;
+};
+
+export type DiarizeAllResult = {
+ diarized: number;
+ alreadyExists: number;
+ noAudio: number;
+ skipped: number;
+ failed: number;
+};
+
+export async function diarizeAll(
+ opts: DiarizeAllOptions,
+): Promise<DiarizeAllResult> {
+ const log = opts.onLog ?? ((m: string) => console.log(m));
+ const settings = getSettings().diarization;
+ const result: DiarizeAllResult = {
+ diarized: 0,
+ alreadyExists: 0,
+ noAudio: 0,
+ skipped: 0,
+ failed: 0,
+ };
+
+ if (!settings.enabled) {
+ log("Diarization is disabled in settings — nothing to do.");
+ return result;
+ }
+ if (!settings.segModel || !settings.embModel) {
+ log(
+ "Diarization has no segmentation/embedding model configured — nothing to do.",
+ );
+ return result;
+ }
+
+ const limit = pLimit(opts.concurrency ?? settings.concurrency);
+ const wanted = new Set(opts.channelSlugs ?? []);
+ const channels = (await listChannelStatsFromDisk(opts.paths)).filter(
+ (ch) => wanted.size === 0 || wanted.has(ch.slug),
+ );
+ const transcribedOnly = opts.transcribedOnly !== false;
+
+ for (const ch of channels) {
+ if (opts.signal?.aborted) break;
+ const dataDir = path.join(opts.paths.channelsDir, ch.slug, "data");
+ const videoIds = await readdir(dataDir).catch(() => [] as string[]);
+ let diarized = 0;
+ let already = 0;
+ let noAudio = 0;
+ let skipped = 0;
+ let failed = 0;
+
+ await Promise.all(
+ videoIds.map((id) =>
+ limit(async () => {
+ if (opts.signal?.aborted) return;
+ const videoDir = path.join(dataDir, id);
+ if (transcribedOnly) {
+ const entries = await readdir(videoDir).catch(
+ () => [] as string[],
+ );
+ if (!entries.includes("transcript.json")) {
+ skipped++;
+ return;
+ }
+ }
+ try {
+ const outcome = await diarizeOneVideo({
+ paths: opts.paths,
+ videoDir,
+ videoId: id,
+ settings,
+ force: opts.force,
+ onLog: log,
+ signal: opts.signal,
+ });
+ if (outcome === "diarized") diarized++;
+ else if (outcome === "already-exists") already++;
+ else if (outcome === "no-audio") noAudio++;
+ else if (outcome === "failed") failed++;
+ else skipped++;
+ } catch (err) {
+ failed++;
+ log(` ! ${ch.slug}/${id}: ${(err as Error).message}`);
+ }
+ }),
+ ),
+ );
+
+ // Only report channels that had something to do — the corpus is 66 channels
+ // and most have no retained audio at all.
+ if (diarized + already + failed > 0) {
+ log(
+ ` ${ch.slug}: diarized=${diarized} already=${already} no-audio=${noAudio} skipped=${skipped} failed=${failed}`,
+ );
+ }
+ result.diarized += diarized;
+ result.alreadyExists += already;
+ result.noAudio += noAudio;
+ result.skipped += skipped;
+ result.failed += failed;
+ }
+
+ log(
+ `Done. diarized=${result.diarized} already=${result.alreadyExists} ` +
+ `no-audio=${result.noAudio} skipped=${result.skipped} failed=${result.failed}`,
+ );
+ return result;
+}
diff --git a/common/controller/diarizeOne.ts b/common/controller/diarizeOne.ts
@@ -0,0 +1,146 @@
+// Run speaker diarization for ONE video and write its diarization.json sidecar.
+//
+// The whole point of this lane is timing, not quality: audio is the one input
+// that goes away (cleanAudioFromTranscribed deletes it once a video is
+// transcribed), so the speaker ranges have to be captured while it is still on
+// disk. Everything downstream — attribution, LLM speaker naming, viewer badges,
+// quote filtering — can be redone from the saved JSON at any time.
+//
+// Never throws for an expected condition. Callers run this from a post-transcribe
+// hook and from a corpus-wide backfill, and in both places a diarization problem
+// must degrade to a skip, never fail the surrounding work.
+
+import path from "node:path";
+import fs from "fs-extra";
+import { execa } from "execa";
+import type { Paths } from "../lib/paths";
+import { getSettings } from "../lib/settings";
+import type { DiarizationSettings } from "../lib/settings";
+import { DIARIZATION_FILENAME } from "../lib/diarization";
+import { hasDiarization } from "../lib/diarization-server";
+import { findSourceMedia, isRealAudioFile } from "../lib/videoStatus";
+import { resolveSavedVideo } from "../lib/savedVideo-server";
+
+const { readdir } = fs;
+
+const AUDIO_PREFERENCE = ["audio.mp3", "audio.m4a", "audio.opus"];
+
+export type DiarizeOneOptions = {
+ paths: Paths;
+ videoDir: string;
+ videoId: string;
+ // Overrides settings.diarization. The backfill passes one resolved copy rather
+ // than re-reading settings per video.
+ settings?: DiarizationSettings;
+ // Redo even when a valid sidecar already exists (an engine/threshold change).
+ force?: boolean;
+ onLog?: (msg: string) => void;
+ signal?: AbortSignal;
+};
+
+export type DiarizeOneOutcome =
+ // Wrote (or rewrote) diarization.json.
+ | "diarized"
+ // A valid sidecar already existed and force was not set.
+ | "already-exists"
+ // The lane is switched off in settings.
+ | "disabled"
+ // Nothing to diarize from: no audio and no persisted source container.
+ | "no-audio"
+ // Configured but unusable — no model paths. Reported distinctly from a
+ // failure so an operator can tell "not set up" from "tried and broke".
+ | "not-configured"
+ // The engine ran and failed. Logged, never thrown.
+ | "failed";
+
+// Resolve something diarizable in the video dir: extracted audio first, then a
+// persisted source container, then the saved-video store. Mirrors
+// transcribeOne's resolver — ffmpeg reads any of them.
+export async function resolveDiarizableMedia(
+ videoDir: string,
+): Promise<string | null> {
+ const entries = await readdir(videoDir).catch(() => [] as string[]);
+ const candidates = entries.filter(isRealAudioFile);
+ if (candidates.length > 0) {
+ for (const preferred of AUDIO_PREFERENCE) {
+ if (candidates.includes(preferred)) return preferred;
+ }
+ return [...candidates].sort()[0];
+ }
+ const container = findSourceMedia(entries);
+ if (container) return container;
+ const saved = await resolveSavedVideo(videoDir);
+ if (saved) return path.relative(videoDir, saved);
+ return null;
+}
+
+export async function diarizeOneVideo(
+ opts: DiarizeOneOptions,
+): Promise<DiarizeOneOutcome> {
+ const log = opts.onLog ?? ((m: string) => console.log(m));
+ const cfg = opts.settings ?? getSettings().diarization;
+
+ if (!cfg.enabled) return "disabled";
+ if (!cfg.segModel || !cfg.embModel) {
+ log(
+ `Diarize ${opts.videoId} skipped: no segmentation/embedding model configured.`,
+ );
+ return "not-configured";
+ }
+ if (!opts.force && (await hasDiarization(opts.videoDir))) {
+ return "already-exists";
+ }
+
+ const media = await resolveDiarizableMedia(opts.videoDir);
+ if (!media) return "no-audio";
+
+ const start = Date.now();
+ log(`Diarize ${opts.videoId} start (${media})`);
+ try {
+ const child = execa(
+ opts.paths.diarizeBin,
+ [
+ "--output",
+ DIARIZATION_FILENAME,
+ "--video-id",
+ opts.videoId,
+ "--seg",
+ cfg.segModel,
+ "--emb",
+ cfg.embModel,
+ "--threshold",
+ String(cfg.threshold),
+ "--threads",
+ String(cfg.threads),
+ "--python",
+ cfg.python,
+ media,
+ ],
+ {
+ cwd: opts.videoDir,
+ cancelSignal: opts.signal,
+ all: true,
+ buffer: false,
+ },
+ );
+ child.all?.on("data", (c: Buffer) => log(c.toString("utf8").trimEnd()));
+ await child;
+ } catch (err) {
+ log(`Diarize ${opts.videoId} failed: ${(err as Error).message}`);
+ return "failed";
+ }
+
+ // Trust the file, not the exit code: an engine that exits 0 without producing
+ // a parseable sidecar has NOT diarized this video, and saying otherwise is
+ // what would let the cleanup sweep delete the audio.
+ if (!(await hasDiarization(opts.videoDir))) {
+ log(
+ `Diarize ${opts.videoId} produced no valid ${DIARIZATION_FILENAME} — treating as failed.`,
+ );
+ return "failed";
+ }
+ log(
+ `Diarize ${opts.videoId} done in ${((Date.now() - start) / 1000).toFixed(1)}s`,
+ );
+ return "diarized";
+}
diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts
@@ -8,6 +8,8 @@ import { detectTranscriptFormat } from "../lib/whisper";
import { getWorkerPool } from "../jobs/workerPool";
import type { TaskTracker } from "../jobs/taskHooks";
import { normalizeTranscript } from "./normalizeTranscript";
+import { diarizeOneVideo } from "./diarizeOne";
+import { getSettings } from "../lib/settings";
import { pingRemoteHealth, transcribeViaRemote } from "./remoteTranscribe";
import { TranscribeError } from "./transcribeError";
import { findSourceMedia, isRealAudioFile } from "../lib/videoStatus";
@@ -252,6 +254,42 @@ export async function transcribeOneVideo(
}
}
+ // Speaker diarization, while the audio is still on disk. THIS is the only
+ // moment it can happen for free: cleanAudioFromTranscribed deletes the audio
+ // once a video is transcribed, and nothing downstream can reconstruct it.
+ //
+ // Off by default and deliberately so — diarization is slower than the
+ // transcription it follows, so the batch workflow leaves this off and uses the
+ // backfill pass instead (see DiarizationSettings.inlineAfterTranscribe). When
+ // it IS on, it stays non-fatal: the established idiom here is a try/catch that
+ // logs and continues, and a diarization problem must never fail a
+ // transcription that already succeeded.
+ //
+ // The "paused" early returns above correctly bypass this — there is no
+ // transcript yet, so there is nothing to diarize alongside.
+ const diarization = getSettings().diarization;
+ if (diarization.enabled && diarization.inlineAfterTranscribe) {
+ try {
+ const outcome = await diarizeOneVideo({
+ paths: opts.paths,
+ videoDir: opts.videoDir,
+ videoId: opts.videoId,
+ settings: diarization,
+ onLog: log,
+ signal: opts.signal,
+ });
+ if (outcome === "failed") {
+ log(
+ `Warning: diarization failed for ${opts.videoId}; audio is retained for a later pass.`,
+ );
+ }
+ } catch (err) {
+ log(
+ `Warning: diarization threw for ${opts.videoId}: ${(err as Error).message}`,
+ );
+ }
+ }
+
return "transcribed";
}
diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts
@@ -186,6 +186,20 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
bookmarkable: true,
queueKeyStrategy: "custom",
},
+ // Speaker-diarization backfill over a channel's retained audio. The capture
+ // lane's catch-all: it picks up everything the post-transcribe hook missed
+ // (videos transcribed before the feature, or with the inline hook off — which
+ // is the recommended way to run a large batch). Drainable, because it is
+ // CPU-hours of work an operator will want to stop without losing what it has
+ // already written, and bookmarkable, because it re-derives its work-list from
+ // disk on every run and so replays correctly with nothing remembered.
+ "diarize-channel": {
+ kind: "diarize-channel",
+ label: "Diarize speakers",
+ drainable: true,
+ bookmarkable: true,
+ queueKeyStrategy: "custom",
+ },
"check-kept-deleted": {
kind: "check-kept-deleted",
label: "Check kept videos",
diff --git a/common/jobs/shutdownCancel.ts b/common/jobs/shutdownCancel.ts
@@ -0,0 +1,63 @@
+// SIGTERM every child we still own when the server shuts down gracefully.
+//
+// WHY THIS IS ITS OWN MODULE, and not inline in editor/instrumentation.ts:
+// `register()` runs in every runtime, and Next statically analyses that file for
+// the EDGE bundle. Node APIs there (`process.once`, `process.kill`,
+// `process.listenerCount`) are flagged as unsupported and the file fails to
+// compile — the `NEXT_RUNTIME !== "nodejs"` early return happens at runtime, and
+// the bundler's static scan does not know about it. Keeping the Node-only code
+// behind a lazy `await import()` inside that guard is the idiom instrumentation.ts
+// already uses for the heartbeat, worker pool and digest sweep.
+//
+// WHAT IT FIXES: a spawned yt-dlp/whisper/parakeet child otherwise outlives the
+// server. Route handlers run in a Next-spawned worker, so when that worker dies
+// the child reparents and keeps burning CPU. The registry holds the only handle
+// that can reach it — runManagedFunction sets `record.abortController` but never
+// `record.child` (controller/../jobs/streamCommand.ts), so `cancel()`'s
+// `child.kill()` branch is a no-op for those jobs and `abort()` is the only kill
+// path. That handle dies with the process unless we use it first.
+//
+// HONEST LIMIT: graceful shutdown only. A SIGKILLed server still orphans its
+// children because nothing in-process gets to run — that case is covered by the
+// fixtures' own watchdog (editor/e2e/fixtures/bin/_watchdog.mjs) in tests, and is
+// simply unavoidable for a real `kill -9` in production.
+
+// Set once the handlers are armed so a second register() — or a second signal —
+// cannot re-enter. register() is documented as running once per server instance,
+// but this is cheap and keeps shutdown re-entrancy-free.
+let shutdownArmed = false;
+
+export function armShutdownCancel(): void {
+ if (shutdownArmed) return;
+ shutdownArmed = true;
+ for (const signal of ["SIGTERM", "SIGINT"] as const) {
+ process.once(signal, () => {
+ try {
+ // Read the global directly: getRegistry() CONSTRUCTS on read, and
+ // building a registry during shutdown just to cancel nothing would be
+ // worse than useless. No registry means nothing ever ran.
+ const registry = globalThis.__yttJobRegistry__;
+ for (const job of registry?.list() ?? []) {
+ if (job.status !== "running" && job.status !== "queued") continue;
+ try {
+ registry?.cancel(job.id);
+ } catch {
+ /* one wedged job must not block the rest, or the exit */
+ }
+ }
+ } catch {
+ /* shutdown must proceed even if the registry is in a bad state */
+ }
+ // No await, no delay: abort() fires execa's cancelSignal inline, so the
+ // SIGTERMs are already delivered by the time we get here. This runs on
+ // every ordinary dev restart, so it must not add latency.
+ //
+ // Re-raise only if nobody else is listening. `once` has already removed
+ // this listener, so a remaining count means another handler (Next's own)
+ // owns the exit — re-raising then would run that handler a second time.
+ if (process.listenerCount(signal) === 0) {
+ process.kill(process.pid, signal);
+ }
+ });
+ }
+}
diff --git a/common/lib/diarization-server.ts b/common/lib/diarization-server.ts
@@ -0,0 +1,47 @@
+import path from "node:path";
+import { readFile, rename, writeFile } from "node:fs/promises";
+import {
+ DIARIZATION_FILENAME,
+ type DiarizationRecord,
+} from "./diarization";
+
+export function diarizationPath(videoDir: string): string {
+ return path.join(videoDir, DIARIZATION_FILENAME);
+}
+
+export async function loadDiarization(
+ videoDir: string,
+): Promise<DiarizationRecord | null> {
+ try {
+ const raw = await readFile(diarizationPath(videoDir), "utf8");
+ const parsed = JSON.parse(raw) as Partial<DiarizationRecord>;
+ if (
+ typeof parsed?.videoId === "string" &&
+ typeof parsed.generatedAt === "string" &&
+ Array.isArray(parsed.turns)
+ ) {
+ return parsed as DiarizationRecord;
+ }
+ return null;
+ } catch {
+ return null;
+ }
+}
+
+// Presence of a VALID sidecar. This is the predicate the cleanup guard keys
+// off, so it deliberately treats a malformed file as absent: a half-written
+// diarization.json must not be what convinces the sweep it is safe to delete
+// the only copy of the audio.
+export async function hasDiarization(videoDir: string): Promise<boolean> {
+ return (await loadDiarization(videoDir)) !== null;
+}
+
+export async function writeDiarization(
+ videoDir: string,
+ record: DiarizationRecord,
+): Promise<void> {
+ const file = diarizationPath(videoDir);
+ const tmp = `${file}.tmp-${process.pid}`;
+ await writeFile(tmp, JSON.stringify(record) + "\n");
+ await rename(tmp, file);
+}
diff --git a/common/lib/diarization.test.ts b/common/lib/diarization.test.ts
@@ -0,0 +1,124 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import path from "node:path";
+import os from "node:os";
+import { mkdtemp, mkdir, writeFile, rm } from "node:fs/promises";
+import {
+ DIARIZATION_FILENAME,
+ diarizationSpeechSeconds,
+ type DiarizationRecord,
+} from "./diarization";
+import {
+ hasDiarization,
+ loadDiarization,
+ writeDiarization,
+} from "./diarization-server";
+import { readVideoFiles, DIARIZATION_FILENAME as STATUS_FILENAME } from "./videoStatus";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test common/lib/diarization.test.ts
+
+async function tmpDir() {
+ return mkdtemp(path.join(os.tmpdir(), "diarization-test-"));
+}
+
+function record(over: Partial<DiarizationRecord> = {}): DiarizationRecord {
+ return {
+ videoId: "vid1",
+ generatedAt: "2026-08-07T00:00:00.000Z",
+ speakers: 2,
+ turns: [
+ { start: 0, end: 4, speaker: 0 },
+ { start: 4, end: 10, speaker: 1 },
+ ],
+ engine: { engine: "sherpa-onnx", threshold: 0.5 },
+ ...over,
+ };
+}
+
+test("the sidecar filename is not claimed by SUB_FILE_RE as a subtitle track", () => {
+ // The regex is /^transcript\.([^.]+)\.([^.]+)$/. A name like
+ // transcript.diarization.json would be indexed as a `diarization`-language
+ // subtitle track — which is exactly the trap this constant exists to avoid.
+ assert.equal(DIARIZATION_FILENAME, "diarization.json");
+ assert.equal(STATUS_FILENAME, DIARIZATION_FILENAME);
+ assert.ok(!/^transcript\.([^.]+)\.([^.]+)$/.test(DIARIZATION_FILENAME));
+});
+
+test("write then load round-trips a record", async () => {
+ const dir = await tmpDir();
+ try {
+ await writeDiarization(dir, record());
+ const loaded = await loadDiarization(dir);
+ assert.ok(loaded);
+ assert.equal(loaded.videoId, "vid1");
+ assert.equal(loaded.turns.length, 2);
+ assert.equal(loaded.engine.engine, "sherpa-onnx");
+ assert.equal(await hasDiarization(dir), true);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+});
+
+test("a missing sidecar reads as absent, not as an error", async () => {
+ const dir = await tmpDir();
+ try {
+ assert.equal(await loadDiarization(dir), null);
+ assert.equal(await hasDiarization(dir), false);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+});
+
+// THE ONE THAT MATTERS. hasDiarization gates deletion of the audio, and audio
+// cannot be regenerated. A truncated or malformed file must read as "not
+// diarized" so the sweep keeps the audio — the failure mode to avoid is a
+// half-written file convincing cleanup that the work is done.
+test("a malformed or truncated sidecar reads as absent so the audio is kept", async () => {
+ const dir = await tmpDir();
+ try {
+ await writeFile(path.join(dir, DIARIZATION_FILENAME), "{ not json");
+ assert.equal(await hasDiarization(dir), false);
+
+ // Parseable JSON, but missing the fields that make it a record.
+ await writeFile(path.join(dir, DIARIZATION_FILENAME), '{"videoId":"vid1"}');
+ assert.equal(await hasDiarization(dir), false);
+
+ // Parseable, has ids, but no turns array.
+ await writeFile(
+ path.join(dir, DIARIZATION_FILENAME),
+ '{"videoId":"vid1","generatedAt":"2026-08-07T00:00:00.000Z"}',
+ );
+ assert.equal(await hasDiarization(dir), false);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+});
+
+test("readVideoFiles reports the sidecar's presence", async () => {
+ const dir = await tmpDir();
+ try {
+ await mkdir(dir, { recursive: true });
+ let files = await readVideoFiles(dir);
+ assert.equal(files.hasDiarization, false);
+
+ await writeDiarization(dir, record());
+ files = await readVideoFiles(dir);
+ assert.equal(files.hasDiarization, true);
+ } finally {
+ await rm(dir, { recursive: true, force: true });
+ }
+});
+
+test("speech seconds sums turns, and overlapping turns are not clamped", () => {
+ assert.equal(diarizationSpeechSeconds(record()), 10);
+ // Two people talking at once is real; reporting 12s of speech over 10s of
+ // audio is the honest answer, not a bug to hide.
+ const overlapping = record({
+ turns: [
+ { start: 0, end: 10, speaker: 0 },
+ { start: 8, end: 10, speaker: 1 },
+ ],
+ });
+ assert.equal(diarizationSpeechSeconds(overlapping), 12);
+});
diff --git a/common/lib/diarization.ts b/common/lib/diarization.ts
@@ -0,0 +1,68 @@
+// Client-safe types and constants for the per-video speaker-diarization
+// sidecar. Mirrors transcribeOutcome.ts: server-only I/O lives in
+// diarization-server.ts.
+//
+// WHY A SIDECAR AND NOT A FIELD ON `Cue`: speaker ranges are deliberately kept
+// out of the cue stream. Putting them inside cues would bump SCHEMA_VERSION
+// (controller/buildIndex.ts), invalidate the IndexedDB transcript cache
+// (components/transcriptStore.ts), force a re-emit of every transcripts page
+// (the SHA-1 skip in buildIndex.ts would miss everywhere), push against the
+// 25 MB archive cap, and inflate the digest token budget (lib/digestPrompt.ts
+// assumes ~10 tokens/cue). A sidecar costs none of that, and attribution can be
+// recomputed from it at any time.
+//
+// WHY THE FILENAME MATTERS: never name a sidecar `transcript.<x>.<y>` —
+// SUB_FILE_RE in lib/videoStatus.ts claims any such file as a subtitle track.
+// The constant lives beside the other video-dir filenames in videoStatus.ts.
+
+// One contiguous stretch of audio attributed to a single (anonymous) speaker
+// cluster. Times are seconds from the start of the audio. `speaker` is a
+// cluster INDEX, not an identity — naming clusters is a later, redoable pass.
+export type DiarizationTurn = {
+ start: number;
+ end: number;
+ speaker: number;
+};
+
+// What produced a given diarization.json. Recorded so a later attribution pass
+// can tell whether a file is worth re-running: swapping the engine, the
+// segmentation model, or the clustering threshold all change the output, and
+// only an explicit record makes that legible after the audio is gone.
+export type DiarizationEngine = {
+ // Wrapper/engine id, e.g. "sherpa-onnx".
+ engine: string;
+ // Model identifiers, free-form per engine (basename is enough — the full path
+ // is machine-specific and would make records non-portable across shards).
+ segmentationModel?: string;
+ embeddingModel?: string;
+ // Engine/library version string, when the engine reports one.
+ version?: string;
+ // Clustering threshold used. The single most consequential knob: it decides
+ // how many speakers come out, and re-running with a different one is the most
+ // likely reason to regenerate.
+ threshold?: number;
+};
+
+export type DiarizationRecord = {
+ videoId: string;
+ // ISO 8601, set when the sidecar is finalized.
+ generatedAt: string;
+ // Audio duration in seconds as the engine saw it.
+ audioSeconds?: number;
+ // Wall-clock diarization time in milliseconds.
+ durationMs?: number;
+ // Distinct speaker clusters found. Denormalized from `turns` so callers can
+ // bucket/report without walking every turn.
+ speakers: number;
+ turns: DiarizationTurn[];
+ engine: DiarizationEngine;
+};
+
+export const DIARIZATION_FILENAME = "diarization.json";
+
+// Total attributed speech, in seconds. Turns may overlap (two people talking at
+// once), so this can exceed the audio duration — that is a real signal, not a
+// bug, and callers that want wall-time coverage should merge ranges first.
+export function diarizationSpeechSeconds(record: DiarizationRecord): number {
+ return record.turns.reduce((sum, t) => sum + Math.max(0, t.end - t.start), 0);
+}
diff --git a/common/lib/paths.ts b/common/lib/paths.ts
@@ -125,6 +125,10 @@ export type Paths = {
parakeetBin: string;
parakeetCliBin: string;
parakeetModel: string;
+ // Speaker-diarization wrapper script (scripts/diarize.mjs), invoked the way
+ // parakeetBin is. Overridable via DIARIZE_BIN — which is also how the e2e
+ // suite swaps in a fake diarizer.
+ diarizeBin: string;
// Base URL of the local ollama server, the local-GPU digest lane
// (common/lib/digestApps.ts posts to `${ollamaUrl}/api/chat`). The FIRST
// URL-valued entry in Paths, so it is normalized here the way
@@ -225,6 +229,12 @@ export function getPaths(): Paths {
path.join(monorepoRoot, "scripts", "parakeet-stitch.mjs"),
parakeetCliBin: process.env.PARAKEET_CLI ?? "parakeet-cli",
parakeetModel: process.env.PARAKEET_MODEL ?? "",
+ // Speaker-diarization wrapper, same shape as parakeetBin: a script we own,
+ // overridable by env so the engine behind it can be swapped (sherpa-onnx
+ // today, pyannote later) without touching any caller.
+ diarizeBin:
+ process.env.DIARIZE_BIN ??
+ path.join(monorepoRoot, "scripts", "diarize.mjs"),
ollamaUrl: (process.env.OLLAMA_URL ?? "http://127.0.0.1:11434").replace(
/\/+$/,
"",
diff --git a/common/lib/settings.ts b/common/lib/settings.ts
@@ -199,6 +199,69 @@ export type SiteSettings = {
// AI digest generation (chapters + topic tags over the existing transcripts).
// Local-first: the metered lane is off by default. See DigestSettings.
digest: DigestSettings;
+ // Speaker diarization captured right after transcription, while the audio is
+ // still on disk. OFF by default. See DiarizationSettings.
+ diarization: DiarizationSettings;
+};
+
+// Configuration for the speaker-diarization capture lane.
+//
+// This exists because AUDIO IS THE ONE PERISHABLE INPUT in the pipeline:
+// cleanAudioFromTranscribed deletes it once a video is transcribed, so
+// diarization has to happen while the audio is still there or not at all. The
+// capture half is deliberately all that ships here — attribution, LLM speaker
+// naming, viewer badges and quote filtering can all be redone later from the
+// saved JSON, whereas the audio cannot.
+export type DiarizationSettings = {
+ // Master switch. OFF by default so a transcription batch can start before this
+ // lands, with diarization backfilled over the retained audio afterwards.
+ //
+ // Turning it ON also arms the cleanup guard: the Clean-audio sweep stops
+ // deleting audio for a transcribed video that has no diarization.json yet.
+ // That is the point — it is what keeps the perishable input alive long enough
+ // to be captured — but it means enabling this holds disk.
+ enabled: boolean;
+ // Run diarization inline in the post-transcribe hook.
+ //
+ // OFF by default, and that default is a MEASURED decision, not caution.
+ // Measured on this box: GPU transcription runs at 221 s/audio-hour (16.3x
+ // realtime, over 3,602 real videos), CPU diarization at ~500-680 s/audio-hour.
+ // Diarization is therefore ~2-3x SLOWER than the transcription it follows, so
+ // running it inline drops whole-pipeline throughput by roughly 3-4x and leaves
+ // the GPU idle while the CPU catches up.
+ //
+ // The intended sequence for a large batch is the opposite: leave this off, let
+ // the batch transcribe at full GPU speed with `enabled` holding the audio, and
+ // diarize afterwards with the backfill pass. Turn it on for steady state, once
+ // the arrival rate is a few videos a day rather than a corpus.
+ inlineAfterTranscribe: boolean;
+ // Clustering threshold — the single most consequential knob, since it decides
+ // how many speakers come out. Larger merges more aggressively.
+ //
+ // The default is 0.9, NOT sherpa-onnx's own 0.5, and that is measured on this
+ // corpus. On a 6-minute excerpt of a two-person interview (known ground truth:
+ // 2 speakers), sherpa's default produced 22 clusters; 0.9 produced 6, with the
+ // top two at 40%/40% of talk time — recognizably the two hosts. Sweep on the
+ // same clip: 0.4→23, 0.5→22, 0.6→17, 0.7→12, 0.8→10, 0.9→6.
+ //
+ // It still over-splits, which is why this is a capture lane and not an answer:
+ // the turns are recorded with the threshold that produced them, so a later
+ // attribution pass can re-cluster or re-run without needing the audio back.
+ threshold: number;
+ // Engine threads per diarize run.
+ threads: number;
+ // Python interpreter for the default sherpa-onnx engine. sherpa-onnx ships
+ // wheels only up to cp313, and this box's system python is 3.14 — so this
+ // usually points at a dedicated venv rather than `python3`.
+ python: string;
+ // ONNX model paths for the default engine. Empty = the lane cannot run, which
+ // is reported as a skip rather than a failure.
+ segModel: string;
+ embModel: string;
+ // How many diarize runs may execute at once in the backfill pass. Kept low by
+ // default: diarization is CPU-bound and competes with GPU feeding and the
+ // digest sweep for the same 8 threads.
+ concurrency: number;
};
// Configuration for the derived-corpus digest layer. Local-first by decision:
@@ -812,6 +875,49 @@ function defaults(): SiteSettings {
// Must be listed here or the allowlist loop in getSettings() drops the key
// entirely and the whole section is never read from disk.
digest: defaultDigest(),
+ diarization: defaultDiarization(),
+ };
+}
+
+export function defaultDiarization(): DiarizationSettings {
+ return {
+ // OFF. Capture is opt-in: turning it on makes the cleanup sweep start
+ // refusing to delete audio for transcribed-but-undiarized videos, which is
+ // correct but is a disk-pressure decision an operator should make.
+ enabled: false,
+ // OFF — see DiarizationSettings.inlineAfterTranscribe. Diarization is slower
+ // than the transcription it would follow, so inline is the exception.
+ inlineAfterTranscribe: false,
+ // Measured, not sherpa-onnx's 0.5 — see DiarizationSettings.threshold.
+ threshold: 0.9,
+ threads: 4,
+ python: "python3",
+ segModel: "",
+ embModel: "",
+ concurrency: 1,
+ };
+}
+
+export function sanitizeDiarization(value: unknown): DiarizationSettings {
+ const d = defaultDiarization();
+ if (!value || typeof value !== "object") return d;
+ const r = value as Record<string, unknown>;
+ const str = (v: unknown, fallback: string) =>
+ typeof v === "string" && v.trim() ? v.trim() : fallback;
+ return {
+ enabled: r.enabled === true,
+ inlineAfterTranscribe: r.inlineAfterTranscribe === true,
+ threshold:
+ typeof r.threshold === "number" &&
+ Number.isFinite(r.threshold) &&
+ r.threshold > 0
+ ? r.threshold
+ : d.threshold,
+ threads: clampPositiveInt(r.threads, d.threads, 64),
+ python: str(r.python, d.python),
+ segModel: str(r.segModel, d.segModel),
+ embModel: str(r.embModel, d.embModel),
+ concurrency: clampPositiveInt(r.concurrency, d.concurrency, 16),
};
}
@@ -1030,6 +1136,7 @@ export function getSettings(): SiteSettings {
merged.savedVideoBackup = sanitizeSavedVideoBackup(merged.savedVideoBackup);
merged.buildPipeline = sanitizeBuildPipeline(merged.buildPipeline);
merged.digest = sanitizeDigest(merged.digest);
+ merged.diarization = sanitizeDiarization(merged.diarization);
// Workers. When the file predates the worker model (no `workers` key),
// synthesize a default list from the (now-settled) active app + per-app
// configs so existing installs behave identically. Otherwise sanitize the
@@ -1225,6 +1332,7 @@ export async function writeSettings(next: SiteSettings): Promise<void> {
savedVideoBackup: sanitizeSavedVideoBackup(next.savedVideoBackup),
buildPipeline: sanitizeBuildPipeline(next.buildPipeline),
digest: sanitizeDigest(next.digest),
+ diarization: sanitizeDiarization(next.diarization),
};
const tmp = `${file}.tmp-${process.pid}`;
await fs.promises.writeFile(tmp, JSON.stringify(merged, null, 2) + "\n");
diff --git a/common/lib/videoStatus.ts b/common/lib/videoStatus.ts
@@ -14,6 +14,11 @@ export type VideoFiles = {
hasNonCanonicalVtt: boolean;
hasWhisper: boolean;
hasCuesJson: boolean;
+ // A speaker-diarization sidecar is present. Existence only — the cleanup
+ // guard in cleanAudioFromTranscribed validates the CONTENT before letting the
+ // sweep delete audio, because a half-written file must not read as "done".
+ // This cheaper flag is what the snapshot's cleanup buckets use.
+ hasDiarization: boolean;
isUntranscribable: boolean;
audioFiles: string[];
// yt-dlp leaves audio.<ext>.part on disk when a download is interrupted.
@@ -34,6 +39,10 @@ export const CUES_JSON_FILENAME = "transcript.cues.json";
export const LIVE_CHAT_FILENAME = "transcript.live_chat.json";
export const LIVE_CHAT_CUES_FILENAME = "live_chat.cues.json";
export const META_FILENAME = "metadata.info.json";
+// Speaker-diarization sidecar. Deliberately NOT named transcript.<x>.<y>:
+// SUB_FILE_RE below would claim any such file as a subtitle track. See
+// lib/diarization.ts for the record shape.
+export const DIARIZATION_FILENAME = "diarization.json";
// A finalized media output under data/<id>/audio.<ext>. Excludes yt-dlp partials
// (.part) and audio-check snapshots (.part.good/.part.testing), the metadata
@@ -204,6 +213,7 @@ export async function readVideoFiles(
);
const hasWhisper = entries.includes(WHISPER_FILENAME);
const hasCuesJson = entries.includes(CUES_JSON_FILENAME);
+ const hasDiarization = entries.includes(DIARIZATION_FILENAME);
const audioFiles = entries.filter(isRealAudioFile);
const partAudioFiles = entries.filter(isPartAudioFile);
let isUntranscribable = false;
@@ -229,6 +239,7 @@ export async function readVideoFiles(
hasNonCanonicalVtt,
hasWhisper,
hasCuesJson,
+ hasDiarization,
isUntranscribable,
audioFiles,
partAudioFiles,
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **Speaker diarization can now be captured while the audio still exists, and Clean audio will no longer delete audio out from under it.** Audio is the one input in this pipeline that goes away: the Clean-audio sweep removes it as soon as a video has a transcript, and nothing downstream can reconstruct it. Everything built on top of speaker turns — attribution, naming the speakers, badges, quote filtering — can be redone from a saved file at any time, but the turns themselves can only be extracted from audio. So this ships the perishable half only: a new **Diarization** section in Settings, off by default, that captures speaker-labelled time ranges to `diarization.json` beside each transcript, and a **Diarize speakers** button in every channel's Cleanup stage that backfills over whatever audio is still on disk. Turning capture on also arms a guard: the sweep now refuses to delete audio for a transcribed video that has no sidecar yet, reports it as `awaiting diarization`, and drops that video out of the channel's "Est. reclaim" so the estimate matches what the sweep will actually take. That hold is the point — it is what keeps the audio alive long enough to be captured — but it does mean enabling this stops reclaiming disk until diarization catches up, and the setting says so. **Inline-after-transcribe is a separate switch, and it is off by default for a measured reason.** Transcription on this box runs at 221 s per audio-hour on the GPU — 16.3× real time, averaged over 3,602 real videos from their own `transcribe-outcome.json` records — while diarization runs on the CPU at roughly 500–680, because no diarization model has been ported to ggml and neither ONNX nor PyTorch has a Vulkan compute path on Linux. Diarization is therefore 2–3× slower than the transcription it would follow, so running it inline makes the whole pipeline 3–4× slower and idles the GPU. The intended shape for a large batch is the opposite: turn capture on so the audio is held, leave inline off so transcription runs at full speed, then backfill. The clustering threshold defaults to **0.9**, not sherpa-onnx's own 0.5, and that is also measured: on a six-minute excerpt of a known two-person interview, 0.5 produced 22 speaker clusters and 0.9 produced 6, with the top two at 40%/40% of talk time — recognizably the two hosts. It still over-splits, which is why every sidecar records the engine, both model names, the library version and the threshold that produced it: over-splitting is the recoverable direction, and a later pass can re-cluster or re-run selectively without ever needing the audio back. The engine sits behind a path (`DIARIZE_BIN` → `scripts/diarize.mjs`, driving sherpa-onnx through `scripts/diarize-sherpa.py`) exactly as the parakeet wrapper does, so swapping in pyannote later is a settings change rather than a code change. A failing diarizer never fails a transcription, and never counts as "done" — the sweep keeps that video's audio. Full spike results, including the disk arithmetic and what is still unmeasured, are in `plans/diarization-spike-results.md`.
- **The monitor widget's builder is a floorplan you rearrange, not a list of checkboxes.** The widget rendered its nine strips in an order hardcoded in JSX — controls, last sync, scheduler, disk, cleanable total, needs cleaning, workers, active jobs, needs work — and there was no order field in its config at all, so reference data you glance at once (the two cleanable strips) sat permanently above the thing you were actually watching. The builder couldn't help, because it wasn't spatial: 22 checkboxes in three fieldsets, nothing in the control surface corresponding to where anything landed. `/widget/builder` is now a scale model of the widget. Each section is a draggable card carrying a small CSS-drawn likeness of what that strip really renders — four dots for Workers, a part-filled bar for Active jobs, a meter with the low-disk notch for Disk, stacked rows for the two channel lists — laid out in the columns and order it will appear in, with the switched-off ones in a tray underneath. Drag it with the mouse or use the ↑/↓/◀/▶ buttons on every card; nothing here is mouse-only, and a section is one checkbox that changes place rather than a control that vanishes when you turn it off. **Columns are new** — up to three, each header showing the number that actually decides whether columns are viable (the per-column pixel width at the current preview size, flagged when it falls under 150px), alongside two wider presets, **Wide** 640×260 and **Panel** 720×420, because columns need width to be worth having. They are honored at every size by default, since an arrangement you made on purpose shouldn't quietly undo itself; **Stack columns when narrow** opts a widget you intend to resize into collapsing back to a single stack, on a container query measured against the widget's own box rather than the viewport (which, embedded in an iframe, is the host page's). Per-section display options moved onto the section that owns them — the batch bar, heading progress and ETA on the Active jobs card, worker names on Workers, the Sync flags on Controls — leaving only the genuinely global ones (section titles, hide-when-idle, the gear, absolute times) in a row of their own. The in-widget gear gets the same model in a shape that fits its ~320px overlay: one row per section in render order with move buttons and a column picker, driven by the same registry and the same placement functions as the board rather than a second copy of them. **Every link written before this still works, and still renders in the original order.** The arrangement rides in one new `l=` param (`/widget?l=ctl.disk.wk.jobs-cln.act` — `.` between sections, `-` between columns, both characters a URL leaves literal where a comma would come back as `%2C`), it is omitted entirely whenever the arrangement is the one the visibility flags already imply, and an inherited layout is normalized against those flags on arrival: unknown codes and repeats are dropped, a section since switched off is dropped, and a section switched **on** but unlisted is appended rather than silently disappearing — which is what keeps an already-copied link working after you enable something new from the widget's own gear. Underneath, the nine hardcoded conditionals became a section registry that drives both configuration surfaces, so adding a tenth strip is now a table entry (there is one waiting: `/api/widget/sync` already polls a digest field that nothing renders). The builder's remembered form state resets once, since a saved entry from before this has no layout to carry forward. See `editor/app/widget/lib/{sections,placement,config}.ts`, `editor/app/widget/builder/components/{LayoutBoard,SectionCard,SectionMini}.tsx`, and the normalization tests in `editor/app/widget/lib/placement.test.ts`.
- **No video a channel has ever listed can be lost again, and a truncated listing can no longer destroy one.** Until now the only record that a video belonged to a channel was the stored playlist file, which the deep pass overwrites wholesale. A video that appeared in a listing but was never downloaded had no folder on disk and no metadata of its own, so the playlist was the *only* place its URL lived — when it dropped out of the listing it was erased with no trace it had ever existed, and no way to even attempt a direct-link recovery. Each channel now keeps a **roster**: every video id it has ever been seen to contain, with the URL it was seen at, added to and never pruned. It is built the first time a channel syncs, from the playlist and the folders on disk *before* anything is rewritten, so the upgrade itself is the protection.
- **A listing that comes back suspiciously small is no longer believed on the first try.** The previous guard only refused a *completely* empty listing; a fetch that returned 100 of 10,795 entries sailed straight through and destroyed the other 10,695. A fresh listing that has lost more than a tenth of its entries — and more than 25 of them, so ordinary churn on a small channel doesn't trip it — is now treated as suspect: the stored video list and the missing-video flags are left exactly as they were, and the sweep is retried on the very next sync instead of after a full day. If a second enumeration reports a similar count it is accepted and acted on, because a real mass deletion repeats and a transient blip does not. No dialog to answer and nothing to override; at worst a genuine deletion lands one cadence later. The threshold is a new **Full-sweep shrink guard** setting (0 turns it off; an empty listing is always refused).
diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts
@@ -4,6 +4,39 @@ import { resetSnapshotScheduler } from "yt-dlp-transcript-common/jobs/snapshotSc
export const dynamic = "force-dynamic";
+// Cancel everything still live BEFORE the registry is dropped.
+//
+// The registry holds the only reachable AbortController for an in-flight child.
+// yt-dlp-style jobs run via runManagedFunction, which sets record.abortController
+// but never record.child (common/jobs/streamCommand.ts:306-319) — so abort() is
+// their ONLY kill path, and the moment the singleton is nulled the child is
+// unkillable by the app and reparents to init when the Next worker dies.
+//
+// Read the global DIRECTLY rather than calling getRegistry(): per this repo's
+// Next-16 lazy-singleton rule, getRegistry() *constructs* on read, so using it
+// here would build a fresh empty registry and dutifully cancel nothing. Reading
+// the (already globally-declared) var also means this route imports no registry
+// module, so it cannot construct one as an import side-effect either.
+function cancelLiveJobs() {
+ const registry = globalThis.__yttJobRegistry__;
+ if (!registry) return;
+ for (const job of registry.list()) {
+ if (job.status !== "running" && job.status !== "queued") continue;
+ try {
+ // Synchronous by design: abort() fires execa's cancelSignal inline, so
+ // SIGTERM is delivered before this route returns. execa's default
+ // forceKillAfterDelay escalates to SIGKILL, so we need not await the child.
+ registry.cancel(job.id);
+ // Then free the scheduler slot even if the record never settles (a child
+ // that ignores SIGTERM would otherwise leave a phantom busy slot behind
+ // for the next spec). Both calls are idempotent.
+ registry.forceRelease(job.id);
+ } catch {
+ /* one wedged job must not stop us cancelling the rest */
+ }
+ }
+}
+
// E2E test harness only. Mounted unconditionally so it's reachable from the
// dev:test script; the editor is intended for localhost use, not deployment.
export async function POST() {
@@ -18,6 +51,11 @@ function invalidate() {
// Clear any pending debounced snapshot regen FIRST (also clears its timer) so
// it can't fire against the about-to-be-reset registry mid-spec.
resetSnapshotScheduler();
+ // Kill live children while we still have the handle that can kill them. This
+ // serves the same intent the wipe below was added for (4057f76: "tests don't
+ // observe stale jobs from a prior spec") rather than fighting it — a job whose
+ // process is gone is a great deal less observable than one that isn't.
+ cancelLiveJobs();
// Reset the in-memory job registry too so tests don't observe stale jobs
// from a prior spec in the same dev-server lifetime.
// eslint-disable-next-line @typescript-eslint/no-explicit-any
diff --git a/editor/app/channels/[slug]/components/stages/CleanupStage.tsx b/editor/app/channels/[slug]/components/stages/CleanupStage.tsx
@@ -9,6 +9,7 @@ import {
checkKeptDeletedAction,
cleanAudioAction,
cleanExtraAudioFormatsAction,
+ diarizeChannelAction,
purgeSupersededAutoSubsAction,
removeWrongFormatAudioAction,
} from "../../whisperActions";
@@ -38,6 +39,10 @@ type Props = {
savedCount: number;
savedBytes: number;
backupConfigured: boolean;
+ // Whether the diarization capture lane is armed. When it is, the Clean-audio
+ // sweep holds audio back for any transcribed video without a sidecar, so the
+ // backfill button that clears that hold belongs next to it.
+ diarizationEnabled: boolean;
};
export function CleanupStage({
@@ -55,6 +60,7 @@ export function CleanupStage({
savedCount,
savedBytes,
backupConfigured,
+ diarizationEnabled,
}: Props) {
const defaultQueueKey = `channel:${slug}`;
const [cleanQueue, setCleanQueue] = useState(defaultQueueKey);
@@ -72,6 +78,30 @@ export function CleanupStage({
defaultQueueKey={defaultQueueKey}
/>
)}
+ {diarizationEnabled && (
+ <div className="flex flex-col gap-2">
+ <Heading
+ title="Diarize speakers"
+ desc="Capture who is speaking when, to diarization.json beside each transcript, for every video here that still has audio. This is the half of the decision that has a deadline: Clean audio below deletes the audio, and nothing can reconstruct it afterwards. While diarization is on, that sweep holds back any transcribed video without a sidecar — so running this is what releases the audio it is protecting."
+ />
+ <StreamActionLog
+ trigger={() => diarizeChannelAction(slug, cleanQueue)}
+ cancelAction={cancelJobAction}
+ buttonLabel="Diarize speakers"
+ runningLabel="Diarizing…"
+ label="Diarize speakers"
+ extraControls={
+ <QueueControl
+ value={cleanQueue}
+ onChange={setCleanQueue}
+ defaultQueueKey={defaultQueueKey}
+ existingQueues={existingQueues}
+ actionLabel="Diarize speakers"
+ />
+ }
+ />
+ </div>
+ )}
<div className="flex flex-col gap-2">
<Heading
title="Clean audio for transcribed videos"
diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx
@@ -424,6 +424,7 @@ export default async function ChannelDetailPage({
savedCount={savedTotals.count}
savedBytes={savedTotals.bytes}
backupConfigured={backupConfigured}
+ diarizationEnabled={settings.diarization.enabled}
/>
),
diagnostics: (
diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts
@@ -25,6 +25,7 @@ import {
} from "yt-dlp-transcript-common/controller/failedTranscodings";
import { clearFailedTranscriptions } from "yt-dlp-transcript-common/controller/failedTranscriptions";
import { cleanAudioFromTranscribed } from "yt-dlp-transcript-common/controller/cleanAudioFromTranscribed";
+import { diarizeAll } from "yt-dlp-transcript-common/controller/diarizeAll";
import { purgeSupersededAutoSubs } from "yt-dlp-transcript-common/controller/purgeSupersededAutoSubs";
import { checkKeptDeleted } from "yt-dlp-transcript-common/controller/checkKeptDeleted";
import { cleanExtraAudioFormats } from "yt-dlp-transcript-common/controller/cleanExtraAudioFormats";
@@ -422,6 +423,38 @@ export async function cleanAudioAction(
});
}
+// Capture speaker turns for this channel's videos that still have audio and no
+// diarization.json. Runs on the channel queue alongside the Clean-audio sweep,
+// deliberately: the two are the paired halves of the same decision — this one
+// captures the perishable input, that one destroys it, and letting them
+// interleave is what the cleanup guard exists to make safe.
+export async function diarizeChannelAction(
+ slug: string,
+ queueKey?: string,
+): Promise<StreamActionResult> {
+ const paths = getPaths();
+ return runManagedFunction({
+ kind: "diarize-channel",
+ queueKey: resolveQueueKey(channelQueueKey(slug), queueKey),
+ paths,
+ channelSlug: slug,
+ spec: { kind: "diarize-channel", slug, params: { queueKey } },
+ fn: async (onLog, signal) => {
+ const result = await diarizeAll({
+ paths,
+ channelSlugs: [slug],
+ onLog,
+ signal,
+ });
+ onLog(
+ `Diarize ${slug}: ${result.diarized} diarized, ${result.alreadyExists} already done, ` +
+ `${result.noAudio} without audio, ${result.failed} failed.`,
+ );
+ revalidatePath(`/channels/${slug}`);
+ },
+ });
+}
+
// Delete the YouTube auto-caption VTTs that our own transcript superseded (the
// supersededAutoSubs bucket). Manual only — never auto-queued — and the single
// irreversible step in the lane, so it lives next to the Clean-audio sweep and
diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts
@@ -24,6 +24,7 @@ import {
checkKeptDeletedAction,
cleanAudioAction,
cleanExtraAudioFormatsAction,
+ diarizeChannelAction,
clearFailedTranscodingsAction,
clearFailedTranscriptionsAction,
removeFailedTranscodingsAction,
@@ -266,6 +267,10 @@ export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = {
const { queueKey } = params(spec);
return cleanAudioAction(spec.slug, queueKey);
},
+ "diarize-channel": (spec) => {
+ const { queueKey } = params(spec);
+ return diarizeChannelAction(spec.slug, queueKey);
+ },
"check-kept-deleted": (spec) => {
const { queueKey } = params(spec);
return checkKeptDeletedAction(spec.slug, queueKey);
diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts
@@ -321,6 +321,37 @@ export async function saveSettingsAction(
}
: dD;
+ // Diarization. Same hidden-marker discipline as the digest block above, and
+ // for the same reason: unchecked checkboxes are absent from a FormData, so
+ // without the marker any unrelated save would read `enabled` as false and
+ // silently disarm the capture lane — which would let the next Clean-audio
+ // sweep delete audio that was being held for diarization. Spread the CURRENT
+ // block first so fields this form does not render survive untouched.
+ const dDiar = getSettings().diarization;
+ const diarizationFormPresent =
+ formData.get("diarizationFormPresent") === "1";
+ const num = (key: string, fallback: number) => {
+ const n = Number.parseFloat(String(formData.get(key) ?? "").trim());
+ return Number.isFinite(n) ? n : fallback;
+ };
+ const diarizationSettings: SiteSettings["diarization"] =
+ diarizationFormPresent
+ ? {
+ ...dDiar,
+ enabled: formData.get("diarizationEnabled") === "on",
+ inlineAfterTranscribe:
+ formData.get("diarizationInlineAfterTranscribe") === "on",
+ threshold: num("diarizationThreshold", dDiar.threshold),
+ threads: num("diarizationThreads", dDiar.threads),
+ concurrency: num("diarizationConcurrency", dDiar.concurrency),
+ python:
+ String(formData.get("diarizationPython") ?? "").trim() ||
+ dDiar.python,
+ segModel: String(formData.get("diarizationSegModel") ?? "").trim(),
+ embModel: String(formData.get("diarizationEmbModel") ?? "").trim(),
+ }
+ : dDiar;
+
const dB = defaultBuildPipeline();
const buildModeRaw = String(formData.get("buildMode") ?? "").trim();
const buildPipeline = {
@@ -375,6 +406,7 @@ export async function saveSettingsAction(
// Same: the Digest section of this form owns these fields, but an unrelated
// save must not reset them (and must never silently flip remoteEnabled on).
digest: digestSettings,
+ diarization: diarizationSettings,
};
try {
await writeSettings(next);
diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx
@@ -690,6 +690,114 @@ export function SettingsForm({ initial, apps, digestApps }: Props) {
/>
</fieldset>
<fieldset className="flex flex-col gap-3 border border-border rounded p-3">
+ <legend className="px-1 text-sm font-medium">Diarization</legend>
+ {/*
+ Same marker discipline as Digest above — and it matters MORE here.
+ Without it an unrelated save would read `diarizationEnabled` as false,
+ disarming the cleanup guard, and the next Clean-audio sweep would
+ delete audio that was being held for diarization. That deletion is
+ not recoverable.
+ */}
+ <input
+ type="hidden"
+ name="diarizationFormPresent"
+ value="1"
+ readOnly
+ />
+ <p className="text-xs text-muted-foreground">
+ Who is speaking when, captured to <code>diarization.json</code> beside
+ each transcript. This runs on <strong>audio</strong>, and audio is the
+ one input here that goes away — the Clean-audio sweep deletes it once a
+ video is transcribed. So this is a one-shot window per video: the
+ speaker turns can be captured now or not at all, while everything built
+ on top of them can be redone from the saved file at any time.
+ </p>
+ <label className="flex items-start gap-2 text-sm">
+ <input
+ type="checkbox"
+ name="diarizationEnabled"
+ defaultChecked={initial.diarization.enabled}
+ className="mt-1"
+ />
+ <span className="flex flex-col gap-1">
+ <span className="font-medium">Capture speaker diarization</span>
+ <span className="text-xs text-muted-foreground">
+ Turning this on also <strong>holds disk</strong>: the Clean-audio
+ sweep stops deleting audio for a transcribed video until its
+ sidecar exists, and each channel's "Est. reclaim"
+ drops to match. That hold is the point — it is what keeps the
+ perishable input alive long enough to capture — but it means free
+ space stops being reclaimed until diarization catches up.
+ </span>
+ </span>
+ </label>
+ <label className="flex items-start gap-2 text-sm">
+ <input
+ type="checkbox"
+ name="diarizationInlineAfterTranscribe"
+ defaultChecked={initial.diarization.inlineAfterTranscribe}
+ className="mt-1"
+ />
+ <span className="flex flex-col gap-1">
+ <span className="font-medium">
+ Diarize inline, right after each transcription
+ </span>
+ <span className="text-xs text-muted-foreground">
+ Off by default, from measurement rather than caution. On this box
+ transcription runs at ~221 s per audio-hour on the GPU (16× real
+ time, over 3,602 videos); diarization runs on the CPU at roughly
+ 500–680, because no diarization model has been ported to ggml and
+ neither ONNX nor PyTorch has a Vulkan path on Linux. Inline
+ therefore makes the whole pipeline ~3–4× slower and idles the GPU.
+ For a large batch, leave this off, let transcription run at full
+ speed with the box above holding the audio, and catch up
+ afterwards with <strong>Diarize speakers</strong> on each channel.
+ Turn it on for steady state.
+ </span>
+ </span>
+ </label>
+ <Field
+ label="Segmentation model (ONNX)"
+ name="diarizationSegModel"
+ defaultValue={initial.diarization.segModel}
+ hint="Absolute path to a pyannote segmentation-3.0 ONNX model. Empty means the lane is not configured and every run reports a skip rather than a failure."
+ />
+ <Field
+ label="Speaker-embedding model (ONNX)"
+ name="diarizationEmbModel"
+ defaultValue={initial.diarization.embModel}
+ hint="Absolute path to a speaker-embedding ONNX model (e.g. NeMo TitaNet)."
+ />
+ <Field
+ label="Python interpreter"
+ name="diarizationPython"
+ defaultValue={initial.diarization.python}
+ hint="Interpreter with sherpa-onnx installed. sherpa-onnx ships wheels only up to CPython 3.13, so on a 3.14 system this must point at a dedicated venv."
+ />
+ <Field
+ label="Clustering threshold"
+ name="diarizationThreshold"
+ defaultValue={String(initial.diarization.threshold)}
+ type="number"
+ step="0.01"
+ hint="The knob that decides how many speakers come out; larger merges more. It is recorded in every sidecar, so changing it is what lets a later pass tell which files are worth regenerating."
+ />
+ <Field
+ label="Engine threads"
+ name="diarizationThreads"
+ defaultValue={String(initial.diarization.threads)}
+ type="number"
+ hint="Threads per diarize run."
+ />
+ <Field
+ label="Backfill concurrency"
+ name="diarizationConcurrency"
+ defaultValue={String(initial.diarization.concurrency)}
+ type="number"
+ hint="Videos diarized at once during a backfill. Default 1: this is CPU-bound work competing with GPU feeding and the digest sweep for the same 8 threads."
+ />
+ </fieldset>
+ <fieldset className="flex flex-col gap-3 border border-border rounded p-3">
<legend className="px-1 text-sm font-medium">Social links</legend>
<p className="text-xs text-muted-foreground">
Default social links shown in every site's footer. Each site can
diff --git a/editor/e2e/diarization.spec.ts b/editor/e2e/diarization.spec.ts
@@ -0,0 +1,312 @@
+import { writeFile } from "node:fs/promises";
+import { test, expect } from "@playwright/test";
+import {
+ pathExists,
+ readJson,
+ resetData,
+ resolvePath,
+ writeSettings,
+ generateReport,
+} from "./helpers";
+import { baseUrl } from "./baseUrl";
+
+// Speaker-diarization capture. The point of this lane is TIMING, not quality:
+// audio is deleted once a video is transcribed, so the speaker turns get
+// captured while the audio is still on disk or never. These specs pin the three
+// behaviours that make that safe.
+//
+// The engine is faked via DIARIZE_BIN (e2e/fixtures/bin/fake-diarize.mjs), wired
+// in the dev:test / start:test scripts the way every other fake binary is.
+
+const SLUG = "test-transcribe";
+const SNAPSHOT_REL = `test-transcripts/channels/${SLUG}/snapshot.json`;
+
+function dataRel(videoId: string, file: string): string {
+ return `test-transcripts/channels/${SLUG}/data/${videoId}/${file}`;
+}
+
+async function seedTranscript(videoId: string): Promise<void> {
+ await writeFile(
+ resolvePath(dataRel(videoId, "transcript.json")),
+ '{"transcription":[]}\n',
+ );
+}
+
+async function seedDiarization(videoId: string): Promise<void> {
+ await writeFile(
+ resolvePath(dataRel(videoId, "diarization.json")),
+ JSON.stringify({
+ videoId,
+ generatedAt: "2026-08-07T00:00:00.000Z",
+ speakers: 1,
+ turns: [{ start: 0, end: 10, speaker: 0 }],
+ engine: { engine: "fake-diarize", threshold: 0.5 },
+ }) + "\n",
+ );
+}
+
+const BASE_SETTINGS = {
+ adminTitle: "Test Admin",
+ maxTranscriptPageBytes: 8388608,
+ sleepBetweenDownloadsSeconds: 0,
+ minFreeDiskGB: 0,
+ verifyAvailabilityBeforeClean: false,
+ syncScheduler: { fullSweepIntervalMinutes: 0 },
+};
+
+function diarizationSettings(over: Record<string, unknown> = {}) {
+ return {
+ ...BASE_SETTINGS,
+ diarization: {
+ enabled: true,
+ inlineAfterTranscribe: true,
+ threshold: 0.5,
+ threads: 1,
+ python: "python3",
+ // Any non-empty pair: the fake engine never reads them, but an empty pair
+ // is reported as "not configured" and the lane never runs.
+ segModel: "/dev/null",
+ embModel: "/dev/null",
+ concurrency: 1,
+ ...over,
+ },
+ };
+}
+
+// These drive a real transcribe batch over three videos and then a cleanup
+// sweep, which does not fit the 30s default on a loaded box.
+const SLOW = 120_000;
+
+// (a) A transcription produces diarization.json.
+test("transcribing a video captures its speaker turns", async ({ page }) => {
+ test.setTimeout(SLOW);
+ await resetData("one-transcribe-channel-with-audio");
+ await writeSettings(diarizationSettings());
+
+ await page.goto(`/channels/${SLUG}`);
+ await page.getByRole("button", { name: "Transcribe stage summary" }).click();
+ await page
+ .getByRole("button", { name: "Transcribe missing", exact: true })
+ .click();
+
+ await expect
+ .poll(async () => pathExists(dataRel("vidA", "diarization.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+
+ const record = await readJson<{
+ videoId: string;
+ speakers: number;
+ turns: { start: number; end: number; speaker: number }[];
+ engine: { engine: string; threshold: number };
+ }>(dataRel("vidA", "diarization.json"));
+
+ expect(record.videoId).toBe("vidA");
+ expect(record.turns.length).toBeGreaterThan(0);
+ expect(record.speakers).toBe(2);
+ // Provenance is the whole reason this is a record and not a bare array: a
+ // later attribution pass has to be able to tell what produced a given file.
+ expect(record.engine.engine).toBe("fake-diarize");
+ expect(record.engine.threshold).toBe(0.5);
+
+ // Let the batch finish before the test ends. Without this it leaks a running
+ // transcribe job into the NEXT spec, whose resetData then races it — the
+ // runner re-creates files in the tree the walk is deleting. (resetData's own
+ // maxRetries comment documents this hazard; the fix belongs here, at the
+ // source, rather than relying on the retry.) vidC is the fake diarizer's
+ // designated failure, so it gets a transcript but never a sidecar.
+ await expect
+ .poll(async () => pathExists(dataRel("vidB", "diarization.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+ await expect
+ .poll(async () => pathExists(dataRel("vidC", "transcript.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+});
+
+// (b) A failing diarizer must not fail the transcription.
+test("a failing diarizer does not fail the transcription", async ({ page }) => {
+ test.setTimeout(SLOW);
+ // The fake engine exits non-zero for vidC (see fake-diarize.mjs). vidA and
+ // vidB succeed in the same run, so this also pins that one video's failure
+ // does not take the batch down with it.
+ await resetData("one-transcribe-channel-with-audio");
+ await writeSettings(diarizationSettings());
+
+ await page.goto(`/channels/${SLUG}`);
+ await page.getByRole("button", { name: "Transcribe stage summary" }).click();
+ await page
+ .getByRole("button", { name: "Transcribe missing", exact: true })
+ .click();
+
+ // The transcript still lands for the video whose diarization failed —
+ // diarization is strictly best-effort and must never fail a transcription
+ // that already succeeded.
+ await expect
+ .poll(async () => pathExists(dataRel("vidC", "transcript.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+ // Its neighbours were diarized normally. Waiting for BOTH also settles the
+ // batch, so this spec does not leak a running job into the next one and the
+ // cleanup below sees a stable tree.
+ await expect
+ .poll(async () => pathExists(dataRel("vidA", "diarization.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+ await expect
+ .poll(async () => pathExists(dataRel("vidB", "diarization.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+
+ // No sidecar for vidC, so its audio stays eligible for a later pass rather
+ // than being silently treated as done.
+ expect(await pathExists(dataRel("vidC", "diarization.json"))).toBe(false);
+
+ // And the guard holds vidC's audio while releasing the diarized ones.
+ await page.getByRole("button", { name: "Cleanup stage summary" }).click();
+ await page.getByRole("button", { name: "Clean audio", exact: true }).click();
+ await expect(page.getByLabel("Clean audio output")).toContainText(
+ "awaiting diarization",
+ { timeout: 30_000 },
+ );
+ expect(await pathExists(dataRel("vidC", "audio.m4a"))).toBe(true);
+});
+
+// (c) THE REGRESSION THAT WOULD SILENTLY DESTROY DATA.
+// Cleanup must refuse to delete audio for a transcribed-but-not-yet-diarized
+// video. Without this guard an async diarize pass races the sweep and loses the
+// only copy of the audio, permanently and with no trace.
+test("cleanup refuses to delete audio for a transcribed-but-undiarized video", async ({
+ page,
+}) => {
+ test.setTimeout(SLOW);
+ await resetData("one-transcribe-channel-with-audio");
+ // Capture ON, inline OFF — the recommended batch shape, and precisely the
+ // window in which the race exists.
+ await writeSettings(diarizationSettings({ inlineAfterTranscribe: false }));
+ await seedTranscript("vidA");
+ await seedTranscript("vidB");
+ // vidA has been diarized; vidB has not.
+ await seedDiarization("vidA");
+ await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {});
+
+ // The snapshot's cleanup bucket must agree with the sweep, or "Est. reclaim"
+ // promises space the sweep is going to refuse to take.
+ //
+ // Poll-with-reload rather than read once: the channel page serves a PERSISTED
+ // snapshot and regenerates on a debounce, and the settings cache this reads
+ // `diarization.enabled` from is invalidated by a fire-and-forget fetch. A
+ // single read can therefore catch a snapshot built from the pre-writeSettings
+ // defaults (diarization off), in which case the guard exclusion has not been
+ // applied yet and vidB is still in the bucket.
+ await generateReport(page, SLUG);
+ await expect
+ .poll(
+ async () => {
+ await page.goto(`/channels/${SLUG}`);
+ const snap = await readJson<{
+ buckets: { transcribedWithAudio?: string[] };
+ }>(SNAPSHOT_REL);
+ return snap.buckets.transcribedWithAudio ?? [];
+ },
+ { timeout: 60_000 },
+ )
+ .toEqual(["vidA"]);
+
+ await page.getByRole("button", { name: "Cleanup stage summary" }).click();
+ await page.getByRole("button", { name: "Clean audio", exact: true }).click();
+ await expect(page.getByLabel("Clean audio output")).toContainText(
+ "awaiting diarization",
+ { timeout: 30_000 },
+ );
+
+ // vidA's audio is reclaimed; vidB's is held.
+ expect(await pathExists(dataRel("vidA", "audio.m4a"))).toBe(false);
+ expect(await pathExists(dataRel("vidB", "audio.m4a"))).toBe(true);
+});
+
+// The backfill loop end to end — this is the workflow a large batch actually
+// uses: transcribe at full GPU speed with the hook off, let the guard hold the
+// audio, then catch up. It is also the only path that exercises diarizeAll.
+test("the Diarize speakers backfill clears the hold and releases the audio", async ({
+ page,
+}) => {
+ test.setTimeout(SLOW);
+ await resetData("one-transcribe-channel-with-audio");
+ await writeSettings(diarizationSettings({ inlineAfterTranscribe: false }));
+ await seedTranscript("vidA");
+ await seedTranscript("vidB");
+ await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {});
+
+ // Nothing is diarized yet, so the sweep must take nothing.
+ await generateReport(page, SLUG);
+ await page.goto(`/channels/${SLUG}`);
+ await page.getByRole("button", { name: "Cleanup stage summary" }).click();
+ await page.getByRole("button", { name: "Clean audio", exact: true }).click();
+ await expect(page.getByLabel("Clean audio output")).toContainText(
+ "awaiting diarization",
+ { timeout: 30_000 },
+ );
+ expect(await pathExists(dataRel("vidA", "audio.m4a"))).toBe(true);
+ expect(await pathExists(dataRel("vidB", "audio.m4a"))).toBe(true);
+
+ // Backfill.
+ await page
+ .getByRole("button", { name: "Diarize speakers", exact: true })
+ .click();
+ await expect
+ .poll(async () => pathExists(dataRel("vidA", "diarization.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+ await expect
+ .poll(async () => pathExists(dataRel("vidB", "diarization.json")), {
+ timeout: 60_000,
+ })
+ .toBe(true);
+
+ // With the sidecars in place the hold clears and the audio is reclaimable —
+ // the whole point of the guard being transient rather than a protection.
+ await generateReport(page, SLUG);
+ await page.goto(`/channels/${SLUG}`);
+ await page.getByRole("button", { name: "Cleanup stage summary" }).click();
+ await page.getByRole("button", { name: "Clean audio", exact: true }).click();
+ await expect(page.getByLabel("Clean audio output")).toContainText(
+ "Cleaned 2 audio file",
+ { timeout: 30_000 },
+ );
+ expect(await pathExists(dataRel("vidA", "audio.m4a"))).toBe(false);
+ expect(await pathExists(dataRel("vidB", "audio.m4a"))).toBe(false);
+});
+
+// The guard is opt-in, and turning it off must restore the old behaviour
+// exactly — otherwise enabling diarization would be a one-way door for disk.
+test("with diarization disabled, cleanup deletes undiarized audio as before", async ({
+ page,
+}) => {
+ test.setTimeout(SLOW);
+ await resetData("one-transcribe-channel-with-audio");
+ await writeSettings(diarizationSettings({ enabled: false }));
+ await seedTranscript("vidA");
+ await seedTranscript("vidB");
+ await fetch(`${baseUrl}/api/test/invalidate-cache`).catch(() => {});
+
+ await generateReport(page, SLUG);
+ await page.goto(`/channels/${SLUG}`);
+ await page.getByRole("button", { name: "Cleanup stage summary" }).click();
+ await page.getByRole("button", { name: "Clean audio", exact: true }).click();
+ await expect(page.getByLabel("Clean audio output")).toContainText(
+ "Cleaned 2 audio file",
+ { timeout: 30_000 },
+ );
+
+ expect(await pathExists(dataRel("vidA", "audio.m4a"))).toBe(false);
+ expect(await pathExists(dataRel("vidB", "audio.m4a"))).toBe(false);
+});
diff --git a/editor/e2e/fixtureProcs.ts b/editor/e2e/fixtureProcs.ts
@@ -0,0 +1,121 @@
+// Find (and reap) e2e fixture-binary processes that outlived the run that
+// spawned them.
+//
+// Why this needs to be strict: these helpers kill processes, so a loose match is
+// a footgun aimed at the developer's own machine. `pgrep -f fixtures/bin` is NOT
+// safe — it matches any shell whose command line merely *mentions* the path,
+// including the very command running the sweep. Verified on this box:
+//
+// a real fixture: argv[0]="node" argv[1]="<...>/e2e/fixtures/bin/fake-ytdlp.mjs"
+// a shell that (path appears in argv[2], the -c script)
+// mentions it: argv[0]="/usr/bin/zsh"
+//
+// So we match on argv[1] specifically — the script path node was handed via the
+// fixtures' `#!/usr/bin/env node` shebang — and require it to sit under this
+// checkout's absolute fixtures/bin directory. A shell can never satisfy that.
+//
+// Linux-only (reads /proc). That matches where this suite runs; on any other
+// platform the scan returns nothing rather than guessing with `ps`, so a sweep
+// degrades to a no-op instead of killing something it misidentified.
+
+import { readFile, readdir } from "node:fs/promises";
+import { dirname, resolve, basename } from "node:path";
+import { fileURLToPath } from "node:url";
+
+const here = dirname(fileURLToPath(import.meta.url));
+
+// Absolute, normalised, with a trailing separator so a startsWith() test cannot
+// match a sibling directory whose name merely shares the prefix.
+export const fixtureBinDir = resolve(here, "fixtures", "bin") + "/";
+
+export type FixtureProc = { pid: number; script: string; argv: string[] };
+
+async function readCmdline(pid: string): Promise<string[] | null> {
+ try {
+ const raw = await readFile(`/proc/${pid}/cmdline`, "utf8");
+ // cmdline is NUL-separated with a trailing NUL; drop the empty tail.
+ const argv = raw.split("\0").filter((s) => s.length > 0);
+ return argv.length > 0 ? argv : null;
+ } catch {
+ // The process exited between readdir and readFile, or it isn't ours to
+ // inspect. Either way it is not a stray we can act on.
+ return null;
+ }
+}
+
+// Every live process whose argv[1] is a script inside this checkout's
+// e2e/fixtures/bin. Never includes the calling process.
+export async function listFixtureProcesses(): Promise<FixtureProc[]> {
+ if (process.platform !== "linux") return [];
+ let entries: string[];
+ try {
+ entries = await readdir("/proc");
+ } catch {
+ return [];
+ }
+ const found: FixtureProc[] = [];
+ for (const entry of entries) {
+ if (!/^\d+$/.test(entry)) continue;
+ const pid = Number(entry);
+ if (pid === process.pid) continue;
+ const argv = await readCmdline(entry);
+ if (!argv) continue;
+ // argv[0] is the interpreter (the shebang resolves to node), argv[1] the
+ // fixture script. Requiring both is what makes this safe to kill on.
+ if (!basename(argv[0]).startsWith("node")) continue;
+ const script = argv[1];
+ if (!script || !resolve(script).startsWith(fixtureBinDir)) continue;
+ found.push({ pid, script: resolve(script), argv });
+ }
+ return found;
+}
+
+// SIGTERM every stray, then SIGKILL whatever is still alive after `graceMs`
+// (a fixture may deliberately ignore SIGTERM — fake-parakeet-stitch does, in
+// `hangterm` dirs). Returns what it killed so callers can log it: a silent
+// reaper would hide a regression in the layers that are supposed to prevent
+// strays in the first place.
+export async function killFixtureProcesses(
+ graceMs = 2_000,
+): Promise<FixtureProc[]> {
+ const strays = await listFixtureProcesses();
+ if (strays.length === 0) return [];
+ for (const p of strays) {
+ try {
+ process.kill(p.pid, "SIGTERM");
+ } catch {
+ /* already gone */
+ }
+ }
+ await new Promise((r) => setTimeout(r, graceMs));
+ for (const p of strays) {
+ try {
+ process.kill(p.pid, "SIGKILL");
+ } catch {
+ /* exited on SIGTERM, as it should have */
+ }
+ }
+ return strays;
+}
+
+// The sweep both the globalSetup and globalTeardown hooks run.
+//
+// It LOGS whatever it reaps, deliberately and loudly. A silent reaper would hide
+// the thing we actually care about: strays are supposed to be impossible now
+// (jobs are cancelled before the registry is dropped, on graceful shutdown, and
+// each fixture self-terminates when orphaned). Anything this finds is a
+// regression in one of those three layers, not a routine cleanup.
+export async function sweepFixtureProcesses(phase: "setup" | "teardown") {
+ const killed = await killFixtureProcesses();
+ if (killed.length === 0) return;
+ const names = killed
+ .map((p) => `${basename(p.script)}(${p.pid})`)
+ .join(", ");
+ const cause =
+ phase === "setup"
+ ? "left behind by an EARLIER run (a killed run cannot run its own teardown)"
+ : "survived this run — the anti-orphan layers did not hold";
+ console.warn(
+ `[fixture-sweep:${phase}] reaped ${killed.length} stray fixture process(es) ${cause}: ${names}`,
+ );
+}
diff --git a/editor/e2e/fixtures/bin/_watchdog.mjs b/editor/e2e/fixtures/bin/_watchdog.mjs
@@ -0,0 +1,56 @@
+// Self-termination for the e2e fake binaries: no fixture should be able to
+// outlive the run that spawned it, whatever happens to the server above it.
+//
+// The app-side fixes (cancelling live jobs before the registry is wiped, and on
+// graceful shutdown) cover every case where something is still alive to send a
+// signal. This covers the case where nothing is: a SIGKILLed dev server can't
+// run shutdown code, so its children reparent and keep spinning. A full suite
+// once left 27 of these alive at ~29% of a core each — enough to saturate the
+// box, push load to 46, and turn a 30-minute suite with 1 failure into a
+// 1.3-hour suite with 13 failures that all looked like real regressions.
+//
+// Two independent triggers, because neither alone is sufficient:
+//
+// 1. Reparenting — our parent died. Checked as "ppid changed since startup",
+// NOT "ppid === 1": on a systemd-user box (this one) a child subreaper
+// adopts orphans, so the pid-1 test would silently never fire.
+// 2. A max-lifetime backstop, for the case where the parent is alive but
+// wedged. The longest legitimate fixture life is 30s
+// (fake-parakeet-stitch's completion backstop, and fake-ytdlp --test-slow,
+// which jobs-retry/jobs-reorder/queues specs rely on), so the default 120s
+// leaves 4x margin. Override with FIXTURE_MAX_LIFETIME_MS if a future
+// fixture legitimately needs longer.
+//
+// BOTH timers are unref'd. That is the load-bearing detail: a ref'd timer would
+// hold the event loop open and keep an otherwise-finished fixture alive — this
+// module would then create exactly the leak it exists to prevent.
+
+const POLL_MS = 1_000;
+const DEFAULT_MAX_LIFETIME_MS = 120_000;
+
+export function installFixtureWatchdog() {
+ const label = process.argv[1]?.split("/").pop() ?? "fixture";
+ const initialPpid = process.ppid;
+
+ const poll = setInterval(() => {
+ // process.ppid is re-read from the OS on each access, so this observes the
+ // reparent rather than a value cached at startup.
+ if (process.ppid !== initialPpid) {
+ process.stderr.write(
+ `[watchdog] ${label}: parent ${initialPpid} died (now ${process.ppid}) — exiting rather than orphaning\n`,
+ );
+ process.exit(1);
+ }
+ }, POLL_MS);
+ poll.unref();
+
+ const maxLifetimeMs =
+ Number(process.env.FIXTURE_MAX_LIFETIME_MS) || DEFAULT_MAX_LIFETIME_MS;
+ const backstop = setTimeout(() => {
+ process.stderr.write(
+ `[watchdog] ${label}: exceeded ${maxLifetimeMs}ms lifetime — exiting\n`,
+ );
+ process.exit(1);
+ }, maxLifetimeMs);
+ backstop.unref();
+}
diff --git a/editor/e2e/fixtures/bin/fake-chough.mjs b/editor/e2e/fixtures/bin/fake-chough.mjs
@@ -6,6 +6,10 @@
// { duration_seconds, chunks, text, chunk_data: [{ start_time, end_time, text }] }
// with times in SECONDS. Run from cwd == the video dir (set by runWhisperBatch).
import { writeFile } from "node:fs/promises";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const argv = process.argv.slice(2);
function arg(flag) {
diff --git a/editor/e2e/fixtures/bin/fake-claude.mjs b/editor/e2e/fixtures/bin/fake-claude.mjs
@@ -16,6 +16,10 @@
// rather than on a schema.
import { readFileSync } from "node:fs";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const argv = process.argv.slice(2);
function flag(name) {
diff --git a/editor/e2e/fixtures/bin/fake-diarize.mjs b/editor/e2e/fixtures/bin/fake-diarize.mjs
@@ -0,0 +1,81 @@
+#!/usr/bin/env node
+// E2E fake for scripts/diarize.mjs (the speaker-diarization wrapper). Mimics the
+// args controller/diarizeOne.ts builds:
+// diarize.mjs --output diarization.json --video-id <id> --seg <m> --emb <m>
+// --threshold <f> --threads <n> --python <p> <audioFile>
+// Output contract: write a DiarizationRecord (common/lib/diarization.ts) to
+// EXACTLY the --output path. Run from cwd == the video dir.
+//
+// Two modes by video dir (cwd), mirroring how fake-parakeet-stitch keys off the
+// dir name:
+// * default: write a plausible two-speaker record and exit 0.
+// * vidC, or any dir containing "diarizefail": exit non-zero WITHOUT writing
+// anything. This is the case that must not fail the surrounding
+// transcription, and must leave the video ineligible for the Clean-audio
+// sweep — audio is unrecoverable, so "the diarizer broke" has to read as
+// "keep the audio", never as "done".
+//
+// vidC is the designated failing id because it is already in the shared
+// one-transcribe-channel-with-audio fixture's archive; a spec can exercise the
+// failure path without renaming dirs into ids the batch would never pick up.
+import { writeFile } from "node:fs/promises";
+import path from "node:path";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
+
+const argv = process.argv.slice(2);
+function arg(flag) {
+ const i = argv.indexOf(flag);
+ return i < 0 ? undefined : argv[i + 1];
+}
+
+const out = arg("--output") ?? arg("-o");
+const videoId = arg("--video-id") ?? path.basename(process.cwd());
+const seg = arg("--seg");
+const emb = arg("--emb");
+const threshold = Number(arg("--threshold") ?? "0.5");
+const audio = argv[argv.length - 1];
+
+if (!out) {
+ process.stderr.write("[fake-diarize] missing --output\n");
+ process.exit(2);
+}
+
+const cwd = process.cwd();
+if (
+ cwd.toLowerCase().includes("diarizefail") ||
+ path.basename(cwd) === "vidC"
+) {
+ process.stderr.write(`[fake-diarize] synthetic engine failure for ${videoId}\n`);
+ process.exit(3);
+}
+
+const turns = [
+ { start: 0, end: 4.5, speaker: 0 },
+ { start: 4.5, end: 7.25, speaker: 1 },
+ { start: 7.25, end: 10, speaker: 0 },
+];
+
+await writeFile(
+ out,
+ JSON.stringify({
+ videoId,
+ generatedAt: new Date().toISOString(),
+ audioSeconds: 10,
+ durationMs: 1,
+ speakers: new Set(turns.map((t) => t.speaker)).size,
+ turns,
+ engine: {
+ engine: "fake-diarize",
+ ...(seg ? { segmentationModel: path.basename(seg) } : {}),
+ ...(emb ? { embeddingModel: path.basename(emb) } : {}),
+ version: "e2e",
+ threshold,
+ },
+ }) + "\n",
+);
+process.stderr.write(
+ `[fake-diarize] ${turns.length} turns for ${videoId} (${audio}) -> ${out}\n`,
+);
diff --git a/editor/e2e/fixtures/bin/fake-ffmpeg.mjs b/editor/e2e/fixtures/bin/fake-ffmpeg.mjs
@@ -9,6 +9,10 @@
// same w.r.t. corruption detection but doesn't actually write
// output.
import { readFile, writeFile } from "node:fs/promises";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const CORRUPT_MARKER = "__CORRUPT__";
// When the probed file contains this marker, the probe (and transcode) sleeps
diff --git a/editor/e2e/fixtures/bin/fake-ffprobe.mjs b/editor/e2e/fixtures/bin/fake-ffprobe.mjs
@@ -5,6 +5,10 @@
// file (`__DUR=<seconds>__`). A file with no marker reports a large duration so
// the guard never falsely trips on it (matches a full-length download).
import { readFile } from "node:fs/promises";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const argv = process.argv.slice(2);
// The input file is the only non-flag argument (and not a value of -show_entries
diff --git a/editor/e2e/fixtures/bin/fake-gallery-dl.mjs b/editor/e2e/fixtures/bin/fake-gallery-dl.mjs
@@ -9,6 +9,10 @@
//
// Deterministic by design: the same tweets every run, so a re-run proves the
// incremental (no-duplicate) guarantee.
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const args = process.argv.slice(2);
diff --git a/editor/e2e/fixtures/bin/fake-parakeet-stitch.mjs b/editor/e2e/fixtures/bin/fake-parakeet-stitch.mjs
@@ -25,6 +25,10 @@
import { mkdir, writeFile, rm } from "node:fs/promises";
import { existsSync } from "node:fs";
import path from "node:path";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const argv = process.argv.slice(2);
function arg(flag) {
diff --git a/editor/e2e/fixtures/bin/fake-whisper.mjs b/editor/e2e/fixtures/bin/fake-whisper.mjs
@@ -4,6 +4,10 @@
// Writes <tmpBase>.json next to cwd (which runWhisperBatch sets to the
// video dir). Format matches what common/lib/whisper.ts parses.
import { writeFile } from "node:fs/promises";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
const argv = process.argv.slice(2);
function arg(flag) {
diff --git a/editor/e2e/fixtures/bin/fake-ytdlp.mjs b/editor/e2e/fixtures/bin/fake-ytdlp.mjs
@@ -15,6 +15,10 @@
import { mkdir, writeFile, readFile, appendFile, stat, rename } from "node:fs/promises";
import { existsSync, openSync, writeSync, closeSync } from "node:fs";
import path from "node:path";
+import { installFixtureWatchdog } from "./_watchdog.mjs";
+
+// Never outlive the run that spawned us — see _watchdog.mjs.
+installFixtureWatchdog();
// Used in audio-check scenarios. The fake-ffmpeg companion treats files
// containing this string as malformed when probing.
diff --git a/editor/e2e/globalSetup.ts b/editor/e2e/globalSetup.ts
@@ -0,0 +1,17 @@
+// Reap fixture strays left by an EARLIER run before this one starts.
+//
+// A run killed with SIGKILL (Ctrl-C twice, `timeout`, an OOM) never gets to run
+// its own globalTeardown, so its strays are still burning CPU when the next run
+// begins. That is precisely how 27 of them once accumulated and turned a
+// 30-minute suite into a 1.3-hour one whose 13 failures were pure contention.
+//
+// ORDERING CAVEAT: Playwright starts the `webServer` entries BEFORE globalSetup
+// (the same constraint export/e2e-2origin/globalSetup.ts documents at :151-152).
+// So this cannot clear a stale dev server ahead of the new one booting — it only
+// reaps fixture children, which is what actually eats the CPU.
+
+import { sweepFixtureProcesses } from "./fixtureProcs";
+
+export default async function globalSetup() {
+ await sweepFixtureProcesses("setup");
+}
diff --git a/editor/e2e/globalTeardown.ts b/editor/e2e/globalTeardown.ts
@@ -0,0 +1,17 @@
+// Last line of defence: nothing this suite spawned may outlive it.
+//
+// The three layers above this one should make strays impossible — live jobs are
+// cancelled before /api/test/invalidate-cache drops the registry, the server
+// cancels them again on graceful shutdown, and each fake binary self-terminates
+// when it is orphaned. This hook exists because "should be impossible" is not
+// the same as "is", and because the failure mode is so expensive: orphaned
+// fixtures do not announce themselves, they just quietly make the next run slow
+// and flaky in ways that read like real regressions.
+//
+// Anything reaped here is logged as a regression, not as routine cleanup.
+
+import { sweepFixtureProcesses } from "./fixtureProcs";
+
+export default async function globalTeardown() {
+ await sweepFixtureProcesses("teardown");
+}
diff --git a/editor/e2e/no-orphan-fixtures.spec.ts b/editor/e2e/no-orphan-fixtures.spec.ts
@@ -0,0 +1,97 @@
+// Regression guard: a fixture child must never outlive the job that owns it.
+//
+// The bug this pins down: /api/test/invalidate-cache used to drop the job
+// registry (globalThis.__yttJobRegistry__ = undefined) without cancelling
+// anything first. For yt-dlp-style jobs the registry's AbortController is the
+// ONLY kill path — runManagedFunction sets record.abortController but never
+// record.child (common/jobs/streamCommand.ts:306-319), so cancel()'s
+// child.kill() branch is a no-op for them. Dropping the registry therefore made
+// every in-flight child unkillable by the app.
+//
+// That route is called by resetData() and writeSettings() — 303 resetData()
+// calls across 83 of 87 spec files, many from inside test bodies while a job is
+// running. So the leak fired on essentially every full run, not (as was long
+// assumed) only on interrupted ones. A full run left 27 strays alive, each
+// spinning ~29% of a core; the suite went from 30 min / 1 failure to 1.3 h / 13
+// failures, and all 13 read like real regressions in unrelated specs.
+//
+// A test suite that lies about what broke is the actual cost, which is why this
+// is worth a dedicated spec.
+
+import { mkdir, writeFile } from "node:fs/promises";
+import { test, expect } from "@playwright/test";
+import { resetData, resolvePath, writeSettings, generateReport } from "./helpers";
+import { listFixtureProcesses } from "./fixtureProcs";
+
+async function makeTranscribeChannel(slug: string, ids: string[]) {
+ const root = resolvePath(`test-transcripts/channels/${slug}`);
+ await mkdir(root, { recursive: true });
+ await writeFile(
+ `${root}/config.json`,
+ JSON.stringify({
+ handling: "transcribe",
+ name: slug,
+ url: "https://odysee.com/@example",
+ audioFormat: "mp3",
+ }),
+ );
+ for (const id of ids) {
+ await mkdir(`${root}/data/${id}`, { recursive: true });
+ await writeFile(`${root}/data/${id}/audio.mp3`, `fake audio ${id}\n`);
+ }
+}
+
+test("resetData cancels in-flight fixture children instead of orphaning them", async ({
+ page,
+}) => {
+ test.setTimeout(90_000);
+ await resetData("empty");
+
+ // parakeet is the right lever here: for a "slowop" dir the fake wrapper waits
+ // on a 30s backstop before finishing (fixtures/bin/fake-parakeet-stitch.mjs).
+ // That is a wide enough window that a stray is unambiguous — a fixture with a
+ // sub-second sleep would "pass" simply by exiting on its own.
+ await writeSettings({
+ workers: [
+ {
+ id: "gpu",
+ name: "GPU parakeet",
+ kind: "local",
+ enabled: true,
+ priority: 0,
+ appId: "parakeet",
+ config: {},
+ },
+ ],
+ });
+ await makeTranscribeChannel("orphan-chan", ["slowoporphan1"]);
+
+ await generateReport(page, "orphan-chan");
+ await page.goto("/channels/orphan-chan");
+ await page.getByRole("button", { name: "Transcribe missing" }).click();
+
+ // Wait for the child to actually exist before trying to strand it, otherwise
+ // the assertion below could pass because nothing had spawned yet.
+ await expect
+ .poll(async () => (await listFixtureProcesses()).length, {
+ timeout: 20_000,
+ message: "expected a fixture child to be running",
+ })
+ .toBeGreaterThan(0);
+
+ // The stranding event: this hits /api/test/invalidate-cache, which wipes the
+ // registry singletons. Pre-fix, the running child survived this with nothing
+ // left holding a reference that could kill it.
+ await resetData("empty");
+
+ // Generous window: the app sends SIGTERM immediately (abort() fires execa's
+ // cancelSignal synchronously) and execa escalates to SIGKILL after 5s, so a
+ // correctly-cancelled child is gone well inside this. A child left orphaned
+ // would still be sleeping on its 30s backstop when this expires.
+ await expect
+ .poll(async () => (await listFixtureProcesses()).length, {
+ timeout: 15_000,
+ message: "fixture child outlived the registry wipe — it was orphaned",
+ })
+ .toBe(0);
+});
diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts
@@ -9,6 +9,23 @@ export async function register() {
// so guard the import — and run nothing on Edge.
if (process.env.NEXT_RUNTIME !== "nodejs") return;
+ // Cancel in-flight children on a graceful shutdown. Armed FIRST, before any of
+ // the runners below: a failure while starting those must not leave the server
+ // without its only way to reap the children they spawn.
+ //
+ // Must be a lazy import like the rest — the module uses process.once/kill/
+ // listenerCount, and inlining it here fails the Edge bundle's static Node-API
+ // scan even though the guard above means it never runs there. See
+ // common/jobs/shutdownCancel.ts.
+ try {
+ const { armShutdownCancel } = await import(
+ "yt-dlp-transcript-common/jobs/shutdownCancel"
+ );
+ armShutdownCancel();
+ } catch {
+ /* failing to arm the reaper must not block server readiness */
+ }
+
// Lazy import inside the guard keeps server-only code out of the Edge bundle.
// startSyncHeartbeat only arms a timer (no synchronous tick), so it never
// blocks the server from becoming ready.
diff --git a/editor/package.json b/editor/package.json
@@ -5,8 +5,8 @@
"type": "module",
"scripts": {
"dev": "next dev --port ${EDITOR_PORT:-3001}",
- "dev:test": "WORKER_TOKEN=test-worker-token TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json EDITOR_CHANGELOG_FILE=$(pwd)/test-changelog.md YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs GALLERY_DL_BIN=$(pwd)/e2e/fixtures/bin/fake-gallery-dl.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs FFPROBE_BIN=$(pwd)/e2e/fixtures/bin/fake-ffprobe.mjs OLLAMA_URL=http://127.0.0.1:${OLLAMA_STUB_PORT:-11435} CLAUDE_BIN=$(pwd)/e2e/fixtures/bin/fake-claude.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 AUDIO_CHECK_INTERVAL_FLOOR_MS_OVERRIDE=50 AUDIO_CHECK_RECOVER_STEP_MS_OVERRIDE=100 AUDIO_CHECK_RECOVER_AFTER_OVERRIDE=2 next dev --port ${PORT:-3011}",
- "start:test": "WORKER_TOKEN=test-worker-token TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json EDITOR_CHANGELOG_FILE=$(pwd)/test-changelog.md YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs GALLERY_DL_BIN=$(pwd)/e2e/fixtures/bin/fake-gallery-dl.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs FFPROBE_BIN=$(pwd)/e2e/fixtures/bin/fake-ffprobe.mjs OLLAMA_URL=http://127.0.0.1:${OLLAMA_STUB_PORT:-11435} CLAUDE_BIN=$(pwd)/e2e/fixtures/bin/fake-claude.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 AUDIO_CHECK_INTERVAL_FLOOR_MS_OVERRIDE=50 AUDIO_CHECK_RECOVER_STEP_MS_OVERRIDE=100 AUDIO_CHECK_RECOVER_AFTER_OVERRIDE=2 next start --port ${PORT:-3011}",
+ "dev:test": "WORKER_TOKEN=test-worker-token TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json EDITOR_CHANGELOG_FILE=$(pwd)/test-changelog.md YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs GALLERY_DL_BIN=$(pwd)/e2e/fixtures/bin/fake-gallery-dl.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null DIARIZE_BIN=$(pwd)/e2e/fixtures/bin/fake-diarize.mjs FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs FFPROBE_BIN=$(pwd)/e2e/fixtures/bin/fake-ffprobe.mjs OLLAMA_URL=http://127.0.0.1:${OLLAMA_STUB_PORT:-11435} CLAUDE_BIN=$(pwd)/e2e/fixtures/bin/fake-claude.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 AUDIO_CHECK_INTERVAL_FLOOR_MS_OVERRIDE=50 AUDIO_CHECK_RECOVER_STEP_MS_OVERRIDE=100 AUDIO_CHECK_RECOVER_AFTER_OVERRIDE=2 next dev --port ${PORT:-3011}",
+ "start:test": "WORKER_TOKEN=test-worker-token TRANSCRIPTS_DIR=$(pwd)/test-transcripts EXPORT_PUBLIC_DIR=$(pwd)/test-transcripts/.export-public SETTINGS_FILE=$(pwd)/test-settings.json EDITOR_CHANGELOG_FILE=$(pwd)/test-changelog.md YTDLP_BIN=$(pwd)/e2e/fixtures/bin/fake-ytdlp.mjs GALLERY_DL_BIN=$(pwd)/e2e/fixtures/bin/fake-gallery-dl.mjs WHISPER_BIN=$(pwd)/e2e/fixtures/bin/fake-whisper.mjs WHISPER_MODEL=/dev/null CHOUGH_BIN=$(pwd)/e2e/fixtures/bin/fake-chough.mjs CHOUGH_MODEL=/dev/null PARAKEET_STITCH_BIN=$(pwd)/e2e/fixtures/bin/fake-parakeet-stitch.mjs PARAKEET_CLI=/dev/null PARAKEET_MODEL=/dev/null DIARIZE_BIN=$(pwd)/e2e/fixtures/bin/fake-diarize.mjs FFMPEG_BIN=$(pwd)/e2e/fixtures/bin/fake-ffmpeg.mjs FFPROBE_BIN=$(pwd)/e2e/fixtures/bin/fake-ffprobe.mjs OLLAMA_URL=http://127.0.0.1:${OLLAMA_STUB_PORT:-11435} CLAUDE_BIN=$(pwd)/e2e/fixtures/bin/fake-claude.mjs AUDIO_CHECK_INTERVAL_MS_OVERRIDE=300 AUDIO_CHECK_SIZE_GATE_OVERRIDE=4096 AUDIO_CHECK_INTERVAL_FLOOR_MS_OVERRIDE=50 AUDIO_CHECK_RECOVER_STEP_MS_OVERRIDE=100 AUDIO_CHECK_RECOVER_AFTER_OVERRIDE=2 next start --port ${PORT:-3011}",
"build": "next build",
"start": "next start --port ${EDITOR_PORT:-3001}",
"lint": "eslint",
diff --git a/editor/playwright.config.ts b/editor/playwright.config.ts
@@ -38,6 +38,12 @@ const exportSitesDir = path.resolve(
export default defineConfig({
testDir: "./e2e",
timeout: 30_000,
+ // Reap fixture binaries that outlived their run — before this one starts (a
+ // SIGKILLed run never gets to clean up after itself) and again after it ends.
+ // Both log anything they find, because by construction they should find
+ // nothing. See e2e/fixtureProcs.ts.
+ globalSetup: "./e2e/globalSetup.ts",
+ globalTeardown: "./e2e/globalTeardown.ts",
// The sharded runner sets CI=true for `reuseExistingServer: !CI` below, not
// for retries, and passes an explicit --retries=0 so its failures stay
// comparable to a serial run's. See scripts/run-sharded-e2e.mjs.
diff --git a/plans/STATE.md b/plans/STATE.md
@@ -172,7 +172,7 @@ borrowed digest is labelled as borrowed.
| 6 · Ollama `/ask` provider | not started | **No dependencies.** |
| 7 · Tags + chat highlights | not started | Tag generation is now VERIFIED WORKING (it had never run once — 0 of 109 sidecars). 11 videos, 0 warnings, metrics in `digest-validate.ts`, e2e in `digest.spec.ts`. Nothing consumes it yet. Known gap for this phase: only 2.4% cross-video tag reuse, and `antisemitism`/`anti-semitism` are two tags — the vocabulary does not converge, so a consumer needs normalization. |
| 8 · Visibility policy | not started | |
-| 9 · Attribution + quote filtering | not started | |
+| 9 · Attribution + quote filtering | **capture half landed** | Only the perishable part: `diarization.json` per video, off by default, plus the cleanup guard that stops the sweep deleting audio for a transcribed-but-undiarized video. Attribution, LLM naming, badges, quote filtering, the text-only lane and snapshot buckets are all still not started — deliberately, because they can be redone from the saved JSON and the audio cannot. Spike results in `plans/diarization-spike-results.md`. |
| 10 · Lead with the derived corpus | not started | |
| 11a · Review queue | **minimum landed** | Total failures now persist warnings (they used to persist none), `buckets.digestWarnings` + a `digest_warnings` filter + an /actionable section, and duplicate-cluster confirm/reject. Approve/dismiss state deferred — needs new persistence. |
| 11b · Viewer feedback | not started | |
@@ -251,6 +251,49 @@ parallel task), Phase 0 transcription benchmark (a separate bottleneck).
## Decisions log
+**Diarization is 2–3× SLOWER than the transcription it follows, so it does not run
+inline (2026-08-07).** Measured, not assumed. GPU transcription came free from the
+corpus's own `transcribe-outcome.json` records — **3,602 videos, 10,790 audio-hours,
+221 s/audio-hour, 16.3× real time**. sherpa-onnx diarization on the CPU managed
+~500 s/audio-hour at best and 500–700 typically. The plan's decision gate said to
+prefer inline over a parallel CPU lane if throughput couldn't keep pace; it can't,
+under every reading of the data. But inline has a cost the gate didn't price — it
+makes the whole pipeline 3–4× slower and idles the GPU, which on the 4,100
+audio-hours of pending retained audio is ~10.5 days versus ~42. So what shipped is
+**both, switchable**: `diarization.enabled` arms the cleanup guard (which is what
+actually protects the audio, independent of *when* diarization runs), and
+`inlineAfterTranscribe` — default **off** — controls the post-transcribe hook. Run a
+big batch with capture on and inline off, then backfill.
+
+**The GPU is not the blocker and the card is fine — no diarization model has been
+ported to ggml (2026-08-07).** parakeet.cpp is built `GGML_VULKAN=ON` and ships
+`libggml-vulkan.so`, which is why ASR hits 16.3× on the RX 6600 XT. Diarization
+can't join it: the models ship as PyTorch (CUDA/ROCm) or ONNX
+(CUDA/ROCm/MIGraphX/OpenVINO/DirectML), and neither runtime has a Vulkan compute
+path on Linux. Ecosystem gap, not hardware limit. ROCm
+(`HSA_OVERRIDE_GFX_VERSION=10.3.0`, gfx1032) was deliberately **not** installed.
+
+**sherpa-onnx's default clustering threshold is wrong for this corpus; ship 0.9
+(2026-08-07).** On a 6-minute excerpt of a known two-person interview, the default
+0.5 produced **22** speaker clusters. Sweep: 0.4→23, 0.5→22, 0.6→17, 0.7→12,
+0.8→10, **0.9→6**, with the top two clusters at 40%/40% of talk time — the two
+hosts. It still over-splits, and that is the acceptable direction: merging clusters
+later is solvable, un-splitting a wrongly-merged one would need the audio back. Every
+sidecar records engine, both model names, version and threshold so a later pass can
+re-run selectively. Talk-time distribution is far more trustworthy than the speaker
+count — the dominant cluster was reliably the host in every file tested.
+
+**Audio deletion was never automatic, and the deadline is disk, not the transcribe
+job (2026-08-07).** `PLAN.md` implied a hook; `cleanAudioFromTranscribed` has exactly
+one call site, an operator-triggered action. What makes this urgent is that
+**161.5 GB of the 387 GB `transcripts/` tree is retained audio (836 files) on a disk
+that is 97% full with 45 GB free** — and only **82 of those 836 are transcribed**, so
+the one-shot window is genuinely still open for 754 videos. Unguarded, a parallel
+diarize lane would accumulate ~340 MB of un-diarized audio per wall hour and fill the
+disk in ~5.5 days. The guard doesn't remove that pressure; it makes it visible
+(`awaiting diarization` in the sweep log, and excluded from "Est. reclaim") instead of
+silently destructive.
+
Recorded with reasons, because these are exactly what a cold agent would otherwise
relitigate. Earlier entries (fabric deferred, ollama-direct as the structured default,
Claude Code on a second queue key, digests in their own page tree, per-section provenance,
diff --git a/plans/diarization-spike-results.md b/plans/diarization-spike-results.md
@@ -0,0 +1,318 @@
+# Diarization spike — results
+
+**Date:** 2026-08-07
+**Engine:** sherpa-onnx 1.13.4 (Python), pyannote segmentation-3.0 ONNX +
+NeMo TitaNet-small speaker embeddings, CPU.
+**Box:** i7-6700 (4c/8t, 2015), 16 GB RAM, RX 6600 XT (gfx1032).
+
+Phase A of the "Capture speaker diarization on new transcriptions" plan, which
+required evaluating on real audio before committing because *"quality here
+cannot be judged from code."* Both unknowns it named — throughput on this box,
+and quality on this content — are answered below, and both changed the design.
+
+---
+
+## READ THIS FIRST: the box was contended throughout
+
+The plan says to measure the box before timing anything. That warning earned its
+place. `/proc/loadavg` ranged from **9 to 49** during these runs, and the same
+6-minute clip took anywhere from **498 to 6340 s/audio-hour** depending on when
+it landed.
+
+Two separate sources, and the second one is a trap worth writing down:
+
+1. Another project's dev server and Playwright suite, running concurrently — the
+ hazard the plan already warned about.
+2. **Leaked `fake-ytdlp.mjs` fixture processes from this repo's own e2e runs.**
+ 27 of them had accumulated, each burning ~29% of a core for over two hours,
+ after e2e runs were interrupted. On a 4-core box that alone is most of the
+ machine. Killing them took load from 46 to 2.7.
+
+So "check the load before timing" is not enough on its own — check *what* is
+producing it. A leaked fixture pool looks exactly like someone else's job, and
+it silently degraded a full e2e run from 45 minutes to 1.3 hours and turned
+4 failures into 13.
+
+Every throughput number below is therefore a **contended** number and should be
+read as an upper bound on cost, not a clean measurement. The quality numbers are
+unaffected: clustering output does not depend on how busy the CPU is.
+
+A clean re-measurement is queued (`quiet-measure.sh` waits for load < 8). The
+design conclusion does not depend on it — see "Why the gate fires either way".
+
+---
+
+## 1. Throughput
+
+### The comparison that matters
+
+GPU transcription was measured for free from the corpus's own
+`transcribe-outcome.json` records — **3,602 videos, 10,790 audio-hours**, real
+production runs, not a benchmark:
+
+| Lane | Rate | Realtime factor |
+| --- | --- | --- |
+| parakeet.cpp ASR, Vulkan, RX 6600 XT | **221 s/audio-hour** | 16.3× |
+| sherpa-onnx diarization, CPU, best observed | ~500 s/audio-hour | ~7× |
+| sherpa-onnx diarization, CPU, typical contended | 500–700 s/audio-hour | 5–7× |
+| sherpa-onnx diarization, CPU, heavily contended | up to 6340 s/audio-hour | 0.6× |
+
+**Diarization is 2–3× slower than the transcription it would follow, at best.**
+
+### Why the gate fires either way
+
+The plan's decision gate: *"if throughput can't keep pace with GPU transcription,
+prefer inline-after-transcribe over a parallel CPU lane."*
+
+Even the single fastest observation (498 s/audio-hour, taken while load was
+lowest) is **2.3× slower** than the GPU's long-run production average. The gate
+fires under every reading of the data, so the design does not hinge on the
+pending clean measurement.
+
+But the gate's own remedy needs one correction for this corpus. Inline was
+preferred because a parallel lane lets audio pile up. That is true, and the disk
+numbers are alarming — see §3 — but inline has a cost the gate does not price:
+it makes the whole pipeline ~3–4× slower and idles the GPU while the CPU works.
+With **749 videos / 4,100 audio-hours** of retained audio pending transcription,
+that is the difference between ~10.5 days and ~42 days of wall clock.
+
+**So what shipped is both, switchable, with inline OFF by default:**
+
+- `diarization.enabled` arms the **cleanup guard** — the sweep stops deleting
+ audio for a transcribed-but-undiarized video. This is what actually protects
+ the perishable input, and it is independent of *when* diarization runs.
+- `diarization.inlineAfterTranscribe` (default **off**) runs it in the
+ post-transcribe hook, for steady state.
+- The `diarize-channel` job backfills over retained audio at full CPU speed
+ without holding up the GPU.
+
+For the batch that is about to run, the intended sequence is: enable capture
+(holds the audio), leave inline off (batch runs at full GPU speed), then backfill.
+
+### Cost of the backfill
+
+4,338 audio-hours of retained audio exist today (836 files, 161.5 GB).
+
+| Rate | Backfill wall time |
+| --- | --- |
+| 500 s/audio-hour | ~25 days |
+| 676 s/audio-hour | ~34 days |
+
+Single-threaded-lane figures. This is the same order as the digest sweep
+(~25–55 days) and competes with it for the same 8 threads.
+
+### Peak RSS
+
+482 MB on a 13-minute file. Memory is dominated by holding decoded audio as
+float32 at 16 kHz — ~230 MB per audio-hour — so a 5-hour VOD needs ~1.2 GB
+resident. On a 16 GB box with ~1 GB free and 7 GB of swap already in use, the
+long tail (the corpus has 8-hour VODs) is a real constraint. Not addressed here;
+noted for whoever runs the backfill.
+
+---
+
+## 2. Quality
+
+### The default threshold is wrong for this content
+
+sherpa-onnx defaults to a clustering threshold of 0.5. On a 6-minute excerpt of
+**MommaOcco/v1cnNSjEmZk — a two-person interview with @ProtonJon, so ground
+truth is exactly 2 speakers** — that produced **22 clusters**.
+
+Sweep on the same clip:
+
+| threshold | speakers | turns | top-4 talk-time share |
+| --- | --- | --- | --- |
+| 0.4 | 23 | 62 | 38%, 31%, 13%, 3% |
+| 0.5 (sherpa default) | 22 | 64 | 38%, 31%, 17%, 3% |
+| 0.6 | 17 | 64 | 40%, 34%, 18%, 3% |
+| 0.7 | 12 | 65 | 39%, 38%, 19%, 1% |
+| 0.8 | 10 | 64 | 39%, 39%, 20%, 1% |
+| **0.9** | **6** | 61 | **40%, 40%, 20%, 0%** |
+
+**The shipped default is 0.9.** At 0.9 the top two clusters sit at 40%/40% —
+recognizably the two hosts — with a 20% third cluster and a negligible tail.
+The *shape* is right even where the count is not.
+
+Corroboration: on a 30-second excerpt of the same interview the wrapper returned
+exactly **2 speakers**, the correct answer. Over-splitting grows with duration,
+which is the expected failure mode for agglomerative clustering over a long
+recording.
+
+### Confirmed on a full-length real file, not just an excerpt
+
+The end-to-end verification run — `ObviousRises-rumble/v6z1o2g`, the whole
+12.8-minute reaction video, through the shipped `scripts/diarize.mjs` at the
+shipped default of 0.9:
+
+| | th=0.5 | th=0.9 |
+| --- | --- | --- |
+| speakers | 29 | **13** |
+| turns | 59 | 51 |
+| dominant cluster | 70.7% | **73%** |
+
+The tail collapses (5%, 3%, 2%, 2% after the top two) while the host's share
+holds. This matters because the excerpt sweep could have been an artifact of
+6-minute clips; it is not — the same improvement shows on the full file.
+
+Cost of that run: 274 s of engine time for 12.8 minutes of audio, i.e.
+1290 s/audio-hour, at load ~30 and only 248% of the available 400% CPU. Another
+contended number, and a reminder of how much contention costs here.
+
+### Other content types
+
+| Clip | Content | th=0.9 result |
+| --- | --- | --- |
+| ObviousRises-rumble/v6z1o2g (6 min) | reaction, plays third-party clips | 8 speakers, top 60% (the host) |
+| kirsche/LEu6R1kn7Gs (6 min @ 1h in) | solo stream | 3 speakers @ th=0.4, top 60%/25%/15% |
+
+The reaction clip's 8 clusters are not obviously wrong — it genuinely contains
+several third-party voices. The solo stream returning 3 clusters at th=0.4 is
+over-split; higher thresholds were not reached for it before the run was cut.
+
+The full 13-minute reaction video at th=0.5 produced **29 speakers across 59
+turns**, with cluster IDs running non-contiguously up to 57 — near-zero merging.
+One cluster held 70.7% of talk time (the host), which is the correct shape buried
+in noise. That single result is what motivated the whole threshold sweep.
+
+### Verdict on quality
+
+`PLAN.md:442-443` set the expectation: *"this misfires on rapid back-and-forth,
+and auto-caption channels have no speaker turns at all."* That holds. What the
+spike adds:
+
+- **Talk-time distribution is more trustworthy than speaker count.** The
+ dominant cluster is reliably the host across every file tested. A downstream
+ pass that asks "which cluster is the main speaker" will do much better than one
+ that trusts the cluster count.
+- **Over-splitting is the failure mode, not under-splitting.** That is the
+ benign direction: merging clusters later is a solvable problem, and the turns
+ and their boundaries are recorded either way. Splitting a cluster that was
+ wrongly merged would need the audio back.
+- **It is worth capturing.** Not because the output is good enough to show a
+ user — it is not — but because the boundaries and the talk-time structure are
+ real, they are recoverable into something better, and they are unobtainable
+ once the audio is gone.
+
+---
+
+## 3. Disk — the pressure the guard makes visible
+
+At the time of writing: **45 GB free, 97% full**, `transcripts/` at 387 GB, of
+which **161.5 GB is retained audio** across 836 files (42%).
+
+Audio averages **33 MB per audio-hour**. If diarization ran as an unguarded
+parallel CPU lane while the GPU transcribed, the backlog would grow at roughly
+(16.3 − 6) ≈ 10 audio-hours of un-diarized audio per wall hour, or **~340 MB per
+wall hour** — filling the remaining 45 GB in about **5.5 days**.
+
+The cleanup guard does not remove that pressure. It makes it **visible and
+non-destructive** instead of silent and permanent: the sweep reports
+`awaiting diarization` and the channel's "Est. reclaim" drops to match, rather
+than the audio quietly disappearing before its diarize job ran.
+
+Also worth recording: **only 82 of the 836 retained-audio videos are
+transcribed.** The other 754 are the pending batch. So the one-shot window this
+plan exists to catch is genuinely still open.
+
+---
+
+## 4. Why not the GPU
+
+`PLAN.md` assumed the card might be the blocker. It is not, and neither is the
+GPU available.
+
+- parakeet.cpp is built `GGML_VULKAN=ON` and ships `libggml-vulkan.so`, so ASR
+ genuinely runs on the RX 6600 XT. That is why transcription hits 16.3×.
+- **No diarization model has been ported to ggml.** They ship as PyTorch
+ (CUDA/ROCm) or ONNX (CUDA/ROCm/MIGraphX/OpenVINO/DirectML). Neither runtime has
+ a Vulkan compute path on Linux.
+
+This is an ecosystem gap, not a hardware limit. The ROCm fallback
+(`HSA_OVERRIDE_GFX_VERSION=10.3.0` for this gfx1032 card) remains available but
+was **not** installed — the plan says not to do that speculatively, and the CPU
+path is sufficient for a backfill that is not on the critical path.
+
+Environment note: the system Python is **3.14.6**, and sherpa-onnx publishes
+wheels only up to cp313. The spike used a dedicated `uv`-managed 3.13 venv, which
+is why `diarization.python` is a configurable path rather than `python3`.
+
+---
+
+## 5. Reproducing
+
+```
+uv venv --python 3.13 diarize-env
+uv pip install --python ./diarize-env/bin/python sherpa-onnx numpy
+# models
+curl -sSLO https://github.com/k2-fsa/sherpa-onnx/releases/download/speaker-segmentation-models/sherpa-onnx-pyannote-segmentation-3-0.tar.bz2
+tar xf sherpa-onnx-pyannote-segmentation-3-0.tar.bz2
+curl -sSL -o titanet.onnx https://github.com/k2-fsa/sherpa-onnx/releases/download/speaker-recongition-models/nemo_en_titanet_small.onnx
+
+scripts/diarize.mjs --output diarization.json --video-id <id> \
+ --python <venv>/bin/python \
+ --seg <...>/sherpa-onnx-pyannote-segmentation-3-0/model.onnx \
+ --emb <...>/titanet.onnx --threshold 0.9 <audioFile>
+```
+
+## 6. End-to-end verification on real data
+
+`ObviousRises-rumble/v6z1o2g` — a real corpus video with a transcript and audio
+still on disk — run through the shipped wrapper, with the cleanup guard's
+decision evaluated read-only before and after (never by running the real sweep,
+which would delete the audio):
+
+```
+ before after
+transcript true true
+audio files audio.mp3 audio.mp3
+diarization false true
+sweep takes it, diarization OFF true true
+sweep takes it, diarization ON false --> true
+```
+
+That is the whole contract in four lines: with capture armed the sweep refuses
+the video until the sidecar exists, and releases it the moment it does. The
+`diarization OFF` column is unchanged in both states, which is what makes
+enabling the feature reversible rather than a one-way door.
+
+## 7. Verification
+
+On a genuinely quiet box (load 2.7, after clearing the leaked fixtures):
+
+| Check | Result |
+| --- | --- |
+| `common` unit suite | **510 / 511** |
+| Editor e2e, full | **433 / 434** in 30.1 min |
+| Diarization e2e (5 specs) | 10/10 across two repeats, and green in-suite |
+| `tsc --noEmit` in common / editor / export | clean |
+| `pnpm build` (editor) | clean — e2e runs in dev mode and never prerenders |
+
+The two residual failures are both **pre-existing and reproduced independently
+of this change**:
+
+- `digestPlan.test.ts › the corpus chunk census reproduces 191,116` — the real
+ corpus has grown to **194,053** chunks. Reproduced identically on a pristine
+ `HEAD` worktree pointed at the same corpus, so it is not this change. The
+ pinned number needs re-baselining by whoever owns the digest plan; it is
+ deliberately left alone here.
+- `jobs-batch-tasks-drain › hard Cancel during a drain` — passes 6/6 in
+ isolation. The rotating contention tail.
+
+One real regression was found and fixed during verification: appending
+`, 0 awaiting diarization` to the cleanup sweep's skip breakdown broke an exact
+string assertion in `pre-clean-availability.spec`. The clause is now emitted only
+when the count is non-zero, so the message is byte-identical for any install with
+the capture lane off.
+
+## 8. What is still open
+
+- **The clean throughput number.** Queued behind a load<8 wait. Every figure
+ here is contended.
+- **Thread scaling.** 2 vs 4 vs 8 threads on a 4c/8t part was not measured;
+ queued with the above.
+- **Peak RSS on a genuinely long file.** Extrapolated (~230 MB per audio-hour),
+ not measured on an 8-hour VOD.
+- **Whether a better embedding model closes the over-splitting gap.** Only
+ TitaNet-small was tried. 3D-Speaker / WeSpeaker English models are the obvious
+ next comparison, and swapping one is a settings change, not a code change.
diff --git a/scripts/diarize-sherpa.py b/scripts/diarize-sherpa.py
@@ -0,0 +1,113 @@
+#!/usr/bin/env python3
+"""sherpa-onnx speaker-diarization engine for scripts/diarize.mjs.
+
+Kept deliberately dumb: decode audio, run pyannote segmentation + a speaker
+embedding model, cluster, print raw turns as JSON on stdout. All policy
+(filenames, provenance shape, atomic writes) lives in the .mjs wrapper, so
+swapping this engine for pyannote later means replacing one file.
+
+stdout is JSON ONLY. Progress goes to stderr so the wrapper can stream it.
+
+Usage:
+ diarize-sherpa.py --seg <model.onnx> --emb <model.onnx> [--threshold F]
+ [--threads N] [--ffmpeg PATH] <audioFile>
+
+Output: {"turns":[{"start":s,"end":s,"speaker":i}], "audioSeconds":f,
+ "version":"...", "sampleRate":n}
+"""
+import argparse, json, subprocess, sys, time
+
+
+def main():
+ ap = argparse.ArgumentParser()
+ ap.add_argument("audio")
+ ap.add_argument("--seg", required=True)
+ ap.add_argument("--emb", required=True)
+ ap.add_argument("--threshold", type=float, default=0.5)
+ ap.add_argument("--threads", type=int, default=4)
+ ap.add_argument("--num-speakers", type=int, default=-1)
+ ap.add_argument("--min-duration-on", type=float, default=0.3)
+ ap.add_argument("--min-duration-off", type=float, default=0.5)
+ ap.add_argument("--ffmpeg", default="ffmpeg")
+ args = ap.parse_args()
+
+ try:
+ import numpy as np
+ import sherpa_onnx
+ except ImportError as e:
+ print(f"diarize-sherpa: missing dependency: {e}", file=sys.stderr)
+ return 3
+
+ cfg = sherpa_onnx.OfflineSpeakerDiarizationConfig(
+ segmentation=sherpa_onnx.OfflineSpeakerSegmentationModelConfig(
+ pyannote=sherpa_onnx.OfflineSpeakerSegmentationPyannoteModelConfig(
+ model=args.seg
+ ),
+ num_threads=args.threads,
+ ),
+ embedding=sherpa_onnx.SpeakerEmbeddingExtractorConfig(
+ model=args.emb, num_threads=args.threads
+ ),
+ clustering=sherpa_onnx.FastClusteringConfig(
+ num_clusters=args.num_speakers, threshold=args.threshold
+ ),
+ min_duration_on=args.min_duration_on,
+ min_duration_off=args.min_duration_off,
+ )
+ if not cfg.validate():
+ print("diarize-sherpa: invalid engine config", file=sys.stderr)
+ return 2
+
+ sd = sherpa_onnx.OfflineSpeakerDiarization(cfg)
+
+ # Decode to the rate the models expect. ffmpeg reads any container, which is
+ # what lets this run against a persisted source video as well as audio.mp3.
+ proc = subprocess.run(
+ [args.ffmpeg, "-v", "error", "-i", args.audio, "-f", "f32le",
+ "-ac", "1", "-ar", str(sd.sample_rate), "-"],
+ stdout=subprocess.PIPE, stderr=subprocess.PIPE,
+ )
+ if proc.returncode != 0:
+ sys.stderr.write(proc.stderr.decode("utf8", "replace"))
+ print(f"diarize-sherpa: ffmpeg failed ({proc.returncode})", file=sys.stderr)
+ return 4
+ samples = np.frombuffer(proc.stdout, dtype=np.float32)
+ if samples.size == 0:
+ print("diarize-sherpa: decoded zero samples", file=sys.stderr)
+ return 5
+
+ audio_seconds = samples.size / sd.sample_rate
+ print(
+ f"diarize-sherpa: {audio_seconds/60:.1f} min decoded, diarizing…",
+ file=sys.stderr, flush=True,
+ )
+
+ t0 = time.time()
+ result = sd.process(samples).sort_by_start_time()
+ elapsed = time.time() - t0
+
+ turns = [
+ {"start": round(s.start, 3), "end": round(s.end, 3), "speaker": s.speaker}
+ for s in result
+ ]
+ speakers = len({t["speaker"] for t in turns})
+ print(
+ f"diarize-sherpa: {len(turns)} turns, {speakers} speakers in "
+ f"{elapsed:.1f}s ({elapsed/max(audio_seconds/3600,1e-9):.0f} s/audio-hour)",
+ file=sys.stderr, flush=True,
+ )
+
+ json.dump(
+ {
+ "turns": turns,
+ "audioSeconds": round(audio_seconds, 3),
+ "sampleRate": sd.sample_rate,
+ "version": getattr(sherpa_onnx, "__version__", None),
+ },
+ sys.stdout,
+ )
+ return 0
+
+
+if __name__ == "__main__":
+ sys.exit(main())
diff --git a/scripts/diarize.mjs b/scripts/diarize.mjs
@@ -0,0 +1,173 @@
+#!/usr/bin/env node
+// Standalone speaker diarization: audio in, diarization.json out.
+//
+// Wraps a diarization engine the same way scripts/parakeet-stitch.mjs wraps
+// parakeet-cli, and for the same reason: the app only ever knows a PATH
+// (paths.diarizeBin / DIARIZE_BIN), so the engine behind it can be swapped
+// without touching a caller. Today that engine is sherpa-onnx via
+// scripts/diarize-sherpa.py (ONNX, CPU); pyannote would be a drop-in
+// replacement for the --engine command.
+//
+// WHY CPU AND NOT THE GPU: parakeet.cpp is built GGML_VULKAN=ON so ASR really
+// does run on the RX 6600 XT, but no diarization model has been ported to ggml.
+// They ship as PyTorch (CUDA/ROCm) or ONNX (CUDA/ROCm/MIGraphX/OpenVINO/
+// DirectML), and neither runtime has a Vulkan compute path on Linux. That is an
+// ecosystem gap, not a hardware limit.
+//
+// Dual-use, by design:
+// * In the app: controller/diarizeOne.ts runs this with cwd == the video dir
+// and an explicit --output, then the caller treats the file as the artifact.
+// * On the CLI: run it directly. With no --output it prints JSON to stdout.
+//
+// Output is a DiarizationRecord (see common/lib/diarization.ts) — speaker-
+// labelled time ranges plus engine provenance, so a later attribution pass can
+// tell what produced a given file and re-run selectively.
+//
+// Usage:
+// diarize.mjs [options] <audioFile>
+//
+// Options (env fallback in parens):
+// --output <f> write the record here (default: stdout)
+// --video-id <id> recorded in the sidecar (default: basename of cwd)
+// --engine <cmd> engine command (DIARIZE_ENGINE_CMD)
+// --python <path> python for the default engine (DIARIZE_PYTHON, "python3")
+// --seg <path> segmentation model (DIARIZE_SEG_MODEL) [required]
+// --emb <path> speaker-embedding model (DIARIZE_EMB_MODEL) [required]
+// --threshold <f> clustering threshold (DIARIZE_THRESHOLD, 0.5)
+// --threads <n> engine threads (DIARIZE_THREADS, 4)
+// --ffmpeg <path> ffmpeg binary (FFMPEG_BIN, "ffmpeg")
+// -h, --help
+
+import { spawn } from "node:child_process";
+import { rename, writeFile } from "node:fs/promises";
+import path from "node:path";
+
+const HELP = `diarize.mjs [options] <audioFile>
+
+Runs speaker diarization and writes a DiarizationRecord JSON sidecar.
+See the header of this file for the full option list.`;
+
+function fail(msg, code = 2) {
+ process.stderr.write(`diarize: ${msg}\n`);
+ process.exit(code);
+}
+
+const argv = process.argv.slice(2);
+if (argv.includes("-h") || argv.includes("--help")) {
+ process.stdout.write(HELP + "\n");
+ process.exit(0);
+}
+
+function arg(flag, fallback) {
+ const i = argv.indexOf(flag);
+ return i < 0 || i === argv.length - 1 ? fallback : argv[i + 1];
+}
+
+const positional = [];
+for (let i = 0; i < argv.length; i++) {
+ if (argv[i].startsWith("-")) {
+ i++; // every flag here takes a value
+ continue;
+ }
+ positional.push(argv[i]);
+}
+
+const audio = positional[0];
+if (!audio) fail("missing <audioFile>");
+
+const output = arg("--output", arg("-o", undefined));
+const videoId = arg("--video-id", path.basename(process.cwd()));
+const python = arg("--python", process.env.DIARIZE_PYTHON ?? "python3");
+const seg = arg("--seg", process.env.DIARIZE_SEG_MODEL ?? "");
+const emb = arg("--emb", process.env.DIARIZE_EMB_MODEL ?? "");
+const threshold = Number(
+ arg("--threshold", process.env.DIARIZE_THRESHOLD ?? "0.5"),
+);
+const threads = Number(arg("--threads", process.env.DIARIZE_THREADS ?? "4"));
+const ffmpeg = arg("--ffmpeg", process.env.FFMPEG_BIN ?? "ffmpeg");
+
+const engineCmd = arg("--engine", process.env.DIARIZE_ENGINE_CMD ?? "");
+
+let cmd;
+let args;
+if (engineCmd) {
+ const parts = engineCmd.split(" ").filter(Boolean);
+ cmd = parts[0];
+ args = [...parts.slice(1), audio];
+} else {
+ if (!seg) fail("missing --seg (or DIARIZE_SEG_MODEL)");
+ if (!emb) fail("missing --emb (or DIARIZE_EMB_MODEL)");
+ cmd = python;
+ args = [
+ path.join(import.meta.dirname, "diarize-sherpa.py"),
+ "--seg", seg,
+ "--emb", emb,
+ "--threshold", String(threshold),
+ "--threads", String(threads),
+ "--ffmpeg", ffmpeg,
+ audio,
+ ];
+}
+
+const started = Date.now();
+const child = spawn(cmd, args, { stdio: ["ignore", "pipe", "inherit"] });
+
+let stdout = "";
+child.stdout.setEncoding("utf8");
+child.stdout.on("data", (d) => {
+ stdout += d;
+});
+
+child.on("error", (err) => fail(`failed to run ${cmd}: ${err.message}`, 3));
+
+child.on("close", async (code) => {
+ if (code !== 0) fail(`engine exited ${code}`, code ?? 1);
+
+ let raw;
+ try {
+ raw = JSON.parse(stdout);
+ } catch {
+ fail("engine produced no parseable JSON on stdout");
+ }
+ if (!Array.isArray(raw?.turns)) fail("engine output has no `turns` array");
+
+ const turns = raw.turns.map((t) => ({
+ start: Number(t.start),
+ end: Number(t.end),
+ speaker: Number(t.speaker),
+ }));
+
+ const record = {
+ videoId,
+ generatedAt: new Date().toISOString(),
+ ...(Number.isFinite(raw.audioSeconds)
+ ? { audioSeconds: raw.audioSeconds }
+ : {}),
+ durationMs: Date.now() - started,
+ speakers: new Set(turns.map((t) => t.speaker)).size,
+ turns,
+ engine: {
+ engine: engineCmd ? path.basename(cmd) : "sherpa-onnx",
+ // basename only: absolute paths are machine-specific and the corpus is
+ // rsynced between shards, so a full path would make records non-portable.
+ ...(seg ? { segmentationModel: path.basename(seg) } : {}),
+ ...(emb ? { embeddingModel: path.basename(emb) } : {}),
+ ...(raw.version ? { version: String(raw.version) } : {}),
+ threshold,
+ },
+ };
+
+ const json = JSON.stringify(record) + "\n";
+ if (!output) {
+ process.stdout.write(json);
+ return;
+ }
+ // Atomic: a half-written sidecar must never be what convinces the cleanup
+ // sweep it is safe to delete the only copy of the audio.
+ const tmp = `${output}.tmp-${process.pid}`;
+ await writeFile(tmp, json);
+ await rename(tmp, output);
+ process.stderr.write(
+ `diarize: ${record.turns.length} turns, ${record.speakers} speakers -> ${output}\n`,
+ );
+});