// Run speaker diarization for ONE video and write its diarization.json sidecar. // // The whole point of this lane is timing, not quality: audio is the one input // that goes away (cleanAudioFromTranscribed deletes it once a video is // transcribed), so the speaker ranges have to be captured while it is still on // disk. Everything downstream — attribution, LLM speaker naming, viewer badges, // quote filtering — can be redone from the saved JSON at any time. // // Never throws for an expected condition. Callers run this from a post-transcribe // hook and from a corpus-wide backfill, and in both places a diarization problem // must degrade to a skip, never fail the surrounding work. import path from "node:path"; import fs from "fs-extra"; import { execa } from "execa"; import type { Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; import type { DiarizationSettings } from "../lib/settings"; import { DIARIZATION_FILENAME, SORTFORMER_DIARIZATION_ENGINE, diarizationTarget, isDiarizationFresh, } from "../lib/diarization"; import { hasDiarization, loadDiarization } from "../lib/diarization-server"; import { findSourceMedia, isRealAudioFile } from "../lib/videoStatus"; import { pickPreferredAudio } from "../lib/mediaFiles"; import { resolveSavedVideo } from "../lib/savedVideo-server"; const { readdir } = fs; export type DiarizeOneOptions = { paths: Paths; videoDir: string; videoId: string; // Overrides settings.diarization. The backfill passes one resolved copy rather // than re-reading settings per video. settings?: DiarizationSettings; // Redo even when a valid sidecar already exists (an engine/threshold change). force?: boolean; onLog?: (msg: string) => void; signal?: AbortSignal; }; export type DiarizeOneOutcome = // Wrote (or rewrote) diarization.json. | "diarized" // A valid sidecar already existed AT THE CURRENT IDENTITY (engine, models, // threshold) and force was not set. A sidecar from a different identity is // stale, not this — it gets regenerated. | "already-exists" // The lane is switched off in settings. | "disabled" // Nothing to diarize from: no audio and no persisted source container. | "no-audio" // Configured but unusable — no model paths. Reported distinctly from a // failure so an operator can tell "not set up" from "tried and broke". | "not-configured" // The engine ran and failed. Logged, never thrown. | "failed"; // Resolve something diarizable in the video dir: extracted audio first, then a // persisted source container, then the saved-video store. Mirrors // transcribeOne's resolver — ffmpeg reads any of them. export async function resolveDiarizableMedia( videoDir: string, ): Promise { const entries = await readdir(videoDir).catch(() => [] as string[]); const preferred = pickPreferredAudio(entries.filter(isRealAudioFile)); if (preferred) return preferred; const container = findSourceMedia(entries); if (container) return container; const saved = await resolveSavedVideo(videoDir); if (saved) return path.relative(videoDir, saved); return null; } export async function diarizeOneVideo( opts: DiarizeOneOptions, ): Promise { const log = opts.onLog ?? ((m: string) => console.log(m)); const cfg = opts.settings ?? getSettings().diarization; if (!cfg.enabled) return "disabled"; // Each engine has its own idea of "configured", and reporting the wrong one // sends an operator to fix a setting the selected engine never reads. const usingSortformer = cfg.engine === SORTFORMER_DIARIZATION_ENGINE; if (usingSortformer) { if (!cfg.sortformerBin || !cfg.sortformerModel) { log( `Diarize ${opts.videoId} skipped: sortformer engine selected but no binary/model configured (run scripts/build-sortformer.sh).`, ); return "not-configured"; } } else if (!cfg.segModel || !cfg.embModel) { log( `Diarize ${opts.videoId} skipped: no segmentation/embedding model configured.`, ); return "not-configured"; } // FRESHNESS, not existence. This used to short-circuit on `hasDiarization` // alone, which meant the two most likely reasons to re-run — a different // clustering threshold, a different model — left the entire corpus looking // done, and the only way to redo anything was `force` over everything. The // record carries its provenance precisely so this comparison is possible; see // isDiarizationFresh for what is compared and why `version` is not. if (!opts.force) { const existing = await loadDiarization(opts.videoDir); if (existing && isDiarizationFresh(existing, diarizationTarget(cfg))) { return "already-exists"; } } const media = await resolveDiarizableMedia(opts.videoDir); if (!media) return "no-audio"; const start = Date.now(); log(`Diarize ${opts.videoId} start (${media})`); try { const engineArgs = usingSortformer ? [ "--engine-kind", SORTFORMER_DIARIZATION_ENGINE, "--sortformer-bin", cfg.sortformerBin, "--sortformer-model", cfg.sortformerModel, "--backend", cfg.backend, "--threads", String(cfg.threads), // No --ffprobe: sortformer never windows, so nothing needs to ask how // long the recording is. Its memory is O(1) in duration. "--ffmpeg", opts.paths.ffmpegBin, ] : [ "--seg", cfg.segModel, "--emb", cfg.embModel, "--threshold", String(cfg.threshold), "--threads", 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, ]; const child = execa( opts.paths.diarizeBin, [ "--output", DIARIZATION_FILENAME, "--video-id", opts.videoId, ...engineArgs, media, ], { cwd: opts.videoDir, cancelSignal: opts.signal, all: true, buffer: false, }, ); child.all?.on("data", (c: Buffer) => log(c.toString("utf8").trimEnd())); await child; } catch (err) { log(`Diarize ${opts.videoId} failed: ${(err as Error).message}`); return "failed"; } // Trust the file, not the exit code: an engine that exits 0 without producing // a parseable sidecar has NOT diarized this video, and saying otherwise is // what would let the cleanup sweep delete the audio. if (!(await hasDiarization(opts.videoDir))) { log( `Diarize ${opts.videoId} produced no valid ${DIARIZATION_FILENAME} — treating as failed.`, ); return "failed"; } log( `Diarize ${opts.videoId} done in ${((Date.now() - start) / 1000).toFixed(1)}s`, ); return "diarized"; }