Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit 29d99aa8374c22265797a4ab2e8501e0555c22a3
parent 9b78c15f806b842ec587cb77ec4c32cad31d8698
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 10 Aug 2026 14:08:48 -0400

Give diarization a GPU engine, and a reason to prefer it

Speaker capture has always over-split. On a 13-minute reaction video with
one host, sherpa-onnx at its tuned threshold finds 13 speakers; on the
corpus's worst file it finds 35, with the dominant cluster at 48.9% and a
tail of noise. The turns are real but the identities are not, and the
threshold sweep that produced the 0.9 default only moved the number.

Sortformer decides speaker turns end-to-end instead of clustering
embeddings afterwards, so it does not have that failure mode: 4 speakers
on both files, agreeing with sherpa on the dominant speaker's share to
within half a point (73.5% vs 73.1%) while collapsing the invented tail.
It has no threshold to tune and a hard ceiling of 4 speakers.

The note in diarize.mjs saying no diarization model had been ported to
ggml is now out of date, and that is what makes this possible: the
Sortformer port links against ggml ALONE, so it inherits the same Vulkan
backend parakeet already uses on the RX 6600 XT.

Measured on this box, same file, idle:

  sherpa-onnx CPU      104.7s   492 s/audio-hour   391% CPU   503 MB
  sortformer Vulkan    190.3s   894 s/audio-hour    97% CPU   559 MB
  sortformer CPU (6t)  277.7s  1305 s/audio-hour   274% CPU  4.84 GB

sherpa is FASTER. This is a quality trade, not a speed one, and the GPU
is simply the cheapest way to run the engine that has the quality. Two
equivalences were verified and both held exactly: the CPU and Vulkan
backends produce byte-identical turns, and piping mp3 through ffmpeg into
the streaming stdin path gives the same 66 turns as running the binary on
a decoded WAV.

Memory is the other win. sherpa's clustering holds an O(n^2) matrix over
segments, which is why diarize-sherpa.py has a windowing/centroid-
reclustering design and why a duration cap exists. Sortformer's state is
a fixed-size speaker cache, and the driver streams raw PCM from ffmpeg
rather than decoding a file, so an 8-hour VOD costs the same 558 MB as a
6-minute clip. No windows, no cap, no ~900 MB temp WAV on a 98% disk.

Engine selection is a setting, and switching it restates the freshness
identity so the other engine's sidecars become work the backfill lane
offers to redo -- the two disagree about how many speakers exist, and a
corpus half-captured by each is not one corpus. The DEFAULT engine still
asserts nothing, because settings cannot predict which binary a --engine
override runs and the e2e fake engine would regenerate forever.

The lane yields the card to transcription (~4.4 GB of an 8 GB card that
parakeet and ollama also want), reusing the declared contendsFor rule
rather than a second copy of it. A `weight > 0` guaranteed share is
ignored while GPU capture is enabled: a share of VRAM is not a slower
run, it is an OOM.

sherpa-onnx stays the default and stays supported, for machines with no
usable GPU.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

Diffstat:
M.diarize/README.md | 114++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
M.gitignore | 4++++
Mcommon/controller/backfillBatch.test.ts | 32++++++++++++++++++++++++++++++++
Mcommon/controller/backfillBatch.ts | 32++++++++++++++++++++++++++++++--
Mcommon/controller/diarizeOne.ts | 70+++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Mcommon/lib/backfillKinds.test.ts | 120+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/backfillKinds.ts | 42++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/diarization.ts | 80++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mcommon/lib/settings.ts | 60+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Meditor/CHANGELOG.md | 4++++
Ascripts/build-sortformer.sh | 164+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Ascripts/diarize-sortformer.mjs | 156+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mscripts/diarize.mjs | 93+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------
Ascripts/sortformer/diarize-file.cpp | 268+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
14 files changed, 1185 insertions(+), 54 deletions(-)

