commit 77bba40a5ec8b6b22a1ea02a4d3ed932e68820e6
parent f08840763fb5fa396f379cb7067941918e3885a0
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 8 Aug 2026 23:59:43 -0400
Window the long files, and scan what is already on disk
PART 3 -- the point of the whole branch. Two independent memory sinks, both
measured, both closed by the same change:
1. sherpa's clustering holds an O(n^2) pairwise matrix over speech SEGMENTS.
Density varies 40x here, which is why duration never predicted the kills:
a sparse 7h42m file finished while a dense 6h12m one died.
2. subprocess.run(stdout=PIPE) buffered the entire decoded stream (1.84 GB for
8h, peaking near 2x while the buffer grows).
Windowing divides per-window n by the window count, so the matrix falls by its
SQUARE, and one window is decoded at a time -- into a PREALLOCATED array via
readinto, so even the per-window buffer never doubles.
Cross-window identity is the hard part and it is built, not hand-waved: each
window's local speakers are reduced to an averaged, L2-normalized
SpeakerEmbeddingExtractor centroid over their longest turns, and the centroids
are clustered globally with FastClustering at the SAME threshold. Both APIs
verified live against the installed 1.13.4 (dim 192; FastClustering is callable
and exposes no named methods).
GATES, in the order the plan set them:
(a) ObviousRises-rumble/v6z1o2g re-runs BYTE-IDENTICAL with windowing not
triggered -- 51 turns, 13 speakers, every non-volatile field equal, no
`windowing` key. The whole-file branch is the original code untouched.
(b) hasanabi/2834644359 (6h12m), which the kernel killed at 10.5 GB,
COMPLETES: 3,433 turns, 96 speakers, 46 min, peak RSS ~2.0 GB sampled
from /proc/<pid>/status (GNU time is not installed here).
(c) windowed vs whole-file on a file that survives both: speech total -0.1%,
identical wall-clock coverage, identity MORE concentrated (top-1 88.2% vs
73.1%, 6 speakers vs 13) rather than fragmented at seams.
THE GATE CAUGHT A REAL DEFECT THAT READING THE CODE DID NOT. The first seam
merge joined only CONSECUTIVE turns in the sorted list. That is not the same as
unioning a speaker's own intervals: two local speakers can map to one global
speaker and their turns then interleave with other people's, leaving the same
speaker overlapping THEMSELVES -- 35 self-overlaps and 24 sub-gap splits on the
6h12m file, double-counting talk time in every downstream consumer. Now a
per-speaker interval union, re-verified at 0 and 0. Cross-speaker overlap is
preserved: that is genuine overlapping speech and the whole-file path emits it.
Windowing is recorded on the record but is NOT in the freshness identity --
top-level, outside `engine`. The target is built from settings alone and cannot
know a video's duration, so a window field there would mark every sidecar on
disk stale for work that is unchanged. diarizeOne now passes --ffmpeg/--ffprobe
explicitly rather than relying on env inheritance, because ffprobe is what
decides whether to window at all.
The cap from the previous commit now defaults to 0. It was a stopgap for this.
PART 4 -- the probes existed and were standalone; nothing had ever run them
against files ALREADY on disk. Tiered because tier 2 is a full decode: tier 1
reads the container (milliseconds), tier 2 decodes only files whose duration
disagrees with the metadata, and is opt-in. Unresolved is `unknown`, never
`corrupt` -- verifyBeforeClean's rule. Deletes nothing.
It deliberately scans every MEDIA-EXTENSION file rather than only recognized
ones, because the file it was written for is precisely one the fixed predicates
exclude. First corpus run, 953 files / 68 channels / 83 s:
* the known 4.04 GiB source-media.temp.mp4 (moov atom not found)
* 3 ZERO-BYTE audio.mp3 in leaflit-rumble -- these satisfy the name-based
predicate, so destinationExists reads them as downloaded and the pipeline
never re-fetches; one sits beside a real audio.aac the mp3-first read
preference would pass over
* 39 transcript.m4a in community-notes (1.17 GiB) -- audio under the SUBTITLE
template, invisible in both directions (SUB_FILE_RE needs two segments)
* 7 duration mismatches, up to 5h, awaiting the opt-in decode
Report is versioned + atomic at the corpus root with PER-CHANNEL totals, so a
per-channel re-scan replaces one entry without the headline describing a scan it
did not do. Review decisions live in a sibling overrides file, because the report
is regenerated wholesale. Job kinds registered properly.
Verified: common 625/625; tsc clean in common and editor; pnpm build clean;
editor e2e 449/451 (the two are races in auto-queue picks and drain-cancel,
neither touching this work -- recordPick fires at DISPATCH, so that spec asserts
transcripts exist before the unit can have written them).
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
21 files changed, 1870 insertions(+), 364 deletions(-)
diff --git a/common/controller/backfillBatch.test.ts b/common/controller/backfillBatch.test.ts
@@ -185,6 +185,8 @@ test("the duration cap round-trips, and 0 means off", () => {
sanitizeDiarization({ maxAudioHours: -2 }).maxAudioHours,
defaultDiarization().maxAudioHours,
);
- // The shipped default: every video that OOM-killed this box was over 6h.
- assert.equal(defaultDiarization().maxAudioHours, 4);
+ // The shipped default is OFF: windowed diarization removed the OOM the cap
+ // was a stopgap for. This is asserted rather than assumed because turning it
+ // back on by accident would silently stop diarizing long videos.
+ assert.equal(defaultDiarization().maxAudioHours, 0);
});
diff --git a/common/controller/diarizeOne.ts b/common/controller/diarizeOne.ts
@@ -123,6 +123,14 @@ export async function diarizeOneVideo(
String(cfg.threads),
"--python",
cfg.python,
+ // Passed explicitly rather than left to env inheritance. ffprobe is what
+ // decides whether a file is long enough to window, and a wrapper that
+ // silently fell back to a bare "ffprobe" on PATH would quietly take the
+ // whole-file path on exactly the long recordings windowing exists for.
+ "--ffmpeg",
+ opts.paths.ffmpegBin,
+ "--ffprobe",
+ opts.paths.ffprobeBin,
media,
],
{
diff --git a/common/controller/scanCorruptMedia.ts b/common/controller/scanCorruptMedia.ts
@@ -0,0 +1,458 @@
+// Scan media already on disk for corruption. Reports; deletes nothing.
+//
+// ALMOST NONE OF THIS IS NEW CODE, and that is the point. The probes have
+// existed since the audio-checked download landed and are already standalone
+// (bin, file, signal) -> result functions with no download coupling. They were
+// simply never run against files that are ALREADY here, which is precisely the
+// gap: a container that arrived corrupt, or was truncated by a disk that filled,
+// is invisible until something tries to read it — and by then the audio may be
+// the only copy.
+//
+// TIERED, BECAUSE TIER 2 IS A FULL DECODE. probeAudioStream transcodes the whole
+// file; run over every audio-hour on disk that is hours of CPU, so it has to be
+// budgeted like the digest sweep rather than like a stat() walk.
+//
+// Tier 1 (cheap, always): ask ffprobe for the container's duration. Costs
+// milliseconds per file and catches today's known case outright — a
+// truncated mp4 fails with "moov atom not found" immediately.
+// Tier 2 (full decode, opt-in): only for files that PASS tier 1 but whose
+// duration disagrees with metadata.info.json. This tier is not optional for
+// correctness: ffmpegStreamProbe.ts records that a truncated container exits
+// 0 and classifies as `partial`, so comparing durations is the only way to
+// catch truncation at all.
+//
+// The rule from verifyBeforeClean.ts is kept: anything that cannot be resolved
+// is reported `unknown` and never acted on. A missing ffprobe makes this scan
+// say "I could not tell", not "these files are corrupt".
+//
+// DELETION IS A SEPARATE, EXPLICIT CLICK. The house convention here is report
+// first, delete second — not a dry-run flag on a deleting command.
+
+import path from "node:path";
+import { mkdir, readdir, readFile, rename, stat, writeFile } from "node:fs/promises";
+import { execa } from "execa";
+import type { Paths } from "../lib/paths";
+import {
+ MEDIA_EXTS,
+ isRealAudioFile,
+ isSourceMediaFile,
+} from "../lib/mediaFiles";
+import { META_FILENAME } from "../lib/videoStatus";
+import { probeAudioStream } from "../ytdlp/ffmpegStreamProbe";
+import {
+ DEFAULT_DURATION_TOLERANCE_RATIO,
+ DEFAULT_DURATION_TOLERANCE_SEC,
+ MEDIA_SCAN_FILENAME,
+ MEDIA_SCAN_OVERRIDES_FILENAME,
+ MEDIA_SCAN_OVERRIDES_VERSION,
+ MEDIA_SCAN_REPORT_VERSION,
+ sanitizeMediaScanOverrides,
+ type MediaScanChannelTotals,
+ type MediaScanFinding,
+ type MediaScanOverrides,
+ type MediaScanReport,
+ type MediaScanVerdict,
+} from "../lib/mediaScan";
+
+const MEDIA_EXT_SET = new Set(MEDIA_EXTS);
+
+// ffprobe exited non-zero. Whether that means "this file is broken" or "I could
+// not run properly" is the difference between a finding and a false accusation,
+// so it is decided from what ffprobe actually said rather than from the exit
+// code alone. Anything unmatched falls through to `unknown`.
+const UNREADABLE_PATTERNS = [
+ /moov atom not found/i,
+ /Invalid data found when processing input/i,
+ /End of file/i,
+ /Invalid argument/i,
+ /could not find codec parameters/i,
+ /Format .* detected only with low score/i,
+];
+
+export type ScanCorruptMediaOptions = {
+ paths: Paths;
+ // Channels to scan. Absent means every channel.
+ channels?: string[];
+ // Run the tier-2 full decode on files whose duration is suspect. OFF by
+ // default: it is a real transcode per file.
+ deepProbe?: boolean;
+ toleranceSeconds?: number;
+ toleranceRatio?: number;
+ onLog?: (msg: string) => void;
+ onProgress?: (done: number, total: number) => void;
+ signal?: AbortSignal;
+};
+
+function emptyTotals(scannedAt: string): MediaScanChannelTotals {
+ return {
+ videosScanned: 0,
+ filesScanned: 0,
+ ok: 0,
+ unreadable: 0,
+ truncated: 0,
+ stray: 0,
+ unknown: 0,
+ deepProbed: 0,
+ bytesAtRisk: 0,
+ scannedAt,
+ };
+}
+
+function fileExt(name: string): string {
+ const dot = name.lastIndexOf(".");
+ return dot >= 0 ? name.slice(dot + 1).toLowerCase() : "";
+}
+
+// Every file in a video dir that CLAIMS to be media by its extension —
+// deliberately wider than isRealAudioFile/isSourceMediaFile. Those predicates
+// answer "is this one of our finished outputs?", and the answer for the corrupt
+// file this scan was written for (source-media.temp.mp4, a yt-dlp postprocessor
+// scratch, 4 GiB, unreadable) is NO. Scanning only recognized files would miss
+// exactly the population most likely to be broken.
+function mediaCandidates(entries: string[]): string[] {
+ return entries.filter((e) => MEDIA_EXT_SET.has(fileExt(e))).sort();
+}
+
+type ContainerProbe = {
+ seconds: number | null;
+ verdict: "ok" | "unreadable" | "unknown";
+ detail?: string;
+};
+
+// Tier 1. Deliberately NOT probeMediaDurationSec: that helper collapses "the
+// container is unreadable" and "I could not run ffprobe" into the same null,
+// which is the one distinction this scan exists to make.
+async function probeContainer(
+ ffprobeBin: string,
+ file: string,
+ signal?: AbortSignal,
+): Promise<ContainerProbe> {
+ const args = [
+ "-v",
+ "error",
+ "-show_entries",
+ "format=duration",
+ "-of",
+ "default=noprint_wrappers=1:nokey=1",
+ file,
+ ];
+ let result;
+ try {
+ result = await execa(ffprobeBin, args, {
+ cancelSignal: signal,
+ reject: false,
+ });
+ } catch (err) {
+ // ffprobe itself could not be run (missing binary, cancelled). Not a verdict
+ // about the file.
+ return {
+ seconds: null,
+ verdict: "unknown",
+ detail: `ffprobe could not run: ${(err as Error).message}`,
+ };
+ }
+ const stderr = String(result.stderr ?? "").trim();
+ const firstLine = stderr.split("\n").find((l) => l.trim()) ?? "";
+ if (result.exitCode !== 0) {
+ const unreadable = UNREADABLE_PATTERNS.some((re) => re.test(stderr));
+ return {
+ seconds: null,
+ verdict: unreadable ? "unreadable" : "unknown",
+ detail: firstLine || `ffprobe exited ${result.exitCode}`,
+ };
+ }
+ const sec = Number.parseFloat(String(result.stdout).trim());
+ if (!Number.isFinite(sec) || sec <= 0) {
+ return {
+ seconds: null,
+ verdict: "unknown",
+ detail: firstLine || "ffprobe reported no usable duration",
+ };
+ }
+ return { seconds: sec, verdict: "ok" };
+}
+
+async function readMetadataDuration(videoDir: string): Promise<number | null> {
+ try {
+ const raw = await readFile(path.join(videoDir, META_FILENAME), "utf8");
+ const d = (JSON.parse(raw) as { duration?: unknown }).duration;
+ return typeof d === "number" && Number.isFinite(d) && d > 0 ? d : null;
+ } catch {
+ return null;
+ }
+}
+
+export async function scanCorruptMedia(
+ opts: ScanCorruptMediaOptions,
+): Promise<MediaScanReport> {
+ const log = opts.onLog ?? ((m: string) => console.log(m));
+ const deepProbe = opts.deepProbe === true;
+ const toleranceSeconds = opts.toleranceSeconds ?? DEFAULT_DURATION_TOLERANCE_SEC;
+ const toleranceRatio = opts.toleranceRatio ?? DEFAULT_DURATION_TOLERANCE_RATIO;
+ const startedAt = Date.now();
+ const scannedAt = new Date().toISOString();
+
+ const allChannels = (
+ await readdir(opts.paths.channelsDir).catch(() => [] as string[])
+ ).sort();
+ const channels =
+ opts.channels && opts.channels.length > 0
+ ? allChannels.filter((c) => opts.channels!.includes(c))
+ : allChannels;
+
+ log(
+ `Media scan over ${channels.length} channel(s), ` +
+ (deepProbe
+ ? "with the full-decode probe for suspect durations."
+ : "container read only (the full-decode probe is off)."),
+ );
+
+ const findings: MediaScanFinding[] = [];
+ const totals: Record<string, MediaScanChannelTotals> = {};
+ let done = 0;
+
+ for (const channelSlug of channels) {
+ if (opts.signal?.aborted) {
+ log("Cancelled.");
+ break;
+ }
+ const dataDir = path.join(opts.paths.channelsDir, channelSlug, "data");
+ const ids = (await readdir(dataDir).catch(() => [] as string[])).sort();
+ const t = emptyTotals(scannedAt);
+
+ for (const videoId of ids) {
+ if (opts.signal?.aborted) break;
+ const videoDir = path.join(dataDir, videoId);
+ const entries = await readdir(videoDir).catch(() => [] as string[]);
+ const candidates = mediaCandidates(entries);
+ if (candidates.length === 0) continue;
+ t.videosScanned++;
+ let metadataSeconds: number | null | undefined;
+
+ for (const file of candidates) {
+ if (opts.signal?.aborted) break;
+ const full = path.join(videoDir, file);
+ const bytes = await stat(full)
+ .then((s) => s.size)
+ .catch(() => 0);
+ t.filesScanned++;
+
+ const recognized = isRealAudioFile(file) || isSourceMediaFile(file);
+ const push = (
+ verdict: Exclude<MediaScanVerdict, "ok">,
+ tier: 1 | 2,
+ containerSeconds: number | null,
+ detail?: string,
+ ) => {
+ t[verdict]++;
+ t.bytesAtRisk += bytes;
+ findings.push({
+ slug: `${channelSlug}/${videoId}`,
+ channelSlug,
+ videoId,
+ file,
+ bytes,
+ verdict,
+ tier,
+ containerSeconds,
+ metadataSeconds: metadataSeconds ?? null,
+ ...(detail ? { detail: detail.slice(0, 300) } : {}),
+ });
+ };
+
+ // Tier 1.
+ const probe = await probeContainer(
+ opts.paths.ffprobeBin,
+ full,
+ opts.signal,
+ );
+ if (probe.verdict === "unreadable") {
+ push("unreadable", 1, null, probe.detail);
+ continue;
+ }
+ if (probe.verdict === "unknown") {
+ push("unknown", 1, null, probe.detail);
+ continue;
+ }
+
+ // Readable. An unrecognized name is reported as a leftover rather than
+ // as damage — it decodes fine, it just is not one of ours.
+ if (!recognized) {
+ push(
+ "stray",
+ 1,
+ probe.seconds,
+ "Readable media the app does not recognize as one of its outputs (scratch or temp leftover).",
+ );
+ continue;
+ }
+
+ if (metadataSeconds === undefined) {
+ metadataSeconds = await readMetadataDuration(videoDir);
+ }
+ const expected = metadataSeconds;
+ const actual = probe.seconds ?? 0;
+ const suspect =
+ expected !== null &&
+ expected - actual >
+ Math.max(toleranceSeconds, expected * toleranceRatio);
+
+ if (!suspect) {
+ t.ok++;
+ continue;
+ }
+
+ // Tier 2, and ONLY here. Without the deep probe this stays a suspicion,
+ // reported honestly as one rather than promoted to a verdict the cheap
+ // tier cannot support.
+ if (!deepProbe) {
+ push(
+ "unknown",
+ 1,
+ probe.seconds,
+ `Container is ${Math.round(expected - actual)}s shorter than the metadata; run the full-decode probe to confirm.`,
+ );
+ continue;
+ }
+ t.deepProbed++;
+ const deep = await probeAudioStream({
+ ffmpegBin: opts.paths.ffmpegBin,
+ file: full,
+ signal: opts.signal ?? new AbortController().signal,
+ });
+ if (deep.verdict === "malformed") {
+ push("unreadable", 2, probe.seconds, deep.stderr.split("\n")[0]);
+ } else {
+ // `clean` AND `partial` both land here: ffmpeg exits 0 on a truncated
+ // container, so the duration gap — not ffmpeg's verdict — is what says
+ // the file is short.
+ push(
+ "truncated",
+ 2,
+ probe.seconds,
+ `Decodes, but ${Math.round(expected - actual)}s shorter than the metadata (ffmpeg: ${deep.verdict}).`,
+ );
+ }
+ }
+ done++;
+ if (done % 500 === 0) opts.onProgress?.(done, done);
+ }
+ totals[channelSlug] = t;
+ if (t.filesScanned > 0) {
+ log(
+ `${channelSlug}: ${t.filesScanned} file(s) in ${t.videosScanned} video dir(s) — ` +
+ `${t.ok} ok, ${t.unreadable} unreadable, ${t.truncated} truncated, ` +
+ `${t.stray} stray, ${t.unknown} unknown.`,
+ );
+ }
+ }
+
+ const report = await mergeAndWriteReport(opts.paths, {
+ version: MEDIA_SCAN_REPORT_VERSION,
+ generatedAt: scannedAt,
+ runConfig: { deepProbe, toleranceSeconds, toleranceRatio },
+ channels: totals,
+ findings,
+ });
+
+ const flagged = report.findings.length;
+ log(
+ `Done: ${flagged} finding(s) across the scanned channels → ${MEDIA_SCAN_FILENAME} ` +
+ `in ${Math.round((Date.now() - startedAt) / 1000)}s. Nothing was deleted.`,
+ );
+ return report;
+}
+
+// A per-channel re-scan must not wipe every other channel's result, and a
+// corpus-wide scan must not leave a stale channel behind. So: replace the
+// scanned channels' totals and findings, keep the rest.
+async function mergeAndWriteReport(
+ paths: Paths,
+ fresh: MediaScanReport,
+): Promise<MediaScanReport> {
+ const prior = await readMediaScanReport(paths);
+ const scanned = new Set(Object.keys(fresh.channels));
+ const channels = { ...(prior?.channels ?? {}) };
+ for (const [slug, totals] of Object.entries(fresh.channels)) {
+ channels[slug] = totals;
+ }
+ const kept = (prior?.findings ?? []).filter(
+ (f) => !scanned.has(f.channelSlug),
+ );
+ const merged: MediaScanReport = {
+ ...fresh,
+ channels,
+ findings: [...kept, ...fresh.findings].sort(
+ (a, b) =>
+ b.bytes - a.bytes || a.slug.localeCompare(b.slug) ||
+ a.file.localeCompare(b.file),
+ ),
+ };
+ const outPath = path.join(paths.transcriptsDir, MEDIA_SCAN_FILENAME);
+ await mkdir(path.dirname(outPath), { recursive: true });
+ const tmp = `${outPath}.tmp-${process.pid}`;
+ await writeFile(tmp, JSON.stringify(merged));
+ await rename(tmp, outPath);
+ return merged;
+}
+
+export async function readMediaScanReport(
+ paths: Paths,
+): Promise<MediaScanReport | null> {
+ try {
+ const raw = await readFile(
+ path.join(paths.transcriptsDir, MEDIA_SCAN_FILENAME),
+ "utf8",
+ );
+ const parsed = JSON.parse(raw) as MediaScanReport;
+ if (typeof parsed.version !== "number") return null;
+ if (!Array.isArray(parsed.findings)) return null;
+ return parsed;
+ } catch {
+ return null;
+ }
+}
+
+export function mediaScanOverridesPath(paths: Paths): string {
+ return path.join(paths.transcriptsDir, MEDIA_SCAN_OVERRIDES_FILENAME);
+}
+
+export async function readMediaScanOverrides(
+ paths: Paths,
+): Promise<MediaScanOverrides> {
+ try {
+ const raw = await readFile(mediaScanOverridesPath(paths), "utf8");
+ return sanitizeMediaScanOverrides(JSON.parse(raw));
+ } catch {
+ return sanitizeMediaScanOverrides(null);
+ }
+}
+
+// Record one finding as reviewed (or clear it), preserving every other decision.
+// Read-modify-write with the same atomic tmp+rename the report uses.
+export async function updateMediaScanOverride(
+ paths: Paths,
+ key: string,
+ patch: { reviewed: boolean; note?: string },
+): Promise<MediaScanOverrides> {
+ const current = await readMediaScanOverrides(paths);
+ const reviewed = { ...current.reviewed };
+ if (patch.reviewed) {
+ reviewed[key] = {
+ reviewedAt: new Date().toISOString(),
+ ...(patch.note ? { note: patch.note } : {}),
+ };
+ } else {
+ delete reviewed[key];
+ }
+ const out: MediaScanOverrides = {
+ version: MEDIA_SCAN_OVERRIDES_VERSION,
+ reviewed,
+ };
+ const file = mediaScanOverridesPath(paths);
+ await mkdir(path.dirname(file), { recursive: true });
+ const tmp = `${file}.tmp-${process.pid}`;
+ await writeFile(tmp, JSON.stringify(out, null, 2) + "\n");
+ await rename(tmp, file);
+ return out;
+}
diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts
@@ -223,6 +223,29 @@ const JOB_KINDS: Record<string, JobKindMeta> = {
bookmarkable: true,
queueKeyStrategy: "custom",
},
+ // The corpus-wide corrupt-media scan and its per-channel twin. Registered
+ // properly, unlike check-availability / refresh-report / detect-duplicates,
+ // which have no entry here and fall back to their raw machine kind in the UI —
+ // that absence is a gap, not a pattern worth copying.
+ "scan-media": {
+ kind: "scan-media",
+ label: "Scan media for corruption",
+ // Reports only; there is nothing in flight to let finish.
+ drainable: false,
+ bookmarkable: false,
+ // Local disk work with nothing to serialize against — see the queueKey ""
+ // escape hatch refresh-report and detect-duplicates use.
+ queueKeyStrategy: "parallel",
+ defaultTier: "background",
+ },
+ "scan-media-channel": {
+ kind: "scan-media-channel",
+ label: "Scan channel media",
+ drainable: false,
+ bookmarkable: false,
+ queueKeyStrategy: "parallel",
+ defaultTier: "background",
+ },
"check-kept-deleted": {
kind: "check-kept-deleted",
label: "Check kept videos",
diff --git a/common/lib/backfillKinds.test.ts b/common/lib/backfillKinds.test.ts
@@ -61,6 +61,14 @@ function settingsWithDiarization(
};
}
+// The duration cap is OFF by default now that windowing exists, so a test of the
+// cap has to ask for one. Stated here once rather than inline, so it is obvious
+// that every OTHER test in this file runs with the cap disabled — which is the
+// shipped configuration.
+function settingsWithCap(hours: number): SiteSettings {
+ return settingsWithDiarization({ maxAudioHours: hours });
+}
+
// A video dir built from a description of what is on disk. `sidecar` is written
// verbatim so a MALFORMED file can be tested — that case is not hypothetical,
// it is what a crash mid-write leaves behind.
@@ -281,19 +289,19 @@ test("deferred: over the cap, with the input right there", async () => {
// point of a third state: this is not missing-input (nothing to work from) and
// not missing (work to do); it is work deliberately not attempted.
assert.equal(
- await classify({ audio: true, durationSec: 5 * 3600 }),
+ await classify({ audio: true, durationSec: 5 * 3600 }, settingsWithCap(4)),
"deferred",
);
});
test("under the cap is ordinary missing work", async () => {
assert.equal(
- await classify({ audio: true, durationSec: 3 * 3600 }),
+ await classify({ audio: true, durationSec: 3 * 3600 }, settingsWithCap(4)),
"missing",
);
// Exactly at the cap is not over it.
assert.equal(
- await classify({ audio: true, durationSec: 4 * 3600 }),
+ await classify({ audio: true, durationSec: 4 * 3600 }, settingsWithCap(4)),
"missing",
);
});
@@ -302,13 +310,16 @@ test("unknown duration is NOT deferred", async () => {
// No metadata.info.json at all. Deferring here would quietly remove a video
// from the work list on the strength of a file that could not be read, and it
// would break every fixture in this file that predates the cap.
- assert.equal(await classify({ audio: true }), "missing");
+ assert.equal(await classify({ audio: true }, settingsWithCap(4)), "missing");
// Present but useless — the same answer, for the same reason.
- assert.equal(await classify({ audio: true, durationSec: 0 }), "missing");
+ assert.equal(
+ await classify({ audio: true, durationSec: 0 }, settingsWithCap(4)),
+ "missing",
+ );
});
-test("maxAudioHours 0 turns the cap off", async () => {
- // Where this setting goes once windowed diarization lands.
+test("maxAudioHours 0 turns the cap off — and 0 is the shipped default", async () => {
+ // Where this setting went once windowed diarization landed.
assert.equal(
await classify(
{ audio: true, durationSec: 12 * 3600 },
@@ -323,11 +334,10 @@ test("the cap never overrules a sidecar that is already there", async () => {
// attempt, not what counts as done. Otherwise raising the cap would look like
// work appearing and lowering it would look like work being undone.
assert.equal(
- await classify({
- audio: true,
- durationSec: 12 * 3600,
- sidecar: sidecar(CURRENT),
- }),
+ await classify(
+ { audio: true, durationSec: 12 * 3600, sidecar: sidecar(CURRENT) },
+ settingsWithCap(4),
+ ),
"present",
);
});
@@ -336,7 +346,10 @@ test("the cap does not resurrect a video whose input is gone", async () => {
// missing-input is checked FIRST. A 12-hour video with nothing to diarize from
// is unreachable, not deferred — deferred promises "we could do this if you
// raised the cap", and that would be a lie here.
- assert.equal(await classify({ durationSec: 12 * 3600 }), "missing-input");
+ assert.equal(
+ await classify({ durationSec: 12 * 3600 }, settingsWithCap(4)),
+ "missing-input",
+ );
});
test("a malformed sidecar reads as absent, never as done", async () => {
diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts
@@ -278,8 +278,10 @@ async function hasDiarizableInput(
// Is this video longer than the diarization duration cap?
//
-// THE CAP IS A STOPGAP FOR AN OOM, and this function is where its two honest
-// limitations live. First, duration is a PROXY: the memory blowup is O(n^2) in
+// THE CAP IS OFF BY DEFAULT NOW — windowed diarization removed the OOM it
+// existed for. It remains because a smaller machine, or a recording longer than
+// anything measured here, may still want it. This function is where its two
+// honest limitations live. First, duration is a PROXY: the memory blowup is O(n^2) in
// speech-SEGMENT count, and turn density varies 40x across this corpus, so a
// sparse 7h42m video is cheaper than a dense 6h12m one. Duration is used anyway
// because it is the only predictor available from metadata already on disk, for
diff --git a/common/lib/diarization.ts b/common/lib/diarization.ts
@@ -56,6 +56,22 @@ export type DiarizationRecord = {
speakers: number;
turns: DiarizationTurn[];
engine: DiarizationEngine;
+ // Present only when the engine processed the audio in windows rather than
+ // whole (long recordings — see scripts/diarize-sherpa.py for why).
+ //
+ // DELIBERATELY OUTSIDE `engine`, AND DELIBERATELY NOT COMPARED. The freshness
+ // target below is built from settings alone and cannot know a given video's
+ // duration, so a window field in the identity would mark every sidecar on disk
+ // stale the day windowing shipped — for work that is unchanged. Keeping it out
+ // here, rather than merely forgetting to compare it, is what stops a later
+ // edit to isDiarizationFresh from sweeping it in by accident. Recorded because
+ // it genuinely describes how the answer was produced, and a re-clustering pass
+ // may want to know a seam existed.
+ windowing?: {
+ windows: number;
+ windowSeconds: number;
+ overlapSeconds: number;
+ };
};
export const DIARIZATION_FILENAME = "diarization.json";
diff --git a/common/lib/mediaScan.test.ts b/common/lib/mediaScan.test.ts
@@ -0,0 +1,130 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import {
+ MEDIA_SCAN_OVERRIDES_VERSION,
+ mediaScanKey,
+ mediaScanTotals,
+ sanitizeMediaScanOverrides,
+ unreviewedFindings,
+ type MediaScanFinding,
+ type MediaScanReport,
+} from "./mediaScan";
+
+function finding(over: Partial<MediaScanFinding> = {}): MediaScanFinding {
+ return {
+ slug: "chan/vid",
+ channelSlug: "chan",
+ videoId: "vid",
+ file: "audio.mp3",
+ bytes: 100,
+ verdict: "unreadable",
+ tier: 1,
+ containerSeconds: null,
+ metadataSeconds: null,
+ ...over,
+ };
+}
+
+function report(over: Partial<MediaScanReport> = {}): MediaScanReport {
+ return {
+ version: 1,
+ generatedAt: "2026-08-08T00:00:00.000Z",
+ runConfig: { deepProbe: false, toleranceSeconds: 30, toleranceRatio: 0.02 },
+ channels: {},
+ findings: [],
+ ...over,
+ };
+}
+
+const CHANNEL_TOTALS = {
+ videosScanned: 2,
+ filesScanned: 3,
+ ok: 2,
+ unreadable: 1,
+ truncated: 0,
+ stray: 0,
+ unknown: 0,
+ deepProbed: 0,
+ bytesAtRisk: 500,
+ scannedAt: "2026-08-08T00:00:00.000Z",
+};
+
+test("corpus totals are the SUM of the per-channel entries", () => {
+ // Derived rather than stored, which is what lets a per-channel re-scan replace
+ // one entry without the headline claiming to describe a scan it did not do.
+ const r = report({
+ channels: {
+ a: { ...CHANNEL_TOTALS },
+ b: {
+ ...CHANNEL_TOTALS,
+ videosScanned: 5,
+ unreadable: 0,
+ truncated: 2,
+ bytesAtRisk: 700,
+ scannedAt: "2026-08-09T00:00:00.000Z",
+ },
+ },
+ });
+ const t = mediaScanTotals(r);
+ assert.equal(t.videosScanned, 7);
+ assert.equal(t.filesScanned, 6);
+ assert.equal(t.unreadable, 1);
+ assert.equal(t.truncated, 2);
+ assert.equal(t.bytesAtRisk, 1200);
+ // The freshest channel scan is what the report is "as of".
+ assert.equal(t.scannedAt, "2026-08-09T00:00:00.000Z");
+});
+
+test("an empty report totals to zero rather than throwing", () => {
+ assert.equal(mediaScanTotals(report()).filesScanned, 0);
+ // A report written by an older build has no `channels` at all.
+ assert.equal(
+ mediaScanTotals({ ...report(), channels: undefined as never }).filesScanned,
+ 0,
+ );
+});
+
+test("a reviewed finding drops out of the work list", () => {
+ const f = finding();
+ const r = report({ findings: [f, finding({ file: "source-media.mp4" })] });
+ const overrides = sanitizeMediaScanOverrides({
+ reviewed: { [mediaScanKey(f)]: { reviewedAt: "2026-08-08T00:00:00.000Z" } },
+ });
+ const left = unreviewedFindings(r, overrides);
+ assert.equal(left.length, 1);
+ assert.equal(left[0].file, "source-media.mp4");
+});
+
+test("the key is per FILE, not per video", () => {
+ // A video dir can hold both a good container and a corrupt scratch copy;
+ // reviewing one must not silently dismiss the other.
+ assert.notEqual(
+ mediaScanKey({ slug: "chan/vid", file: "source-media.mp4" }),
+ mediaScanKey({ slug: "chan/vid", file: "source-media.temp.mp4" }),
+ );
+});
+
+test("a malformed overrides file reads as no decisions, never throws", () => {
+ // Hand-editable, so it will eventually be hand-broken. The page still has to
+ // render.
+ for (const junk of [null, undefined, 7, "{}", { reviewed: 3 }, {}]) {
+ const o = sanitizeMediaScanOverrides(junk);
+ assert.deepEqual(o.reviewed, {});
+ assert.equal(o.version, MEDIA_SCAN_OVERRIDES_VERSION);
+ }
+ // Entries missing the timestamp are dropped, the rest survive.
+ const mixed = sanitizeMediaScanOverrides({
+ reviewed: {
+ good: { reviewedAt: "2026-08-08T00:00:00.000Z", note: "known clip" },
+ bad: { note: "no timestamp" },
+ alsoBad: 5,
+ },
+ });
+ assert.deepEqual(Object.keys(mixed.reviewed), ["good"]);
+ assert.equal(mixed.reviewed.good.note, "known clip");
+});
+
+test("no report means no findings, not a crash", () => {
+ assert.deepEqual(unreviewedFindings(null, null), []);
+ assert.equal(unreviewedFindings(report({ findings: [finding()] }), null).length, 1);
+});
diff --git a/common/lib/mediaScan.ts b/common/lib/mediaScan.ts
@@ -0,0 +1,187 @@
+// Client-safe types and constants for the corrupt-media scan. The walk itself
+// lives in controller/scanCorruptMedia.ts; this file is what a page or a widget
+// may import.
+//
+// MODELLED ON lib/duplicates.ts, deliberately: a versioned report JSON written
+// atomically at the corpus root, findings that mutate nothing, and human review
+// decisions kept in a SIBLING overrides file rather than in the report — because
+// the report is regenerated wholesale by every run and a decision recorded
+// inside it would be destroyed by the next one.
+
+export const MEDIA_SCAN_REPORT_VERSION = 1;
+export const MEDIA_SCAN_FILENAME = "media-scan.json";
+export const MEDIA_SCAN_OVERRIDES_FILENAME = "media-scan.overrides.json";
+export const MEDIA_SCAN_OVERRIDES_VERSION = 1;
+
+// How far apart the container's own duration and the metadata's may be before
+// the file is treated as suspect. Generous: a container rounds, and yt-dlp's
+// metadata duration comes from the site rather than from the bytes.
+export const DEFAULT_DURATION_TOLERANCE_SEC = 30;
+// A relative tolerance too, because 30s means something very different on a
+// 40-second clip than on an eight-hour stream.
+export const DEFAULT_DURATION_TOLERANCE_RATIO = 0.02;
+
+// What the scan concluded about one file.
+//
+// ok — readable, and its length agrees with what we know.
+// unreadable — ffprobe RAN and could not parse the container. This is the
+// today's-known-case verdict: a truncated mp4 fails with
+// "moov atom not found" in milliseconds.
+// truncated — readable, but materially shorter than the metadata says. Only
+// reachable with the deep probe, because ffmpeg exits 0 on a
+// truncated container (see ffmpegStreamProbe.ts) so a duration
+// comparison is the ONLY way to catch it.
+// stray — a media file the app does not recognize as one of its own
+// outputs: a leftover scratch or temp file. Not a corruption
+// verdict at all, and reported separately for that reason — but
+// it is how the 4 GiB source-media.temp.mp4 becomes visible.
+// unknown — the scan could not reach a conclusion. ffprobe missing, a
+// probe error, no duration to compare against. NEVER acted on.
+export type MediaScanVerdict =
+ | "ok"
+ | "unreadable"
+ | "truncated"
+ | "stray"
+ | "unknown";
+
+// The verdicts a human is being asked to look at. `ok` files are not written to
+// the report at all — the report is a work list, not an inventory.
+export const MEDIA_SCAN_FINDING_VERDICTS: ReadonlyArray<MediaScanVerdict> = [
+ "unreadable",
+ "truncated",
+ "stray",
+ "unknown",
+];
+
+export type MediaScanFinding = {
+ // `${channelSlug}/${videoId}`, the slug the rest of the app uses.
+ slug: string;
+ channelSlug: string;
+ videoId: string;
+ // Filename within the video dir.
+ file: string;
+ bytes: number;
+ verdict: MediaScanVerdict;
+ // Which tier reached the verdict. 1 is the cheap container read every file
+ // gets; 2 is the full decode, which is opt-in and narrowed.
+ tier: 1 | 2;
+ // The container's own duration, and what metadata.info.json claims. Either may
+ // be null — that is what an `unknown` usually means.
+ containerSeconds: number | null;
+ metadataSeconds: number | null;
+ // First meaningful line of probe output. The evidence, so a verdict can be
+ // argued with rather than merely trusted.
+ detail?: string;
+};
+
+export type MediaScanChannelTotals = {
+ videosScanned: number;
+ filesScanned: number;
+ ok: number;
+ unreadable: number;
+ truncated: number;
+ stray: number;
+ unknown: number;
+ // Files that got the tier-2 full decode.
+ deepProbed: number;
+ bytesAtRisk: number;
+ scannedAt: string;
+};
+
+export type MediaScanReport = {
+ version: number;
+ generatedAt: string;
+ runConfig: {
+ deepProbe: boolean;
+ toleranceSeconds: number;
+ toleranceRatio: number;
+ };
+ // PER CHANNEL, not one flat set of totals. A per-channel re-scan replaces one
+ // entry and its findings; the corpus-wide numbers are then the sum of these,
+ // so a partial run can never leave the headline claiming to describe a scan it
+ // did not do.
+ channels: Record<string, MediaScanChannelTotals>;
+ findings: MediaScanFinding[];
+};
+
+// Corpus-wide totals, derived rather than stored — see MediaScanReport.channels.
+export function mediaScanTotals(report: MediaScanReport): MediaScanChannelTotals {
+ const out: MediaScanChannelTotals = {
+ videosScanned: 0,
+ filesScanned: 0,
+ ok: 0,
+ unreadable: 0,
+ truncated: 0,
+ stray: 0,
+ unknown: 0,
+ deepProbed: 0,
+ bytesAtRisk: 0,
+ scannedAt: "",
+ };
+ for (const t of Object.values(report.channels ?? {})) {
+ out.videosScanned += t.videosScanned;
+ out.filesScanned += t.filesScanned;
+ out.ok += t.ok;
+ out.unreadable += t.unreadable;
+ out.truncated += t.truncated;
+ out.stray += t.stray;
+ out.unknown += t.unknown;
+ out.deepProbed += t.deepProbed;
+ out.bytesAtRisk += t.bytesAtRisk;
+ if (t.scannedAt > out.scannedAt) out.scannedAt = t.scannedAt;
+ }
+ return out;
+}
+
+// A finding's stable identity across runs, for the overrides file.
+export function mediaScanKey(f: {
+ slug: string;
+ file: string;
+}): string {
+ return `${f.slug}/${f.file}`;
+}
+
+// "Looked at it; it is fine." Kept OUT of the report for the reason stated in
+// the header.
+export type MediaScanOverride = {
+ reviewedAt: string;
+ note?: string;
+};
+
+export type MediaScanOverrides = {
+ version: number;
+ reviewed: Record<string, MediaScanOverride>;
+};
+
+// Never throws. A malformed overrides file reads as "no decisions recorded", so
+// the page still renders rather than 500ing on a hand-edited file.
+export function sanitizeMediaScanOverrides(value: unknown): MediaScanOverrides {
+ const empty: MediaScanOverrides = {
+ version: MEDIA_SCAN_OVERRIDES_VERSION,
+ reviewed: {},
+ };
+ if (!value || typeof value !== "object") return empty;
+ const raw = value as { reviewed?: unknown };
+ if (!raw.reviewed || typeof raw.reviewed !== "object") return empty;
+ const reviewed: Record<string, MediaScanOverride> = {};
+ for (const [key, v] of Object.entries(raw.reviewed as object)) {
+ if (!v || typeof v !== "object") continue;
+ const o = v as { reviewedAt?: unknown; note?: unknown };
+ if (typeof o.reviewedAt !== "string") continue;
+ reviewed[key] = {
+ reviewedAt: o.reviewedAt,
+ ...(typeof o.note === "string" && o.note ? { note: o.note } : {}),
+ };
+ }
+ return { version: MEDIA_SCAN_OVERRIDES_VERSION, reviewed };
+}
+
+// Findings a human has not already dismissed. What every surface should count.
+export function unreviewedFindings(
+ report: MediaScanReport | null,
+ overrides: MediaScanOverrides | null,
+): MediaScanFinding[] {
+ if (!report) return [];
+ const reviewed = overrides?.reviewed ?? {};
+ return report.findings.filter((f) => !reviewed[mediaScanKey(f)]);
+}
diff --git a/common/lib/settings.ts b/common/lib/settings.ts
@@ -1101,10 +1101,13 @@ export function defaultDiarization(): DiarizationSettings {
segModel: "",
embModel: "",
concurrency: 1,
- // 4 hours. Every video that OOM-killed this box was over 6h; the shortest
- // was 6h12m. 4 leaves headroom for the density variance duration cannot see.
- // Temporary — see DiarizationSettings.maxAudioHours.
- maxAudioHours: 4,
+ // OFF, because windowing made it unnecessary — which is what it was always
+ // for. It shipped at 4 hours as a stopgap while long recordings were being
+ // OOM-killed; the engine now processes them in windows and the 6h12m file
+ // that was being killed at 10.5 GB completes at ~2 GB. The knob stays, and
+ // stays honest about what it does, for a machine smaller than this one or a
+ // recording longer than anything measured here.
+ maxAudioHours: 0,
};
}
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,8 @@
# Changelog
## [Unreleased]
+- **Speaker diarization now works on very long recordings instead of dying on them.** The engine compares every detected speaker turn against every other one, so its memory grows with the *square* of how many turns there are — on a dense six-hour stream that reached 10.5 GB and the kernel killed it, after about forty minutes of work, having produced nothing. Six of ten videos over six hours died this way. A second, quieter problem sat underneath it: the audio was decoded into memory **all at once**, which is another 1.8 GB for an eight-hour file before any analysis starts. **Long recordings are now processed in 45-minute windows**, which fixes both at once rather than trading one off against the other — each window has a fraction of the turns, so the comparison shrinks by the *square* of that fraction, and only one window's audio is ever in memory. Measured on the six-hour stream that was being killed: **peak memory 2.0 GB against 10.5 GB, and it completes** — 3,433 speaker turns out of a file that previously produced none. The hard part of windowing is not the windows, it is that the same person appears in several of them and has to be recognized as one speaker rather than nine; each window's speakers are therefore reduced to a voice fingerprint and those are matched across the whole recording, with turns rejoined across the seams. **Short recordings are untouched — the same code path, the same results, byte for byte** — which is checked against an existing file and is what keeps everything already produced reproducible. Windowing is recorded on the file it produces but is deliberately not part of what makes a result look out-of-date, so this does not mark a single existing diarization stale. With this in place the **Max audio hours** limit added above now defaults to off.
+- **A new scan checks the media already on disk for corruption, and it found some.** The integrity probes have existed since audio-checked downloads landed, but they only ever ran on files as they arrived — nothing had ever asked whether the files sitting here are still readable. They are not all readable. On this corpus the first run took **83 seconds** across 953 media files and turned up: a **4.0 GiB** unreadable video container left behind by a download that never finished (on a disk at 99% full); **three zero-byte audio files** whose videos consequently look "already downloaded" to the pipeline and so are never fetched again, one of them sitting next to a perfectly good file the app would have passed over in its favour; **39 audio files under a name the app does not recognize**, so 1.2 GiB of audio that is present but invisible; and **7 files whose length disagrees with their metadata** by anything up to five hours, which is what a truncated download looks like. The scan is **tiered on purpose**: the default pass only reads each file's header, which costs milliseconds, and the thorough pass — a full decode — is opt-in and runs only on files the cheap pass found reason to doubt, because running it over everything would be hours of processor time. Anything it cannot resolve is reported as *"could not tell"* and never as a fault. **It deletes nothing.** Findings appear on Actionable and per-channel under Diagnostics, each with the evidence for the verdict, and can be marked reviewed to dismiss them; deleting a file remains the separate, deliberate click it already was.
- **Very long recordings are now set aside instead of being attempted and killed.** Speaker diarization on a multi-hour stream can exhaust this machine's memory and be killed by the kernel partway through — six of ten videos over six hours died that way, each after about forty minutes of work, producing nothing. The cause is not length as such: the engine holds a comparison of every detected speaker turn against every other, so its memory grows with the *square* of how many turns there are, and turn density varies fortyfold across this corpus (a sparse seven-hour reaction video finished; a dense six-hour stream did not). Length is simply the only predictor available for free, from information already on disk, before spending the forty minutes. So there is now a **Max audio hours** setting (default **4**, `0` turns it off), and a video over it is reported as **deferred**: a third number, shown on the channel's Backfill card and in its summary line, and **never added to the work still to do**. That last part is the point — a limit that quietly shrank the backlog would let a capped corpus report itself as finished. Deferred is not a failure and not "no input"; it is work deliberately not attempted, and raising the limit is all it takes to ask for it. **This is a stopgap and it is meant to be removed:** windowing the engine's work is the real fix, and when that lands the default goes to 0. Note that a backfill already running when this shipped keeps the old behaviour until it is stopped and started again.
- **A file called `audio.en-orig.vtt` is a subtitle, and the cleanup sweep used to think it was audio.** Every place in the app that asked "is this a media file?" answered by looking at the *start* of the name — anything beginning with `audio.` counted — and then subtracting a list of known exceptions. A list of exceptions is only ever as good as the last thing someone remembered to add to it, and it had fallen behind: a leaked subtitle track (`audio.en-orig.vtt`) matched nothing on the list, and a live-chat download fragment (`audio.live_chat.json.part-Frag114`) slipped past the `.part` rule because it doesn't *end* in `.part`. Since the sweep that deletes finished audio shares that same question, it could delete a subtitle believing it was audio — a small amount of real, unrecoverable data loss, in a sweep with no dry run and no trash. The same gap on the video-container side was worse: yt-dlp's own scratch file `source-media.temp.mp4` read as a finished video, and that is the one file the app *writes* from — it would extract audio from the corrupt scratch copy and then move it into the saved-video store as the permanent archive copy, orphaning the real container. **The rule is now the shape of the name rather than a list of exceptions**: a finished media file is `audio.<ext>` or `source-media.<ext>` with exactly one extension and nothing between, and `<ext>` has to be a format we actually recognize. That one rule rejects all four junk files, and it also rejects something an extension check alone would have let back in — the app's own half-written transcode output, `audio.tmp-12345.mp3`, which has a perfectly ordinary audio extension and whose deletion would race the transcode still writing it. Three knock-on fixes ride along: the "remove extra audio formats" sweep and the transcode-retry pass each carried their own private, staler copy of the old exception list (so the former could delete that subtitle too, and the latter could hand ffmpeg a `.vtt` file as "source audio"), and the audio-integrity checker could mistake a live-chat fragment for a resumable download. All four now ask the same question in the same place. **Everything moves in the safe direction:** strictly fewer files are ever deleted, videos whose only "audio" was one of these fakes now report honestly as having no audio instead of as a failed transcription, and such a video is no longer mistaken for one that has already been downloaded — so it gets fetched. See `common/lib/mediaFiles.ts` (new, with tests seeded from the real files this corpus holds) and `common/lib/videoStatus.ts`.
- **Speaker-attribution turns that land outside their own chunk are now counted instead of quietly dropped.** The text-only attribution lane asks the model about one slice of a transcript at a time and discards any answer whose timestamp points somewhere outside that slice — correctly, since a mark in another slice's territory would fight that slice's own answer. But it discarded them silently, with no count anywhere, and a validation run measured the drop rate at **30% — 520 of 1,705 turns**. That is the difference between "this lane is working" and "this lane is throwing away one answer in three", and nothing on disk or in the log could tell you which. The discard is now recorded on the record itself (one entry per chunk, carrying the count, so a high rate doesn't bury every other warning) and summarized in the job log as a rate. Behaviour is otherwise unchanged and the lane remains switched off by default; this is instrumentation, not a change of policy.
diff --git a/editor/app/actionable/actions.ts b/editor/app/actionable/actions.ts
@@ -15,6 +15,11 @@ import {
readDuplicateReport,
updateDuplicateOverride,
} from "yt-dlp-transcript-common/controller/duplicateShorts";
+import {
+ scanCorruptMedia,
+ updateMediaScanOverride,
+} from "yt-dlp-transcript-common/controller/scanCorruptMedia";
+import { mediaScanTotals } from "yt-dlp-transcript-common/lib/mediaScan";
import { shareClusterFromCanonical } from "yt-dlp-transcript-common/controller/digestSharing";
import {
clearIncompleteTranscriptsAction,
@@ -247,3 +252,69 @@ export async function redownloadAllIncompleteTranscriptsAction(): Promise<Global
revalidatePath("/actionable");
return { ok: true, channels: slugs.length, affected: queued };
}
+
+
+// ---------------------------------------------------------------------------
+// Corrupt-media scan
+// ---------------------------------------------------------------------------
+
+export type RunMediaScanResult =
+ | { ok: true; findings: number; filesScanned: number }
+ | { ok: false; error: string };
+
+// Scan media already on disk. Reports; deletes nothing — deletion stays the
+// separate, explicit per-file click that already exists on the video page.
+//
+// `deepProbe` turns on the tier-2 full decode for files whose duration looks
+// wrong. Off by default because that tier is a real transcode per file: it has
+// to be budgeted like the digest sweep, not like a stat() walk.
+export async function runMediaScanAction(
+ opts: { channels?: string[]; deepProbe?: boolean } = {},
+): Promise<RunMediaScanResult> {
+ const paths = getPaths();
+ let findings = 0;
+ let filesScanned = 0;
+ const perChannel = (opts.channels?.length ?? 0) > 0;
+ const result = await runManagedFunction({
+ kind: perChannel ? "scan-media-channel" : "scan-media",
+ // Empty queueKey: local read-only disk work, no reason to wait behind
+ // sync/download work (see runDuplicateDetectionAction).
+ queueKey: "",
+ paths,
+ fn: async (onLog, signal) => {
+ const report = await scanCorruptMedia({
+ paths,
+ channels: opts.channels,
+ deepProbe: opts.deepProbe,
+ onLog,
+ signal,
+ });
+ findings = report.findings.length;
+ filesScanned = mediaScanTotals(report).filesScanned;
+ },
+ });
+ if (!result.ok) return { ok: false, error: result.error };
+ await drainStream(result.stream);
+ revalidatePath("/actionable");
+ return { ok: true, findings, filesScanned };
+}
+
+export type ReviewMediaFindingResult =
+ | { ok: true }
+ | { ok: false; error: string };
+
+// "Looked at it; it is fine." Recorded in a sibling overrides file, because the
+// report is regenerated wholesale by every scan and a decision written into it
+// would be destroyed by the next run.
+export async function reviewMediaFindingAction(
+ key: string,
+ reviewed: boolean,
+): Promise<ReviewMediaFindingResult> {
+ try {
+ await updateMediaScanOverride(getPaths(), key, { reviewed });
+ revalidatePath("/actionable");
+ return { ok: true };
+ } catch (e) {
+ return { ok: false, error: (e as Error).message };
+ }
+}
diff --git a/editor/app/actionable/components/MediaScanFindingRow.tsx b/editor/app/actionable/components/MediaScanFindingRow.tsx
@@ -0,0 +1,121 @@
+"use client";
+
+import { useState, useTransition } from "react";
+import Link from "next/link";
+import { formatBytes } from "yt-dlp-transcript-common/lib/format";
+import {
+ mediaScanKey,
+ type MediaScanFinding,
+} from "yt-dlp-transcript-common/lib/mediaScan";
+import { reviewMediaFindingAction } from "../actions";
+
+const VERDICT_LABEL: Record<MediaScanFinding["verdict"], string> = {
+ ok: "ok",
+ unreadable: "unreadable",
+ truncated: "truncated",
+ stray: "stray file",
+ unknown: "could not tell",
+};
+
+const VERDICT_HELP: Record<MediaScanFinding["verdict"], string> = {
+ ok: "",
+ unreadable:
+ "ffprobe could not parse the container at all. This file cannot be transcribed or diarized from.",
+ truncated:
+ "Decodes, but is materially shorter than the metadata says — the download was cut off.",
+ stray:
+ "Readable media that is not one of the app's own outputs: a leftover scratch or temp file taking up disk.",
+ unknown:
+ "The scan could not reach a conclusion. Nothing is claimed about this file.",
+};
+
+export function MediaScanFindingRow({
+ finding,
+ reviewed,
+}: {
+ finding: MediaScanFinding;
+ reviewed: boolean;
+}) {
+ const [pending, startTransition] = useTransition();
+ const [error, setError] = useState<string | null>(null);
+ const key = mediaScanKey(finding);
+
+ function toggle() {
+ setError(null);
+ startTransition(async () => {
+ const r = await reviewMediaFindingAction(key, !reviewed);
+ if (!r.ok) setError(r.error);
+ });
+ }
+
+ return (
+ <li
+ aria-label={`media finding ${key}`}
+ data-verdict={finding.verdict}
+ className="border border-border rounded-md p-3 flex flex-col gap-1"
+ >
+ <div className="flex items-center gap-2 flex-wrap text-xs">
+ <span className="rounded px-2 py-0.5 border border-border font-medium">
+ {VERDICT_LABEL[finding.verdict]}
+ </span>
+ <span className="text-muted-foreground">
+ tier {finding.tier} · {formatBytes(finding.bytes)}
+ </span>
+ {reviewed && (
+ <span className="text-muted-foreground">· marked reviewed</span>
+ )}
+ </div>
+ <p className="text-sm">
+ <Link
+ href={`/channels/${encodeURIComponent(finding.channelSlug)}/videos/${encodeURIComponent(finding.videoId)}`}
+ className="underline underline-offset-2"
+ >
+ {finding.slug}
+ </Link>{" "}
+ <code className="text-xs">{finding.file}</code>
+ </p>
+ <p className="text-xs text-muted-foreground">
+ {VERDICT_HELP[finding.verdict]}
+ </p>
+ {(finding.containerSeconds !== null ||
+ finding.metadataSeconds !== null) && (
+ <p className="text-xs text-muted-foreground">
+ container{" "}
+ {finding.containerSeconds === null
+ ? "unreadable"
+ : `${Math.round(finding.containerSeconds)}s`}{" "}
+ · metadata{" "}
+ {finding.metadataSeconds === null
+ ? "unknown"
+ : `${Math.round(finding.metadataSeconds)}s`}
+ </p>
+ )}
+ {finding.detail && (
+ <p className="text-xs text-muted-foreground font-mono break-all">
+ {finding.detail}
+ </p>
+ )}
+ <div className="flex items-center gap-2">
+ {/* Deletion is deliberately NOT here. The convention in this app is
+ report first, delete second: the per-file delete already exists on
+ the video page, behind its own confirm, and a scan that could delete
+ what it flagged would be a very different and much more dangerous
+ thing than a scan. */}
+ <button
+ type="button"
+ onClick={toggle}
+ disabled={pending}
+ aria-label={`${reviewed ? "unreview" : "review"} ${key}`}
+ className="px-2 py-1 rounded border border-border text-xs hover:bg-muted disabled:opacity-50"
+ >
+ {reviewed ? "Un-review" : "Mark reviewed"}
+ </button>
+ {error && (
+ <span role="alert" className="text-xs text-destructive">
+ {error}
+ </span>
+ )}
+ </div>
+ </li>
+ );
+}
diff --git a/editor/app/actionable/components/RunMediaScanButton.tsx b/editor/app/actionable/components/RunMediaScanButton.tsx
@@ -0,0 +1,93 @@
+"use client";
+
+import { useState } from "react";
+import {
+ runMediaScanAction,
+ type RunMediaScanResult,
+} from "../actions";
+
+type Status =
+ | { kind: "idle" }
+ | { kind: "running" }
+ | { kind: "done"; result: RunMediaScanResult }
+ | { kind: "error"; message: string };
+
+// The depth choice is the whole reason this has a select rather than one button.
+// The container read costs milliseconds per file; the full decode is a real
+// transcode of every suspect file, so it has to be asked for rather than
+// happening because someone clicked the obvious thing.
+type Depth = "quick" | "deep";
+
+export function RunMediaScanButton({ slug }: { slug?: string }) {
+ const [depth, setDepth] = useState<Depth>("quick");
+ const [status, setStatus] = useState<Status>({ kind: "idle" });
+
+ async function handleClick() {
+ setStatus({ kind: "running" });
+ try {
+ const result = await runMediaScanAction({
+ ...(slug ? { channels: [slug] } : {}),
+ deepProbe: depth === "deep",
+ });
+ setStatus({ kind: "done", result });
+ } catch (e) {
+ setStatus({ kind: "error", message: (e as Error).message });
+ }
+ }
+
+ const running = status.kind === "running";
+ return (
+ <div className="flex items-center gap-2 flex-wrap">
+ <label className="sr-only" htmlFor={`media-scan-depth-${slug ?? "all"}`}>
+ Scan depth
+ </label>
+ <select
+ id={`media-scan-depth-${slug ?? "all"}`}
+ aria-label="media scan depth"
+ value={depth}
+ onChange={(e) => setDepth(e.target.value as Depth)}
+ disabled={running}
+ className="text-sm rounded-md border border-border bg-transparent px-2 py-2"
+ >
+ <option value="quick">Quick (read containers)</option>
+ <option value="deep">Deep (decode suspect files)</option>
+ </select>
+ <button
+ type="button"
+ onClick={handleClick}
+ disabled={running}
+ aria-label={slug ? "scan channel media" : "scan media for corruption"}
+ className="px-3 py-2 rounded-md bg-primary text-primary-foreground text-sm font-medium hover:opacity-90 disabled:opacity-50"
+ >
+ {running ? "Scanning…" : slug ? "Scan this channel" : "Scan media"}
+ </button>
+ {status.kind === "done" && status.result.ok && (
+ <span
+ aria-label="media scan result"
+ className="text-xs text-muted-foreground"
+ >
+ {status.result.filesScanned.toLocaleString()} file(s) checked ·{" "}
+ {status.result.findings.toLocaleString()} finding(s)
+ </span>
+ )}
+ {status.kind === "done" && !status.result.ok && (
+ <span
+ role="alert"
+ aria-label="media scan error"
+ className="text-xs text-destructive"
+ >
+ {status.result.error}
+ </span>
+ )}
+ {status.kind === "error" && (
+ <span
+ role="alert"
+ aria-label="media scan error"
+ className="text-xs text-destructive"
+ >
+ {status.message}
+ </span>
+ )}
+ </div>
+ );
+}
diff --git a/editor/app/actionable/lib/loadActionable.ts b/editor/app/actionable/lib/loadActionable.ts
@@ -16,6 +16,14 @@ import type {
DuplicateOverrides,
DuplicateReport,
} from "yt-dlp-transcript-common/lib/duplicates";
+import type {
+ MediaScanOverrides,
+ MediaScanReport,
+} from "yt-dlp-transcript-common/lib/mediaScan";
+import {
+ readMediaScanOverrides,
+ readMediaScanReport,
+} from "yt-dlp-transcript-common/controller/scanCorruptMedia";
export type ActionableRow = {
channel: ChannelBrief;
@@ -43,6 +51,12 @@ export type ActionableSummary = {
// review UI cannot show what has already been decided without it, and a
// review queue that forgets its own answers re-asks every question.
duplicateOverrides: DuplicateOverrides;
+ // The corrupt-media scan, and the "reviewed, this one's fine" decisions kept
+ // beside it. Null when the scan has never been run — which is NOT the same as
+ // "nothing is wrong", and the section says so rather than rendering an
+ // all-clear it has no evidence for.
+ mediaScan: MediaScanReport | null;
+ mediaScanOverrides: MediaScanOverrides;
};
export function isStaleOrMissing(row: ActionableRow): boolean {
@@ -157,11 +171,14 @@ export function actionableCleanExtraFormatsBytes(row: ActionableRow): number {
export async function loadActionableSummary(
paths: Paths,
): Promise<ActionableSummary> {
- const [channels, duplicates, duplicateOverrides] = await Promise.all([
- getChannelBriefs(paths),
- getDuplicateReport(paths),
- readDuplicateOverrides(paths),
- ]);
+ const [channels, duplicates, duplicateOverrides, mediaScan, mediaScanOverrides] =
+ await Promise.all([
+ getChannelBriefs(paths),
+ getDuplicateReport(paths),
+ readDuplicateOverrides(paths),
+ readMediaScanReport(paths),
+ readMediaScanOverrides(paths),
+ ]);
// The brief already read the snapshot; this used to read each one a second
// time on top of a full corpus walk.
const rows: ActionableRow[] = channels.map((channel) => ({
@@ -245,6 +262,8 @@ export async function loadActionableSummary(
backfill,
duplicates,
duplicateOverrides,
+ mediaScan,
+ mediaScanOverrides,
};
}
diff --git a/editor/app/actionable/page.tsx b/editor/app/actionable/page.tsx
@@ -22,6 +22,15 @@ import { InlineActionButton } from "./components/InlineActionButton";
import { FixAllIncompleteButton } from "./components/FixAllIncompleteButton";
import { RefreshAllReportsButton } from "./components/RefreshAllReportsButton";
import { RunDuplicateDetectionButton } from "./components/RunDuplicateDetectionButton";
+import { RunMediaScanButton } from "./components/RunMediaScanButton";
+import { MediaScanFindingRow } from "./components/MediaScanFindingRow";
+import {
+ mediaScanKey,
+ mediaScanTotals,
+ unreviewedFindings,
+ type MediaScanOverrides,
+ type MediaScanReport,
+} from "yt-dlp-transcript-common/lib/mediaScan";
import { DuplicateClusterReview } from "./components/DuplicateClusterReview";
import type {
DuplicateCluster,
@@ -315,10 +324,82 @@ export default async function ActionablePage() {
report={summary.duplicates}
overrides={summary.duplicateOverrides}
/>
+ <MediaScanSection
+ report={summary.mediaScan}
+ overrides={summary.mediaScanOverrides}
+ />
</div>
);
}
+// Corrupt media already on disk. Flag-only, exactly like the duplicates section
+// above it — the probes have existed since the audio-checked download landed;
+// what was missing was ever running them against files that are already here.
+function MediaScanSection({
+ report,
+ overrides,
+}: {
+ report: MediaScanReport | null;
+ overrides: MediaScanOverrides;
+}) {
+ const findings = report?.findings ?? [];
+ const outstanding = unreviewedFindings(report, overrides);
+ const totals = report ? mediaScanTotals(report) : null;
+ return (
+ <section aria-label="media scan" className="flex flex-col gap-2">
+ <div className="flex items-center justify-between flex-wrap gap-2">
+ <h2 className="text-lg font-semibold">
+ Media integrity{outstanding.length > 0 ? ` (${outstanding.length})` : ""}
+ </h2>
+ <RunMediaScanButton />
+ </div>
+ <p className="text-sm text-muted-foreground">
+ Checks the audio and video files still on disk for corruption or
+ truncation. Flag-only — nothing here deletes anything.
+ {totals && totals.scannedAt ? (
+ <>
+ {" "}
+ Last scan{" "}
+ <time dateTime={totals.scannedAt}>
+ {new Date(totals.scannedAt).toLocaleString()}
+ </time>{" "}
+ · {totals.filesScanned.toLocaleString()} file(s) across{" "}
+ {totals.videosScanned.toLocaleString()} video dir(s)
+ {totals.deepProbed > 0
+ ? `, ${totals.deepProbed.toLocaleString()} fully decoded`
+ : ""}
+ .
+ </>
+ ) : (
+ " Never run."
+ )}
+ </p>
+ {findings.length === 0 ? (
+ <p
+ aria-label="media scan empty"
+ className="text-sm text-muted-foreground border border-dashed border-border rounded p-4"
+ >
+ {/* "Never run" is NOT "all clear", and saying so would be the one way
+ this section could actively mislead. */}
+ {report
+ ? "Every file the scan read was intact."
+ : "Run a scan to check the media on disk."}
+ </p>
+ ) : (
+ <ul className="flex flex-col gap-3">
+ {findings.map((f) => (
+ <MediaScanFindingRow
+ key={mediaScanKey(f)}
+ finding={f}
+ reviewed={Boolean(overrides.reviewed[mediaScanKey(f)])}
+ />
+ ))}
+ </ul>
+ )}
+ </section>
+ );
+}
+
function DuplicatesSection({
report,
overrides,
diff --git a/editor/app/channels/[slug]/components/stages/DiagnosticsStage.tsx b/editor/app/channels/[slug]/components/stages/DiagnosticsStage.tsx
@@ -28,6 +28,7 @@ import {
} from "../../../../components/ShardControl";
import { RetryBucketControl } from "../RetryBucketControl";
import { VideoIdList } from "../VideoIdList";
+import { RunMediaScanButton } from "../../../../actionable/components/RunMediaScanButton";
function parseConcurrency(s: string): number | undefined {
const trimmed = s.trim();
@@ -131,6 +132,13 @@ export function DiagnosticsStage({
</div>
<div className="flex flex-col gap-2">
<Heading
+ title="Media integrity"
+ desc="Check this channel's audio and video files on disk for corruption or truncation. Reports only — nothing is deleted."
+ />
+ <RunMediaScanButton slug={slug} />
+ </div>
+ <div className="flex flex-col gap-2">
+ <Heading
title="Inspect a video"
desc="Jump to a specific video by ID."
/>
diff --git a/plans/FACTS.md b/plans/FACTS.md
@@ -1812,3 +1812,120 @@ later on a bigger box or with a segment-bounded clusterer.
`diarization.concurrency: 1` is now doubly load-bearing: two 11 GB transients
would not merely be tight, they would be unschedulable.
+
+## Windowed diarization (2026-08-08) — the OOM wall is REMOVED, with numbers
+
+### The two memory sinks, and why windowing closes both
+
+1. **`sherpa_onnx` clustering is O(n²) in SPEECH-SEGMENT count**, not in duration. Segment
+ density varies ~40x across this corpus, which is why duration never predicted the kills:
+ a sparse 7h42m file finished while a dense 6h12m one died.
+2. **The whole-file decode buffered the entire stream** — `subprocess.run(stdout=PIPE)`
+ accumulates one `bytes` object, ~1.84 GB for an 8h file, peaking near 2x that while the
+ buffer grows.
+
+Windowing divides per-window `n` by the window count, so the matrix falls by its SQUARE, and
+one window is decoded at a time. The window decode also writes straight into a preallocated
+`np.empty(n, float32)` via `readinto`, so even the per-window buffer never doubles.
+
+### Measured, on the file that was being killed
+
+`hasanabi/2834644359`, 6.21 h (22,368 s), 9 windows of 45 min with 30 s overlap:
+
+| | before | after |
+| --- | --- | --- |
+| outcome | **OOM-killed** (kernel, `anon-rss 10,512,464 kB`) | **completed** |
+| peak RSS | 10.5 GB | **~2.0 GB** (`VmHWM` 2,061,008 kB) |
+| result | nothing | 3,433 turns, 96 speakers, 2,774 s (442 s/audio-hour) |
+
+Peak RSS is sampled from `/proc/<pid>/status` — **GNU `time` is not installed on this box**.
+Two kernel kills are in the journal at 18:39 and 19:21 on 2026-08-08, both at ~10.5 GB
+anon-rss, which is what the O(n²) estimate brackets (30k segments = 6.7 GiB, 40k = 11.9 GiB).
+
+### Reproducibility gate — short files are byte-identical
+
+`ObviousRises-rumble/v6z1o2g` re-ran through the new script with windowing not triggered:
+**51 turns, 13 speakers, every non-volatile field equal** to the sidecar on disk, and no
+`windowing` key. The whole-file branch is the original code, unchanged, and `--window-minutes 0`
+forces it. This is what keeps the 51 sidecars already written reproducible.
+
+### Quality gate — compared on talk-time, not speaker count
+
+Windowed vs whole-file on the same file (4-minute windows, i.e. a HARSHER seam test than the
+shipped 45-minute default):
+
+| | whole-file | windowed |
+| --- | --- | --- |
+| total speech | 713.3 s | 712.7 s (**−0.1%**) |
+| wall-clock covered | 711.3 s, 0.03–765.97 | 711.4 s, 0.03–765.98 |
+| speakers | 13 | 6 |
+| top-1 / top-2 talk share | 73.1% / 84.3% | 88.2% / 94.0% |
+
+Windowing's risk is fragmenting one person across seams. The opposite happened: no speech is
+lost, coverage is identical, and identity is MORE concentrated — consistent with the known
+over-splitting at threshold 0.9, since centroids averaged over a speaker's longest turns are
+a cleaner signal than sherpa's per-segment clustering.
+
+### THE BUG THE GATE CAUGHT, which reading the code did not
+
+The first seam merge joined only **consecutive** turns in the sorted list. That is not
+equivalent to merging a speaker's own intervals: two different LOCAL speakers can map to one
+GLOBAL speaker, and their turns then interleave with other people's. On the 6h12m file that
+left **35 same-speaker self-overlaps and 24 sub-0.5s splits** — talk time double-counted in
+every downstream consumer. Fixed to a per-speaker interval union; re-verified at 0 and 0.
+Cross-speaker overlap is deliberately preserved (genuine overlapping speech; the whole-file
+path emits it too — 702 such pairs on this file).
+
+### API facts, verified live against the installed sherpa-onnx 1.13.4
+
+- `SpeakerEmbeddingExtractor(cfg)` — `.dim` is **192**; `.create_stream()` →
+ `accept_waveform(sr, samples)` / `input_finished()`, then `.is_ready(stream)` and
+ `.compute(stream)`.
+- `FastClustering(FastClusteringConfig(num_clusters, threshold))` is **callable**: it takes an
+ arbitrary `(n, 192)` float32 array and returns a list of labels. It exposes no named
+ methods (`dir()` is empty) — call the instance.
+- `ffmpeg -ss N -t D -i <file> -f f32le -ac 1 -ar 16000 -` is sample-exact on this corpus's
+ files, so `-ss` BEFORE `-i` (input seek) is safe for window boundaries.
+
+### Windowing is recorded but is NOT part of the freshness identity
+
+`DiarizationRecord.windowing` is top-level, deliberately outside `engine`. The freshness
+target is built from settings alone and cannot know a video's duration, so a window field in
+the identity would mark every sidecar on disk stale the day windowing shipped. Same treatment
+`version` already gets, for the same reason.
+
+### The cap is now OFF by default
+
+`settings.diarization.maxAudioHours` shipped at 4 h as the stopgap and defaults to **0** now
+that windowing exists. Residual risk, stated rather than hidden: the engine decides whether to
+window from **ffprobe**, and if ffprobe cannot measure the file it falls back to the
+whole-file path (preserving determinism). A multi-hour file whose duration cannot be measured
+would therefore take the old path; setting `maxAudioHours` again is the backstop.
+
+## Corrupt-media scan (2026-08-08) — first corpus run, and it found four things
+
+Tier 1 only (ffprobe container read), **953 media files across 68 channels in 83 s**:
+
+| verdict | files | size |
+| --- | ---: | ---: |
+| unreadable | 4 | 4.04 GiB |
+| stray (unrecognized name) | 39 | 1.17 GiB |
+| unknown (duration mismatch) | 7 | 0.86 GiB |
+
+- **`nuxanor-kick/fd7bd973-…/source-media.temp.mp4`, 4.04 GiB** — `moov atom not found`. The
+ known case, and the only `source-media.*` in the corpus. Already transcribed; disk at 99%.
+- **3 ZERO-BYTE `audio.mp3`** in `leaflit-rumble` (`v6v9mq3`, `v6v9o25`, `v6v9r2r`), none
+ transcribed. This is the more interesting find: a 0-byte file satisfies the NAME-based
+ `isRealAudioFile`, so `destinationExists` reads these videos as already downloaded and the
+ pipeline never re-fetches them. `v6v9r2r` also holds a real `audio.aac`, which
+ `AUDIO_READ_PREFERENCE` (mp3 first) would pass over in favour of the empty file.
+- **39 `transcript.m4a` in `community-notes`, 1.17 GiB** — audio written under the SUBTITLE
+ output template. Invisible to the app in both directions: not audio (wrong base name), and
+ not a subtitle either, because `SUB_FILE_RE` requires `transcript.<x>.<y>` — two segments.
+- **7 duration mismatches**, up to 5 h (`piratesoftware/yTj0ciOOZys` is 18,315 s short).
+ Reported as `unknown`, never as `corrupt`: `ffmpegStreamProbe` records that a truncated
+ container exits 0, so only the opt-in tier-2 decode can settle these.
+
+**The scan deliberately looks at every MEDIA-EXTENSION file, not only recognized ones.** The
+file it was written for is precisely one the fixed predicates exclude — scanning only
+recognized outputs would have missed it.
diff --git a/plans/STATE.md b/plans/STATE.md
@@ -3,316 +3,70 @@
The working memory for the local-AI derived-corpus work. Rewritten at the end of every
session, before context is cleared. See [`README.md`](README.md) for the protocol.
-**Last updated:** 2026-08-08 (later) — **diarization is runnable again, and the text-only
-lane came back from the dead on a schema change.** Two tracks, no shared resource.
-
-**TRACK B — diarization environment rebuilt, and it is byte-identical.** The 2026-08-07
-spike shipped `scripts/diarize.mjs` working but its venv was ephemeral, so the code was
-shipped and unrunnable. Now durable at **`.diarize/`**: `uv`-managed Python 3.13.11,
-**sherpa-onnx 1.13.4**, pyannote segmentation-3.0 + TitaNet-small (148 MB, gitignored;
-`.diarize/README.md` is tracked and carries the rebuild recipe). Re-running the wrapper on
-`ObviousRises-rumble/v6z1o2g` to a **temp path** produced **byte-identical turns** (51
-turns, 13 speakers) against the sidecar the spike made — so the pipeline is deterministic
-and the rebuilt env is provably the same one. Nothing was written to the corpus.
-It also closed an item the spike left open: **585 s/audio-hour on a quiet box** (124.5 s
-for 766 s of audio, 381% CPU, loadavg 2.7), inside the spike's 500–700 estimate, so **no
-conclusion changes** — still ~2.6× slower than the GPU's 221 s/audio-hour ASR.
-
-**Re-measured reach, and it corrects the pilot by an order of magnitude:** of 78,938 video
-dirs, **893** still have audio but only **87** are also transcribed, so the diarized
-backfill reaches **0.115%**, not the ~1.1% the pilot's survey implied. The number that
-matters is the other one: **806 videos have audio and are NOT yet transcribed**, and audio
-is deleted at transcribe time. That is the irreversible loss, and it is still preventable.
-
-**AWAITING YOUR GO-AHEAD (nothing done):** writing
-`settings.diarization.{python,segModel,embModel}` — a change to your `settings.json`, so I
-did not make it. Separately, arming `inlineAfterTranscribe` is a policy call with a real
-cost (~3–4× slower per transcription), presented with numbers, not flipped by me.
-
-**TRACK A — the schema fix was right, and the per-turn lane is still DEAD.** Round 2 passed
-all four thresholds on 2 videos; round 3 ran the frozen 8-video sample to settle it
-(40 chunk + 8 cast calls, same engine, on a genuinely quiet box — loadavg **0.16** vs round
-2's 6.3). New harness `common/bin/attribution-bakeoff.ts` **writes no sidecar** and checks
-that itself (`attribution.json` **9 before, 9 after**).
-
-| criterion | bar | round 1 | round 2 (n=2) | round 3 (n=8) |
-| --- | ---: | ---: | ---: | ---: |
-| new labels/chunk | ≤ 0.3 | 14.70 | 0.29 | **0.27** PASS |
-| transcript-text labels | 0 | rampant | 0% | **0%** PASS |
-| off-cast labels | 0 | n/a | 0 | **0** PASS |
-| s/chunk engine | < 25 | 66.2 | 24.1 | **28.6 FAIL** |
-
-**Keep both halves of this.** The closed enum did exactly what it was designed to do and
-never once failed — 0 off-cast labels across ~1,705 turns, 0% transcript text. The
-14.70 → 0.27 collapse is real. **But cost fails** (round 2's 24.1 was a small-sample
-artifact; output scales with video length), projecting **63 days** against a digest lane
-that already beats it at 11.2 s/chunk.
-
-**And the passing headline was hiding a lane that produces almost nothing usable.** Of
-1,705 emitted turns: 338 accepted, 847 `Unknown`, and **520 discarded as out-of-range** —
-**20% usable**. `outOfRangeTurns` was wired after round 2 and live for the first time here;
-it is the round's most valuable finding, because the shipped `findTurns` drops those turns
-on the same condition **without counting them**. That is a second failure, independent of
-the label space, invisible in rounds 1 and 2. Meanwhile **4 of 8 videos are single-label at
-100% talk time and a 5th produced nothing**, so the 0.27 is a mean over videos that mostly
-emitted one label. At most 2 of 8 gave defensible multi-speaker attribution.
-
-**Pass A is better than pass B but is NOT clean, correcting round 2's read.** Round 2 said
-"every cast name checked against the transcript, none hallucinated". More precisely: names
-are grounded in the transcript **or the title** — `Chrissie Mayr` scores confidence 1.00 and
-appears 0× in the transcript, only in the title. Worse, **pass A names TOPICS as speakers**:
-on the 8-hour Hasan VOD the cast omits the host entirely and offers `Tony`/`Raymond` (55 and
-28 mentions, both discussed not speaking), which pass B then made 35%/65% of the video.
-`Twitter` was offered as a person elsewhere. At 5.1 s / 1 call per video, cast-only is
-**4.4 days** corpus-wide and remains the most interesting survivor — but it is a
-**promising candidate, not a shippable one**, and needs its own pre-stated bar. The round-1
-decision table's "pass A sane ⇒ ship cast-only" was written before anyone had seen pass A's
-output on 8 videos.
-
-**DECIDED: do not adopt the closed-cast lane.** `attributeOne.ts` is unchanged and
-`ATTRIBUTION_PROMPT_VERSION` is **not bumped**, so the 9 sidecars on disk stay valid. The
-diarized lane remains the one worth having.
-
-**CORRECTION MADE TO FACTS.md (protocol: correct in place, note it here).** Two claims in
-the 2026-08-08 pilot write-up were wrong and are now marked corrected:
-1. **The text lane's collapse was blamed on the `knownSpeakers` roster.** It is not. Per
- `round1-bakeoff.json`'s `labelFirstChunk`: on `destiny/5nmDzKB23OU` **15 of 16 labels
- first appear in chunk 0**, where `findTurns` passes no roster at all; on
- `rekietalaw/EsZhaCfc8HQ` the collapse is **intermittent** (new labels by chunk: 2, 17,
- 3, 15, 2, 2, **48**, 1 — chunk 7 recovers after chunk 6 catastrophes), which prompt
- length cannot explain. The real cause is **schema-adherence collapse on a long prompt**:
- dozens of labels sit at exactly 60 chars, the schema's own `maxLength`. That reframed the
- fix from "make the roster work" to "close the label space" — and that fix then worked.
-2. **"~1.1% of the legacy corpus"** for the diarized lane overstated reach 10×; audio is
- necessary but not sufficient, a transcript is also needed. Corrected to 87 videos, 0.115%.
-
-**Previously (2026-08-08, earlier):** **the attribution pilot ran, and the text-only lane is
-REJECTED on measurement.** Stage 2 of PLAN.md, run headless under `tsx` via the new
-`common/bin/attribution-pilot.ts` — no editor, no runner armed, `settings.json` untouched
-(the lane is enabled through `attributeOneVideo`'s in-memory `settings` override). 8 videos
-/ 40 chunks of the FROZEN digest bake-off sample, so the per-chunk cost is directly
-comparable to that bake-off's on the same eight videos.
-
-**It fails both pre-stated thresholds, and not narrowly.** Cross-chunk identity: **14.70
-new labels per chunk** on the multi-chunk videos against a "≲ 0.3 healthy, ≳ 1 do not
-sweep" bar (worst: 346 labels on one 13-chunk video). Cost: **66.2 s/chunk engine against
-the digest lane's 11.2**, projecting **147–159 days, not 25–55**. Engine and wall are
-within 7%, so this is decode volume, not contention.
-
-**The mechanism is worse than churn and is the thing to remember: on multi-chunk videos
-`qwen2.5:7b` stops emitting labels and puts VERBATIM TRANSCRIPT TEXT in the `speaker`
-field**, truncated at the schema's 60-char cap. The schema constrains a label's LENGTH,
-not its nature, and `isUselessSpeakerLabel` only blocks "Speaker N" — neither guard can
-catch it. Single-chunk videos in the same run produced clean names (`Jeremy`, `Donald
-Trump`, `Stephen Colbert`), so the failure appears once the prompt carries a
-`knownSpeakers` roster and grows. That is precisely the cross-chunk identity mechanism the
-lane exists for, which is why this is a lane result and not a prompt nit.
-**[CORRECTED 2026-08-08 — the roster sentence above is WRONG; the collapse starts in chunk
-0 where no roster exists, and is intermittent. See the correction at the top of this file
-and the marked block in FACTS.md. The closed-cast schema then fixed it.]**
-
-**The diarized lane's one-call claim HOLDS** (n=1, the corpus's only `diarization.json`):
-1 call, 9 clusters named, confidences 0.30–0.90, 13.4 s wall, and the labels are genuinely
-good — `Mr. Obvious` at 73.8% of talk time plus named guests. Combined with the survey
-finding that only **865** of 78,917 video dirs still have audio, the diarized lane is the
-going-forward story and **making diarization runnable is now the highest-value next step**
-(sherpa-onnx uninstalled; system python 3.14.6 vs sherpa's cp313 wheels; pyenv has 3.12.12
-and `uv` is available).
-**[CORRECTED/DONE 2026-08-08 — "865 → ~1.1% reach" overstates by 10×: audio is necessary
-but a transcript is too, so the lane reaches 87 videos (0.115%). And sherpa-onnx IS now
-installed at `.diarize/`. See the top of this file.]** Full numbers: FACTS.md, "Attribution pilot" (2026-08-08). Raw
-artifacts in `plans/attribution-pilot/`. 9 sidecars were written into the real corpus;
-undo is `find transcripts/channels -name attribution.json -delete`.
-
-**Still open, deliberately:** the auto-caption quality sample (`quarteringvlogs`, 13
-chunks, ~15 min) was NOT run — it can only confirm a lane that already failed
-catastrophically, and the GPU was reclaimed. If the text-only lane is ever retried on a
-better instruction-follower, run it then; the harness takes `--sample channel:<slug>`.
-
-**Previously:** 2026-08-07 — **attribution is built, both lanes, and NEITHER is armed.
-Three measurements taken today invert PLAN.md's stated sequencing, and the reason
-matters.** On disk: **122** `ai-digest.json` of 73,367 transcribed (**0.17%**), **1**
-`diarization.json`, and `settings.digest.sweepEnabled` is **false**. PLAN.md says "run the
-text-only pass over the whole backlog immediately, then upgrade selectively" — but
-text-only attribution must read the transcript to find speaker changes at all, so it costs
-roughly the digest sweep's chunk count (**~194,000 model calls, another 25–55 GPU-days**),
-and it would spend that *before the digest sweep has finished 0.2% of its own*. Both lanes
-contend for the same 8 GB card and `backfillBatch`'s yield watches only the transcription
-lane, so nothing would arbitrate between them. The diarized lane is the opposite:
-`diarization.json` already carries globally-consistent cluster indices, so the model only
-names N clusters from samples — **about one call per video** — which is both cheaper and
-better, and its only problem is input coverage, which is exactly what the backfill lane
-exists to fix. So: **build both, pilot both, arm neither corpus-wide.** A measured pilot
-picks the default, the way the digest layer was done (pilot → bake-off → measured defaults
-→ sweep), which is the one pattern here that has worked.
-
-**Also landed:** the backfill lane itself is committed (`51b87bc`), and it now has three
-registered kinds instead of one — which is what turned it from a wrapper into an
-abstraction. Registering the two attribution kinds lit the channel stage card,
-`/actionable`, the dashboard instrument and the widget strip with **zero UI changes**,
-because all four iterate `snapshot.backfill[kindId]`. That was the claim the registry was
-justified by, and it held.
-
-**Recorded so it is not re-proposed: attribution cannot be a digest section.** Folding a
-`speakers` section into the digest prompt looks like it halves the GPU bill, and at 0.17%
-coverage invalidating the existing 122 digests would cost nothing. It does not work.
-Digests are generated **chunk-local** (the measured default at `PROMPT_VERSION` 2), so each
-chunk is labelled with no knowledge of the others — and the one property attribution needs
-above all is that speaker 0 in chunk 1 is the same person in chunk 30. Cross-chunk identity
-is the hard part of the text-only lane and the reason it is its own artifact.
-
-**Previously:** 2026-07-30 — **the cost model was wrong in its UNITS, and fixing
-that dissolved the conflict between three measurements.** A chunk (one model call)
-is the unit of work; chunk density varies 4× across the corpus, so
-seconds-per-audio-hour is not a unit. Re-priced, 27 / 90 / 151 s/audio-hour all
-agree to within ~2% and the corpus is **191,116 chunks → ~25 days idle to ~55
-contended, not 81**. The ~6–8 h bake-off replicate matrix was descoped as
-unnecessary: the census cost zero GPU time and settled the model choice too. Tags
-were generated for the first time ever (they had NEVER run) and cost 5–15%, not
-100%. Two shipped bugs fixed. The contention measurement was attempted and
-**aborted on evidence** — see below; it is the one thing still open.
-
-**Previously:** 2026-07-29 — **the backfill is at the starting line.** A corpus-wide sweep
-can now be started from one control, survives a server restart, can be paused without being
-lost, yields the GPU to transcription, and reports coverage and an ETA in the unit the work
-is actually priced in. Stages A–D of the "cheapest-first" plan are all in; see the decisions
-below for what the measurements changed.
-
-**The headline is a negative result, and it should be read before planning any more cost
-work.** The plan's own premise — shrink the work-list via duplicate sharing — does not pay.
-Measured by `common/bin/digest-plan.ts`:
-
-| `--near` | clusters | to generate | saved by sharing |
-| --- | --- | --- | --- |
-| 0.6 (was) | 2,846 | 76,804 h | 0.5 d |
-| 0.45 | 7,110 | 75,851 h | 1.5 d |
-| **0.35 (now)** | **7,434** | **75,638 h** | **1.7 d** |
-| 0.25 | 7,493 | 75,616 h | 1.8 d |
-
-(The `sweep days` column this table used to carry has been dropped: it was
-audio-hours × 90 s, and that basis is retracted — see the correction below. The
-current threshold reprices to **185,483 chunks → 53.5 days** at the measured
-24.9 s/chunk, with **5,633 chunks / 1.6 days** avoided by sharing. The
-conclusion about sharing is unchanged, because it is a RATIO and both bases
-agree it is ~2%.)
-
-Loosening the threshold as far as it goes buys under 2 days. The plan assumed mirrors
-skew long; **they skew short**. Cluster members are 19% of the corpus by video count and
-**4.3% of its audio-hours** — the sweep is dominated by unclustered long-form VODs, and
-`HasanAbiVODs3` alone (8,331 audio-hours) outweighs every mirror in the corpus combined.
-Only ~55% of mirrors pass the alignment gate, so half the nominal saving is refused anyway.
-
-**CORRECTION (2026-07-29, later the same day): the paragraph that stood here was
-wrong, and it was my own.** It claimed GPU arbitration was "worth ~56 days" and
-"about fifty times every duplicate lever put together", from 90 s/audio-hour
-measured against 27 s projected. That 2.8× is not a contention penalty. It is
-worse than confounded — it compares the validation run's **wall** clock against
-round 2's **engine** time, which are not the same quantity, across samples with
-different length mixes.
-
-Re-priced per chunk (the unit a model call is actually billed in), the same three
-measurements agree to within ~2% and the factor decomposes exactly: **1.49× from
-the validation sample's length mix × 2.22× per-chunk cost = 3.30×**, against the
-3.33× observed. The honest corpus range is **~25 days idle to ~55 contended**, and
-the ~81-day headline was wrong on the validation run's own data. See
-`controller/digestPlan.ts` for the census and the arithmetic.
-
-**Restated as a hypothesis with a named test.** The remaining 2.22× per-chunk gap
-has three parts: idle engine time, contended engine time, and per-chunk non-engine
-wall inflation in production (dominated by the yield gate idling at its 3 s poll).
-No bake-off run can see the third — a standalone `tsx` process reads `busy: false`
-because `getWorkerPool()` is a `globalThis` singleton. The test is a scoped
-production sweep run **twice, differing only in `digest.yieldToCpuWorkers`**, with
-wall/engine computed from `.jobs/<id>.meta.json` against the summed per-call lines
-in `.jobs/<id>.log`. On an idle box the third term is provably ~zero: the
-`teamrcn` job ran 41.578 s wall against 41.4 s of summed engine time — **1.004×**
-— and it did all the sidecar work a bake-off skips. Under contention it is
-unbounded, which is the point of measuring it.
-
-`controller/digestYield.ts` is implemented, and **it shipped with a bug**: it
-yielded for any busy `kind === "local"` worker, but this box runs one GPU worker
-beside two pinned to `device: "cpu"`, so at `parallelTranscriptions: 2` the digest
-lane stopped dead for transcription competing for zero GPU shaders. Fixed and
-made configurable (`digest.yieldToCpuWorkers`, default off). That fix is plausibly
-a bigger throughput lever than the model choice, and it is a logic error rather
-than physics.
-
-**The contention measurement was ATTEMPTED AND ABORTED, and it is still open.**
-Not descoped — attempted, aborted on evidence, and left with its tooling in place.
-The box could not support it: an unrelated containerised Playwright workload for
-another project was running (and respawning), and over the 90 s of sampling
-`plans/tools/aborted-run-2026-07-30.csv` records **swap going 9.6 → 18.4 GB of
-19 GB**, `MemAvailable` down to 1.7 GB, load average 3.4 → 17.0, and finally the
-editor dev server dying mid-request with 28 s response times. A run under those
-conditions measures the swap, which this document already warns about. **A number
-from it would have been worse than no number**, since it would have gone into
-FACTS.md and been planned against — the exact failure this session spent its time
-correcting.
-
-What the attempt DID confirm in production, none of which is timing-dependent:
-
-- **`device` reaches the yield.** `GET /api/workers` reports `device: "cpu"` for
- both CPU workers and no `device` for the GPU one, from the same `summary()` call
- `transcriptionActivity()` reads. That was the missing plumbing in the 1b fix.
-- **The sweep resumes from settings alone**, again: a restart re-launched
- `digest-sweep` + `digest-channel-local` on the persisted scope with no cursor.
-- **The chunk-priced log is live**: `Pass 1: 1 channel(s), 8 chunks / 1
- audio-hours to generate (~0.0 days at the measured rate, 9.14 chunks/audio-h)`.
-- **Tags generate through the real job path** with 0 warnings.
-- **The per-call breakdown is now IN the job log** — `in 5s wall (load 0.4s,
- prefill 0.1s, decode 3s, engine 3.5s)`. This is the thing that was previously
- unrecoverable: job logs recorded per-call seconds with no way to separate engine
- from wall, and the last run's logs had already rotated away.
-
-**What remains unmeasured: the per-chunk wall/engine multiplier under real GPU
-contention, and the live A/B of `yieldToCpuWorkers`.** The A/B is unit-tested in
-both directions against this box's exact worker set
-(`controller/digestYield.test.ts`), so what is missing is specifically a
-production run with genuinely busy CPU workers — which needs eligible
-transcription work, and the only work available is 2.3–4.6 GB multi-hour VODs that
-would take hours on CPU-only workers.
-
-To finish it, on a box with nothing else running:
-
-1. `plans/tools/sample-contention.sh out.csv` — start the sampler first.
-2. Enable only the two `device: "cpu"` workers (note: `transcripts/.workers/defaults.json`
- OVERRIDES `settings.json`'s per-worker `enabled` at boot — an hour lost to that).
-3. Queue transcription for a channel with media on disk (`cornbreadman` has 7).
-4. Arm a scoped sweep, `sections: ["chapters","tags"]`, and read
- `.jobs/<id>.log` for the yield line and the per-call breakdown.
-5. Repeat with `digest.yieldToCpuWorkers: true`, changing NOTHING else.
-6. Void any cell whose CSV shows `gpu_worker_busy` 0 throughout, swap growing, or
- thermal throttling.
-
-**The sweep was run for real, not just built.** Scoped to `teamrcn` (7 videos,
-0.27 audio-hours) against ollama and the real corpus: it armed with its scope
-persisted, launched the per-channel job, generated 7 digests with correct
-provenance at `promptVersion` 2, refreshed the snapshot and ended clean. A second
-process then resumed **from settings alone**, found nothing to do, launched
-**zero** jobs and exited — the B2 claim proven, since the correct amount of
-rework after a restart is none.
-
-Doing that found two bugs no test would have. Every channel's progress bar would
-have stalled one short forever (`countMissingDigests` counted a transcript-less
-directory the batch correctly drops — 2,987 such videos corpus-wide), and
-stopping a sweep left its channel scope behind, so the next sweep would silently
-inherit a weeks-old list and report itself finished. Both fixed and re-verified.
-
-Throughput datapoint, uncontended: **~151 s per audio-hour** on 2.4-minute
-videos — which looked alarming and is in fact *exactly on model*. Those 7 videos
-are 1 chunk each over 0.27 audio-hours, i.e. 26 chunks/audio-hour against the
-corpus's 2.47, and 151/26 = **5.9 s per chunk** on an idle box. Priced in chunks
-it is the cheapest datapoint anyone has taken here, not the most expensive. This
-is the clearest single illustration of why s/audio-hour is not a usable unit.
-
-**Previous entry —** 2026-07-29 — **Phases 2 and 3: the digest layer SHIPS.** The 102 digests
-already on disk now travel generation → LMDB → a shared `/digests/<slug>/` page tree →
-compose → the viewer, where a reader can open a Digest panel, click a chapter and seek to it.
-Proven end to end against the real corpus, not fixtures: `community-notes/v2cywen` renders its
-20 real chapters (first at 01:09:48 — the validation run's worst coverage gap) and the
-playhead follows a click. The control is hidden on the 99.9% of videos with no digest, and a
-borrowed digest is labelled as borrowed.
-
+**Last updated:** 2026-08-08 (latest) — **the OOM wall is gone, and it was removed rather
+than worked around.** Branch `feat/diarization-oom-wall`, four parts, all verified.
+
+**THE WALL: windowed diarization.** 6 of 10 videos over 6 hours were being killed by the
+kernel mid-run, ~40 minutes of CPU each, producing nothing. Two independent memory sinks,
+both now closed by the same change: `sherpa_onnx` clustering holds an O(n²) pairwise matrix
+over speech segments, and `subprocess.run(stdout=PIPE)` buffered the whole decoded stream.
+Long files are now processed in 45-minute windows (`--window-after-minutes 90`, `0` disables
+windowing) with per-window decode straight into a preallocated array. Cross-window identity
+is real, not hand-waved: each window's local speakers are reduced to an averaged, normalized
+`SpeakerEmbeddingExtractor` centroid and those are clustered globally with `FastClustering`
+at the same threshold. Both APIs verified live against the installed 1.13.4.
+
+| gate | result |
+| --- | --- |
+| (a) short file unchanged | **PASS** — `ObviousRises-rumble/v6z1o2g` re-ran byte-identical: 51 turns, 13 speakers, every non-volatile field equal, no `windowing` key |
+| (b) the file that OOM'd | **PASS** — `hasanabi/2834644359` (6h12m) completes: 3,433 turns, 96 speakers, 46 min, **peak RSS ~2.0 GB vs the 10.5 GB kernel kill** (`/proc/<pid>/status` VmHWM; GNU `time` is not installed) |
+| (c) quality | **PASS** — windowed vs whole-file on `v6z1o2g`: speech total **−0.1%**, identical wall-clock coverage, identity MORE concentrated (top-1 88.2% vs 73.1%, 6 speakers vs 13) rather than fragmented |
+
+**A REAL DEFECT WAS FOUND BY CHECKING, NOT BY READING.** The first seam merge joined only
+CONSECUTIVE turns in the sorted list, which silently leaves the same speaker overlapping
+THEMSELVES when two local speakers map to one global one — 35 self-overlaps and 24 sub-gap
+splits on the 6h12m file, which would double-count talk time in every downstream consumer.
+Now a per-speaker interval union; re-verified at 0 and 0. Cross-speaker overlap is left
+alone, because that is genuine overlapping speech the whole-file path emits too.
+
+**The cap shipped and is already switched off.** `settings.diarization.maxAudioHours` landed
+at 4h as the stopgap and now defaults to **0** — windowing removed the thing it was for. It
+stays as a knob for a smaller machine. Over-cap videos report as a new `deferred`
+`BackfillState`: a THIRD number, never summed into `reachableBackfillWork`, so a capped
+corpus can never read as finished. The dangerous line was `backfillBatch`'s candidate pull,
+which handled the states it knew and FELL THROUGH TO DISPATCHING everything else, with no
+exhaustiveness check anywhere in the repo. It is now a pure `candidateAction()` with a
+`never` default — a compile error instead of a silently-diarized six-hour video.
+
+**The detection bugs are fixed, and the fix is a shape, not a longer denylist.**
+`isRealAudioFile` / `isSourceMediaFile` / `isPartAudioFile` are anchored single-segment
+matches over `AUDIO_EXTS ∪ VIDEO_CONTAINER_EXTS`, in a new pure `common/lib/mediaFiles.ts`
+(`videoStatus.ts` re-exports, so no import site changed). An extension allowlist alone would
+NOT have been enough — this repo's own `audio.tmp-<pid>.<fmt>` carries a valid audio
+extension. Three more sites carried private copies of the old denylist and are folded in.
+Safety is provable: `audioFilesToRemove` is monotone in the predicate, so the delete set can
+only shrink.
+
+**THE MEDIA SCAN FOUND MORE THAN IT WAS WRITTEN FOR.** 953 media files across 68 channels in
+**83 seconds**, tier 1 only:
+- the known 4.0 GiB unreadable `nuxanor-kick/.../source-media.temp.mp4` (`moov atom not found`)
+- **3 zero-byte `audio.mp3`** in `leaflit-rumble` (`v6v9mq3`, `v6v9o25`, `v6v9r2r`) — these
+ read as "downloaded", so the pipeline never re-fetches them, and `v6v9r2r` has a real
+ `audio.aac` beside it that the mp3-first read preference would pass over. Newly visible.
+- **39 `transcript.m4a` in `community-notes`, 1.17 GiB** — audio under a name nothing in the
+ app recognizes (not audio, not a subtitle: `SUB_FILE_RE` needs two segments). Present but
+ invisible.
+- **7 duration mismatches** (up to 5h) awaiting the opt-in tier-2 decode.
+
+**STILL YOUR CALL, nothing done:** deleting anything the scan found; whether
+`community-notes`' 39 `transcript.m4a` files should be renamed into place (that is a WRITE to
+the corpus and it is not mine to make); and re-running the 4 previously-`failed` diarizations
+that were really `no-audio`.
+
+**Coordination note:** `settings.backfill.enabled` was set to false (the same field the Pause
+button writes) while `backfillKinds.ts` / `backfillBatch.ts` were edited, and restored after.
+A backfill job that was already running keeps the OLD code until it is stopped and started
+again — server modules hot-reload under `next dev`, but a job in flight holds its closure.
---
## Phase status
diff --git a/scripts/diarize-sherpa.py b/scripts/diarize-sherpa.py
@@ -8,15 +8,338 @@ swapping this engine for pyannote later means replacing one file.
stdout is JSON ONLY. Progress goes to stderr so the wrapper can stream it.
+WINDOWING, AND WHY IT EXISTS
+----------------------------
+Long recordings used to be killed by the kernel partway through. Two independent
+memory sinks, both measured rather than guessed:
+
+1. sherpa's clustering holds a pairwise distance matrix over speech-segment
+ embeddings -- O(n^2) in SEGMENT count. At 30k segments that is 6.7 GiB and at
+ 40k it is 11.9 GiB, which brackets the 10.6 GB and 9.6 GB peaks observed
+ before this. Note it is segment count, not duration: speaker-turn density
+ varies 40x across this corpus, which is why a sparse 7h42m file survived
+ while a dense 6h12m one did not.
+2. the whole-file decode below buffers the ENTIRE decoded stream as one bytes
+ object -- 1.84 GB for an 8h file, and Python's incremental buffer growth
+ peaks near twice that.
+
+Windowing fixes both at once, which is why it is the real answer and not a
+mitigation: per-window n falls by the window count so the matrix falls by its
+SQUARE, and decoding one window at a time removes the whole-file buffer
+entirely. A 45-minute window holds 165 MiB against 1.72 GiB for a whole 8h file.
+
+The hard part is that speaker 2 in window 1 and speaker 5 in window 2 may be the
+same person, and a diarization that cannot say so is not worth having. So each
+window's local speakers are reduced to an embedding centroid and the centroids
+are clustered globally, with the same threshold, to recover identity across
+seams.
+
+SHORT FILES TAKE THE EXACT ORIGINAL CODE PATH -- same decode call, same
+sd.process, same speaker ids. That is a hard requirement, not tidiness: it is
+what keeps every sidecar already on disk reproducible.
+
Usage:
diarize-sherpa.py --seg <model.onnx> --emb <model.onnx> [--threshold F]
- [--threads N] [--ffmpeg PATH] <audioFile>
+ [--threads N] [--ffmpeg PATH] [--ffprobe PATH]
+ [--window-minutes M] [--window-after-minutes M]
+ <audioFile>
Output: {"turns":[{"start":s,"end":s,"speaker":i}], "audioSeconds":f,
- "version":"...", "sampleRate":n}
+ "version":"...", "sampleRate":n, "windowing":{...}|null}
"""
import argparse, json, subprocess, sys, time
+# How much of a speaker's audio to embed when building its per-window centroid.
+# Whole turns are used up to a cap: TitaNet needs a second or two to say anything
+# useful, and past ~10s per excerpt more audio stops changing the embedding.
+EMB_MAX_SEGMENTS = 5
+EMB_MAX_SEC = 10.0
+EMB_MIN_SEC = 0.5
+# Two consecutive turns of the same GLOBAL speaker separated by less than this
+# are joined. Matches the segmentation default for min_duration_off, so a seam
+# does not leave a visible split that the same audio processed whole would not
+# have had.
+MERGE_GAP_SEC = 0.5
+
+
+def decode_whole_file(np, ffmpeg, audio, sample_rate):
+ """The ORIGINAL decode, unchanged. See the module docstring: short files must
+ take this path byte-for-byte so existing sidecars stay reproducible."""
+ proc = subprocess.run(
+ [ffmpeg, "-v", "error", "-i", audio, "-f", "f32le",
+ "-ac", "1", "-ar", str(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 None
+ return np.frombuffer(proc.stdout, dtype=np.float32)
+
+
+def decode_window(np, ffmpeg, audio, sample_rate, start_sec, dur_sec):
+ """Decode ONE window into a PREALLOCATED array.
+
+ The allocation is the point. subprocess.run(stdout=PIPE) accumulates the
+ whole stream and grows its buffer as it goes, so peak memory is roughly
+ twice the final size; here the exact sample count is known in advance from
+ the window length, so ffmpeg writes straight into one array that is never
+ resized. Returns the array trimmed to what actually arrived -- the last
+ window is short, and a truncated file is shorter than its container claims.
+ """
+ n = int(round(dur_sec * sample_rate))
+ buf = np.empty(n, dtype=np.float32)
+ view = memoryview(buf).cast("B")
+ args = [ffmpeg, "-v", "error",
+ # -ss BEFORE -i is an input seek: ffmpeg does not decode and discard
+ # the preceding hours, which is what makes a late window as cheap as
+ # an early one. Verified sample-exact on this corpus's files.
+ "-ss", f"{start_sec:.6f}", "-t", f"{dur_sec:.6f}",
+ "-i", audio, "-f", "f32le",
+ "-ac", "1", "-ar", str(sample_rate), "-"]
+ proc = subprocess.Popen(args, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
+ filled = 0
+ try:
+ while filled < view.nbytes:
+ got = proc.stdout.readinto(view[filled:])
+ if not got:
+ break
+ filled += got
+ # ffmpeg can round a window up by a frame. Drained, not buffered.
+ while proc.stdout.read(65536):
+ pass
+ finally:
+ proc.stdout.close()
+ err = proc.stderr.read().decode("utf8", "replace")
+ proc.stderr.close()
+ rc = proc.wait()
+ if rc != 0:
+ sys.stderr.write(err)
+ raise RuntimeError(f"ffmpeg failed ({rc}) decoding window at {start_sec:.1f}s")
+ return buf[: filled // 4]
+
+
+def probe_duration(ffprobe, audio):
+ """Container duration in seconds, or None when it cannot be measured.
+
+ None is load-bearing: it sends the run down the whole-file path, i.e. exactly
+ what this script did before windowing existed. Guessing a duration and
+ windowing on it would be the one way to make an unreadable file produce a
+ WRONG answer rather than the old one.
+ """
+ try:
+ p = subprocess.run(
+ [ffprobe, "-v", "error", "-show_entries", "format=duration",
+ "-of", "default=noprint_wrappers=1:nokey=1", audio],
+ stdout=subprocess.PIPE, stderr=subprocess.PIPE,
+ )
+ if p.returncode != 0:
+ return None
+ sec = float(p.stdout.decode("utf8", "replace").strip())
+ return sec if sec > 0 else None
+ except Exception:
+ return None
+
+
+def plan_windows(total_sec, window_sec, overlap_sec):
+ """[(start, dur, own_start, own_end)] covering total_sec.
+
+ Windows OVERLAP so the segmentation model has context either side of a seam,
+ but each second of the recording is OWNED by exactly one window -- the
+ ownership boundaries meet in the middle of the overlap. A turn is kept by the
+ window that owns its midpoint, so the overlap buys context without producing
+ the same speech twice.
+ """
+ step = window_sec - overlap_sec
+ spans = []
+ start = 0.0
+ while True:
+ dur = min(window_sec, total_sec - start)
+ spans.append((start, dur))
+ if start + dur >= total_sec - 1e-6:
+ break
+ start += step
+ out = []
+ for i, (s, d) in enumerate(spans):
+ own_s = s if i == 0 else s + overlap_sec / 2.0
+ own_e = s + d if i == len(spans) - 1 else s + d - overlap_sec / 2.0
+ out.append((s, d, own_s, own_e))
+ return out
+
+
+def speaker_centroid(np, extractor, samples, sample_rate, segments):
+ """One L2-normalized embedding for a local speaker, or None.
+
+ Averaged over the speaker's LONGEST turns rather than computed from one:
+ a single turn can be a two-word interjection, and an identity built from that
+ is what makes a speaker split across a seam.
+ """
+ picked = sorted(segments, key=lambda s: s[1] - s[0], reverse=True)
+ picked = picked[:EMB_MAX_SEGMENTS]
+ vecs = []
+ for start, end, _ in picked:
+ end = min(end, start + EMB_MAX_SEC)
+ if end - start < EMB_MIN_SEC:
+ continue
+ a = int(start * sample_rate)
+ b = min(int(end * sample_rate), samples.size)
+ if b - a < int(EMB_MIN_SEC * sample_rate):
+ continue
+ stream = extractor.create_stream()
+ stream.accept_waveform(sample_rate, samples[a:b])
+ stream.input_finished()
+ if not extractor.is_ready(stream):
+ continue
+ v = np.asarray(extractor.compute(stream), dtype=np.float32)
+ norm = float(np.linalg.norm(v))
+ if norm > 0:
+ vecs.append(v / norm)
+ if not vecs:
+ return None
+ mean = np.mean(np.stack(vecs), axis=0)
+ norm = float(np.linalg.norm(mean))
+ return mean / norm if norm > 0 else None
+
+
+def merge_and_renumber(turns):
+ """Absolute-time turns -> final output.
+
+ Two jobs, both seam repair. Renumbering by first appearance, because global
+ cluster labels come out of the clusterer in no useful order and the rest of
+ this app reads speaker 0 as "the first voice heard". Then unioning each
+ speaker's own intervals, so a speaker who talked straight through a seam is
+ not reported as having stopped and started.
+
+ PER SPEAKER, NOT PAIRWISE DOWN THE SORTED LIST. Merging only CONSECUTIVE
+ turns looks equivalent and is not: two different local speakers can map to
+ the same global one, and their turns then interleave with — or overlap —
+ other people's. A pairwise pass leaves those unjoined, and worse leaves the
+ same speaker overlapping THEMSELVES, which double-counts their talk time in
+ every downstream consumer. Measured on a 6h12m VOD before this was fixed: 35
+ self-overlaps and 24 sub-gap splits that a pairwise merge could not see.
+ Cross-speaker overlap is left alone — that is genuine overlapping speech and
+ the whole-file path emits it too.
+ """
+ order = {}
+ for t in sorted(turns, key=lambda t: (t["start"], t["end"])):
+ if t["speaker"] not in order:
+ order[t["speaker"]] = len(order)
+
+ by_speaker = {}
+ for t in turns:
+ by_speaker.setdefault(order[t["speaker"]], []).append(
+ (t["start"], t["end"])
+ )
+
+ merged = []
+ for speaker, spans in by_speaker.items():
+ spans.sort()
+ cur_start, cur_end = spans[0]
+ for start, end in spans[1:]:
+ if start - cur_end <= MERGE_GAP_SEC:
+ cur_end = max(cur_end, end)
+ else:
+ merged.append(
+ {"start": cur_start, "end": cur_end, "speaker": speaker}
+ )
+ cur_start, cur_end = start, end
+ merged.append({"start": cur_start, "end": cur_end, "speaker": speaker})
+
+ merged.sort(key=lambda t: (t["start"], t["end"], t["speaker"]))
+ return merged
+
+
+def run_windowed(np, sherpa_onnx, sd, args, total_sec, window_sec, overlap_sec):
+ """Window -> local diarization -> per-speaker centroid -> global clustering."""
+ sample_rate = sd.sample_rate
+ windows = plan_windows(total_sec, window_sec, overlap_sec)
+ print(
+ f"diarize-sherpa: {total_sec/3600:.2f} h -> {len(windows)} window(s) of "
+ f"{window_sec/60:.0f} min ({overlap_sec:.0f}s overlap), windowed mode",
+ file=sys.stderr, flush=True,
+ )
+
+ extractor = sherpa_onnx.SpeakerEmbeddingExtractor(
+ sherpa_onnx.SpeakerEmbeddingExtractorConfig(
+ model=args.emb, num_threads=args.threads
+ )
+ )
+
+ kept = [] # {"start","end","key"} in ABSOLUTE seconds
+ centroids = [] # parallel to keys
+ keys = [] # (window index, local speaker)
+ decoded_sec = 0.0
+
+ for wi, (ws, dur, own_s, own_e) in enumerate(windows):
+ samples = decode_window(np, args.ffmpeg, args.audio, sample_rate, ws, dur)
+ if samples.size == 0:
+ print(f"diarize-sherpa: window {wi+1} decoded zero samples; skipping",
+ file=sys.stderr, flush=True)
+ continue
+ decoded_sec += samples.size / sample_rate
+ t0 = time.time()
+ result = sd.process(samples).sort_by_start_time()
+ segs = [(s.start, s.end, s.speaker) for s in result]
+
+ # Keep the turns this window OWNS, by midpoint. The last window owns its
+ # own right edge, or the tail of the recording would be dropped.
+ is_last = wi == len(windows) - 1
+ owned_by_speaker = {}
+ for start, end, spk in segs:
+ mid = start + (end - start) / 2.0 + ws
+ if mid < own_s or (mid >= own_e and not is_last):
+ continue
+ owned_by_speaker.setdefault(spk, []).append((start, end, spk))
+ kept.append({"start": ws + start, "end": ws + end, "key": (wi, spk)})
+
+ # Centroids ONLY for speakers with owned turns: a local speaker whose
+ # every turn belongs to the neighbouring window would otherwise add a
+ # cluster that names nothing.
+ for spk, spk_segs in owned_by_speaker.items():
+ c = speaker_centroid(np, extractor, samples, sample_rate, spk_segs)
+ if c is None:
+ continue
+ centroids.append(c)
+ keys.append((wi, spk))
+
+ elapsed = time.time() - t0
+ print(
+ f"diarize-sherpa: window {wi+1}/{len(windows)} "
+ f"({ws/60:.0f}-{(ws+dur)/60:.0f} min): {len(segs)} turns, "
+ f"{len(owned_by_speaker)} local speaker(s) in {elapsed:.1f}s",
+ file=sys.stderr, flush=True,
+ )
+ del samples, result, segs
+
+ if not kept:
+ return [], decoded_sec, len(windows)
+
+ # Cross-window identity. Same clusterer and same threshold as within a
+ # window, so "these two are the same person" means the same thing at both
+ # scales.
+ if centroids:
+ labels = sherpa_onnx.FastClustering(
+ sherpa_onnx.FastClusteringConfig(
+ num_clusters=args.num_speakers, threshold=args.threshold
+ )
+ )(np.stack(centroids))
+ global_of = {key: int(label) for key, label in zip(keys, labels)}
+ else:
+ global_of = {}
+
+ # A local speaker with no usable centroid (all its turns too short to embed)
+ # keeps an identity of its own rather than being merged into someone else or
+ # dropped. Unmatched, but present and honest.
+ spare = (max(global_of.values()) + 1) if global_of else 0
+ for t in kept:
+ if t["key"] not in global_of:
+ global_of[t["key"]] = spare
+ spare += 1
+ t["speaker"] = global_of[t["key"]]
+ del t["key"]
+
+ return merge_and_renumber(kept), decoded_sec, len(windows)
+
def main():
ap = argparse.ArgumentParser()
@@ -29,6 +352,15 @@ def main():
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")
+ ap.add_argument("--ffprobe", default="ffprobe")
+ # Window length. 0 disables windowing entirely and forces the original
+ # whole-file path, whatever the duration.
+ ap.add_argument("--window-minutes", type=float, default=45.0)
+ # Files at or under this take the original whole-file path. Well above any
+ # duration that has ever been a problem here, so windowing stays the
+ # exception rather than quietly becoming the default.
+ ap.add_argument("--window-after-minutes", type=float, default=90.0)
+ ap.add_argument("--window-overlap-seconds", type=float, default=30.0)
args = ap.parse_args()
try:
@@ -60,36 +392,70 @@ def main():
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,
+ window_sec = args.window_minutes * 60.0
+ overlap_sec = max(0.0, min(args.window_overlap_seconds, window_sec / 2.0))
+ # Measured BEFORE decoding anything, because deciding after a whole-file
+ # decode would already have paid the memory this exists to avoid.
+ total_sec = probe_duration(args.ffprobe, args.audio) if window_sec > 0 else None
+ windowed = (
+ window_sec > 0
+ and total_sec is not None
+ and total_sec > args.window_after_minutes * 60.0
)
- 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,
- )
+ if windowed:
+ # Timed WHOLE, decode included: in this mode decoding is interleaved with
+ # the work rather than being a step before it, so there is no honest way
+ # to separate them.
+ t0 = time.time()
+ try:
+ turns, audio_seconds, window_count = run_windowed(
+ np, sherpa_onnx, sd, args, total_sec, window_sec, overlap_sec
+ )
+ except RuntimeError as e:
+ print(f"diarize-sherpa: {e}", file=sys.stderr)
+ return 4
+ elapsed = time.time() - t0
+ windowing = {
+ "windows": window_count,
+ "windowSeconds": round(window_sec, 3),
+ "overlapSeconds": round(overlap_sec, 3),
+ }
+ else:
+ # 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.
+ samples = decode_whole_file(np, args.ffmpeg, args.audio, sd.sample_rate)
+ if samples is None:
+ return 4
+ 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,
+ )
+ # Timed from HERE, excluding the decode, exactly as before windowing
+ # existed — the s/audio-hour figures already recorded in plans/ are
+ # engine time, and a number that quietly changed meaning is worse than no
+ # number.
+ 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
+ ]
+ windowing = None
- t0 = time.time()
- result = sd.process(samples).sort_by_start_time()
- elapsed = time.time() - t0
+ if windowed:
+ turns = [
+ {"start": round(t["start"], 3), "end": round(t["end"], 3),
+ "speaker": int(t["speaker"])}
+ for t in turns
+ ]
- 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 "
@@ -103,6 +469,7 @@ def main():
"audioSeconds": round(audio_seconds, 3),
"sampleRate": sd.sample_rate,
"version": getattr(sherpa_onnx, "__version__", None),
+ "windowing": windowing,
},
sys.stdout,
)
diff --git a/scripts/diarize.mjs b/scripts/diarize.mjs
@@ -36,7 +36,21 @@
// --threshold <f> clustering threshold (DIARIZE_THRESHOLD, 0.5)
// --threads <n> engine threads (DIARIZE_THREADS, 4)
// --ffmpeg <path> ffmpeg binary (FFMPEG_BIN, "ffmpeg")
+// --ffprobe <path> ffprobe binary (FFPROBE_BIN, "ffprobe")
+// --window-minutes <m> window length, 0 = never window
+// (DIARIZE_WINDOW_MINUTES, 45)
+// --window-after-minutes <m> only window files longer than this
+// (DIARIZE_WINDOW_AFTER_MINUTES, 90)
// -h, --help
+//
+// WINDOWING IS RECORDED BUT IS NOT PART OF THE FRESHNESS IDENTITY. It lands on
+// the record as a top-level `windowing` field, deliberately OUTSIDE `engine` —
+// isDiarizationFresh compares engine/models/threshold, and the target it
+// compares against is built from settings alone, which cannot know a given
+// video's duration. A window field in that identity would therefore mark every
+// sidecar on disk stale on the day windowing shipped, for work that is
+// unchanged. Same treatment `version` already gets, and for the same reason.
+// See common/lib/diarization.ts, which says so in as many words.
import { spawn } from "node:child_process";
import { rename, writeFile } from "node:fs/promises";
@@ -85,6 +99,15 @@ const threshold = Number(
);
const threads = Number(arg("--threads", process.env.DIARIZE_THREADS ?? "4"));
const ffmpeg = arg("--ffmpeg", process.env.FFMPEG_BIN ?? "ffmpeg");
+const ffprobe = arg("--ffprobe", process.env.FFPROBE_BIN ?? "ffprobe");
+const windowMinutes = arg(
+ "--window-minutes",
+ process.env.DIARIZE_WINDOW_MINUTES ?? "45",
+);
+const windowAfterMinutes = arg(
+ "--window-after-minutes",
+ process.env.DIARIZE_WINDOW_AFTER_MINUTES ?? "90",
+);
const engineCmd = arg("--engine", process.env.DIARIZE_ENGINE_CMD ?? "");
@@ -105,6 +128,9 @@ if (engineCmd) {
"--threshold", String(threshold),
"--threads", String(threads),
"--ffmpeg", ffmpeg,
+ "--ffprobe", ffprobe,
+ "--window-minutes", String(windowMinutes),
+ "--window-after-minutes", String(windowAfterMinutes),
audio,
];
}
@@ -144,6 +170,11 @@ child.on("close", async (code) => {
? { audioSeconds: raw.audioSeconds }
: {}),
durationMs: Date.now() - started,
+ // Outside `engine` on purpose — see the header. Absent for a whole-file run,
+ // which is what every sidecar written before windowing existed looks like.
+ ...(raw.windowing && typeof raw.windowing === "object"
+ ? { windowing: raw.windowing }
+ : {}),
speakers: new Set(turns.map((t) => t.speaker)).size,
turns,
engine: {