// Corpus-wide diarization backfill: walk every channel/video that still has // audio and no diarization.json, and capture its speaker turns. // // This is what catches everything the post-transcribe hook cannot — videos // transcribed before the feature landed, videos transcribed with the inline hook // off (the recommended way to run a large batch), and anything whose inline // attempt failed. Modelled on normalizeAllTranscripts: pLimit over the channel // walk, with diarizeOneVideo short-circuiting on an existing sidecar so a // re-run does no rework. // // Concurrency defaults to settings.diarization.concurrency (1), NOT the 8 that // normalizeAll uses: normalizing cues is I/O, diarization is CPU-bound and // competes with GPU feeding and the digest sweep for the same 8 threads. import path from "node:path"; import { readdir } from "node:fs/promises"; import pLimit from "p-limit"; import { listChannelStatsFromDisk } from "./channels"; import { diarizeOneVideo } from "./diarizeOne"; import { getSettings } from "../lib/settings"; import type { Paths } from "../lib/paths"; export type DiarizeAllOptions = { paths: Paths; // Restrict to these channel slugs. Empty/omitted = the whole corpus. channelSlugs?: string[]; // Only diarize videos that already have a transcript. ON by default: the // capture lane exists to pair speaker turns with a transcript, and an // untranscribed video's audio is not at risk from the cleanup sweep yet. transcribedOnly?: boolean; force?: boolean; concurrency?: number; onLog?: (msg: string) => void; signal?: AbortSignal; }; export type DiarizeAllResult = { diarized: number; alreadyExists: number; noAudio: number; skipped: number; failed: number; }; export async function diarizeAll( opts: DiarizeAllOptions, ): Promise { const log = opts.onLog ?? ((m: string) => console.log(m)); const settings = getSettings().diarization; const result: DiarizeAllResult = { diarized: 0, alreadyExists: 0, noAudio: 0, skipped: 0, failed: 0, }; if (!settings.enabled) { log("Diarization is disabled in settings — nothing to do."); return result; } if (!settings.segModel || !settings.embModel) { log( "Diarization has no segmentation/embedding model configured — nothing to do.", ); return result; } const limit = pLimit(opts.concurrency ?? settings.concurrency); const wanted = new Set(opts.channelSlugs ?? []); const channels = (await listChannelStatsFromDisk(opts.paths)).filter( (ch) => wanted.size === 0 || wanted.has(ch.slug), ); const transcribedOnly = opts.transcribedOnly !== false; for (const ch of channels) { if (opts.signal?.aborted) break; const dataDir = path.join(opts.paths.channelsDir, ch.slug, "data"); const videoIds = await readdir(dataDir).catch(() => [] as string[]); let diarized = 0; let already = 0; let noAudio = 0; let skipped = 0; let failed = 0; await Promise.all( videoIds.map((id) => limit(async () => { if (opts.signal?.aborted) return; const videoDir = path.join(dataDir, id); if (transcribedOnly) { const entries = await readdir(videoDir).catch( () => [] as string[], ); if (!entries.includes("transcript.json")) { skipped++; return; } } try { const outcome = await diarizeOneVideo({ paths: opts.paths, videoDir, videoId: id, settings, force: opts.force, onLog: log, signal: opts.signal, }); if (outcome === "diarized") diarized++; else if (outcome === "already-exists") already++; else if (outcome === "no-audio") noAudio++; else if (outcome === "failed") failed++; else skipped++; } catch (err) { failed++; log(` ! ${ch.slug}/${id}: ${(err as Error).message}`); } }), ), ); // Only report channels that had something to do — the corpus is 66 channels // and most have no retained audio at all. if (diarized + already + failed > 0) { log( ` ${ch.slug}: diarized=${diarized} already=${already} no-audio=${noAudio} skipped=${skipped} failed=${failed}`, ); } result.diarized += diarized; result.alreadyExists += already; result.noAudio += noAudio; result.skipped += skipped; result.failed += failed; } log( `Done. diarized=${result.diarized} already=${result.alreadyExists} ` + `no-audio=${result.noAudio} skipped=${result.skipped} failed=${result.failed}`, ); return result; }