diff --git a/.diarize/README.md b/.diarize/README.md @@ -1,8 +1,31 @@ # `.diarize/` — the diarization runtime -The Python environment and ONNX models `scripts/diarize.mjs` needs. **Everything -here except this file is gitignored** (`env/` and `models/`, ~148 MB, both -machine-specific). This file is tracked so the environment is reproducible. +The engines `scripts/diarize.mjs` needs. **Everything here except this file is +gitignored** — `env/` and `models/` for sherpa-onnx (~148 MB) and `sortformer/` +for the ggml engine (~530 MB), all machine-specific. This file is tracked so the +environment is reproducible. + +Two engines live here, selected by `settings.diarization.engine`: + +| | `sherpa-onnx` (default) | `sortformer` | +| --- | --- | --- | +| method | pyannote segmentation + TitaNet embeddings, agglomerative clustering | end-to-end streaming Sortformer | +| device | CPU only | Vulkan or CPU | +| speakers on `v6z1o2g` | 13 (ground truth: 1 host + clips) | **4** | +| speakers on the corpus's worst file | 35 | **4** | +| dominant-speaker share | 73.1% | 73.5% | +| throughput | **492–585 s/audio-hour**, ~3.9 cores | 894 s/audio-hour on Vulkan (~1 core), 1305 on tuned CPU | +| memory | ~230 MB/audio-hour, O(n²) in turn count | **558 MB flat**, O(1) in duration | +| knobs | threshold (0.9 here) | none | + +The sherpa range is two clean runs of the SAME file rather than an estimate: +585 s/audio-hour on 2026-08-08 (§ below) and 492 on 2026-08-10, both on an idle +box. Treat sub-20% differences between runs as noise on this hardware. + +sherpa-onnx is FASTER in wall clock. sortformer is chosen for quality: +over-splitting is the failure mode this lane has always had, and an end-to-end +model does not have it. See `plans/diarization-spike-results.md` for the original +CPU-only spike and the measurements that decided the defaults. ## Why it exists at all @@ -87,3 +110,88 @@ when the spike ran it under load ~30. Memory: ~230 MB per audio-hour of decoded float32 audio, so an 8-hour VOD needs ~1.8 GB resident. Untested on the long tail. + +--- + +## The sortformer engine + +### Rebuilding it + +```sh +scripts/build-sortformer.sh # Vulkan (default) +SORTFORMER_BACKEND=cpu scripts/build-sortformer.sh +``` + +Needs `git`, `cmake`, a C++17 compiler, and for Vulkan the loader + headers and +`glslc` (Arch: `vulkan-headers shaderc`). Installing `patchelf` is optional — it +sets `RPATH=$ORIGIN` on the binary; without it the wrapper sets +`LD_LIBRARY_PATH` itself. + +The script clones **openresearchtools/engine at a pinned commit** +(`8bb4928c`, 2026-03-21), stages only `ggml` + `tools/realtime` + our driver, and +builds those. It deliberately uses none of upstream's build system, whose scripts +are PowerShell/CUDA — the Sortformer code links against `ggml` alone (no llama, +no whisper, no Rust), which is what makes that possible. First build is ~10 +minutes, almost all of it compiling Vulkan shaders; rebuilds after a driver edit +are seconds, because `ggml` is staged once per pin. + +The model (471 MB, F32) is downloaded from +`openresearchtools/diar_streaming_sortformer_4spk-v2.1-gguf`. **It is under the +NVIDIA Open Model License, not the engine's MIT** — the GGUF is a conversion of +`nvidia/diar_streaming_sortformer_4spk-v2.1`. + +### Wiring it up + +```json +"diarization": { + "engine": "sortformer", + "backend": "vulkan", + "sortformerBin": "<repo>/.diarize/sortformer/diarize-file", + "sortformerModel": "<repo>/.diarize/sortformer/diar_streaming_sortformer_4spk-v2.1.gguf" +} +``` + +`threshold`, `segModel`, `embModel` and `python` are ignored by this engine and +do not enter its freshness identity. **Switching `engine` marks every sidecar +written by the other one stale**, which is intended — see `diarizationTarget` — +but on the retained audio that is weeks of rework, not a toggle. + +### Checking a rebuild reproduces the record on disk + +Read-only, writes nothing to the corpus: + +```sh +node scripts/diarize.mjs --output /tmp/repro.json --video-id v6z1o2g \ + --engine-kind sortformer \ + --sortformer-bin "$PWD/.diarize/sortformer/diarize-file" \ + --sortformer-model "$PWD/.diarize/sortformer/diar_streaming_sortformer_4spk-v2.1.gguf" \ + --backend vulkan --threads 6 \ + transcripts/channels/ObviousRises-rumble/data/v6z1o2g/audio.mp3 +``` + +Expected on 2026-08-10: **66 turns, 4 speakers, `audioSeconds` 766.101**, in +206.5 s (3.71x realtime). Two independent equivalences were checked and both +held exactly: + +* the **CPU and Vulkan backends produce byte-identical turns**, so the backend is + a throughput choice and never a quality one; and +* piping mp3 through ffmpeg into the streaming stdin path gives the **same 66 + turns** as running the binary directly on a decoded WAV. + +### Notes for whoever touches the driver + +`scripts/sortformer/diarize-file.cpp` is vendored here rather than taken from +upstream because upstream's `llama-realtime-smoke` is a parity tool: it retains +every intermediate matrix, demands PyTorch reference fixtures, and — the fatal +part — drops the event flags, so its JSON mixes 2,033 *preview* re-emissions in +with the 66 real spans. + +Two things in it are non-obvious and both were found the hard way: + +* **stdin is forced back to blocking.** Node hands children non-blocking pipes, + so `read()` returns `EAGAIN` long before EOF; it worked from a shell and failed + under the app. +* **thread count is set explicitly.** The upstream Sortformer path never sets + one, so the CPU backend silently runs at ggml's default of 4. On this 4c/8t box + 6 is the optimum (1305 s/audio-hour) and **8 is worse than 4** (1401) through + oversubscription. diff --git a/.gitignore b/.gitignore @@ -120,3 +120,7 @@ yarn-error.log* /.diarize/env/ /.diarize/models/ /.diarize/logs/ +# Sortformer engine: the built ggml binary + shared objects + the 471 MB GGUF. +# Machine-specific and rebuildable with scripts/build-sortformer.sh. +/.diarize/sortformer/ +/.diarize/.sortformer-build/ diff --git a/common/controller/backfillBatch.test.ts b/common/controller/backfillBatch.test.ts @@ -197,3 +197,35 @@ test("the duration cap round-trips, and 0 means off", () => { // back on by accident would silently stop diarizing long videos. assert.equal(defaultDiarization().maxAudioHours, 0); }); + +test("an unknown engine or backend falls back to the default, never to nothing", () => { + // The DEFAULT ENGINE IS LOAD-BEARING as a fallback, not just as a starting + // point: it is the one every sidecar already on disk matches, so falling back + // to it leaves the corpus fresh. Falling back to sortformer on a typo would + // mark all of it stale and offer weeks of rework. + assert.equal(defaultDiarization().engine, "sherpa-onnx"); + assert.equal(sanitizeDiarization({ engine: "sortformer" }).engine, "sortformer"); + assert.equal( + sanitizeDiarization({ engine: "sortfromer" }).engine, + defaultDiarization().engine, + ); + assert.equal(sanitizeDiarization({ engine: 7 }).engine, defaultDiarization().engine); + + // The backend defaults to the GPU, which is safe only because the lane yields + // the card rather than sharing it — see diarizationLaneFor. + assert.equal(defaultDiarization().backend, "vulkan"); + assert.equal(sanitizeDiarization({ backend: "cpu" }).backend, "cpu"); + assert.equal( + sanitizeDiarization({ backend: "rocm" }).backend, + defaultDiarization().backend, + ); + + // Paths are plain strings and empty means "not configured", which diarizeOne + // reports as a skip rather than a failure. + assert.equal(sanitizeDiarization({}).sortformerBin, ""); + assert.equal(sanitizeDiarization({}).sortformerModel, ""); + assert.equal( + sanitizeDiarization({ sortformerBin: " /opt/diarize-file " }).sortformerBin, + "/opt/diarize-file", + ); +}); diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts @@ -41,6 +41,7 @@ import { readVideoFiles } from "../lib/videoStatus"; import { addBackfillState, emptyBackfillCounts, + laneYieldsToTranscription, reachableBackfillWork, resolveBackfillKinds, type BackfillClassification, @@ -475,7 +476,8 @@ export async function runBackfillBatch( // poll and survives a restart with no boot hook (the downloadsPaused // pattern). Returning 0 makes runPool idle-wait — a hold; returning null // from next() would END the batch, which is not the same thing. - const live = getSettings().backfill; + const liveSettings = getSettings(); + const live = liveSettings.backfill; if (!live.enabled) { if (!yielding) { yielding = true; @@ -497,8 +499,34 @@ export async function runBackfillBatch( // can overlap on the CPU; the benefit is that neither can deadlock the // other. const activity = transcriptionActivity(); + // A GUARANTEED SHARE IS A CPU CONCEPT. `weight > 0` means "keep a slice of + // the cores running even while transcription works", which is reasonable + // when the contended resource is cores and divisible. It is not available + // for VRAM: sortformer on Vulkan holds ~4.4 GB of the same 8 GB card + // parakeet is using, so "a small share" is not a slower run, it is an + // out-of-memory failure of whichever lane allocates second. + // + // So a run that could dispatch GPU work is idle-only whatever the weight + // says. Which kinds those are is DECLARED (laneFor/contendsFor), not + // re-tested here, so the rule stays next to the queue key it belongs with. + // The cost is bluntness — one such kind makes the whole run idle-only, + // including its cheap CPU-bound siblings — and that is the deliberate + // direction to be wrong in, since the alternative is an OOM mid-sweep. + // + // ONLY kinds that declare laneFor are considered, and that is the fix for a + // real regression rather than a nicety. `digest` declares `contendsFor: + // "gpu"` statically, so testing every kind's lane made this lane idle-only + // whenever a digest was in the run — which broke the guarantee that a digest + // runs CONCURRENTLY with a backfill rather than behind it. A kind with a + // static GPU lane already arbitrates itself (digestBatch has its own yield); + // laneFor marks the kinds whose resource changes under them from settings, + // which today is diarization alone, and those are the ones nothing else is + // deciding for. + const gpuBound = kinds.some( + (k) => k.laneFor && laneYieldsToTranscription(k.laneFor(liveSettings)), + ); const limit = backfillLimit({ - weight: live.weight, + weight: gpuBound ? 0 : live.weight, slots: live.concurrency, primaryBusy: activity.busy, }); diff --git a/common/controller/diarizeOne.ts b/common/controller/diarizeOne.ts @@ -18,6 +18,7 @@ import { getSettings } from "../lib/settings"; import type { DiarizationSettings } from "../lib/settings"; import { DIARIZATION_FILENAME, + SORTFORMER_DIARIZATION_ENGINE, diarizationTarget, isDiarizationFresh, } from "../lib/diarization"; @@ -81,7 +82,17 @@ export async function diarizeOneVideo( const cfg = opts.settings ?? getSettings().diarization; if (!cfg.enabled) return "disabled"; - if (!cfg.segModel || !cfg.embModel) { + // 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.`, ); @@ -106,6 +117,44 @@ export async function diarizeOneVideo( 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, [ @@ -113,24 +162,7 @@ export async function diarizeOneVideo( DIARIZATION_FILENAME, "--video-id", opts.videoId, - "--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, + ...engineArgs, media, ], { diff --git a/common/lib/backfillKinds.test.ts b/common/lib/backfillKinds.test.ts @@ -12,6 +12,7 @@ import { operationCatalog, operationLabel, digestLaneFor, + diarizationLaneFor, laneYieldsToTranscription, addBackfillState, emptyBackfillCounts, @@ -40,6 +41,8 @@ import { import { DEFAULT_DIARIZATION_THRESHOLD, DIARIZATION_FILENAME, + SORTFORMER_DIARIZATION_ENGINE, + diarizationTarget, isDiarizationFresh, type DiarizationRecord, } from "./diarization"; @@ -245,6 +248,123 @@ test("a different engine binary is not, by itself, stale", async () => { ); }); +// The engine-as-a-setting case the test above says needs "no new logic". It is +// asserted for sortformer and NOT for the default, and that asymmetry is the +// whole design: the default engine still cannot predict which binary a +// `--engine` override runs, so it must keep asserting nothing. +test("selecting sortformer asserts the engine; selecting the default still does not", () => { + const rec = (engine: DiarizationRecord["engine"]): DiarizationRecord => ({ + videoId: "v", + generatedAt: "now", + speakers: 1, + turns: [], + engine, + }); + const sherpaSidecar = rec({ + engine: "sherpa-onnx", + segmentationModel: "seg-1.onnx", + embeddingModel: "emb-1.onnx", + threshold: DEFAULT_DIARIZATION_THRESHOLD, + }); + const sortformerSidecar = rec({ + engine: SORTFORMER_DIARIZATION_ENGINE, + model: "sortformer-4spk.gguf", + }); + + const sortformerTarget = diarizationTarget({ + engine: SORTFORMER_DIARIZATION_ENGINE, + sortformerModel: "/abs/path/sortformer-4spk.gguf", + // Still configured, because settings carry one set of fields for both + // engines. None of it may leak into the sortformer identity. + segModel: "/abs/seg-1.onnx", + embModel: "/abs/emb-1.onnx", + threshold: DEFAULT_DIARIZATION_THRESHOLD, + }); + + // Basename only, as everywhere else: the corpus is rsynced between shards. + assert.equal(sortformerTarget.model, "sortformer-4spk.gguf"); + assert.equal(sortformerTarget.segmentationModel, undefined); + assert.equal(sortformerTarget.embeddingModel, undefined); + + assert.equal(isDiarizationFresh(sortformerSidecar, sortformerTarget), true); + // The point of switching: every sherpa sidecar becomes work the backfill lane + // will offer to redo, rather than silently staying half a corpus. + assert.equal(isDiarizationFresh(sherpaSidecar, sortformerTarget), false); + + // And back the other way, with no special case needed — a sortformer record + // carries neither segmentation nor embedding model, so it fails the sherpa + // comparison on the models alone. + const sherpaTarget = diarizationTarget({ + engine: "sherpa-onnx", + segModel: "/abs/seg-1.onnx", + embModel: "/abs/emb-1.onnx", + threshold: DEFAULT_DIARIZATION_THRESHOLD, + }); + assert.equal(sherpaTarget.engine, undefined); + assert.equal(isDiarizationFresh(sherpaSidecar, sherpaTarget), true); + assert.equal(isDiarizationFresh(sortformerSidecar, sherpaTarget), false); +}); + +// Sortformer has no clustering step, so the threshold cannot have changed any of +// its turns. Letting it into the identity would regenerate the whole corpus for +// an edit that provably could not affect it. +test("the clustering threshold does not stale a sortformer sidecar", () => { + const sidecar: DiarizationRecord = { + videoId: "v", + generatedAt: "now", + speakers: 1, + turns: [], + engine: { engine: SORTFORMER_DIARIZATION_ENGINE, model: "m.gguf" }, + }; + const at = (threshold: number) => + diarizationTarget({ + engine: SORTFORMER_DIARIZATION_ENGINE, + sortformerModel: "m.gguf", + threshold, + }); + assert.equal(isDiarizationFresh(sidecar, at(DEFAULT_DIARIZATION_THRESHOLD)), true); + assert.equal(isDiarizationFresh(sidecar, at(0.4)), true); + // The model itself IS the identity, though — a different one is a redo. + assert.equal( + isDiarizationFresh( + sidecar, + diarizationTarget({ + engine: SORTFORMER_DIARIZATION_ENGINE, + sortformerModel: "other.gguf", + }), + ), + false, + ); +}); + +// Which resource diarization competes for is a CONFIGURATION outcome, and the +// backfill lane reads it to decide whether a guaranteed share is even available. +// Getting this wrong in the permissive direction puts ~4.4 GB of sortformer next +// to parakeet on an 8 GB card. +test("diarization contends for the GPU only as sortformer on vulkan", () => { + const lane = (engine: "sherpa-onnx" | "sortformer", backend: "vulkan" | "cpu") => + diarizationLaneFor({ engine, backend }); + + assert.equal(lane("sortformer", "vulkan").contendsFor, "gpu"); + assert.equal(laneYieldsToTranscription(lane("sortformer", "vulkan")), true); + + // The same engine on the CPU backend is only after cores. + assert.equal(lane("sortformer", "cpu").contendsFor, "cpu"); + assert.equal(laneYieldsToTranscription(lane("sortformer", "cpu")), false); + + // sherpa-onnx is ONNX/CPU, so a stale `backend: vulkan` left in settings must + // NOT make it claim the card — that would park the lane behind transcription + // for work using no shaders at all, which is the bug digestYield.ts already + // records having hit once. + assert.equal(lane("sherpa-onnx", "vulkan").contendsFor, "cpu"); + assert.equal(laneYieldsToTranscription(lane("sherpa-onnx", "vulkan")), false); + + // The queue key never changes: two diarizations must not run at once whichever + // engine is selected. + assert.equal(lane("sortformer", "vulkan").queueKey, BACKFILL_QUEUE); + assert.equal(lane("sherpa-onnx", "cpu").queueKey, BACKFILL_QUEUE); +}); + // THE COMPATIBILITY RULE, and the reason it is written down: without it, adding // a field to the provenance would mark all 77,000 videos stale at once. test("an absent recorded field compares equal to today's default", async () => { diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts @@ -73,8 +73,11 @@ import { } from "./digest"; import { loadDigest } from "./digest-server"; import { + SORTFORMER_DIARIZATION_ENGINE, diarizationTarget, isDiarizationFresh, + type DiarizationBackend, + type DiarizationEngineId, type DiarizationFreshnessTarget, } from "./diarization"; import { loadDiarization } from "./diarization-server"; @@ -259,7 +262,19 @@ export type BackfillKind = { tier: BackfillCostTier; // Where this operation's work runs. See BackfillLane — the queue key is what // keeps CPU and GPU operations overlapping instead of taking turns. + // + // This is the DECLARED lane, which for a kind with a choice means its default. + // Prefer laneFor() when a live answer is needed. lane: BackfillLane; + // The lane this kind would actually use under the given settings, for the + // kinds whose scarce resource is a configuration choice rather than a fact. + // Diarization is one: sherpa-onnx is CPU-only, while sortformer on the Vulkan + // backend holds ~4.4 GB of the same 8 GB card the transcription engine wants. + // + // Optional because most kinds have no choice, and `lane` is the answer for + // them. Mirrors digestLaneFor, which solved the same problem for the digest + // operation's two lanes. + laneFor?(settings: SiteSettings): BackfillLane; // Ids of other kinds in this table whose output this one consumes. // // PURELY DECLARATIVE. It does not gate anything by itself — a kind still @@ -315,7 +330,12 @@ const diarization: BackfillKind = { // kinds on purpose — two channels' worth of diarization at once just thrashes // cores. How much of the machine it may take is backfillLimit()'s question, // not the queue's. + // The DEFAULT engine's answer, written out rather than derived: settings.ts is + // a TYPE-only import here, and pulling defaultDiarization() in as a value would + // make this module's initialization depend on it at runtime. laneFor is the + // live answer, and diarizationLaneFor is where the rule actually lives. lane: { queueKey: BACKFILL_QUEUE, contendsFor: "cpu" }, + laneFor: (settings) => diarizationLaneFor(settings.diarization), enabled: (settings) => settings.diarization.enabled && !!settings.diarization.segModel && @@ -808,6 +828,28 @@ export function digestLaneFor(appLane: DigestLane): BackfillLane { { queueKey: DIGEST_REMOTE_QUEUE, contendsFor: "network" }; } +// Which resource a diarization run competes for, which follows from the engine +// it is configured with rather than from a separate setting — the same shape as +// digestLaneFor, and stated here so nothing has to re-test the engine id. +// +// The queue key does NOT change with the engine. Diarization serializes against +// itself either way, and giving the GPU variant its own key would only let two +// diarizations run at once — which is precisely what must not happen when each +// holds ~4.4 GB of an 8 GB card. +export function diarizationLaneFor(diarization: { + engine: DiarizationEngineId; + backend: DiarizationBackend; +}): BackfillLane { + return diarization.engine === SORTFORMER_DIARIZATION_ENGINE && + diarization.backend === "vulkan" + ? // Competes with the transcription engine for the same VRAM. Yields. + { queueKey: BACKFILL_QUEUE, contendsFor: "gpu" } + : // sherpa-onnx is ONNX/CPU, and sortformer on the CPU backend is likewise + // only after cores. Contends for CPU, whose share the backfill lane's own + // weight already governs. + { queueKey: BACKFILL_QUEUE, contendsFor: "cpu" }; +} + // Whether a lane must stand aside while transcription is working. One rule, so // the digest lane and any future GPU operation cannot answer it differently. export function laneYieldsToTranscription(lane: BackfillLane): boolean { diff --git a/common/lib/diarization.ts b/common/lib/diarization.ts @@ -35,6 +35,12 @@ export type DiarizationEngine = { // is machine-specific and would make records non-portable across shards). segmentationModel?: string; embeddingModel?: string; + // The single-model engines' model, kept separate from the segmentation/ + // embedding PAIR above rather than overloading it. Sortformer is one GGUF and + // has no embedding model at all, and writing it into `segmentationModel` would + // make a sortformer record compare equal to a sherpa one that happened to share + // a basename. + model?: string; // Engine/library version string, when the engine reports one. version?: string; // Clustering threshold used. The single most consequential knob: it decides @@ -81,6 +87,31 @@ export const DIARIZATION_FILENAME = "diarization.json"; // freshness comparator below and the wrapper agree on one spelling. export const DEFAULT_DIARIZATION_ENGINE = "sherpa-onnx"; +// The ggml/Sortformer engine (scripts/build-sortformer.sh, scripts/ +// diarize-sortformer.mjs). End-to-end rather than clustered: no segmentation + +// embedding pair, no threshold, a hard ceiling of 4 speakers, and — because its +// state is a fixed-size speaker cache — memory that is O(1) in recording length +// rather than O(n^2) in segment count. +export const SORTFORMER_DIARIZATION_ENGINE = "sortformer"; + +export type DiarizationEngineId = + | typeof DEFAULT_DIARIZATION_ENGINE + | typeof SORTFORMER_DIARIZATION_ENGINE; + +// Compute device for engines that have a choice. Only sortformer does; sherpa is +// ONNX/CPU here (no Vulkan compute path on Linux for onnxruntime). +export type DiarizationBackend = "vulkan" | "cpu"; + +export const DIARIZATION_ENGINE_IDS: readonly DiarizationEngineId[] = [ + DEFAULT_DIARIZATION_ENGINE, + SORTFORMER_DIARIZATION_ENGINE, +]; + +export const DIARIZATION_BACKENDS: readonly DiarizationBackend[] = [ + "vulkan", + "cpu", +]; + // The clustering-threshold default, measured on this corpus (see // DiarizationSettings.threshold for the sweep that produced it). It lives here // rather than only in settings.ts because isDiarizationFresh needs it to @@ -93,25 +124,30 @@ export const DEFAULT_DIARIZATION_THRESHOLD = 0.9; // idea: a sidecar is stale when its recorded provenance differs from this, not // when it is old. export type DiarizationFreshnessTarget = { - // OPTIONAL, and compared only when present — which today means never. + // OPTIONAL, and compared only when present — which now means "only when the + // configured engine is not the default". // - // The engine is not a setting: scripts/diarize.mjs records whatever binary - // actually ran (its own default, or the basename of a `--engine` / - // DIARIZE_ENGINE_CMD replacement), and nothing in DiarizationSettings can say - // which that will be. Asserting a hardcoded "sherpa-onnx" here would mark - // every sidecar produced by any other wrapper permanently stale — an infinite - // regeneration loop at ~500-680 s/audio-hour, and one the e2e fake engine - // would trip on its first run. + // The asymmetry is deliberate and load-bearing. scripts/diarize.mjs records + // whatever actually ran, including the basename of a `--engine` / + // DIARIZE_ENGINE_CMD replacement, so asserting a hardcoded "sherpa-onnx" here + // would mark every sidecar from any other wrapper permanently stale — an + // infinite regeneration loop at ~500-680 s/audio-hour, and one the e2e fake + // engine trips on its first run. So the default engine still asserts nothing + // and compares exactly as it always did. // - // So the comparison is restricted to what the configuration can genuinely - // predict: the models and the threshold, which are also the knobs that - // actually change the output. The field stays here, and the comparison stays - // written, so that adding an engine setting later needs no new logic. + // Selecting `sortformer` IS a configuration statement, and there the engine is + // asserted: sherpa sidecars become stale and the backfill lane offers to redo + // them, which is the point — the two engines disagree about how many speakers + // exist, and a corpus half-diarized by each is not one corpus. Switching back + // needs no special case: a sortformer record carries neither segmentation nor + // embedding model, so it fails the sherpa comparison on the models anyway. engine?: string; // Basenames, as DiarizationEngine records them — full paths are // machine-specific and would make every record stale on another shard. segmentationModel?: string; embeddingModel?: string; + // Single-model engines. See DiarizationEngine.model. + model?: string; threshold?: number; }; @@ -132,14 +168,27 @@ function baseName(p: string): string { // guards disagreeing about the identity is how a corpus ends up either // regenerating forever or never. export function diarizationTarget(cfg: { + engine?: DiarizationEngineId; segModel?: string; embModel?: string; + sortformerModel?: string; threshold?: number; }): DiarizationFreshnessTarget { + // Sortformer is end-to-end: one model, no segmentation/embedding pair, and no + // clustering threshold. Comparing sherpa's knobs against it would mark every + // sortformer sidecar stale on a threshold edit that could not have changed a + // single one of its turns. + if (cfg.engine === SORTFORMER_DIARIZATION_ENGINE) { + return { + engine: SORTFORMER_DIARIZATION_ENGINE, + ...(cfg.sortformerModel ? { model: baseName(cfg.sortformerModel) } : {}), + }; + } return { - // `engine` is deliberately NOT set — see DiarizationFreshnessTarget. Settings - // cannot know which binary will run, so claiming to compare it would mark - // every sidecar from any other wrapper permanently stale. + // `engine` is deliberately NOT set for the default engine — see + // DiarizationFreshnessTarget. Settings cannot know which binary a `--engine` + // override will run, so claiming to compare it would mark every sidecar from + // any other wrapper permanently stale. ...(cfg.segModel ? { segmentationModel: baseName(cfg.segModel) } : {}), ...(cfg.embModel ? { embeddingModel: baseName(cfg.embModel) } : {}), threshold: cfg.threshold ?? DEFAULT_DIARIZATION_THRESHOLD, @@ -184,6 +233,7 @@ export function isDiarizationFresh( (target.engine === undefined || e.engine === target.engine) && sameModel(e.segmentationModel, target.segmentationModel) && sameModel(e.embeddingModel, target.embeddingModel) && + sameModel(e.model, target.model) && sameThreshold(e.threshold, target.threshold) ); } diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -24,7 +24,14 @@ import { defaultAutoQueue, sanitizeAutoQueue, } from "../jobs/autoQueuePolicy"; -import { DEFAULT_DIARIZATION_THRESHOLD } from "./diarization"; +import { + DEFAULT_DIARIZATION_ENGINE, + DEFAULT_DIARIZATION_THRESHOLD, + DIARIZATION_BACKENDS, + DIARIZATION_ENGINE_IDS, + type DiarizationBackend, + type DiarizationEngineId, +} from "./diarization"; import { ATTRIBUTION_PROMPT_VERSION } from "./attribution"; import { DEFAULT_COOKIE_MODE, @@ -358,6 +365,33 @@ export type DiarizationSettings = { threshold: number; // Engine threads per diarize run. threads: number; + // Which engine runs. "sherpa-onnx" is the shipped default and what every + // sidecar on disk was produced by; "sortformer" is the ggml engine built by + // scripts/build-sortformer.sh. + // + // CHANGING THIS RESTATES THE FRESHNESS IDENTITY (see diarizationTarget), so + // every sidecar written by the other engine becomes stale and the backfill lane + // offers to redo it. That is intended — the two disagree about how many + // speakers exist, and a corpus half-diarized by each is not one corpus — but on + // the retained audio it is weeks of work, not a toggle. + // + // Why anyone would: on the same file, sherpa at its tuned threshold returns 13 + // speakers and sortformer returns 4, agreeing on the dominant speaker's share + // to within half a point (73.1% vs 73.5%). On the corpus's worst case sherpa + // returns 35 and sortformer 4. Over-splitting is the failure mode this lane has + // always had, and sortformer is end-to-end rather than clustered, so it does + // not have it. The cost is a hard ceiling of 4 speakers and ~1.8x the wall + // clock. + engine: DiarizationEngineId; + // Compute device for the sortformer engine; ignored by sherpa-onnx, which has + // no Vulkan compute path on Linux. + // + // "vulkan" is 1.5x faster than a thread-tuned CPU run (894 vs 1305 + // s/audio-hour, measured on this box) and holds 558 MB resident instead of + // 4.84 GB by keeping weights and activations in VRAM. It also takes ~4.4 GB of + // an 8 GB card, which is why the lane YIELDS to transcription rather than + // sharing — see controller/digestYield.ts. + backend: DiarizationBackend; // Python interpreter for the default sherpa-onnx engine. sherpa-onnx ships // wheels only up to cp313, and this box's system python is 3.14 — so this // usually points at a dedicated venv rather than `python3`. @@ -366,6 +400,11 @@ export type DiarizationSettings = { // is reported as a skip rather than a failure. segModel: string; embModel: string; + // Binary and model for the sortformer engine, both produced by + // scripts/build-sortformer.sh. Empty = that engine cannot run, reported as the + // same "not-configured" skip as an unset segModel/embModel. + sortformerBin: string; + sortformerModel: string; // How many diarize runs may execute at once in the backfill pass. Kept low by // default: diarization is CPU-bound and competes with GPU feeding and the // digest sweep for the same 8 threads. @@ -1111,9 +1150,17 @@ export function defaultDiarization(): DiarizationSettings { // and the comparator from drifting apart. threshold: DEFAULT_DIARIZATION_THRESHOLD, threads: 4, + // The engine every sidecar on disk was produced by. Switching is an explicit + // decision that restates the freshness identity — see DiarizationSettings. + engine: DEFAULT_DIARIZATION_ENGINE, + // Only consulted when engine is "sortformer". Defaulting to the GPU is safe + // because the lane yields the card to transcription rather than sharing it. + backend: "vulkan", python: "python3", segModel: "", embModel: "", + sortformerBin: "", + sortformerModel: "", concurrency: 1, // 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 @@ -1141,9 +1188,20 @@ export function sanitizeDiarization(value: unknown): DiarizationSettings { ? r.threshold : d.threshold, threads: clampPositiveInt(r.threads, d.threads, 64), + // An unknown engine falls back to the default rather than disabling the lane: + // a typo in settings.json must not silently stop diarization, and the default + // is the one every existing sidecar already matches. + engine: DIARIZATION_ENGINE_IDS.includes(r.engine as DiarizationEngineId) + ? (r.engine as DiarizationEngineId) + : d.engine, + backend: DIARIZATION_BACKENDS.includes(r.backend as DiarizationBackend) + ? (r.backend as DiarizationBackend) + : d.backend, python: str(r.python, d.python), segModel: str(r.segModel, d.segModel), embModel: str(r.embModel, d.embModel), + sortformerBin: str(r.sortformerBin, d.sortformerBin), + sortformerModel: str(r.sortformerModel, d.sortformerModel), concurrency: clampPositiveInt(r.concurrency, d.concurrency, 16), // 0 is meaningful here (cap off), so this cannot use clampPositiveInt. // Fractional hours are allowed — the knob is a duration, not a count. diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,10 @@ # Changelog ## [Unreleased] +- **Speaker capture can now run on the graphics card, and stops inventing speakers that were never there.** The old engine works by grouping voices it thinks sound alike, and it splits far too eagerly: on a 13-minute reaction video with one host it found **13 speakers**, and on the worst video in the corpus it found **35**. The new engine decides speaker turns directly instead of grouping them afterwards, and returns **4** in both cases — agreeing with the old one about how much of the video the main speaker talks for (73.5% against 73.1%) while collapsing the invented tail. It has no threshold to tune, and it caps at 4 speakers, so a panel of five will merge two of them rather than split one into twenty. Turn it on with `"engine": "sortformer"` after running `scripts/build-sortformer.sh`; the previous engine stays the default and stays supported for machines with no usable GPU. +- **Speaker capture no longer has a length limit to worry about.** The old engine's memory grew with the square of how many turns a recording contains, which is why long streams had to be processed in windows and why a duration cap existed at all. The new engine holds a fixed amount of memory no matter how long the recording is — 558 MB for a 13-minute video or an 8-hour one — so nothing has to be windowed, capped, or decoded to a temporary file first. +- **Switching speaker-capture engines correctly marks the old results as work to redo.** The two engines disagree about how many speakers exist, so a corpus half-captured by each is not one corpus. Changing the engine now shows every recording captured by the other one as outstanding backfill work, rather than leaving it looking finished. Changing unrelated settings — a clustering threshold the new engine does not even have — no longer marks anything stale. +- **Speaker capture on the GPU stands aside for transcription instead of fighting it for video memory.** It needs about 4.4 GB of the 8 GB card that transcription also uses, so it now holds while transcription is working and resumes the moment the card is free — the same behaviour the digest sweep already had. The "run a guaranteed share alongside" setting is ignored while GPU capture is enabled, because a share of video memory is not a slower run, it is a failure. - **The digest backlog that could never start now has a button that clears it, and honest copy about why.** Nearly 2,000 videos across the corpus have a transcript but no compact transcript file beside it, and the digest generator refuses those — so they sat in a "waiting" count on the Digest card under text that said the problem would resolve on its own. It does not. Nothing writes that file automatically for a channel whose captions are *downloaded* rather than transcribed here, which is why one channel alone accounts for 1,683 of them and was showing about 6% digest coverage. The card now says what is actually wrong, says that nothing will fix it unattended, and offers **Normalize transcripts** right there to fix it for that channel — instead of the corpus-wide button on the Build page that walks all 79,000 video folders. The button appears only when there is something for it to do. - **Skipped-video explanations on the Backfill card now match the reason each one was skipped.** The card had one hardcoded sentence about the speaker-capture length limit and showed it for every kind of skipped video, including ones skipped for completely unrelated reasons. Each kind of work now supplies its own explanation and gets its own line. - **There is now one definition of "digested", and every screen reads it.** The digest layer kept its own private tally of what still needed doing, separate from the one the digest runner actually uses — and the two could disagree by an entire channel. They are now the same number. Two consequences you will see immediately. **Channels whose reports were months old start reporting digest work at all**: eleven of them predated the old tally entirely and had been quietly reading as "nothing to do" (they account for about 1,500 videos, most of them in one channel). And **videos with no usable transcript file are no longer offered as work**: they were being handed to the digest runner, which looked at them and immediately put them back. On the current corpus that is 1,989 videos, and 1,683 of them are in a single channel — so nearly all of that channel's apparent digest backlog was work that could never have started. The overall count goes *down* slightly as a result, from 75,613 to 75,199, which is the counter becoming honest rather than anything being skipped. diff --git a/scripts/build-sortformer.sh b/scripts/build-sortformer.sh @@ -0,0 +1,164 @@ +#!/usr/bin/env bash +# Build the Sortformer diarization engine into .diarize/sortformer/. +# +# WHAT THIS BUILDS, AND WHY SO LITTLE OF IT: openresearchtools/engine is a whole +# local-AI runtime (Rust CLI, in-process llama bridge, PDF/VLM modules, a cluster +# controller). We want exactly one thing from it — the native ggml implementation +# of NVIDIA's streaming Sortformer diarizer in `diarize/addons/overlay`. That code +# links against `ggml` ALONE: no llama, no whisper, no Rust. So this script stages +# just `ggml` + `tools/realtime` + our driver and builds those, which is why it +# needs no part of upstream's build system (whose scripts are PowerShell/CUDA and +# would otherwise be the blocker on Linux). +# +# The upstream commit is PINNED. This is a research repo with no releases, and an +# unpinned clone would silently change what produced every sidecar on disk. +# +# Usage: +# scripts/build-sortformer.sh # Vulkan (default) +# SORTFORMER_BACKEND=cpu scripts/build-sortformer.sh +# SORTFORMER_JOBS=2 scripts/build-sortformer.sh +# +# Requires: git, cmake, a C++17 compiler, and for the Vulkan build the Vulkan +# loader + headers and `glslc` (Arch: vulkan-headers shaderc). + +set -euo pipefail + +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +OUT_DIR="${SORTFORMER_DIR:-$REPO_ROOT/.diarize/sortformer}" +WORK_DIR="${SORTFORMER_BUILD_DIR:-$REPO_ROOT/.diarize/.sortformer-build}" + +# Pinned 2026-08-10. Bumping this changes model execution, so re-run the +# reproducibility check in .diarize/README.md before trusting a new pin. +ENGINE_REPO="https://github.com/openresearchtools/engine" +ENGINE_COMMIT="8bb4928c429b76754c7c7cfbdec0c27419bef684" + +MODEL_URL="https://huggingface.co/openresearchtools/diar_streaming_sortformer_4spk-v2.1-gguf/resolve/main/diar_streaming_sortformer_4spk-v2.1.gguf" +MODEL_NAME="diar_streaming_sortformer_4spk-v2.1.gguf" + +BACKEND="${SORTFORMER_BACKEND:-vulkan}" +JOBS="${SORTFORMER_JOBS:-$(nproc 2>/dev/null || echo 4)}" + +say() { printf '\033[1m==>\033[0m %s\n' "$*"; } + +for tool in git cmake; do + command -v "$tool" >/dev/null || { echo "missing required tool: $tool" >&2; exit 1; } +done + +case "$BACKEND" in + vulkan) + command -v glslc >/dev/null || { + echo "glslc not found — the Vulkan build compiles shaders with it." >&2 + echo "Install shaderc, or build CPU-only with SORTFORMER_BACKEND=cpu." >&2 + exit 1 + } + CMAKE_BACKEND_FLAGS=(-DGGML_VULKAN=ON) + ;; + cpu) + CMAKE_BACKEND_FLAGS=() + ;; + *) + echo "SORTFORMER_BACKEND must be 'vulkan' or 'cpu' (got '$BACKEND')" >&2 + exit 1 + ;; +esac + +mkdir -p "$WORK_DIR" + +# 1. Upstream source, at the pin. +SRC="$WORK_DIR/engine" +if [ -d "$SRC/.git" ] && [ "$(git -C "$SRC" rev-parse HEAD 2>/dev/null)" = "$ENGINE_COMMIT" ]; then + say "engine source already at $ENGINE_COMMIT" +else + say "fetching engine @ ${ENGINE_COMMIT:0:12}" + rm -rf "$SRC" + git init -q "$SRC" + git -C "$SRC" remote add origin "$ENGINE_REPO" + # Blobless partial clone: the full history is ~200 MB and we need one commit. + git -C "$SRC" fetch -q --depth 1 --filter=blob:none origin "$ENGINE_COMMIT" + git -C "$SRC" checkout -q FETCH_HEAD +fi + +# 2. Stage the minimal tree. Layout matters: tools/realtime/CMakeLists.txt refers +# to ../../ggml/include, so ggml must sit beside tools/, exactly as upstream. +# +# ggml is staged ONCE per pin and then left alone. Re-copying it would bump +# every source mtime and force a full rebuild of the Vulkan shaders (~10 min) +# every time only our driver changed. `cp -u` on the small realtime tree keeps +# edits cheap while still picking up a new pin. +STAGE="$WORK_DIR/stage" +mkdir -p "$STAGE/tools" +if [ -f "$STAGE/.ggml-commit" ] && [ "$(cat "$STAGE/.ggml-commit")" = "$ENGINE_COMMIT" ]; then + say "ggml already staged for this pin" +else + say "staging ggml" + rm -rf "$STAGE/ggml" "$STAGE/build" + cp -r "$SRC/third_party/llama.cpp/ggml" "$STAGE/ggml" + echo "$ENGINE_COMMIT" > "$STAGE/.ggml-commit" +fi +say "staging tools/realtime + driver" +mkdir -p "$STAGE/tools/realtime" +cp -ru "$SRC/diarize/addons/overlay/llama.cpp/tools/realtime/." "$STAGE/tools/realtime/" +cp -u "$REPO_ROOT/scripts/sortformer/diarize-file.cpp" "$STAGE/tools/realtime/diarize-file.cpp" + +cat > "$STAGE/CMakeLists.txt" <<'EOF' +cmake_minimum_required(VERSION 3.14) +project(sortformer-diar C CXX) +set(CMAKE_CXX_STANDARD 17) +set(CMAKE_POSITION_INDEPENDENT_CODE ON) +add_subdirectory(ggml) +add_subdirectory(tools/realtime) +EOF + +# Upstream's CMakeLists builds their parity tool; we append our driver rather than +# replacing it, so the staged tree stays a faithful copy plus one target. Guarded +# because the staged tree now persists across runs. +if ! grep -q 'add_executable(diarize-file' "$STAGE/tools/realtime/CMakeLists.txt"; then + cat >> "$STAGE/tools/realtime/CMakeLists.txt" <<'EOF' + +add_executable(diarize-file diarize-file.cpp) +target_compile_features(diarize-file PRIVATE cxx_std_17) +target_link_libraries(diarize-file PRIVATE ${TARGET_LIB}) +EOF +fi + +# 3. Build. +say "building ($BACKEND, -j$JOBS) — the Vulkan shader compile is the slow part" +cmake -S "$STAGE" -B "$STAGE/build" \ + -DCMAKE_BUILD_TYPE=Release \ + -DGGML_NATIVE=ON \ + "${CMAKE_BACKEND_FLAGS[@]}" >/dev/null +cmake --build "$STAGE/build" -j"$JOBS" --target diarize-file + +# 4. Install the binary beside the ggml shared objects it loads, and point its +# RPATH at $ORIGIN so it runs without LD_LIBRARY_PATH from any cwd. +say "installing to $OUT_DIR" +mkdir -p "$OUT_DIR" +cp "$STAGE/build/tools/realtime/diarize-file" "$OUT_DIR/diarize-file" +find "$STAGE/build" -name 'libggml*.so*' -exec cp -a {} "$OUT_DIR/" \; +if command -v patchelf >/dev/null; then + patchelf --set-rpath '$ORIGIN' "$OUT_DIR/diarize-file" +else + echo "note: patchelf not found; scripts/diarize-sortformer.mjs sets LD_LIBRARY_PATH itself" >&2 +fi + +# 5. Model. +if [ -f "$OUT_DIR/$MODEL_NAME" ]; then + say "model already present" +else + say "downloading model (~471 MB)" + curl -fSL --progress-bar -o "$OUT_DIR/$MODEL_NAME.part" "$MODEL_URL" + mv "$OUT_DIR/$MODEL_NAME.part" "$OUT_DIR/$MODEL_NAME" +fi + +say "done" +echo +echo " binary: $OUT_DIR/diarize-file" +echo " model: $OUT_DIR/$MODEL_NAME" +echo +echo "Wire it up in settings.json:" +echo " \"diarization\": {" +echo " \"engine\": \"sortformer\"," +echo " \"sortformerBin\": \"$OUT_DIR/diarize-file\"," +echo " \"sortformerModel\": \"$OUT_DIR/$MODEL_NAME\"," +echo " \"backend\": \"$BACKEND\"" +echo " }" diff --git a/scripts/diarize-sortformer.mjs b/scripts/diarize-sortformer.mjs @@ -0,0 +1,156 @@ +#!/usr/bin/env node +// Sortformer diarization engine for scripts/diarize.mjs. +// +// The exact counterpart of scripts/diarize-sherpa.py, and deliberately the same +// shape: decode audio, run the model, print raw turns as JSON on stdout. All +// policy (filenames, provenance shape, atomic writes) lives in the .mjs wrapper. +// +// stdout is JSON ONLY. Progress goes to stderr so the wrapper can stream it. +// +// WHAT THIS ADDS OVER THE BINARY: audio decoding. scripts/sortformer/ +// diarize-file.cpp only speaks 16 kHz mono PCM, and the corpus holds mp3, m4a, +// webm and whole video containers. ffmpeg does that here, exactly as the sherpa +// engine shells out to ffmpeg for the same reason. +// +// NO TEMPORARY WAV, ON PURPOSE. A decoded 8-hour VOD is ~900 MB, this box's disk +// runs at 98%, and /tmp is frequently a tmpfs — writing one would be the largest +// single risk in this lane. ffmpeg's raw output is piped straight into the +// engine, which consumes it in blocks, so neither process ever holds the file. +// +// Usage (invoked by diarize.mjs; runnable directly for debugging): +// diarize-sortformer.mjs --bin <diarize-file> --model <gguf> [options] <audio> +// +// --bin <path> engine binary [required] +// --model <path> sortformer .gguf [required] +// --backend <b> vulkan | cpu (default vulkan) +// --threads <n> CPU backend threads only (default 6 — see below) +// --ffmpeg <path> ffmpeg binary (FFMPEG_BIN, "ffmpeg") + +import { spawn } from "node:child_process"; +import path from "node:path"; + +const SAMPLE_RATE = 16000; // what the model's frontend expects; not a knob + +function fail(msg, code = 2) { + process.stderr.write(`diarize-sortformer: ${msg}\n`); + process.exit(code); +} + +const argv = process.argv.slice(2); +if (argv.includes("-h") || argv.includes("--help")) { + process.stdout.write( + "diarize-sortformer.mjs --bin <diarize-file> --model <gguf> [options] <audio>\n", + ); + process.exit(0); +} + +function arg(flag, fallback) { + const i = argv.indexOf(flag); + return i < 0 || i === argv.length - 1 ? fallback : argv[i + 1]; +} + +const positional = []; +for (let i = 0; i < argv.length; i++) { + if (argv[i].startsWith("-")) { + i++; // every flag here takes a value + continue; + } + positional.push(argv[i]); +} + +const audio = positional[0]; +if (!audio) fail("missing <audioFile>"); + +const bin = arg("--bin", process.env.SORTFORMER_BIN ?? ""); +const model = arg("--model", process.env.SORTFORMER_MODEL ?? ""); +const backend = String(arg("--backend", "vulkan")).toLowerCase(); +// 6, not 4 and not 8: measured on this box at 1305 s/audio-hour against 1514 for +// ggml's default of 4, and 1401 for 8 — which is WORSE than 4 through +// oversubscription on a 4c/8t part. Ignored by the Vulkan backend. +const threads = Number(arg("--threads", "6")); +const ffmpeg = arg("--ffmpeg", process.env.FFMPEG_BIN ?? "ffmpeg"); + +if (!bin) fail("missing --bin (or SORTFORMER_BIN)"); +if (!model) fail("missing --model (or SORTFORMER_MODEL)"); + +// Settings speak "vulkan"/"cpu"; ggml speaks device names. One translation, here, +// so nothing upstream of this file has to know ggml's spelling. +const ggmlBackend = + backend === "cpu" ? "CPU" : backend === "vulkan" ? "Vulkan0" : null; +if (!ggmlBackend) fail(`--backend must be 'vulkan' or 'cpu' (got '${backend}')`); + +// ffmpeg: anything in -> raw s16le mono 16k on stdout. `-vn` because a source +// container carries a video stream that would otherwise be negotiated first. +const dec = spawn( + ffmpeg, + [ + "-nostdin", + "-v", "error", + "-i", audio, + "-vn", + "-ac", "1", + "-ar", String(SAMPLE_RATE), + "-f", "s16le", + "pipe:1", + ], + { stdio: ["ignore", "pipe", "inherit"] }, +); + +const engine = spawn( + bin, + [ + "--model", model, + "--audio-raw", "-", + "--sample-rate", String(SAMPLE_RATE), + "--backend", ggmlBackend, + "--threads", String(threads), + ], + { + stdio: ["pipe", "pipe", "inherit"], + env: { + ...process.env, + // The build sets RPATH=$ORIGIN when patchelf is available; this covers the + // case where it was not, so the binary finds its ggml .so files either way. + LD_LIBRARY_PATH: [path.dirname(path.resolve(bin)), process.env.LD_LIBRARY_PATH] + .filter(Boolean) + .join(":"), + }, + }, +); + +dec.on("error", (err) => fail(`failed to run ${ffmpeg}: ${err.message}`, 3)); +engine.on("error", (err) => fail(`failed to run ${bin}: ${err.message}`, 3)); + +dec.stdout.pipe(engine.stdin); +// The engine exiting early (bad model, unusable backend) closes the pipe under +// ffmpeg. That is an expected race, not a crash to report. +engine.stdin.on("error", () => dec.kill("SIGTERM")); + +let stdout = ""; +engine.stdout.setEncoding("utf8"); +engine.stdout.on("data", (d) => { + stdout += d; +}); + +let decCode = null; +let engineCode = null; + +function finish() { + if (decCode === null || engineCode === null) return; + // The engine's verdict comes first: it is the one that can fail for a reason + // worth reporting. A non-zero ffmpeg with a healthy engine still means the + // audio was truncated, so neither is allowed to pass silently. + if (engineCode !== 0) fail(`engine exited ${engineCode}`, engineCode ?? 1); + if (decCode !== 0) fail(`ffmpeg exited ${decCode} — audio may be truncated`, decCode ?? 1); + if (!stdout.trim()) fail("engine produced no output"); + process.stdout.write(stdout.endsWith("\n") ? stdout : stdout + "\n"); +} + +dec.on("close", (code) => { + decCode = code ?? 0; + finish(); +}); +engine.on("close", (code) => { + engineCode = code ?? 0; + finish(); +}); diff --git a/scripts/diarize.mjs b/scripts/diarize.mjs @@ -8,11 +8,24 @@ // scripts/diarize-sherpa.py (ONNX, CPU); pyannote would be a drop-in // replacement for the --engine command. // -// WHY CPU AND NOT THE GPU: parakeet.cpp is built GGML_VULKAN=ON so ASR really -// does run on the RX 6600 XT, but no diarization model has been ported to ggml. -// They ship as PyTorch (CUDA/ROCm) or ONNX (CUDA/ROCm/MIGraphX/OpenVINO/ -// DirectML), and neither runtime has a Vulkan compute path on Linux. That is an -// ecosystem gap, not a hardware limit. +// TWO ENGINES, PICKED BY --engine-kind: +// sherpa-onnx (default) scripts/diarize-sherpa.py — ONNX segmentation + +// embeddings, agglomerative clustering, CPU only. +// sortformer scripts/diarize-sortformer.mjs — NVIDIA's streaming +// Sortformer as ggml, end-to-end, Vulkan or CPU. +// +// THE GPU NOTE THIS FILE USED TO CARRY IS NOW OUT OF DATE, and the correction is +// worth keeping: it said no diarization model had been ported to ggml, so the +// RX 6600 XT could not run one while parakeet.cpp happily did ASR on it. That was +// true of PyTorch and ONNX runtimes, and it is no longer true — openresearchtools +// ported Sortformer to ggml, which inherits the same Vulkan backend parakeet +// uses. Measured on this box: 894 s/audio-hour on Vulkan against 1305 for the +// same engine on tuned CPU threads, on ONE core instead of 2.2. +// +// It is still not the fastest lane: sherpa-onnx does the same file in +// 492 s/audio-hour across ~4 cores. Sortformer is chosen for QUALITY (4 speakers +// where sherpa splits into 13, and 4 where sherpa splits into 35), and the GPU is +// chosen because it is the cheapest way to run it. // // Dual-use, by design: // * In the app: controller/diarizeOne.ts runs this with cwd == the video dir @@ -29,7 +42,11 @@ // Options (env fallback in parens): // --output <f> write the record here (default: stdout) // --video-id <id> recorded in the sidecar (default: basename of cwd) +// --engine-kind <k> sherpa-onnx | sortformer (DIARIZE_ENGINE_KIND, sherpa-onnx) // --engine <cmd> engine command (DIARIZE_ENGINE_CMD) +// --sortformer-bin <p> engine binary (SORTFORMER_BIN) [sortformer] +// --sortformer-model <p> .gguf model (SORTFORMER_MODEL) [sortformer] +// --backend <b> vulkan | cpu (default vulkan) [sortformer] // --python <path> python for the default engine (DIARIZE_PYTHON, "python3") // --seg <path> segmentation model (DIARIZE_SEG_MODEL) [required] // --emb <path> speaker-embedding model (DIARIZE_EMB_MODEL) [required] @@ -110,6 +127,24 @@ const windowAfterMinutes = arg( ); const engineCmd = arg("--engine", process.env.DIARIZE_ENGINE_CMD ?? ""); +const engineKind = arg( + "--engine-kind", + process.env.DIARIZE_ENGINE_KIND ?? "sherpa-onnx", +); +const sortformerBin = arg("--sortformer-bin", process.env.SORTFORMER_BIN ?? ""); +const sortformerModel = arg( + "--sortformer-model", + process.env.SORTFORMER_MODEL ?? "", +); +const backend = arg("--backend", "vulkan"); + +if (!engineCmd && engineKind !== "sherpa-onnx" && engineKind !== "sortformer") { + fail(`--engine-kind must be 'sherpa-onnx' or 'sortformer' (got '${engineKind}')`); +} + +// `--engine` still wins over `--engine-kind`. It is the raw escape hatch (and +// what the e2e fake engine uses), so a kind must never quietly override it. +const usingSortformer = !engineCmd && engineKind === "sortformer"; let cmd; let args; @@ -117,6 +152,21 @@ if (engineCmd) { const parts = engineCmd.split(" ").filter(Boolean); cmd = parts[0]; args = [...parts.slice(1), audio]; +} else if (usingSortformer) { + if (!sortformerBin) fail("missing --sortformer-bin (or SORTFORMER_BIN)"); + if (!sortformerModel) fail("missing --sortformer-model (or SORTFORMER_MODEL)"); + // process.execPath, not a bare "node": the app may run under a node that is not + // the one on PATH, and the engine has to be the same runtime as its wrapper. + cmd = process.execPath; + args = [ + path.join(import.meta.dirname, "diarize-sortformer.mjs"), + "--bin", sortformerBin, + "--model", sortformerModel, + "--backend", backend, + "--threads", String(threads), + "--ffmpeg", ffmpeg, + audio, + ]; } else { if (!seg) fail("missing --seg (or DIARIZE_SEG_MODEL)"); if (!emb) fail("missing --emb (or DIARIZE_EMB_MODEL)"); @@ -177,15 +227,30 @@ child.on("close", async (code) => { : {}), speakers: new Set(turns.map((t) => t.speaker)).size, turns, - engine: { - engine: engineCmd ? path.basename(cmd) : "sherpa-onnx", - // basename only: absolute paths are machine-specific and the corpus is - // rsynced between shards, so a full path would make records non-portable. - ...(seg ? { segmentationModel: path.basename(seg) } : {}), - ...(emb ? { embeddingModel: path.basename(emb) } : {}), - ...(raw.version ? { version: String(raw.version) } : {}), - threshold, - }, + // Provenance is per-engine, and the fields a given engine CANNOT have are + // left off rather than written empty. Recording sherpa's threshold on a + // sortformer record would be a lie that isDiarizationFresh then acts on: it + // has no clustering step and no threshold, so a threshold edit must not mark + // its sidecars stale. + engine: usingSortformer + ? { + engine: "sortformer", + model: path.basename(sortformerModel), + // The driver reports its ggml backend here ("sortformer/Vulkan0"), so + // the record says which device produced it. Excluded from the freshness + // identity along with every other `version` — CPU and Vulkan produce + // byte-identical turns, so re-running one on the other buys nothing. + ...(raw.version ? { version: String(raw.version) } : {}), + } + : { + engine: engineCmd ? path.basename(cmd) : "sherpa-onnx", + // basename only: absolute paths are machine-specific and the corpus is + // rsynced between shards, so a full path would make records non-portable. + ...(seg ? { segmentationModel: path.basename(seg) } : {}), + ...(emb ? { embeddingModel: path.basename(emb) } : {}), + ...(raw.version ? { version: String(raw.version) } : {}), + threshold, + }, }; const json = JSON.stringify(record) + "\n"; diff --git a/scripts/sortformer/diarize-file.cpp b/scripts/sortformer/diarize-file.cpp @@ -0,0 +1,268 @@ +// Sortformer diarization engine: 16 kHz mono WAV in, speaker turns out. +// +// Kept deliberately dumb, exactly like scripts/diarize-sherpa.py: read audio, +// run the model, print raw turns as JSON on stdout. All policy (filenames, +// provenance shape, atomic writes) lives in scripts/diarize.mjs, and audio +// decoding lives in scripts/diarize-sortformer.mjs, so this file is only ever +// "run the model". +// +// stdout is JSON ONLY. Anything human goes to stderr. +// +// WHY THIS FILE EXISTS RATHER THAN AN UPSTREAM BINARY: openresearchtools/engine +// ships `llama-realtime-smoke`, a PARITY tool. It constructs the backend with +// capture_debug=true (retaining every intermediate matrix for every step), it +// requires a directory of PyTorch reference fixtures we do not have, and its +// JSON dump drops the event flags — which is fatal, because 2033 of the 2099 +// events it emits for a 12.8-minute file are PREVIEW re-emissions of spans that +// are still growing. Only the 66 non-preview events are the answer. This driver +// is that tool with the debug capture off, the fixtures gone, and the preview +// events filtered. +// +// WHY NO MERGE PASS: the non-preview events are already the postprocessed, +// disjoint spans — running an interval union over them is a no-op (66 in, 66 +// out, verified). Unlike the sherpa engine there is no clustering step and no +// threshold: Sortformer is end-to-end, which is the whole reason it does not +// over-split. +// +// WHY RAW PCM ON STDIN IS THE PRIMARY INPUT: this corpus has 8-hour VODs. A +// decoded 16 kHz mono WAV of one is ~900 MB on a disk that is 98% full, and +// buffering it as float32 is ~460 MB resident. Streaming ffmpeg's output straight +// into push_audio() costs neither: the model's own state is a fixed-size speaker +// cache, so memory becomes O(1) in duration rather than O(n). That is the single +// biggest advantage over the sherpa engine, whose O(n^2) clustering matrix is +// what forced the whole windowing/centroid-reclustering design in +// diarize-sherpa.py. --audio-wav is kept for direct CLI use and testing. + +#include "backend-factory.h" +#include "sortformer/sortformer-backend.h" +#include "stream-manager.h" + +#include <algorithm> +#include <chrono> +#include <cstdint> +#include <cstdio> +#include <cstring> +#include <fstream> +#include <iomanip> +#include <iostream> +#include <memory> +#include <sstream> +#include <stdexcept> +#include <string> +#include <vector> + +// POSIX only, which is all scripts/build-sortformer.sh targets. The upstream +// engine builds on Windows too; this driver deliberately does not try to. +#include <cerrno> +#include <fcntl.h> +#include <unistd.h> + +namespace { + +// 16-bit PCM mono, which is what `ffmpeg -ac 1 -ar 16000 -c:a pcm_s16le` makes. +// Deliberately not a general WAV reader: the adapter always hands us that exact +// format, and quietly accepting anything else would mean silently diarizing +// resampled-wrong audio. +std::vector<float> load_wav_s16_mono(const std::string & path, uint32_t & sample_rate) { + std::ifstream f(path, std::ios::binary); + if (!f) throw std::runtime_error("cannot open wav: " + path); + + char riff[12]; + f.read(riff, 12); + if (std::strncmp(riff, "RIFF", 4) != 0 || std::strncmp(riff + 8, "WAVE", 4) != 0) + throw std::runtime_error("not a RIFF/WAVE file: " + path); + + uint16_t channels = 0, bits = 0; + while (f) { + char id[4]; + uint32_t sz = 0; + f.read(id, 4); + f.read(reinterpret_cast<char *>(&sz), 4); + if (!f) break; + + if (std::strncmp(id, "fmt ", 4) == 0) { + std::vector<char> fmt(sz); + f.read(fmt.data(), sz); + if (sz >= 16) { + std::memcpy(&channels, fmt.data() + 2, 2); + std::memcpy(&sample_rate, fmt.data() + 4, 4); + std::memcpy(&bits, fmt.data() + 14, 2); + } + } else if (std::strncmp(id, "data", 4) == 0) { + if (channels != 1 || bits != 16) + throw std::runtime_error("expected 16-bit mono wav, got " + + std::to_string(channels) + "ch/" + + std::to_string(bits) + "bit"); + const size_t n = sz / 2; + std::vector<int16_t> pcm(n); + f.read(reinterpret_cast<char *>(pcm.data()), sz); + std::vector<float> out(n); + for (size_t i = 0; i < n; ++i) out[i] = static_cast<float>(pcm[i]) / 32768.0f; + return out; + } else { + f.seekg(sz + (sz & 1), std::ios::cur); // chunks are word-aligned + } + } + throw std::runtime_error("no data chunk in wav: " + path); +} + +// The sortformer path never sets a thread count (only voxtral's runtime does), +// so the CPU backend runs at ggml's default of 4. On this box 6 threads is the +// optimum (1305 s/audio-hour) and 8 is WORSE than 4 through oversubscription, so +// this has to be reachable rather than left to the default. +void set_backend_n_threads(ggml_backend_t backend, int n_threads) { + if (backend == nullptr) return; + ggml_backend_dev_t dev = ggml_backend_get_device(backend); + ggml_backend_reg_t reg = dev ? ggml_backend_dev_backend_reg(dev) : nullptr; + if (reg == nullptr) return; + auto * fn = (ggml_backend_set_n_threads_t) + ggml_backend_reg_get_proc_address(reg, "ggml_backend_set_n_threads"); + if (fn != nullptr) fn(backend, n_threads); +} + +// Stream s16le mono PCM from stdin straight into the session. Never holds more +// than one block, so an 8-hour VOD costs the same resident memory as a 6-minute +// clip. Returns the number of samples fed. +size_t feed_raw_stdin(llama::realtime::stream_manager & mgr, int64_t sid, + uint32_t sample_rate, size_t block_samples) { + std::vector<int16_t> pcm(block_samples); + std::vector<float> block(block_samples); + size_t total = 0; + + // FORCE THE PIPE BACK TO BLOCKING. Node sets the pipes it hands a child to + // non-blocking, so a plain read() returns -1/EAGAIN long before EOF and the + // stream looks like an I/O error. It works from a shell and fails under the + // app, which is exactly the sort of difference that gets found in production + // rather than in a test — so it is fixed here, in the binary, rather than + // being made the caller's problem. + const int flags = fcntl(STDIN_FILENO, F_GETFL, 0); + if (flags != -1 && (flags & O_NONBLOCK)) + fcntl(STDIN_FILENO, F_SETFL, flags & ~O_NONBLOCK); + + // read(2) rather than fread: it lets EINTR be retried without the ambiguity + // of a short fread, and a partial read is normal on a pipe rather than an + // error to distinguish from EOF. + while (true) { + size_t filled = 0; + const size_t want = block_samples * sizeof(int16_t); + auto * buf = reinterpret_cast<char *>(pcm.data()); + while (filled < want) { + const ssize_t n = ::read(STDIN_FILENO, buf + filled, want - filled); + if (n == 0) break; // EOF + if (n < 0) { + if (errno == EINTR) continue; + throw std::runtime_error(std::string("read error on stdin: ") + + std::strerror(errno)); + } + filled += static_cast<size_t>(n); + } + const size_t got = filled / sizeof(int16_t); + if (got == 0) break; + for (size_t i = 0; i < got; ++i) block[i] = static_cast<float>(pcm[i]) / 32768.0f; + mgr.push_audio(sid, block.data(), got, sample_rate); + total += got; + if (filled < want) break; // short read means EOF + } + return total; +} + +void usage(const char * argv0) { + std::cerr + << "usage: " << argv0 << " --model <sortformer.gguf> (--audio-raw - | --audio-wav <f>)\n" + << " [--sample-rate 16000] [--backend Vulkan0|CPU]\n" + << " [--threads N] [--feed-ms N]\n\n" + << "--audio-raw - read s16le mono PCM from stdin (streaming, O(1) memory)\n" + << "--audio-wav <f> read a 16 kHz mono 16-bit WAV file\n\n" + << "Prints {\"turns\":[{start,end,speaker}],\"audioSeconds\":N,\"version\":\"...\"}\n" + << "on stdout. Progress and errors go to stderr.\n"; +} + +} // namespace + +int main(int argc, char ** argv) { + try { + std::string gguf, wav, raw, backend = "Vulkan0"; + uint32_t sr = 16000; + int n_threads = 0; + double feed_ms = 100.0; + + for (int i = 1; i < argc; ++i) { + const std::string a = argv[i]; + if (a == "--help" || a == "-h") { usage(argv[0]); return 0; } + if (i + 1 >= argc) throw std::invalid_argument("missing value for " + a); + if (a == "--model") gguf = argv[++i]; + else if (a == "--audio-wav") wav = argv[++i]; + else if (a == "--audio-raw") raw = argv[++i]; + else if (a == "--sample-rate") sr = static_cast<uint32_t>(std::stoul(argv[++i])); + else if (a == "--backend") backend = argv[++i]; + else if (a == "--threads") n_threads = std::stoi(argv[++i]); + else if (a == "--feed-ms") feed_ms = std::stod(argv[++i]); + else throw std::invalid_argument("unknown argument: " + a); + } + if (gguf.empty() || (wav.empty() == raw.empty())) { usage(argv[0]); return 2; } + if (!raw.empty() && raw != "-") + throw std::invalid_argument("--audio-raw only supports '-' (stdin)"); + + // Load the model BEFORE reading stdin, so a bad model path fails before + // ffmpeg has decoded anything. capture_debug=false — see the header. + auto be = std::make_unique<llama::realtime::sortformer_stream_backend>(gguf, backend, false); + auto * be_ptr = be.get(); + if (n_threads > 0) set_backend_n_threads(be_ptr->model().backend(), n_threads); + + llama::realtime::stream_manager mgr; + const int64_t sid = mgr.create_session(std::move(be)); + std::cerr << "sortformer: backend=" << be_ptr->backend_name() + << " input=" << (raw.empty() ? wav : "stdin") << "\n"; + + const auto t0 = std::chrono::steady_clock::now(); + size_t n_samples = 0; + if (!raw.empty()) { + if (sr == 0) throw std::runtime_error("--sample-rate must be non-zero"); + n_samples = feed_raw_stdin( + mgr, sid, sr, + std::max<size_t>(1, static_cast<size_t>((feed_ms / 1000.0) * sr))); + } else { + const auto audio = load_wav_s16_mono(wav, sr); + if (sr == 0) throw std::runtime_error("wav reports a zero sample rate"); + const size_t feed = + std::max<size_t>(1, static_cast<size_t>((feed_ms / 1000.0) * sr)); + for (size_t off = 0; off < audio.size(); off += feed) { + const size_t n = std::min(feed, audio.size() - off); + mgr.push_audio(sid, audio.data() + static_cast<ptrdiff_t>(off), n, sr); + } + n_samples = audio.size(); + } + mgr.flush_session(sid); + const double infer_sec = + std::chrono::duration<double>(std::chrono::steady_clock::now() - t0).count(); + const double audio_sec = static_cast<double>(n_samples) / static_cast<double>(sr); + + // PREVIEW EVENTS ARE NOT THE ANSWER. A streaming span is re-emitted every + // time it grows; only the final, postprocessed commits carry the result. + const auto events = mgr.drain_events(sid, 0); + std::ostringstream turns; + size_t n_turns = 0; + for (const auto & e : events) { + if (e.type != llama::realtime::event_type::speaker_span_commit) continue; + if (e.flags & llama::realtime::event_flag_preview) continue; + if (e.end_sec <= e.begin_sec) continue; + if (n_turns++) turns << ","; + turns << "{\"start\":" << e.begin_sec + << ",\"end\":" << e.end_sec + << ",\"speaker\":" << e.speaker_id << "}"; + } + + std::cerr << "sortformer: " << n_turns << " turns in " + << std::fixed << std::setprecision(1) << infer_sec << "s (" + << std::setprecision(2) << (audio_sec / infer_sec) << "x realtime)\n"; + + std::cout << std::setprecision(6) + << "{\"turns\":[" << turns.str() << "]" + << ",\"audioSeconds\":" << audio_sec + << ",\"version\":\"" << be_ptr->backend_name() << "\"}\n"; + return 0; + } catch (const std::exception & e) { + std::cerr << "sortformer: error: " << e.what() << "\n"; + return 1; + } +}