commit 653bd65287646770448172696fff68aaac176c9c
parent c10085f2ce6453c045461e2d634db341e9125552
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Fri, 7 Aug 2026 23:27:40 -0400
Merge feat/attribution-backfill: put names to the speakers, and make the backfill registry earn its name
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
19 files changed, 3140 insertions(+), 38 deletions(-)
diff --git a/PLAN.md b/PLAN.md
@@ -404,24 +404,48 @@ posture is set here.
### Phase 9 — Attribution (two lanes) + quote filtering
-> **The CAPTURE half landed 2026-08-07** — `diarization.json` per video, off by default,
-> plus the cleanup guard that stops the Clean-audio sweep deleting audio for a
-> transcribed-but-undiarized video. That is deliberately all of it: capture is the only
-> part with a deadline, because audio is the only input that cannot be regenerated.
-> Everything below — `attribution.json`, the text-only lane, `attributionStatus.ts`,
-> snapshot buckets, badges, upgrade jobs, quote filtering — is still to do and can be
-> built from the saved sidecars at any time. Measurements that changed the design (and
-> corrected two assumptions in the text below) are in
+> **CAPTURE AND BOTH ATTRIBUTION LANES LANDED 2026-08-07. Nothing is armed.**
+> `diarization.json` per video with the cleanup guard, then `attribution.json` with both
+> lanes, `attributionStatus.ts`, and the freshness comparators — all off by default, under
+> an off-by-default master switch. What is still to do: `attributionMs` through the build,
+> viewer badges, per-channel export counts, and quote filtering (which additionally
+> depends on Phase 8). Measurements that changed the design are in
> [`plans/diarization-spike-results.md`](plans/diarization-spike-results.md); the
> decisions are in `plans/STATE.md`.
>
-> Two corrections to what this section says: **audio deletion is not automatic** —
-> `cleanAudioFromTranscribed` has one call site, an operator-triggered action, so the
-> deadline is disk pressure, not the transcribe job. And **diarization does not run
-> "right after transcription" by default**: it is 2–3× slower than the transcription it
-> would follow (221 s/audio-hour on the GPU vs ~500–680 on the CPU), so the inline hook
-> ships off and the intended shape for a batch is capture-on, inline-off, backfill after.
-
+> **Four corrections to what this section says.**
+>
+> **1. The sequencing below is inverted.** "Run the text-only pass over the whole backlog
+> immediately, then upgrade selectively" was written before the numbers existed. 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 while the digest sweep has completed **0.17%** of its own (122
+> `ai-digest.json` of 73,367). Both want the same 8 GB card and nothing arbitrates between
+> them. The diarized lane costs about **one call per video** and is better, and its only
+> problem is input coverage — which the backfill lane now fixes. So: both built, neither
+> armed, and a measured pilot picks the default.
+>
+> **2. The upgrade job does not need building.** This section describes a bespoke job that
+> "re-downloads audio, diarizes, attributes, then removes the re-fetched audio in a
+> `finally`". That job now exists generically: `attribution-diarized` reports a text-only
+> record as `missing` work, so the lane, the sweep and the indicators queue the upgrade
+> with no new machinery, and the re-acquire-and-delete-in-a-`finally` half is
+> [`common/controller/backfillReacquire.ts`](common/controller/backfillReacquire.ts).
+>
+> **3. Audio deletion is not automatic** — `cleanAudioFromTranscribed` has one call site,
+> an operator-triggered action, so the deadline is disk pressure, not the transcribe job.
+>
+> **4. Diarization does not run "right after transcription" by default**: it is 2–3×
+> slower than the transcription it would follow (221 s/audio-hour on the GPU vs ~500–680 on
+> the CPU), so the inline hook ships off and the intended shape for a batch is capture-on,
+> inline-off, backfill after.
+>
+> **And one thing NOT to re-propose: attribution cannot be a digest section.** Folding
+> `speakers` into the digest prompt looks like it halves the bill. Digests are generated
+> **chunk-local**, 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.
+>
Largest and least certain. Ship 1–8 first; keep it off by default.
**Blocking constraint:** audio is deleted once a video is transcribed, except for dirs
@@ -435,16 +459,18 @@ Two first-class lanes producing the same artifact at different quality tiers:
| Diarization-assisted | Going forward; on-demand upgrade | audio → speaker turns → LLM labels the turns | `method: "diarized"` |
| Text-only | Legacy videos with no audio | cues alone → LLM segments and labels from content | `method: "text-only"` |
-Neither is a fallback for the other in code. Going forward, run diarization right after
-transcription while audio is still on disk and **before** the cleanup sweep. For legacy, run
-the text-only pass over the whole backlog immediately, then upgrade selectively.
+Neither is a fallback for the other in code. Going forward, capture diarization while audio
+is still on disk and **before** the cleanup sweep. (The original text said to run the
+text-only pass over the whole backlog immediately and upgrade selectively — see correction
+1 above: that is 25–55 GPU-days spent on the worse lane, and the sequencing is inverted.)
- `scripts/diarize.mjs` behind a `DIARIZE_BIN` path, exactly as the parakeet wrapper works.
Keeps the engine swappable: evaluate `sherpa-onnx` (CPU-friendly, no HF token) against
`pyannote` (better, needs token + GPU) **on real audio before committing** — quality here
cannot be judged from code.
-- `attributeCues.ts` writes `attribution.json` with ranges, confidence, method, and engine
- provenance.
+- ~~`attributeCues.ts`~~ **`common/controller/attributeOne.ts`** writes `attribution.json`
+ with ranges, confidence, method, and engine provenance. Both lanes live in it because
+ they share every guard around the artifact; only the middle differs.
- Filter in the page writer alongside Phase 8. **Bias toward dropping on uncertainty** — the
cost is asymmetric. Make the threshold policy-driven so a channel can filter on `diarized`
only and ignore `text-only` labels it doesn't trust.
@@ -452,10 +478,13 @@ the text-only pass over the whole backlog immediately, then upgrade selectively.
(`"none" | "text-only" | "diarized" | "stale"`), new snapshot buckets, a per-video badge,
per-channel counts. Mark a video **ineligible** when availability says deleted/private and
no audio is retained — a distinct state, not an upgrade button that can only fail.
-- Re-processing: per-video re-attribute / upgrade, per-channel bulk over the text-only
+- ~~Re-processing: per-video re-attribute / upgrade, per-channel bulk over the text-only
bucket. The upgrade job re-downloads audio, diarizes, attributes, then removes the
- re-fetched audio **in a `finally`** unless `do-not-clean.json` is present — this job can
- pull gigabytes and a cancelled run must not silently fill the disk.
+ re-fetched audio **in a `finally`**…~~ **Done generically — see correction 2 above.**
+ `attribution-diarized` reports a text-only record as `missing`, so the backfill lane, the
+ corpus sweep and the per-channel button all queue the upgrade already; the
+ fetch-use-delete-in-a-`finally` half is `controller/backfillReacquire.ts`, which honours
+ `do-not-clean.json` and a free-disk floor and has e2e tests for the leak paths.
- `attributionMs` threads through the build the same way `digestMs` does.
- **Set expectations honestly:** this misfires on rapid back-and-forth, and auto-caption
channels have no speaker turns at all with cue boundaries that don't align to them. A
diff --git a/common/controller/attributeOne.ts b/common/controller/attributeOne.ts
@@ -0,0 +1,634 @@
+// Attribute ONE video — put names to the speakers — and write attribution.json
+// beside the transcript. Both lanes live here because they produce the same
+// artifact and share every guard around it; only the middle (how the speakers
+// are found) differs.
+//
+// Pipeline, reusing the declared single source of truth for each step:
+// isCuesJsonFresh — don't attribute a transcript about to be rewritten
+// readNormalizedTranscript
+// loadAttribution — the DOWNGRADE guard, and the freshness skip
+// diarized: loadDiarization -> selectClusterSamples -> ONE call
+// text-only: chunkCuesForContext -> one call per chunk, roster carried forward
+// digest app — ollama (local) or the metered claude lane
+// writeAttribution — tmp + rename
+//
+// TWO KINDS WRITE THIS ONE FILE, so the ordering rule is enforced HERE as well
+// as in the registry's state(). state() is a read of disk taken at pull time; the
+// pool may hold a candidate for minutes afterwards, and in that gap the diarized
+// lane can legitimately have written a better record. A text-only run that
+// checked only what state() saw would then overwrite it. isAttributionDowngrade
+// is checked immediately before the write for exactly that window.
+//
+// PER-CHUNK FAILURE ISOLATION, copied from digestVideo.ts: one 500 from ollama
+// costs that chunk, not the 30-call video. chunks/chunksOk go into provenance so
+// a partial record is legible as partial afterwards.
+
+import type { Paths } from "../lib/paths";
+import { getSettings, type AttributionSettings } from "../lib/settings";
+import type { DigestApp, DigestAppConfig } from "../lib/digestApps";
+import {
+ ATTRIBUTION_PROMPT_VERSION,
+ isAttributionDowngrade,
+ isAttributionFresh,
+ type AttributionFreshnessTarget,
+ type AttributionMethod,
+ type AttributionRecord,
+ type AttributionSegment,
+ type AttributionSpeaker,
+ type AttributionWarning,
+} from "../lib/attribution";
+import { resolveAttributionTarget } from "./attributionTarget";
+import { loadAttribution, writeAttribution } from "../lib/attribution-server";
+import { loadDiarization } from "../lib/diarization-server";
+import type { DiarizationRecord } from "../lib/diarization";
+import {
+ SPEAKER_SYSTEM_PROMPT,
+ TURN_SYSTEM_PROMPT,
+ buildSpeakerPrompt,
+ buildTurnPrompt,
+ isUselessSpeakerLabel,
+ selectClusterSamples,
+ speakerSchema,
+ turnSchema,
+} from "../lib/attributionPrompt";
+import {
+ DIGEST_OVERLAP_CUES,
+ hmsToSeconds,
+ maxCuesForContext,
+ toHms,
+} from "../lib/digestPrompt";
+import { chunkCuesForContext } from "../lib/transcriptWindow";
+import { transcriptToMarkdown } from "../lib/transcriptToMarkdown";
+import { readDigestContext } from "../lib/digestContext-server";
+import type { Cue } from "../lib/vtt";
+import {
+ isCuesJsonFresh,
+ readNormalizedTranscript,
+} from "./normalizeTranscript";
+
+export type AttributeOneOptions = {
+ paths: Paths;
+ videoDir: string;
+ videoId: string;
+ channelSlug: string;
+ method: AttributionMethod;
+ // Overrides settings.attribution. The backfill passes one resolved copy rather
+ // than re-reading settings per video.
+ settings?: AttributionSettings;
+ // Redo even when the recorded identity matches. Never overrides the DOWNGRADE
+ // rule — forcing a regeneration is not the same as asking for a worse record,
+ // and nothing in the UI should be able to request the second by accident.
+ force?: boolean;
+ onLog?: (msg: string) => void;
+ signal?: AbortSignal;
+};
+
+export type AttributeOneOutcome =
+ // Wrote (or rewrote) attribution.json.
+ | "attributed"
+ // A record at the current identity already existed and force was not set.
+ | "already-exists"
+ // The text-only lane found a DIARIZED record. Distinct from already-exists:
+ // nothing is stale, there is simply a better record here and this lane must
+ // not touch it.
+ | "outranked"
+ // The lane, or the feature, is switched off.
+ | "disabled"
+ // Diarized lane only: no diarization.json to name clusters from.
+ | "no-diarization"
+ // No usable transcript — absent, or about to be rewritten (stale cues).
+ | "no-transcript"
+ // The engine produced nothing usable. Logged, never thrown.
+ | "failed";
+
+export async function attributeOneVideo(
+ opts: AttributeOneOptions,
+): Promise<AttributeOneOutcome> {
+ const log = opts.onLog ?? ((m: string) => console.log(m));
+ const cfg = opts.settings ?? getSettings().attribution;
+ const method = opts.method;
+
+ if (!cfg.enabled) return "disabled";
+ if (method === "diarized" ? !cfg.diarizedEnabled : !cfg.textOnlyEnabled) {
+ return "disabled";
+ }
+
+ // Don't attribute a transcript that is about to be rewritten: a stale
+ // cues.json means the raw transcript changed under it, so the attribution
+ // would describe superseded text and then look fresh forever. Same guard, same
+ // reason, as digestVideo's.
+ const { fresh: cuesFresh, cuesPath } = await isCuesJsonFresh(opts.videoDir);
+ if (!cuesFresh) return "no-transcript";
+ const transcript = await readNormalizedTranscript(cuesPath);
+ const cues = transcript?.cues ?? [];
+ if (!transcript || cues.length === 0) return "no-transcript";
+
+ const diarization =
+ method === "diarized" ? await loadDiarization(opts.videoDir) : null;
+ if (method === "diarized" && !diarization) return "no-diarization";
+
+ const existing = await loadAttribution(opts.videoDir);
+ // THE DOWNGRADE GUARD, and it comes FIRST — before the engine is resolved,
+ // before settings are consulted for a model, before anything that could fail
+ // for an unrelated reason. A diarized record is not "stale" for the text lane;
+ // it is simply better, and the text lane has nothing to do here regardless of
+ // prompt version, model, or whether ollama is even reachable. Putting it after
+ // any of that would make the answer depend on things it must not depend on.
+ if (isAttributionDowngrade(existing, method)) return "outranked";
+
+ const { target, app, config, modelRequested } = resolveAttributionTarget(
+ method,
+ cfg,
+ );
+ // The per-video half of the diarized identity. See
+ // AttributionProvenance.diarizationGeneratedAt: a cluster index means nothing
+ // except relative to the run that produced it.
+ const fullTarget: AttributionFreshnessTarget = {
+ ...target,
+ ...(diarization?.generatedAt
+ ? { diarizationGeneratedAt: diarization.generatedAt }
+ : {}),
+ };
+ if (!opts.force && isAttributionFresh(existing, fullTarget)) {
+ return "already-exists";
+ }
+
+ const context = await readDigestContext(opts.paths, opts.channelSlug);
+ const warnings: AttributionWarning[] = [];
+ const started = Date.now();
+ let speakers: AttributionSpeaker[];
+ let segments: AttributionSegment[];
+ let reportedModel = modelRequested;
+ let costUsd = 0;
+ let chunks: number | undefined;
+ let chunksOk: number | undefined;
+
+ const runInput = {
+ config,
+ signal: opts.signal,
+ onLog: opts.onLog,
+ };
+ const title = transcript.title || opts.videoId;
+ const channel = transcript.channel || opts.channelSlug;
+
+ try {
+ if (method === "diarized") {
+ const out = await nameClusters({
+ app,
+ runInput,
+ diarization: diarization!,
+ cues,
+ title,
+ channel,
+ contextNote: context.note,
+ warnings,
+ log,
+ videoId: opts.videoId,
+ });
+ if (!out) return "failed";
+ speakers = out.speakers;
+ segments = out.segments;
+ reportedModel = out.model || reportedModel;
+ costUsd += out.costUsd;
+ } else {
+ const out = await findTurns({
+ app,
+ runInput,
+ transcript,
+ cues,
+ config,
+ title,
+ channel,
+ contextNote: context.note,
+ warnings,
+ log,
+ videoId: opts.videoId,
+ channelSlug: opts.channelSlug,
+ signal: opts.signal,
+ });
+ if (!out) return "failed";
+ speakers = out.speakers;
+ segments = out.segments;
+ reportedModel = out.model || reportedModel;
+ costUsd += out.costUsd;
+ chunks = out.chunks;
+ chunksOk = out.chunksOk;
+ }
+ } catch (err) {
+ // A cancel is not an attribution failure — let it propagate so the batch
+ // stops rather than recording a bogus empty record.
+ if (opts.signal?.aborted) throw err;
+ log(
+ `Attribute ${opts.videoId} failed: ${(err as Error)?.message ?? String(err)}`,
+ );
+ return "failed";
+ }
+
+ if (speakers.length === 0) {
+ // Do NOT write. An empty record carrying the current identity would read as
+ // fresh, and this video would never be retried — the same trap digestVideo
+ // documents at its no-output branch.
+ log(
+ `Attribute ${opts.videoId}: no speaker survived the guards (${warnings.length} warning(s)); leaving the sidecar untouched so it retries.`,
+ );
+ return "failed";
+ }
+
+ // RE-CHECKED, against disk, immediately before writing. state() ran at pull
+ // time and the pool may have held this candidate for minutes; the diarized
+ // lane can have landed a better record in that window. See the header.
+ if (isAttributionDowngrade(await loadAttribution(opts.videoDir), method)) {
+ log(
+ `Attribute ${opts.videoId}: a diarized record landed while this text-only run was in flight — keeping it.`,
+ );
+ return "outranked";
+ }
+
+ const record: AttributionRecord = {
+ videoId: opts.videoId,
+ generatedAt: new Date().toISOString(),
+ speakers,
+ segments,
+ ...(warnings.length > 0 ? { warnings } : {}),
+ provenance: {
+ method,
+ appId: app.id,
+ model: reportedModel,
+ modelRequested,
+ promptVersion: cfg.promptVersion ?? ATTRIBUTION_PROMPT_VERSION,
+ generatedAt: new Date().toISOString(),
+ ...(diarization?.generatedAt
+ ? { diarizationGeneratedAt: diarization.generatedAt }
+ : {}),
+ ...(chunks !== undefined ? { chunks, chunksOk } : {}),
+ ...(app.metered && costUsd > 0 ? { costUsd } : {}),
+ durationMs: Date.now() - started,
+ },
+ };
+ await writeAttribution(opts.videoDir, record);
+ log(
+ `Attribute ${opts.videoId} (${method}): ${speakers.length} speaker(s), ` +
+ `${segments.length} segment(s)` +
+ (chunks !== undefined ? `, ${chunksOk}/${chunks} chunk(s)` : "") +
+ (warnings.length > 0 ? `, ${warnings.length} warning(s)` : "") +
+ (costUsd > 0 ? `, $${costUsd.toFixed(4)}` : "") +
+ ".",
+ );
+ return "attributed";
+}
+
+// ---------------------------------------------------------------------------
+// Diarized lane — ONE call
+// ---------------------------------------------------------------------------
+
+type LaneResult = {
+ speakers: AttributionSpeaker[];
+ segments: AttributionSegment[];
+ model: string;
+ costUsd: number;
+};
+
+async function nameClusters(args: {
+ app: DigestApp;
+ runInput: { config: DigestAppConfig; signal?: AbortSignal; onLog?: (m: string) => void };
+ diarization: DiarizationRecord;
+ cues: Cue[];
+ title: string;
+ channel: string;
+ contextNote?: string;
+ warnings: AttributionWarning[];
+ log: (m: string) => void;
+ videoId: string;
+}): Promise<LaneResult | null> {
+ const clusters = selectClusterSamples(args.diarization.turns, args.cues);
+ if (clusters.length === 0) {
+ args.warnings.push({ code: "empty", detail: "no cluster had usable speech" });
+ return null;
+ }
+ args.log(
+ `Attribute ${args.videoId}: naming ${clusters.length} of ${args.diarization.speakers} cluster(s) in one call.`,
+ );
+ const result = await args.app.run({
+ ...args.runInput,
+ system: SPEAKER_SYSTEM_PROMPT,
+ prompt: buildSpeakerPrompt({
+ title: args.title,
+ channel: args.channel,
+ clusters,
+ ...(args.contextNote ? { contextNote: args.contextNote } : {}),
+ }),
+ schema: speakerSchema(clusters.map((c) => c.cluster)),
+ });
+
+ const named = new Map<number, { label: string; confidence?: number }>();
+ const raw = (result.data as { speakers?: unknown })?.speakers;
+ const allowed = new Set(clusters.map((c) => c.cluster));
+ if (Array.isArray(raw)) {
+ for (const item of raw) {
+ const e = item as { cluster?: unknown; label?: unknown; confidence?: unknown };
+ if (typeof e?.cluster !== "number" || typeof e.label !== "string") continue;
+ // The schema pins this with an enum, and the parser checks it anyway —
+ // the same belt-and-braces the digest parser applies to HMS_RE, for the
+ // same reason: an engine that ignores the schema must not be able to
+ // attribute speech to a cluster that does not exist.
+ if (!allowed.has(e.cluster)) {
+ args.warnings.push({
+ code: "unknown-cluster",
+ detail: `cluster ${e.cluster}`,
+ });
+ continue;
+ }
+ if (isUselessSpeakerLabel(e.label)) continue;
+ if (named.has(e.cluster)) continue;
+ named.set(e.cluster, {
+ label: e.label.trim(),
+ ...(typeof e.confidence === "number"
+ ? { confidence: clamp01(e.confidence) }
+ : {}),
+ });
+ }
+ }
+ if (named.size === 0) {
+ args.warnings.push({ code: "empty", detail: "no usable label" });
+ return null;
+ }
+
+ // Speaker order follows the CLUSTER ranking (heaviest first), not the model's
+ // reply order, so `speakers[0]` is the dominant voice whatever the engine did.
+ const speakers: AttributionSpeaker[] = [];
+ const indexOfCluster = new Map<number, number>();
+ for (const c of clusters) {
+ const hit = named.get(c.cluster);
+ if (!hit) continue;
+ indexOfCluster.set(c.cluster, speakers.length);
+ speakers.push({
+ index: speakers.length,
+ label: hit.label,
+ cluster: c.cluster,
+ ...(hit.confidence !== undefined ? { confidence: hit.confidence } : {}),
+ seconds: c.seconds,
+ });
+ }
+
+ // Segments come straight from the diarizer's turns — the acoustic boundaries
+ // are better than anything a model reading text would propose, and re-deriving
+ // them from the labels would throw that away. Turns whose cluster went unnamed
+ // are simply dropped: unattributed is the correct answer for them.
+ const segments = mergeAdjacent(
+ args.diarization.turns
+ .filter((t) => indexOfCluster.has(t.speaker))
+ .map((t) => ({
+ start: t.start,
+ end: t.end,
+ speaker: indexOfCluster.get(t.speaker)!,
+ }))
+ .sort((a, b) => a.start - b.start),
+ );
+ return {
+ speakers,
+ segments,
+ model: result.model,
+ costUsd: typeof result.costUsd === "number" ? result.costUsd : 0,
+ };
+}
+
+// ---------------------------------------------------------------------------
+// Text-only lane — one call per chunk, roster carried across the seams
+// ---------------------------------------------------------------------------
+
+async function findTurns(args: {
+ app: DigestApp;
+ runInput: { config: DigestAppConfig; signal?: AbortSignal; onLog?: (m: string) => void };
+ transcript: { id: string; title: string; channel?: string; duration?: number };
+ cues: Cue[];
+ config: DigestAppConfig;
+ title: string;
+ channel: string;
+ contextNote?: string;
+ warnings: AttributionWarning[];
+ log: (m: string) => void;
+ videoId: string;
+ channelSlug: string;
+ signal?: AbortSignal;
+}): Promise<(LaneResult & { chunks: number; chunksOk: number }) | null> {
+ // Sized to the CONFIGURED context, exactly as the digest chunker is: a smaller
+ // window with the default chunk size overflows and ollama truncates SILENTLY,
+ // which reads as a bad model rather than a misconfiguration.
+ const maxCues = maxCuesForContext(args.config.numCtx);
+ const chunks = chunkCuesForContext(args.cues, {
+ maxCues,
+ overlapCues: DIGEST_OVERLAP_CUES,
+ });
+ args.log(
+ `Attribute ${args.videoId}: ${args.cues.length} cues → ${chunks.length} chunk(s), text-only.`,
+ );
+
+ // The roster, keyed by a normalized label so "The Host" and "the host" are one
+ // person. Insertion order IS speaker order, so the first voice heard is
+ // speaker 0.
+ const roster = new Map<string, number>();
+ const labels: string[] = [];
+ const marks: { at: number; speaker: number }[] = [];
+ let reportedModel = "";
+ let costUsd = 0;
+ let chunksOk = 0;
+
+ for (let i = 0; i < chunks.length; i++) {
+ args.signal?.throwIfAborted();
+ const chunk = chunks[i];
+ const startSeconds = Math.max(0, Math.floor(chunk[0].start));
+ const endSeconds = Math.max(
+ startSeconds,
+ Math.ceil(chunk[chunk.length - 1].end || chunk[chunk.length - 1].start),
+ );
+ try {
+ const result = await args.app.run({
+ ...args.runInput,
+ system: TURN_SYSTEM_PROMPT,
+ prompt: buildTurnPrompt({
+ title: args.title,
+ channel: args.channel,
+ startSeconds,
+ endSeconds,
+ transcript: renderChunk(args.transcript, chunk),
+ // The roster as it stands after every EARLIER chunk. See
+ // TurnPromptInput.knownSpeakers — this is the cross-chunk identity
+ // mechanism, and the reason the calls are sequential rather than
+ // parallel.
+ ...(labels.length > 0 ? { knownSpeakers: labels.slice() } : {}),
+ ...(args.contextNote ? { contextNote: args.contextNote } : {}),
+ }),
+ schema: turnSchema(),
+ });
+ chunksOk++;
+ reportedModel = result.model || reportedModel;
+ if (typeof result.costUsd === "number") costUsd += result.costUsd;
+
+ const raw = (result.data as { turns?: unknown })?.turns;
+ if (!Array.isArray(raw)) continue;
+ for (const item of raw) {
+ const e = item as { start?: unknown; speaker?: unknown };
+ if (typeof e?.start !== "string" || typeof e.speaker !== "string") continue;
+ // The parser's own copy of the rule the schema pins, for the case an
+ // engine ignores the schema.
+ const at = hmsToSeconds(e.start);
+ if (at === null) {
+ args.warnings.push({
+ code: "bad-timestamp",
+ chunk: i,
+ detail: e.start.slice(0, 40),
+ });
+ continue;
+ }
+ // Out of the chunk's own range is wrong even when the topic is real —
+ // clamped away rather than trusted, since a mark in another chunk's
+ // territory would fight that chunk's own answer.
+ if (at < startSeconds || at > endSeconds) continue;
+ if (isUselessSpeakerLabel(e.speaker)) continue;
+ const key = normalizeLabel(e.speaker);
+ let index = roster.get(key);
+ if (index === undefined) {
+ index = labels.length;
+ roster.set(key, index);
+ labels.push(e.speaker.trim());
+ }
+ marks.push({ at, speaker: index });
+ }
+ } catch (err) {
+ if (args.signal?.aborted) throw err;
+ const message = (err as Error)?.message ?? String(err);
+ args.warnings.push({
+ code: "chunk-failed",
+ chunk: i,
+ detail: message.slice(0, 300),
+ });
+ args.log(
+ `Attribute ${args.videoId}: chunk ${i + 1}/${chunks.length} failed: ${message}`,
+ );
+ }
+ }
+
+ if (labels.length === 0) {
+ return { speakers: [], segments: [], model: reportedModel, costUsd, chunks: chunks.length, chunksOk };
+ }
+
+ // Marks -> segments. Each mark runs until the next one; the last runs to the
+ // end of the transcript. Duplicates at a seam (the overlap means two chunks
+ // see the same stretch) collapse because a mark that does not CHANGE the
+ // speaker is not a boundary.
+ marks.sort((a, b) => a.at - b.at || a.speaker - b.speaker);
+ const lastCue = args.cues[args.cues.length - 1];
+ const videoEnd = Math.max(lastCue.end || lastCue.start, marks[marks.length - 1]?.at ?? 0);
+ const segments: AttributionSegment[] = [];
+ for (let i = 0; i < marks.length; i++) {
+ const m = marks[i];
+ if (segments.length > 0 && segments[segments.length - 1].speaker === m.speaker) {
+ continue;
+ }
+ const end = nextChange(marks, i, m.speaker, videoEnd);
+ // A zero-length segment is a model emitting two marks on the same second;
+ // it is noise, not a turn.
+ if (end > m.at) segments.push({ start: m.at, end, speaker: m.speaker });
+ }
+
+ const seconds = new Map<number, number>();
+ for (const s of segments) {
+ seconds.set(s.speaker, (seconds.get(s.speaker) ?? 0) + (s.end - s.start));
+ }
+ const speakers: AttributionSpeaker[] = labels.map((label, index) => ({
+ index,
+ label,
+ // NO confidence field. The text-only lane has nothing honest to put here:
+ // the model was not asked for one (it is answering "where does the speaker
+ // change", not "how sure are you who this is"), and inventing a number would
+ // be exactly the overclaiming PLAN.md warns against. A consumer that wants
+ // to filter on confidence should filter on `method` instead.
+ seconds: Math.round(seconds.get(index) ?? 0),
+ }));
+ return {
+ speakers,
+ segments,
+ model: reportedModel,
+ costUsd,
+ chunks: chunks.length,
+ chunksOk,
+ };
+}
+
+// Where the speaker next changes after mark `i`, or the end of the video.
+function nextChange(
+ marks: { at: number; speaker: number }[],
+ i: number,
+ speaker: number,
+ videoEnd: number,
+): number {
+ for (let j = i + 1; j < marks.length; j++) {
+ if (marks[j].speaker !== speaker) return marks[j].at;
+ }
+ return videoEnd;
+}
+
+// ---------------------------------------------------------------------------
+
+function clamp01(n: number): number {
+ if (!Number.isFinite(n)) return 0;
+ return Math.min(1, Math.max(0, n));
+}
+
+// "The Host:" and "the host" are the same person. Deliberately conservative —
+// it folds case, surrounding punctuation and internal whitespace, and nothing
+// else. Fuzzier matching ("Jane" == "Jane Doe") would merge two people on a
+// guess, and the asymmetry runs the other way: a split speaker is a visible
+// quality problem, a merged one silently attributes words to the wrong person.
+function normalizeLabel(label: string): string {
+ return label
+ .trim()
+ .replace(/\s+/g, " ")
+ .replace(/^[^\p{L}\p{N}]+|[^\p{L}\p{N}]+$/gu, "")
+ .toLowerCase();
+}
+
+function mergeAdjacent(segments: AttributionSegment[]): AttributionSegment[] {
+ const out: AttributionSegment[] = [];
+ for (const s of segments) {
+ const last = out[out.length - 1];
+ if (last && last.speaker === s.speaker && s.start <= last.end) {
+ last.end = Math.max(last.end, s.end);
+ continue;
+ }
+ out.push({ ...s });
+ }
+ return out;
+}
+
+// Render one chunk as the text the engine sees. transcriptToMarkdown is the
+// declared single source of truth for "transcript -> text for an AI", and
+// toHms — not aiHandoff's hms() — because the schema pattern requires two digits
+// in all three fields and the prompt tells the model to copy its starts from
+// these markers.
+//
+// ABSOLUTE timestamps, unlike the digest's measured-best chunk-local mode. The
+// difference is that a digest chapter is a label for a chunk, while an
+// attribution mark has to be placed on one shared timeline that later chunks
+// keep adding to — re-basing each chunk to 00:00:00 would make every chunk's
+// marks collide with every other chunk's.
+function renderChunk(
+ transcript: { id: string; title: string; channel?: string; duration?: number },
+ cues: Cue[],
+): string {
+ return transcriptToMarkdown(
+ {
+ id: transcript.id,
+ title: transcript.title,
+ channel: transcript.channel,
+ duration: transcript.duration,
+ cues,
+ },
+ {
+ timestamps: true,
+ includeDescription: false,
+ includeTags: false,
+ stampForCue: (_clock, seconds) => toHms(Math.max(0, seconds)),
+ },
+ );
+}
diff --git a/common/controller/attributionTarget.ts b/common/controller/attributionTarget.ts
@@ -0,0 +1,76 @@
+// The ONE place an attribution freshness target is derived — the direct mirror
+// of controller/digestTarget.ts, and split out of attributeOne.ts for the same
+// reason that module exists separately from digestVideo.ts: the COUNTERS need
+// the identity without needing the runner.
+//
+// lib/backfillKinds.ts is imported by controller/channelSnapshot.ts, which
+// classifies every video of every channel and sits on the editor's hot path.
+// Reaching the runner from there drags its whole import graph in with it — the
+// transcript normalizer, the markdown renderer, the digest prompt module and the
+// channel-context reader — none of which a CLASSIFICATION needs. So the target
+// resolver is eager (it is a settings read plus the app registry) and the runner
+// is loaded lazily at the point of actually running something.
+//
+// Honesty about what this is worth: it is a structural argument, not a measured
+// win. It was written while chasing a flaky drain e2e, on the theory that the
+// extra dev-mode module compilation was losing a race. That theory did not
+// survive measurement — the test fails at roughly the same rate with and without
+// this split (3/6 against 2/6 back-to-back), and it fails on a pristine HEAD too.
+// The split is kept because keeping a runner off a per-video hot path is right on
+// its own terms, not because it fixed anything.
+
+import { getSettings, type AttributionSettings } from "../lib/settings";
+import {
+ getDigestApp,
+ type DigestApp,
+ type DigestAppConfig,
+} from "../lib/digestApps";
+import {
+ attributionTarget,
+ type AttributionFreshnessTarget,
+ type AttributionMethod,
+} from "../lib/attribution";
+
+export type ResolvedAttributionTarget = {
+ target: AttributionFreshnessTarget;
+ app: DigestApp;
+ config: DigestAppConfig;
+ // What the config asked for. Distinct from what the engine reports having run
+ // — freshness compares this one (see AttributionProvenance.modelRequested).
+ modelRequested: string;
+};
+
+// The identity the current configuration would produce, minus the per-video
+// half.
+//
+// The diarized lane's identity also includes the generatedAt of the
+// diarization.json it names clusters from, and that is a disk read. Callers that
+// have the value pass it in; the registry's state() adds it only inside the one
+// branch that has already paid for the read.
+export function resolveAttributionTarget(
+ method: AttributionMethod,
+ cfg?: AttributionSettings,
+): ResolvedAttributionTarget {
+ const settings = getSettings();
+ const attribution = cfg ?? settings.attribution;
+ const app = getDigestApp(attribution.appId);
+ // The digest app's OWN config block — the ollama URL, context size and
+ // timeout. Attribution is a digest-app workload; a second copy of that config
+ // would be one more thing to keep in step for no benefit.
+ const config: DigestAppConfig = settings.digest.apps[app.id] ?? {};
+ const modelRequested =
+ attribution.model.trim() || config.model?.trim() || app.defaultModel();
+ return {
+ target: attributionTarget(
+ {
+ appId: app.id,
+ model: modelRequested,
+ promptVersion: attribution.promptVersion,
+ },
+ method,
+ ),
+ app,
+ config,
+ modelRequested,
+ };
+}
diff --git a/common/controller/backfillBatch.ts b/common/controller/backfillBatch.ts
Binary files differ.
diff --git a/common/controller/backfillSweep.ts b/common/controller/backfillSweep.ts
@@ -34,8 +34,12 @@ import { drainStream } from "../jobs/drainStream";
import { makeTaskTracker } from "../jobs/taskHooks";
import { requestChannelSnapshot } from "../jobs/snapshotScheduler";
import { BACKFILL_QUEUE } from "../lib/queueKeys";
-import { resolveBackfillKinds } from "../lib/backfillKinds";
-import { listChannelStatsFromDisk } from "./channels";
+import {
+ reachableBackfillWork,
+ resolveBackfillKinds,
+ type BackfillSnapshotEntry,
+} from "../lib/backfillKinds";
+import { listChannelStatsFromDisk, readChannelSnapshot } from "./channels";
import { countBackfillWork, runBackfillBatch } from "./backfillBatch";
export const BACKFILL_SWEEP_KIND = "backfill-sweep";
@@ -98,6 +102,25 @@ export type BackfillPlanEntry = {
// Reachable only, deliberately: ordering by total remaining would sort the
// corpus by how much of it is unreachable, which on the measured numbers (835
// reachable vs ~76,270 needing a re-download) means sorting by noise.
+//
+// PLANNED OFF SNAPSHOTS, not off a corpus walk. countBackfillWork() is a
+// readdir + readVideoFiles + sidecar read for EVERY video — ~77,000 directory
+// reads — and this function runs once per pass of a sweep that runs for days.
+// The channel snapshots already hold these exact numbers in
+// snapshot.backfill[kind], computed by a report the system regenerates anyway
+// (runBackfillChannelJob calls requestChannelSnapshot when it finishes), so the
+// walk was re-deriving, every hour, a number sitting in a file.
+//
+// The fallback is not a nicety: a channel that has never been snapshotted has no
+// entry, and planning it as zero would silently exclude it from the sweep
+// forever. So an ABSENT snapshot walks; a present one is trusted.
+//
+// Trusting it costs a bounded staleness — a snapshot regenerates on a ~1s
+// debounce after mutations, and worst case a channel is planned against
+// last-pass numbers. That is harmless here for the same reason the sweep stores
+// no cursor: the batch re-derives eligibility from disk on every pull, so a
+// wrong plan costs at most one channel visited in the wrong order, never work
+// done twice or work missed.
export async function buildBackfillSweepPlan(opts: {
paths: Paths;
kindIds?: string[];
@@ -108,20 +131,55 @@ export async function buildBackfillSweepPlan(opts: {
const channels = (await listChannelStatsFromDisk(opts.paths)).filter(
(ch) => wanted.size === 0 || wanted.has(ch.slug),
);
+ // Resolved once: the same filtering resolveBackfillKinds does for the batch,
+ // so the plan counts exactly the kinds the run will act on. A snapshot holds
+ // an entry per kind that was enabled when it was WRITTEN, which may be more
+ // than are in scope now.
+ const kinds = resolveBackfillKinds(getSettings(), opts.kindIds).map(
+ (k) => k.id,
+ );
const plan: BackfillPlanEntry[] = [];
+ let walked = 0;
for (const ch of channels) {
- const counts = await countBackfillWork(opts.paths, ch.slug, opts.kindIds);
+ const snapshot = await readChannelSnapshot(opts.paths, ch.slug);
+ const counts = snapshot?.backfill
+ ? countFromSnapshot(snapshot.backfill, kinds)
+ : ((walked++,
+ await countBackfillWork(opts.paths, ch.slug, opts.kindIds)));
plan.push({
channelSlug: ch.slug,
reachable: counts.reachable,
missingInput: counts.missingInput,
});
}
+ if (walked > 0) {
+ opts.onLog?.(
+ `${walked} channel(s) had no snapshot yet and were counted by walking their videos.`,
+ );
+ }
return plan.sort(
(a, b) => b.reachable - a.reachable || a.channelSlug.localeCompare(b.channelSlug),
);
}
+// The snapshot's per-kind counts, summed over the kinds in scope. `missing` and
+// `stale` are the reachable half; `missingInput` stays its own number and is
+// never added to them — see lib/backfillKinds.ts's header.
+function countFromSnapshot(
+ backfill: Record<string, BackfillSnapshotEntry>,
+ kinds: string[],
+): { reachable: number; missingInput: number } {
+ let reachable = 0;
+ let missingInput = 0;
+ for (const id of kinds) {
+ const entry = backfill[id];
+ if (!entry) continue;
+ reachable += reachableBackfillWork(entry);
+ missingInput += entry.missingInput;
+ }
+ return { reachable, missingInput };
+}
+
// Run ONE channel through the batch, as its own managed job — same kind, same
// queue and same progress metric a hand-clicked run produces, so a sweep is
// inspectable with the tools that already exist rather than being an opaque
@@ -165,6 +223,9 @@ export async function runBackfillChannelJob(opts: {
onLog(
`Backfill ${channelSlug}: ${batch.succeeded} done, ${batch.fresh} already current, ` +
`${batch.failed} failed; ${batch.missingInput} still need their media re-acquired` +
+ (batch.skipped > 0
+ ? `; ${batch.skipped} skipped (their transcript is mid-rewrite)`
+ : "") +
(batch.reacquired > 0
? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up`
: "") +
diff --git a/common/lib/attribution-server.ts b/common/lib/attribution-server.ts
@@ -0,0 +1,52 @@
+import path from "node:path";
+import { readFile, rename, writeFile } from "node:fs/promises";
+import { ATTRIBUTION_FILENAME, type AttributionRecord } from "./attribution";
+
+export function attributionPath(videoDir: string): string {
+ return path.join(videoDir, ATTRIBUTION_FILENAME);
+}
+
+export async function loadAttribution(
+ videoDir: string,
+): Promise<AttributionRecord | null> {
+ try {
+ const raw = await readFile(attributionPath(videoDir), "utf8");
+ const parsed = JSON.parse(raw) as Partial<AttributionRecord>;
+ // Validate the SHAPE, not just the parse. A half-written sidecar must read
+ // as absent everywhere — the same rule diarization-server.ts's
+ // hasDiarization encodes, and here it matters for a second reason: a
+ // malformed file that read as "present" would also read as a DOWNGRADE
+ // guard, and could block the diarized lane from ever writing a real record.
+ if (
+ typeof parsed?.videoId === "string" &&
+ typeof parsed.generatedAt === "string" &&
+ Array.isArray(parsed.speakers) &&
+ Array.isArray(parsed.segments) &&
+ !!parsed.provenance &&
+ typeof parsed.provenance.method === "string"
+ ) {
+ return parsed as AttributionRecord;
+ }
+ return null;
+ } catch {
+ return null;
+ }
+}
+
+export async function hasAttribution(videoDir: string): Promise<boolean> {
+ return (await loadAttribution(videoDir)) !== null;
+}
+
+// tmp + rename, so a crash mid-write leaves the previous record rather than a
+// truncated one. Identical to writeDiarization; the atomicity is what lets
+// loadAttribution treat "parsed but wrong shape" as a real anomaly rather than
+// the expected state of a file being written.
+export async function writeAttribution(
+ videoDir: string,
+ record: AttributionRecord,
+): Promise<void> {
+ const file = attributionPath(videoDir);
+ const tmp = `${file}.tmp-${process.pid}`;
+ await writeFile(tmp, JSON.stringify(record) + "\n");
+ await rename(tmp, file);
+}
diff --git a/common/lib/attribution.test.ts b/common/lib/attribution.test.ts
@@ -0,0 +1,355 @@
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import {
+ ATTRIBUTION_PROMPT_VERSION,
+ attributionTarget,
+ isAttributionDowngrade,
+ isAttributionFresh,
+ attributionSpeechSeconds,
+ type AttributionRecord,
+} from "./attribution";
+import { attributionStatus } from "./attributionStatus";
+import {
+ isUselessSpeakerLabel,
+ selectClusterSamples,
+ speakerSchema,
+ turnSchema,
+} from "./attributionPrompt";
+import {
+ defaultAttribution,
+ sanitizeAttribution,
+} from "./settings";
+import type { DiarizationTurn } from "./diarization";
+import type { Cue } from "./vtt";
+
+// Run with:
+// pnpm --filter yt-dlp-transcript-common exec tsx --test common/lib/attribution.test.ts
+
+const CURRENT = {
+ appId: "ollama-direct",
+ model: "qwen2.5:7b",
+ promptVersion: ATTRIBUTION_PROMPT_VERSION,
+};
+
+function record(
+ over: Partial<AttributionRecord["provenance"]> = {},
+): AttributionRecord {
+ return {
+ videoId: "vid1",
+ generatedAt: "2026-08-07T00:00:00.000Z",
+ speakers: [{ index: 0, label: "Host" }],
+ segments: [{ start: 0, end: 10, speaker: 0 }],
+ provenance: {
+ method: "diarized",
+ appId: CURRENT.appId,
+ model: "qwen2.5:7b-instruct",
+ modelRequested: CURRENT.model,
+ promptVersion: ATTRIBUTION_PROMPT_VERSION,
+ generatedAt: "2026-08-07T00:00:00.000Z",
+ ...over,
+ },
+ };
+}
+
+// ---------------------------------------------------------------------------
+// Freshness
+// ---------------------------------------------------------------------------
+
+test("fresh: the record matches the identity we would produce now", () => {
+ assert.equal(
+ isAttributionFresh(record(), attributionTarget(CURRENT, "diarized")),
+ true,
+ );
+});
+
+test("the REQUESTED model is compared, not the one the engine reported", () => {
+ // "qwen2.5:7b" resolving to "qwen2.5:7b-instruct" is not a model change, and
+ // treating it as one would regenerate the whole corpus on every run. Same rule
+ // as DigestProvenance.modelRequested.
+ assert.equal(
+ isAttributionFresh(record(), attributionTarget(CURRENT, "diarized")),
+ true,
+ );
+ assert.equal(
+ isAttributionFresh(
+ record(),
+ attributionTarget({ ...CURRENT, model: "llama3:8b" }, "diarized"),
+ ),
+ false,
+ );
+});
+
+test("stale: a different app or prompt generation is work", () => {
+ assert.equal(
+ isAttributionFresh(
+ record(),
+ attributionTarget({ ...CURRENT, appId: "claude-code" }, "diarized"),
+ ),
+ false,
+ );
+ assert.equal(
+ isAttributionFresh(
+ record(),
+ attributionTarget(
+ { ...CURRENT, promptVersion: ATTRIBUTION_PROMPT_VERSION + 1 },
+ "diarized",
+ ),
+ ),
+ false,
+ );
+});
+
+// THE COMPATIBILITY RULE. Without it, adding a field to the provenance would
+// mark every record on disk stale at once — the trap lib/digest.ts's sameVariant
+// and lib/diarization.ts's sameThreshold both exist to avoid.
+test("an absent recorded promptVersion compares equal to today's default", () => {
+ const old = record();
+ delete (old.provenance as { promptVersion?: number }).promptVersion;
+ assert.equal(
+ isAttributionFresh(old, attributionTarget(CURRENT, "diarized")),
+ true,
+ );
+ // ...equal to the DEFAULT specifically, not to anything: with a bumped
+ // generation configured, the same record is stale.
+ assert.equal(
+ isAttributionFresh(
+ old,
+ attributionTarget(
+ { ...CURRENT, promptVersion: ATTRIBUTION_PROMPT_VERSION + 1 },
+ "diarized",
+ ),
+ ),
+ false,
+ );
+});
+
+test("a record from the other lane is never fresh for this one", () => {
+ assert.equal(
+ isAttributionFresh(
+ record({ method: "text-only" }),
+ attributionTarget(CURRENT, "diarized"),
+ ),
+ false,
+ );
+ assert.equal(
+ isAttributionFresh(record(), attributionTarget(CURRENT, "text-only")),
+ false,
+ );
+});
+
+// A cluster index means nothing except relative to the diarization run that
+// produced it: re-run at a different threshold and cluster 3 is a different
+// person, or nobody. Without this the names would silently point at new clusters.
+test("re-diarizing invalidates the names that pointed at the old clusters", () => {
+ const named = record({ diarizationGeneratedAt: "2026-08-01T00:00:00.000Z" });
+ assert.equal(
+ isAttributionFresh(named, {
+ ...attributionTarget(CURRENT, "diarized"),
+ diarizationGeneratedAt: "2026-08-01T00:00:00.000Z",
+ }),
+ true,
+ );
+ assert.equal(
+ isAttributionFresh(named, {
+ ...attributionTarget(CURRENT, "diarized"),
+ diarizationGeneratedAt: "2026-08-09T00:00:00.000Z",
+ }),
+ false,
+ );
+ // The text-only lane never reads diarization, so its target carries no such
+ // value and the comparison must not fire at all.
+ assert.equal(
+ attributionTarget(CURRENT, "text-only").diarizationGeneratedAt,
+ undefined,
+ );
+});
+
+test("no record is never fresh", () => {
+ assert.equal(
+ isAttributionFresh(null, attributionTarget(CURRENT, "diarized")),
+ false,
+ );
+});
+
+// ---------------------------------------------------------------------------
+// The ordering rule — the whole safety of two kinds writing one file
+// ---------------------------------------------------------------------------
+
+test("diarized may overwrite text-only; text-only may never overwrite diarized", () => {
+ const textOnly = record({ method: "text-only" });
+ const diarized = record();
+ // The UPGRADE. This is the point of the second kind existing.
+ assert.equal(isAttributionDowngrade(textOnly, "diarized"), false);
+ // The DOWNGRADE. ~30 calls of guessed identity must not replace one call
+ // grounded in acoustic clustering.
+ assert.equal(isAttributionDowngrade(diarized, "text-only"), true);
+ // An EQUAL method is a regeneration (a prompt bump, a model change), not a
+ // downgrade — otherwise nothing could ever be redone.
+ assert.equal(isAttributionDowngrade(diarized, "diarized"), false);
+ assert.equal(isAttributionDowngrade(textOnly, "text-only"), false);
+ // Nothing on disk blocks nothing.
+ assert.equal(isAttributionDowngrade(null, "text-only"), false);
+});
+
+// ---------------------------------------------------------------------------
+// Status
+// ---------------------------------------------------------------------------
+
+test("status keeps the two lanes apart, and only claims stale when it can know", () => {
+ assert.equal(attributionStatus(null), "none");
+ assert.equal(attributionStatus(record()), "diarized");
+ assert.equal(attributionStatus(record({ method: "text-only" })), "text-only");
+ // With a target, staleness is knowable.
+ assert.equal(
+ attributionStatus(
+ record(),
+ attributionTarget({ ...CURRENT, model: "llama3:8b" }, "diarized"),
+ ),
+ "stale",
+ );
+ // WITHOUT one — an exported viewer, which has no settings and no engine —
+ // staleness is unknowable, so the honest answer is the method recorded.
+ assert.equal(attributionStatus(record()), "diarized");
+});
+
+// ---------------------------------------------------------------------------
+// Cluster sampling — what the diarized lane's one call actually costs
+// ---------------------------------------------------------------------------
+
+function cue(start: number, end: number, text: string): Cue {
+ return { start, end, text } as Cue;
+}
+
+test("clusters are ranked by talk time, and the noise floor drops fragments", () => {
+ const turns: DiarizationTurn[] = [
+ { start: 0, end: 100, speaker: 1 },
+ { start: 100, end: 160, speaker: 2 },
+ // 0.5s of crosstalk out of 160.5 — under the 1% floor, and naming it would
+ // spend prompt on nobody. Diarization over-splits; this is the direction
+ // that biases toward dropping on uncertainty.
+ { start: 160, end: 160.5, speaker: 3 },
+ ];
+ const cues = [
+ cue(0, 50, "welcome back to the show"),
+ cue(50, 100, "today we are talking about the filing"),
+ cue(100, 160, "thanks for having me on"),
+ cue(160, 160.5, "mm"),
+ ];
+ const picked = selectClusterSamples(turns, cues);
+ assert.deepEqual(
+ picked.map((c) => c.cluster),
+ [1, 2],
+ );
+ assert.equal(picked[0].seconds, 100);
+ assert.ok(picked[0].share > picked[1].share);
+ assert.ok(picked[0].samples.length > 0);
+ assert.ok(picked[0].samples[0].text.includes("welcome back"));
+});
+
+test("the speaker cap holds however badly the diarizer over-split", () => {
+ // Measured on a real file: 29 clusters before the threshold change, 13 after.
+ // The cap is what keeps one call one call.
+ const turns: DiarizationTurn[] = Array.from({ length: 40 }, (_, i) => ({
+ start: i * 10,
+ end: i * 10 + 10,
+ speaker: i,
+ }));
+ const cues = turns.map((t, i) => cue(t.start, t.end, `line ${i}`));
+ assert.equal(selectClusterSamples(turns, cues).length, 12);
+});
+
+test("a diarization with no speech yields nothing to name", () => {
+ assert.deepEqual(selectClusterSamples([], []), []);
+ assert.deepEqual(
+ selectClusterSamples([{ start: 5, end: 5, speaker: 0 }], []),
+ [],
+ );
+});
+
+// ---------------------------------------------------------------------------
+// Schemas + label guards
+// ---------------------------------------------------------------------------
+
+test("the cluster field is an enum of the clusters that actually exist", () => {
+ // An ENUM, not an integer range: numeric minimum/maximum is widely dropped by
+ // schema->grammar conversion, and a hallucinated cluster index would attribute
+ // speech to a speaker who does not exist.
+ const schema = speakerSchema([0, 3, 7]) as Record<string, never>;
+ const items = (schema.properties as never as Record<string, never>).speakers as never as {
+ items: { properties: { cluster: { enum: number[] } } };
+ maxItems: number;
+ };
+ assert.deepEqual(items.items.properties.cluster.enum, [0, 3, 7]);
+ assert.equal(items.maxItems, 3);
+});
+
+test("the turn schema pins timestamps with the same regex the digest measured", () => {
+ const schema = turnSchema() as Record<string, never>;
+ const turns = (schema.properties as never as Record<string, never>).turns as never as {
+ items: { properties: { start: { pattern: string } } };
+ };
+ // Two digits per field, so ":00:27" and "1:2:3" are rejected by the decoder.
+ assert.equal(turns.items.properties.start.pattern, "^[0-9][0-9]:[0-9][0-9]:[0-9][0-9]$");
+});
+
+test('"Speaker 1" is rejected: the number is already known', () => {
+ assert.equal(isUselessSpeakerLabel("Speaker 1"), true);
+ assert.equal(isUselessSpeakerLabel("speaker"), true);
+ assert.equal(isUselessSpeakerLabel("Unknown"), true);
+ assert.equal(isUselessSpeakerLabel("Voice 2"), true);
+ assert.equal(isUselessSpeakerLabel("a"), true);
+ assert.equal(isUselessSpeakerLabel("Host"), false);
+ assert.equal(isUselessSpeakerLabel("Jane Doe"), false);
+ // Not a false positive: a real name that merely contains the word.
+ assert.equal(isUselessSpeakerLabel("Speaker of the House"), false);
+});
+
+test("speech seconds sum the attributed segments", () => {
+ assert.equal(attributionSpeechSeconds(record()), 10);
+});
+
+// ---------------------------------------------------------------------------
+// Settings
+// ---------------------------------------------------------------------------
+
+test("nothing is armed by default", () => {
+ const d = defaultAttribution();
+ assert.equal(d.enabled, false);
+ // Both lanes off UNDER the master switch, so turning the feature on to look at
+ // it cannot start a ~194,000-call corpus sweep.
+ assert.equal(d.diarizedEnabled, false);
+ assert.equal(d.textOnlyEnabled, false);
+});
+
+test("promptVersion is floored at the shipped constant, never pinned below it", () => {
+ // Pinning freshness to a superseded generation freezes that generation's
+ // output into the corpus, indistinguishable from the current prompt's — the
+ // exact trap digestPrompt.ts's version 1 -> 2 note documents.
+ assert.equal(
+ sanitizeAttribution({ promptVersion: 0 }).promptVersion,
+ ATTRIBUTION_PROMPT_VERSION,
+ );
+ assert.equal(
+ sanitizeAttribution({ promptVersion: -5 }).promptVersion,
+ ATTRIBUTION_PROMPT_VERSION,
+ );
+ // Raising it IS allowed: that is how an operator forces a corpus-wide redo
+ // after a prompt tweak, without a code change.
+ assert.equal(
+ sanitizeAttribution({ promptVersion: ATTRIBUTION_PROMPT_VERSION + 3 })
+ .promptVersion,
+ ATTRIBUTION_PROMPT_VERSION + 3,
+ );
+ assert.equal(
+ sanitizeAttribution({ promptVersion: "nonsense" }).promptVersion,
+ ATTRIBUTION_PROMPT_VERSION,
+ );
+});
+
+test("an empty model override survives sanitization", () => {
+ // Empty is meaningful — "use the app's own model" — so it must not revert to
+ // a default that only looks the same.
+ assert.equal(sanitizeAttribution({ model: " " }).model, "");
+ assert.equal(sanitizeAttribution({ model: " llama3:8b " }).model, "llama3:8b");
+ assert.equal(sanitizeAttribution(null).model, "");
+});
diff --git a/common/lib/attribution.ts b/common/lib/attribution.ts
@@ -0,0 +1,254 @@
+// Client-safe types and constants for the per-video speaker-ATTRIBUTION sidecar
+// — who is speaking, as opposed to diarization.json's anonymous "someone else is
+// speaking now". Mirrors diarization.ts exactly: server-only I/O lives in
+// attribution-server.ts, and the freshness comparator lives here so the runner
+// and the backfill registry share ONE definition of what is stale.
+//
+// TWO LANES PRODUCE THIS FILE, and which one produced a given record is the most
+// consequential field in it:
+//
+// "diarized" — diarization.json already carries globally-consistent cluster
+// indices, so the model only has to NAME N clusters from samples.
+// About one call per video, and the cross-chunk identity problem
+// below does not exist: the engine already solved it acoustically.
+// "text-only" — no diarization.json. The model has to read the transcript to
+// find speaker changes at all, which costs roughly the digest
+// sweep's chunk count (~194,000 calls corpus-wide) and has to
+// re-establish identity at every chunk seam.
+//
+// So the two are not interchangeable outputs of one job — one is cheaper AND
+// better — and the ordering rule below is what keeps them from undoing each
+// other.
+//
+// WHY THIS IS NOT A DIGEST SECTION. Folding a `speakers` section into the digest
+// prompt looks like it halves the GPU bill. It does not work: digests are
+// generated CHUNK-LOCAL (the measured default at digest 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 this is its own artifact. Recorded here so the idea is not re-proposed
+// as an obvious saving.
+//
+// WHY THE FILENAME MATTERS: never name a sidecar `transcript.<x>.<y>` —
+// SUB_FILE_RE in lib/videoStatus.ts claims any such file as a subtitle track.
+
+// One person, as this record names them. `index` is what segments point at, so
+// it is stable within the record regardless of which lane produced it.
+export type AttributionSpeaker = {
+ index: number;
+ // What the model called them: a name when the transcript supports one
+ // ("Jane Doe"), a role when it does not ("Host", "Caller 1").
+ label: string;
+ // The diarization cluster this speaker IS. Set only by the diarized lane —
+ // text-only has no clusters to point at, which is exactly the difference
+ // between the two lanes.
+ cluster?: number;
+ // 0..1, as the model reported it. Kept because the filtering that eventually
+ // consumes this must bias toward DROPPING on uncertainty, and it cannot do
+ // that without a number to threshold on. Text-only is expected to be lower,
+ // especially on auto-caption channels with no speaker turns at all.
+ confidence?: number;
+ // Seconds of attributed speech. Denormalized from `segments` so callers can
+ // report talk-time shares without walking every segment.
+ seconds?: number;
+};
+
+// One contiguous stretch attributed to a speaker. Seconds from the start of the
+// video, the same base diarization.json and the cue stream use.
+export type AttributionSegment = {
+ start: number;
+ end: number;
+ // Index into speakers[].
+ speaker: number;
+};
+
+// A guard rejection, kept for the same reason DigestWarning is: the worst
+// outcome a generator has must leave a trace somewhere other than a job log
+// that rotates.
+export type AttributionWarning = {
+ code: "chunk-failed" | "bad-timestamp" | "unknown-cluster" | "empty";
+ chunk?: number;
+ detail?: string;
+};
+
+export type AttributionMethod = "text-only" | "diarized";
+
+// What produced a given attribution.json — the identity isAttributionFresh
+// compares against what the current configuration WOULD produce now.
+export type AttributionProvenance = {
+ method: AttributionMethod;
+ appId: string;
+ // The model that ACTUALLY ran, as the engine reported it.
+ model: string;
+ // What the config asked for. Freshness compares THIS one, so a request for
+ // "qwen2.5" resolving to "qwen2.5:7b" does not read as a model change and
+ // trigger a needless regeneration (DigestProvenance.modelRequested's rule).
+ modelRequested?: string;
+ promptVersion: number;
+ generatedAt: string;
+ // DIARIZED LANE ONLY: the `generatedAt` of the diarization.json whose cluster
+ // indices these speakers name.
+ //
+ // This is load-bearing, not bookkeeping. A cluster index is meaningful only
+ // relative to one diarization run — re-run it at a different threshold and
+ // cluster 3 is a different person, or nobody. Without this field a
+ // re-diarization would silently leave the old names pointing at new clusters,
+ // and nothing on disk could tell. With it, the naming goes stale exactly when
+ // the thing it names is replaced.
+ diarizationGeneratedAt?: string;
+ // Text-only lane: how the video was sliced and how much of it survived. Same
+ // pair digestVideo records, and for the same reason — one 500 from ollama
+ // costs that chunk, not the 30-call video, and a partial result has to be
+ // legible as partial afterwards.
+ chunks?: number;
+ chunksOk?: number;
+ // Metered lanes only.
+ costUsd?: number;
+ durationMs?: number;
+};
+
+export type AttributionRecord = {
+ videoId: string;
+ // ISO 8601, set when the sidecar is finalized.
+ generatedAt: string;
+ speakers: AttributionSpeaker[];
+ segments: AttributionSegment[];
+ warnings?: AttributionWarning[];
+ provenance: AttributionProvenance;
+};
+
+export const ATTRIBUTION_FILENAME = "attribution.json";
+
+// Bump for ANY change to the prompt text or the JSON schema in
+// attributionPrompt.ts. A record whose recorded promptVersion differs is
+// regenerated; one that matches is skipped, and that skip is what keeps a re-run
+// minutes long over a sample instead of weeks over the corpus.
+//
+// It lives HERE rather than in attributionPrompt.ts so settings.ts can import it
+// for the default without depending on the prompt module — the same arrangement
+// DEFAULT_DIARIZATION_THRESHOLD has, and for the same reason: this module has no
+// dependencies at all, so anything may import it.
+export const ATTRIBUTION_PROMPT_VERSION = 1;
+
+// THE ORDERING RULE, and the whole safety of two backfill kinds writing one file.
+//
+// Both `attribution-text` and `attribution-diarized` write attribution.json. The
+// diarized lane may overwrite a text-only record — that is an UPGRADE, and it is
+// the point of the second kind existing. The text lane must never overwrite a
+// diarized one — that is a DOWNGRADE, and it would silently replace one
+// call-per-video output grounded in acoustic clustering with a ~30-call
+// reconstruction that guesses at identity across chunk seams.
+//
+// Expressed as a rank so both the registry's state() and the runner's own guard
+// read from one definition rather than each spelling the comparison out.
+const METHOD_RANK: Record<AttributionMethod, number> = {
+ "text-only": 1,
+ diarized: 2,
+};
+
+// Would writing `candidate` over `existing` lose information? True only for a
+// strict downgrade — an equal method is a legitimate regeneration (a prompt bump,
+// a model change), not a downgrade.
+export function isAttributionDowngrade(
+ existing: AttributionRecord | null,
+ candidate: AttributionMethod,
+): boolean {
+ if (!existing) return false;
+ const have = existing.provenance?.method;
+ if (!have) return false;
+ return METHOD_RANK[have] > METHOD_RANK[candidate];
+}
+
+// What we WOULD produce for this video now, as an identity. The mirror of
+// DigestFreshnessTarget and DiarizationFreshnessTarget.
+export type AttributionFreshnessTarget = {
+ method: AttributionMethod;
+ appId: string;
+ model: string;
+ promptVersion: number;
+ // Diarized lane only: the generatedAt of the diarization.json on disk right
+ // now. Absent means "do not compare" — which is what the text-only lane wants,
+ // since it never reads diarization at all.
+ diarizationGeneratedAt?: string;
+};
+
+// The identity the current configuration would produce, from a structurally
+// typed config rather than SiteSettings — so this module stays free of a
+// settings import and settings.ts can keep importing the version constant above.
+//
+// ONE definition, and the callers are the runner's short-circuit and the backfill
+// registry's state(). A comparator and the writer it guards disagreeing about
+// the identity is how a corpus ends up either regenerating forever or never.
+export function attributionTarget(
+ cfg: {
+ appId: string;
+ model: string;
+ promptVersion?: number;
+ },
+ method: AttributionMethod,
+ diarizationGeneratedAt?: string,
+): AttributionFreshnessTarget {
+ return {
+ method,
+ appId: cfg.appId,
+ model: cfg.model,
+ promptVersion: cfg.promptVersion ?? ATTRIBUTION_PROMPT_VERSION,
+ ...(method === "diarized" && diarizationGeneratedAt
+ ? { diarizationGeneratedAt }
+ : {}),
+ };
+}
+
+// An absent recorded field is read as TODAY'S DEFAULT. Copied from digest.ts's
+// sameVariant and diarization.ts's sameThreshold for the reason both exist: the
+// compatibility rule has to be explicit, or adding a field later silently
+// invalidates every record on disk.
+function sameVersion(a: number | undefined, b: number | undefined): boolean {
+ return (
+ (a ?? ATTRIBUTION_PROMPT_VERSION) === (b ?? ATTRIBUTION_PROMPT_VERSION)
+ );
+}
+
+// Is this sidecar what the current configuration would produce, FOR THIS LANE?
+//
+// `method` is compared like everything else, which has a consequence worth being
+// explicit about: a diarized record is not "fresh" against a text-only target,
+// and a text-only record is not "fresh" against a diarized one. That is correct
+// for the diarized kind (a text-only record is an upgrade opportunity — work)
+// but would be wrong for the text kind, where a diarized record means there is
+// nothing to do. The text kind therefore checks isAttributionDowngrade FIRST and
+// only asks this question when it is genuinely its own record it is looking at.
+// See lib/backfillKinds.ts.
+export function isAttributionFresh(
+ record: AttributionRecord | null,
+ target: AttributionFreshnessTarget,
+): boolean {
+ if (!record) return false;
+ const p = record.provenance;
+ if (!p) return false;
+ if (p.method !== target.method) return false;
+ if (
+ target.diarizationGeneratedAt !== undefined &&
+ p.diarizationGeneratedAt !== target.diarizationGeneratedAt
+ ) {
+ // The clusters these names point at have been replaced. See
+ // AttributionProvenance.diarizationGeneratedAt.
+ return false;
+ }
+ return (
+ p.appId === target.appId &&
+ (p.modelRequested ?? p.model) === target.model &&
+ sameVersion(p.promptVersion, target.promptVersion)
+ );
+}
+
+// Total attributed speech, in seconds. Unlike diarizationSpeechSeconds this does
+// NOT expect overlap — a segment stream that assigns one speaker per stretch is
+// what both lanes produce — but it is still tolerant of it, because a model that
+// emits overlapping ranges is a quality problem, not a crash.
+export function attributionSpeechSeconds(record: AttributionRecord): number {
+ return record.segments.reduce(
+ (sum, s) => sum + Math.max(0, s.end - s.start),
+ 0,
+ );
+}
diff --git a/common/lib/attributionPrompt.ts b/common/lib/attributionPrompt.ts
@@ -0,0 +1,360 @@
+// The prompts + JSON schemas the two attribution lanes are driven with. Pure —
+// no I/O — so it unit-tests and can be imported anywhere. Mirrors
+// digestPrompt.ts, and reuses its timestamp machinery rather than restating it:
+// HMS_PATTERN, HMS_RE, toHms and hmsToSeconds are the SAME contract, and a
+// second copy of a regex that a measured run proved fixes malformed output is
+// exactly the sort of drift that makes one lane quietly worse than the other.
+//
+// Bump ATTRIBUTION_PROMPT_VERSION (in attribution.ts) for any change here.
+//
+// THE TWO LANES ARE PRICED IN DIFFERENT UNITS, and the prompts are why:
+//
+// diarized — ONE call per video. diarization.json already carries globally
+// consistent cluster indices, so the model is only asked to put a
+// name to each cluster given samples of what it said.
+// text-only — ONE call per CHUNK, because finding the speaker changes at all
+// requires reading the transcript. On this corpus that is ~194,000
+// calls, the same order as the digest sweep.
+//
+// Do not price either in seconds-per-audio-hour. Chunk density varies fourfold
+// across this corpus, which is documented in plans/STATE.md as the unit error
+// behind a retracted throughput headline.
+
+import { HMS_PATTERN, toHms } from "./digestPrompt";
+import type { Cue } from "./vtt";
+import type { DiarizationTurn } from "./diarization";
+
+// How many distinct speakers we will ask a model to name for one video.
+//
+// Not a schema nicety — diarization OVER-SPLITS. Measured on this corpus: a
+// known two-person interview clustered to 6 at the shipped 0.9 threshold, and a
+// full real file produced 13 speakers with the host at 73% of talk time. Naming
+// all of them wastes prompt on clusters that are seconds of crosstalk, and a
+// long enum makes a constrained decoder's job harder for no gain. So: the
+// heaviest N by speech time get named, the rest are simply left unattributed —
+// which is the direction that biases toward dropping on uncertainty.
+export const ATTRIBUTION_MAX_SPEAKERS = 12;
+
+// Clusters under this share of attributed speech are not worth a name; they are
+// overwhelmingly crosstalk or an over-split fragment of a speaker already named.
+export const ATTRIBUTION_MIN_CLUSTER_SHARE = 0.01;
+
+// Utterances shown per cluster, and how much of each. Six short samples spread
+// across the video identifies a recurring speaker far better than one long
+// passage, and keeps the whole naming call inside a default 8k window even at
+// the 12-cluster cap.
+export const ATTRIBUTION_SAMPLES_PER_SPEAKER = 6;
+export const ATTRIBUTION_SAMPLE_MAX_CHARS = 220;
+
+// A label the model may not use. "Speaker 1" is what an unhelpful model falls
+// back to, and it is indistinguishable from the cluster index we already have.
+export const ATTRIBUTION_LABEL_MIN = 2;
+export const ATTRIBUTION_LABEL_MAX = 60;
+
+// ---------------------------------------------------------------------------
+// Diarized lane — name N clusters, one call
+// ---------------------------------------------------------------------------
+
+export type ClusterSample = {
+ cluster: number;
+ // Seconds of speech attributed to this cluster by the diarizer.
+ seconds: number;
+ // Share of all attributed speech, 0..1. Stated in the prompt because talk-time
+ // is the single strongest cue for "this is the host".
+ share: number;
+ // Representative utterances, in video order.
+ samples: { clock: string; text: string }[];
+};
+
+export function speakerSchema(clusters: number[]): Record<string, unknown> {
+ return {
+ type: "object",
+ properties: {
+ speakers: {
+ type: "array",
+ minItems: 1,
+ maxItems: Math.max(1, clusters.length),
+ items: {
+ type: "object",
+ properties: {
+ // An ENUM of the actual cluster indices, not an integer with a
+ // range. Both engines convert a schema to a grammar, and an
+ // enumerated literal alternation is the constraint that converts
+ // reliably everywhere — numeric `minimum`/`maximum` is widely
+ // ignored by grammar conversion, which would put us back to
+ // validating a hallucinated cluster after the fact.
+ cluster: { enum: clusters },
+ label: {
+ type: "string",
+ minLength: ATTRIBUTION_LABEL_MIN,
+ maxLength: ATTRIBUTION_LABEL_MAX,
+ },
+ confidence: { type: "number" },
+ },
+ required: ["cluster", "label", "confidence"],
+ },
+ },
+ },
+ required: ["speakers"],
+ };
+}
+
+export const SPEAKER_SYSTEM_PROMPT = [
+ "You identify who is speaking in a transcript.",
+ "You reply with JSON only, matching the provided schema exactly.",
+ "Every label you write is in ENGLISH.",
+ "You never invent a name: a label is either a name the transcript itself",
+ "states, or a role you can justify from what that speaker says.",
+].join(" ");
+
+export type SpeakerPromptInput = {
+ title: string;
+ channel: string;
+ clusters: ClusterSample[];
+ contextNote?: string;
+};
+
+export function buildSpeakerPrompt(input: SpeakerPromptInput): string {
+ const lines: string[] = [];
+ lines.push(
+ `The transcript of "${input.title}" (${input.channel}) has been split by an audio`,
+ `diarizer into ${input.clusters.length} speaker cluster(s). Each cluster is one voice.`,
+ );
+ lines.push("");
+ lines.push(
+ "Give each cluster a label. Prefer a real name when the transcript states one —",
+ 'someone is introduced, greeted, or names themselves. When it does not, use a role',
+ 'you can justify from the material ("Host", "Interviewer", "Caller"). Do NOT use',
+ '"Speaker 1"-style labels: the number is already known and adds nothing.',
+ );
+ lines.push("");
+ lines.push(
+ "Set `confidence` between 0 and 1 for each label: 0.9 when the transcript names",
+ "that person outright, 0.5 for a role you inferred, 0.2 or less when you are",
+ "guessing. An honest low number is more useful than a confident wrong name.",
+ );
+ if (input.contextNote?.trim()) {
+ lines.push("");
+ lines.push("Context for this channel (use it for names and recurring people):");
+ lines.push(input.contextNote.trim());
+ }
+ for (const c of input.clusters) {
+ lines.push("");
+ lines.push(
+ `Cluster ${c.cluster} — ${Math.round(c.share * 100)}% of the speech (${toHms(c.seconds)} total):`,
+ );
+ for (const s of c.samples) lines.push(` [${s.clock}] ${s.text}`);
+ }
+ return lines.join("\n");
+}
+
+// Pick the clusters worth naming, and representative samples for each.
+//
+// Pure and separate from the runner so the selection can be tested without a
+// model: it is the whole of what the diarized lane costs, and getting it wrong
+// (all samples from one minute, or every three-second fragment named) is the
+// difference between one good call and one useless one.
+export function selectClusterSamples(
+ turns: DiarizationTurn[],
+ cues: Cue[],
+ opts: {
+ maxSpeakers?: number;
+ samplesPerSpeaker?: number;
+ minShare?: number;
+ } = {},
+): ClusterSample[] {
+ const maxSpeakers = opts.maxSpeakers ?? ATTRIBUTION_MAX_SPEAKERS;
+ const samplesPer = opts.samplesPerSpeaker ?? ATTRIBUTION_SAMPLES_PER_SPEAKER;
+ const minShare = opts.minShare ?? ATTRIBUTION_MIN_CLUSTER_SHARE;
+
+ const seconds = new Map<number, number>();
+ for (const t of turns) {
+ const d = Math.max(0, t.end - t.start);
+ seconds.set(t.speaker, (seconds.get(t.speaker) ?? 0) + d);
+ }
+ const total = [...seconds.values()].reduce((a, b) => a + b, 0);
+ if (total <= 0) return [];
+
+ const ranked = [...seconds.entries()]
+ .map(([cluster, s]) => ({ cluster, seconds: s, share: s / total }))
+ .filter((c) => c.share >= minShare)
+ .sort((a, b) => b.seconds - a.seconds || a.cluster - b.cluster)
+ .slice(0, maxSpeakers);
+
+ return ranked.map((c) => {
+ // The cluster's own turns, longest first — a long turn carries enough words
+ // to identify a voice, while the two-second ones are backchannel ("right",
+ // "mm-hm") that identifies nobody.
+ const mine = turns
+ .filter((t) => t.speaker === c.cluster)
+ .sort((a, b) => b.end - b.start - (a.end - a.start));
+ const picked: { clock: string; text: string }[] = [];
+ const used = new Set<number>();
+ for (const turn of mine) {
+ if (picked.length >= samplesPer) break;
+ const text = cuesInRange(cues, turn.start, turn.end, used);
+ if (!text) continue;
+ picked.push({ clock: toHms(turn.start), text });
+ }
+ // Back in video order: a model reading samples chronologically can follow a
+ // conversation, and an introduction almost always comes first.
+ picked.sort((a, b) => a.clock.localeCompare(b.clock));
+ return {
+ cluster: c.cluster,
+ seconds: Math.round(c.seconds),
+ share: c.share,
+ samples: picked,
+ };
+ });
+}
+
+// The cue text overlapping [start, end], collapsed and capped. `used` stops two
+// turns of the same cluster from quoting the same cues back.
+function cuesInRange(
+ cues: Cue[],
+ start: number,
+ end: number,
+ used: Set<number>,
+): string {
+ const parts: string[] = [];
+ for (let i = 0; i < cues.length; i++) {
+ const c = cues[i];
+ if (c.start >= end) break;
+ if ((c.end || c.start) <= start) continue;
+ if (used.has(i)) continue;
+ const text = c.text.trim().replace(/\s+/g, " ");
+ if (!text) continue;
+ used.add(i);
+ parts.push(text);
+ if (parts.join(" ").length >= ATTRIBUTION_SAMPLE_MAX_CHARS) break;
+ }
+ return parts.join(" ").slice(0, ATTRIBUTION_SAMPLE_MAX_CHARS);
+}
+
+// ---------------------------------------------------------------------------
+// Text-only lane — find the speaker changes, one call per chunk
+// ---------------------------------------------------------------------------
+
+export function turnSchema(): Record<string, unknown> {
+ return {
+ type: "object",
+ properties: {
+ turns: {
+ type: "array",
+ minItems: 1,
+ items: {
+ type: "object",
+ properties: {
+ // The SAME pin digestPrompt measured: two digits per field, so
+ // ":00:27" and "1:2:3" are rejected by the decoder itself rather
+ // than by a parser afterwards.
+ start: { type: "string", pattern: HMS_PATTERN },
+ speaker: {
+ type: "string",
+ minLength: ATTRIBUTION_LABEL_MIN,
+ maxLength: ATTRIBUTION_LABEL_MAX,
+ },
+ },
+ required: ["start", "speaker"],
+ },
+ },
+ },
+ required: ["turns"],
+ };
+}
+
+export const TURN_SYSTEM_PROMPT = [
+ "You mark where the speaker changes in a transcript.",
+ "You reply with JSON only, matching the provided schema exactly.",
+ "Every speaker label you write is in ENGLISH.",
+ "You never invent a timestamp: every start you emit is copied from a",
+ "[HH:MM:SS] marker that appears in the transcript you were given.",
+ "You reuse a speaker label you have already been given whenever the same",
+ "person is speaking, because the same person must have the same label",
+ "everywhere in the video.",
+].join(" ");
+
+export type TurnPromptInput = {
+ title: string;
+ channel: string;
+ // The chunk's own range, in seconds. Stated in the prompt — one of the three
+ // changes that eliminated out-of-range output for the digest lane.
+ startSeconds: number;
+ endSeconds: number;
+ // The transcript slice, rendered as `[HH:MM:SS] text` lines.
+ transcript: string;
+ // Labels established by EARLIER chunks of this same video.
+ //
+ // THIS IS THE CROSS-CHUNK IDENTITY MECHANISM, and it is the whole reason the
+ // text-only lane is hard. Each call sees one slice, so without carrying the
+ // roster forward the model would invent a fresh cast every chunk and speaker 0
+ // in chunk 1 would have nothing to do with speaker 0 in chunk 30 — the one
+ // property attribution needs above all others. Carrying it does not GUARANTEE
+ // consistency (a model can still split one person across two labels), which is
+ // why the diarized lane, where the clustering solved this acoustically, is
+ // both cheaper and better.
+ knownSpeakers?: string[];
+ contextNote?: string;
+};
+
+export function buildTurnPrompt(input: TurnPromptInput): string {
+ const from = toHms(input.startSeconds);
+ const to = toHms(input.endSeconds);
+ const lines: string[] = [];
+
+ lines.push(
+ `Below is one section of the transcript of "${input.title}" (${input.channel}).`,
+ );
+ lines.push("");
+ lines.push(`This section covers ${from} to ${to}.`);
+ lines.push(
+ `Mark every point where the SPEAKER changes. EVERY start you emit MUST be between`,
+ `${from} and ${to} inclusive, and MUST be copied from a [HH:MM:SS] marker in the`,
+ "transcript below — do not compute or estimate one.",
+ );
+ lines.push("");
+ lines.push(
+ "Emit one entry for the first speaker of the section and one for each change",
+ "after that. Do not emit an entry per line: a person speaking continuously for",
+ "five minutes is ONE entry.",
+ );
+ if (input.knownSpeakers && input.knownSpeakers.length > 0) {
+ lines.push("");
+ lines.push(
+ "These labels were already used earlier in THIS video. Reuse the exact label",
+ "whenever the same person is speaking — the same person must not get two",
+ "different labels:",
+ );
+ for (const s of input.knownSpeakers) lines.push(` - ${s}`);
+ lines.push(
+ "Introduce a new label only when it is genuinely someone who has not spoken yet.",
+ );
+ }
+ lines.push("");
+ lines.push(
+ "Label people by name when the transcript states one, otherwise by a role you can",
+ 'justify from what they say ("Host", "Guest", "Caller"). Do NOT use "Speaker 1".',
+ );
+ if (input.contextNote?.trim()) {
+ lines.push("");
+ lines.push("Context for this channel (use it for names and recurring people):");
+ lines.push(input.contextNote.trim());
+ }
+ lines.push("");
+ lines.push("Transcript section:");
+ lines.push("");
+ lines.push(input.transcript);
+ return lines.join("\n");
+}
+
+// A label the schema accepted but that carries no information. Rejected by the
+// parser rather than by the schema because a `not`/pattern-negation constraint
+// does not survive grammar conversion — the same reason the digest parser
+// re-checks HMS_RE that the schema already pinned.
+const USELESS_LABEL_RE = /^(speaker|person|voice|unknown)\s*[0-9]*$/i;
+
+export function isUselessSpeakerLabel(label: string): boolean {
+ const trimmed = label.trim();
+ if (trimmed.length < ATTRIBUTION_LABEL_MIN) return true;
+ return USELESS_LABEL_RE.test(trimmed);
+}
diff --git a/common/lib/attributionStatus.ts b/common/lib/attributionStatus.ts
@@ -0,0 +1,56 @@
+// The one-word answer to "what does this video know about who is speaking?",
+// client-safe so a badge, a snapshot bucket and a stage card can all read it
+// without pulling a server module in.
+//
+// FOUR STATES, and the two middle ones are deliberately not collapsed into a
+// single "attributed". The lanes are not interchangeable: "diarized" names
+// acoustic clusters and costs about one model call per video, while "text-only"
+// reconstructs speaker changes from the transcript at ~30 calls and has to
+// re-establish identity at every chunk seam. A UI that showed one word for both
+// would be claiming a confidence the text-only lane has not earned — and PLAN.md
+// is explicit that this misfires on rapid back-and-forth and has nothing to work
+// with on auto-caption channels. The filtering that eventually consumes this is
+// supposed to be able to say "diarized only", which it cannot do against a
+// status that has already thrown the distinction away.
+
+import {
+ isAttributionFresh,
+ type AttributionFreshnessTarget,
+ type AttributionRecord,
+} from "./attribution";
+
+export type AttributionStatus = "none" | "text-only" | "diarized" | "stale";
+
+// `target` is OPTIONAL, and the two cases are genuinely different questions.
+//
+// With a target — a server-side caller that has settings — "stale" is reachable
+// and means the record exists but is not what we would produce now. Without one
+// — an exported viewer, which has no settings and no engine — staleness is
+// unknowable, so the honest answer is the method that IS recorded. Returning
+// "stale" there would be a guess, and returning "none" would hide real work.
+export function attributionStatus(
+ record: AttributionRecord | null,
+ target?: AttributionFreshnessTarget,
+): AttributionStatus {
+ if (!record) return "none";
+ const method = record.provenance?.method;
+ if (method !== "text-only" && method !== "diarized") return "none";
+ if (target && !isAttributionFresh(record, target)) return "stale";
+ return method;
+}
+
+// UI copy, kept beside the type so a second surface cannot invent its own
+// wording for the same state.
+export function attributionStatusLabel(status: AttributionStatus): string {
+ switch (status) {
+ case "diarized":
+ return "Speakers named from audio";
+ case "text-only":
+ // Hedged on purpose. See the header.
+ return "Speakers guessed from the transcript";
+ case "stale":
+ return "Speaker names are out of date";
+ default:
+ return "No speaker names";
+ }
+}
diff --git a/common/lib/backfillKinds.test.ts b/common/lib/backfillKinds.test.ts
@@ -2,7 +2,7 @@ import { test } from "node:test";
import assert from "node:assert/strict";
import path from "node:path";
import os from "node:os";
-import { mkdtemp, mkdir, writeFile, rm } from "node:fs/promises";
+import { mkdtemp, mkdir, writeFile, rm, utimes } from "node:fs/promises";
import {
getBackfillKind,
laneBackfillKinds,
@@ -12,7 +12,20 @@ import {
reachableBackfillWork,
type BackfillClassification,
} from "./backfillKinds";
-import { readVideoFiles, SOURCE_MEDIA_BASENAME } from "./videoStatus";
+import {
+ readVideoFiles,
+ CUES_JSON_FILENAME,
+ META_FILENAME,
+ SOURCE_MEDIA_BASENAME,
+} from "./videoStatus";
+import { CUES_FILE_VERSION } from "../controller/normalizeTranscript";
+import { attributeOneVideo } from "../controller/attributeOne";
+import { OLLAMA_DIGEST_APP_ID } from "./digest";
+import {
+ ATTRIBUTION_FILENAME,
+ ATTRIBUTION_PROMPT_VERSION,
+ type AttributionRecord,
+} from "./attribution";
import {
DEFAULT_DIARIZATION_THRESHOLD,
DIARIZATION_FILENAME,
@@ -304,6 +317,328 @@ test("an unknown or stale kind id in the sweep scope is dropped, not fatal", ()
);
});
+// ---------------------------------------------------------------------------
+// Attribution — the second and third kinds, and the pair that shares one file
+// ---------------------------------------------------------------------------
+
+const attrText = getBackfillKind("attribution-text")!;
+const attrDiarized = getBackfillKind("attribution-diarized")!;
+
+function settingsWithAttribution(
+ over: Partial<SiteSettings["attribution"]> = {},
+): SiteSettings {
+ const s = defaultSiteSettings();
+ return {
+ ...s,
+ attribution: {
+ ...s.attribution,
+ enabled: true,
+ diarizedEnabled: true,
+ textOnlyEnabled: true,
+ appId: OLLAMA_DIGEST_APP_ID,
+ model: "qwen2.5:7b",
+ ...over,
+ },
+ };
+}
+
+function attrSidecar(over: Partial<AttributionRecord["provenance"]> = {}): string {
+ return JSON.stringify({
+ videoId: "vid1",
+ generatedAt: "2026-08-07T00:00:00.000Z",
+ speakers: [{ index: 0, label: "Host" }],
+ segments: [{ start: 0, end: 10, speaker: 0 }],
+ provenance: {
+ method: "diarized",
+ appId: OLLAMA_DIGEST_APP_ID,
+ model: "qwen2.5:7b",
+ modelRequested: "qwen2.5:7b",
+ promptVersion: ATTRIBUTION_PROMPT_VERSION,
+ generatedAt: "2026-08-07T00:00:00.000Z",
+ ...over,
+ },
+ } satisfies AttributionRecord);
+}
+
+// Same fixture discipline as above — a real directory, because state() is
+// defined as "read from disk, never stored".
+async function attrFixture(opts: {
+ transcript?: boolean;
+ diarization?: string;
+ attribution?: string;
+}): Promise<{ dir: string; cleanup: () => Promise<void> }> {
+ const root = await mkdtemp(path.join(os.tmpdir(), "backfill-attr-"));
+ const dir = path.join(root, "vid1");
+ await mkdir(dir, { recursive: true });
+ if (opts.transcript !== false) {
+ await writeFile(
+ path.join(dir, "transcript.json"),
+ JSON.stringify({ transcription: [{ text: "hi" }] }),
+ );
+ }
+ if (opts.diarization !== undefined) {
+ await writeFile(path.join(dir, DIARIZATION_FILENAME), opts.diarization);
+ }
+ if (opts.attribution !== undefined) {
+ await writeFile(path.join(dir, ATTRIBUTION_FILENAME), opts.attribution);
+ }
+ return { dir, cleanup: () => rm(root, { recursive: true, force: true }) };
+}
+
+async function classifyAttr(
+ kind: typeof attrText,
+ dirOpts: Parameters<typeof attrFixture>[0],
+ settings: SiteSettings = settingsWithAttribution(),
+): Promise<BackfillClassification> {
+ const { dir, cleanup } = await attrFixture(dirOpts);
+ try {
+ const files = await readVideoFiles(dir, { checkUntranscribable: true });
+ return await kind.state({
+ videoDir: dir,
+ videoId: "vid1",
+ files,
+ target: kind.resolveTarget(settings),
+ });
+ } finally {
+ await cleanup();
+ }
+}
+
+// A diarization sidecar at a known generatedAt, so the "re-diarizing invalidates
+// the naming" case has something to compare against.
+const DIARIZED_AT = "2026-08-06T00:00:00.000Z";
+function diarizationSidecar(generatedAt = DIARIZED_AT): string {
+ return JSON.stringify({
+ videoId: "vid1",
+ generatedAt,
+ speakers: 2,
+ turns: [{ start: 0, end: 10, speaker: 0 }],
+ engine: { engine: "sherpa-onnx" },
+ } satisfies DiarizationRecord);
+}
+
+test("attribution-text reaches every transcribed video and never reports missing-input", async () => {
+ // The claim that prices this lane: its input is the cue stream, which every
+ // transcribed video has. That is why it can reach the whole corpus, and why
+ // running it over the whole corpus costs ~194,000 model calls.
+ assert.equal(await classifyAttr(attrText, {}), "missing");
+ assert.equal(
+ await classifyAttr(attrText, { transcript: false }),
+ "not-applicable",
+ );
+});
+
+test("attribution-text: fresh, stale, and a malformed file that reads as absent", async () => {
+ assert.equal(
+ await classifyAttr(attrText, {
+ attribution: attrSidecar({ method: "text-only" }),
+ }),
+ "present",
+ );
+ assert.equal(
+ await classifyAttr(attrText, {
+ attribution: attrSidecar({ method: "text-only", model: "llama3:8b", modelRequested: "llama3:8b" }),
+ }),
+ "stale",
+ );
+ // A half-written sidecar must read as absent — and here that matters twice
+ // over, because a malformed file that read as "present" could masquerade as a
+ // diarized record and block this lane forever.
+ assert.equal(
+ await classifyAttr(attrText, { attribution: "{ not json" }),
+ "missing",
+ );
+});
+
+// THE DOWNGRADE RULE, in the counter. Reporting a diarized record as work would
+// put this lane in a loop of "attempt, refuse, still outstanding" on every pass
+// of a multi-day sweep.
+test("attribution-text has nothing to do where a diarized record exists", async () => {
+ assert.equal(
+ await classifyAttr(attrText, { attribution: attrSidecar() }),
+ "present",
+ );
+ // Not even when that diarized record is itself stale — it is not this lane's
+ // record to redo.
+ assert.equal(
+ await classifyAttr(attrText, {
+ attribution: attrSidecar({ model: "llama3:8b", modelRequested: "llama3:8b" }),
+ }),
+ "present",
+ );
+});
+
+test("attribution-diarized: no diarization.json is missing-input, not missing", async () => {
+ // ~73,000 videos on this corpus, against a handful reachable. The whole reason
+ // the four-state split exists — a single "remaining" number here would put
+ // every channel at the top of every list forever.
+ assert.equal(await classifyAttr(attrDiarized, {}), "missing-input");
+ assert.equal(
+ await classifyAttr(attrDiarized, { attribution: attrSidecar() }),
+ "missing-input",
+ );
+});
+
+test("attribution-diarized treats a text-only record as the upgrade queue", async () => {
+ // `missing`, not `stale`: the diarized record this kind is responsible for
+ // genuinely was never made. This is the whole of PLAN.md's bespoke "upgrade
+ // job", and it falls out of the registry rather than needing new machinery.
+ assert.equal(
+ await classifyAttr(attrDiarized, {
+ diarization: diarizationSidecar(),
+ attribution: attrSidecar({ method: "text-only" }),
+ }),
+ "missing",
+ );
+ assert.equal(
+ await classifyAttr(attrDiarized, { diarization: diarizationSidecar() }),
+ "missing",
+ );
+});
+
+test("attribution-diarized: present, and stale when the clusters underneath change", async () => {
+ assert.equal(
+ await classifyAttr(attrDiarized, {
+ diarization: diarizationSidecar(),
+ attribution: attrSidecar({ diarizationGeneratedAt: DIARIZED_AT }),
+ }),
+ "present",
+ );
+ // Re-diarized since. Cluster 3 is now a different person, or nobody, so the
+ // names that pointed at it are work again.
+ assert.equal(
+ await classifyAttr(attrDiarized, {
+ diarization: diarizationSidecar("2026-08-09T00:00:00.000Z"),
+ attribution: attrSidecar({ diarizationGeneratedAt: DIARIZED_AT }),
+ }),
+ "stale",
+ );
+ // A model change is stale the ordinary way too.
+ assert.equal(
+ await classifyAttr(attrDiarized, {
+ diarization: diarizationSidecar(),
+ attribution: attrSidecar({
+ diarizationGeneratedAt: DIARIZED_AT,
+ model: "llama3:8b",
+ modelRequested: "llama3:8b",
+ }),
+ }),
+ "stale",
+ );
+});
+
+// Capture may legitimately be switched off after a run: a sidecar on disk is a
+// perfectly good input, and refusing to name it would strand exactly the work
+// the capture lane exists to protect.
+test("attribution-diarized does not require diarization CAPTURE to still be on", async () => {
+ const s = settingsWithAttribution();
+ assert.equal(s.diarization.enabled, false);
+ assert.equal(attrDiarized.enabled(s), true);
+ assert.equal(
+ await classifyAttr(
+ attrDiarized,
+ { diarization: diarizationSidecar() },
+ s,
+ ),
+ "missing",
+ );
+});
+
+test("each attribution lane is gated separately, under one master switch", () => {
+ const off = defaultSiteSettings();
+ assert.equal(attrText.enabled(off), false);
+ assert.equal(attrDiarized.enabled(off), false);
+ // Turning the feature on must not by itself arm a ~194,000-call sweep.
+ const onlyMaster = settingsWithAttribution({
+ diarizedEnabled: false,
+ textOnlyEnabled: false,
+ });
+ assert.equal(attrText.enabled(onlyMaster), false);
+ assert.equal(attrDiarized.enabled(onlyMaster), false);
+ assert.equal(
+ attrText.enabled(settingsWithAttribution({ textOnlyEnabled: false })),
+ false,
+ );
+ assert.equal(
+ attrDiarized.enabled(settingsWithAttribution({ diarizedEnabled: false })),
+ false,
+ );
+ // Both lanes on. Diarization CAPTURE is still off in this fixture, so two —
+ // which is also the point: an attribution backfill does not need the capture
+ // lane armed to have work.
+ assert.equal(laneBackfillKinds(settingsWithAttribution()).length, 2);
+});
+
+// The cheap, better lane must get to a video first: on a video that has
+// diarization, running the text lane first would spend ~30 model calls producing
+// a record the diarized lane then replaces.
+test("the diarized lane is ordered ahead of the text-only lane", () => {
+ assert.deepEqual(
+ resolveBackfillKinds(settingsWithAttribution(), []).map((k) => k.id),
+ ["attribution-diarized", "attribution-text"],
+ );
+});
+
+// THE RUNNER'S OWN GUARD, and it is not the same test as the counter's. state()
+// runs at pull time; the pool can hold a candidate for minutes afterwards, and
+// the diarized lane can land a better record in that window. This asserts the
+// refusal happens with NO engine call at all — no settings resolved, no model
+// chosen, nothing that could fail for an unrelated reason.
+test("the text lane refuses to downgrade a diarized record, without calling an engine", async () => {
+ const { dir, cleanup } = await attrFixture({
+ attribution: attrSidecar(),
+ });
+ try {
+ // A full transcript fixture, so the run gets PAST the cues guard and the
+ // refusal is genuinely the downgrade rule rather than a missing transcript.
+ await writeTranscriptFixture(dir);
+ const outcome = await attributeOneVideo({
+ // A bogus engine: reaching it at all is the failure this test is looking
+ // for. The guard fires before anything resolves an app.
+ paths: { channelsDir: dir } as never,
+ videoDir: dir,
+ videoId: "vid1",
+ channelSlug: "chan",
+ method: "text-only",
+ settings: {
+ ...settingsWithAttribution().attribution,
+ appId: "no-such-engine",
+ },
+ // Even FORCED. Forcing a regeneration is not the same as asking for a
+ // worse record, and nothing in the UI should be able to request the second
+ // by accident.
+ force: true,
+ onLog: () => {},
+ });
+ assert.equal(outcome, "outranked");
+ } finally {
+ await cleanup();
+ }
+});
+
+// mtimes have to ascend: isCuesJsonFresh compares cues.json against the metadata
+// and the raw transcript, and a cues.json older than either means the transcript
+// changed underneath and must not be attributed.
+async function writeTranscriptFixture(dir: string): Promise<void> {
+ await writeFile(
+ path.join(dir, META_FILENAME),
+ JSON.stringify({ id: "vid1", title: "A video" }),
+ );
+ await writeFile(
+ path.join(dir, CUES_JSON_FILENAME),
+ JSON.stringify({
+ version: CUES_FILE_VERSION,
+ id: "vid1",
+ title: "A video",
+ cues: [{ start: 0, end: 5, text: "hello" }],
+ }),
+ );
+ const base = Date.now() / 1000;
+ await utimes(path.join(dir, META_FILENAME), base, base);
+ await utimes(path.join(dir, "transcript.json"), base, base);
+ await utimes(path.join(dir, CUES_JSON_FILENAME), base + 10, base + 10);
+}
+
// The accounting rule the whole feature turns on: reachable work and
// needs-re-acquiring are never added together.
test("counts keep reachable work and needs-re-acquiring apart", () => {
diff --git a/common/lib/backfillKinds.ts b/common/lib/backfillKinds.ts
@@ -36,7 +36,20 @@
// the plan named. No client component imports it — the UI is handed plain
// numbers off the channel snapshot, and labels as props.
//
-// WHY ONLY ONE ENTRY TODAY. The other catch-up mechanisms in the repo do not fit
+// THREE ENTRIES, AND THE SECOND PAIR IS WHAT MAKES THIS AN ABSTRACTION. A
+// registry with one entry is a wrapper: nothing proved that "a feature declares
+// what it needs and the system supplies the lane, the share and the indicator"
+// was true. Attribution is the test of it, and it passed — registering
+// `attribution-diarized` and `attribution-text` lit the channel stage card,
+// /actionable, the dashboard instrument and the widget strip with ZERO UI
+// changes, because all four iterate snapshot.backfill[kindId].
+//
+// It also exercised the parts of the shape that one entry could not:
+// `missing-input` for something other than audio (diarization.json, of which
+// this corpus has one), and two kinds writing the SAME FILE at different quality
+// tiers — see the ordering rule above the attribution entries.
+//
+// WHAT IS STILL NOT HERE. The other catch-up mechanisms in the repo do not fit
// this per-video probe, and forcing them in would make the table lie:
// controller/backfillAvailability.ts is CHANNEL-scoped (one JSON map, folded
// into sync at runYtdlp.ts), and controller/normalizeAll.ts has no recorded
@@ -53,6 +66,14 @@ import {
} from "./diarization";
import { loadDiarization } from "./diarization-server";
import {
+ ATTRIBUTION_FILENAME,
+ isAttributionDowngrade,
+ isAttributionFresh,
+ type AttributionFreshnessTarget,
+ type AttributionMethod,
+} from "./attribution";
+import { loadAttribution } from "./attribution-server";
+import {
findSourceMedia,
isVideoTranscribed,
type VideoFiles,
@@ -60,6 +81,16 @@ import {
import { SAVED_VIDEO_POINTER_FILENAME } from "./savedVideo";
import { resolveSavedVideo } from "./savedVideo-server";
import { diarizeOneVideo } from "../controller/diarizeOne";
+// The TARGET resolver only — a settings read plus the digest app registry.
+// attributeOne is loaded LAZILY inside run() below: this module is imported by
+// controller/channelSnapshot.ts, which classifies every video of every channel,
+// so its eager import graph sits on the editor's hot path, and a classification
+// needs none of the runner's (the transcript normalizer, the markdown renderer,
+// the digest prompt module, the channel-context reader). See
+// controller/attributionTarget.ts for what this is and is not worth — it is a
+// structural argument, not a measured speedup.
+import { resolveAttributionTarget } from "../controller/attributionTarget";
+import type { AttributeOneOutcome } from "../controller/attributeOne";
// What a video's relationship to a backfill is, right now, read from disk.
//
@@ -109,6 +140,12 @@ export type BackfillRunOutcome =
| "done"
| "already-present"
| "missing-input"
+ // Nothing to do for a reason state() could not see WITHOUT a per-video file
+ // read — the case that matters is a transcript whose cues.json is stale, which
+ // costs a read of the source file and cannot be paid 77,000 times per pass.
+ // Counted separately from `failed` so a transient condition that resolves
+ // itself does not report as a broken engine.
+ | "skipped"
| "not-configured"
| "disabled"
| "failed";
@@ -210,8 +247,193 @@ async function hasDiarizableInput(
return (await resolveSavedVideo(videoDir)) !== null;
}
+// ---------------------------------------------------------------------------
+// Attribution — the second and third entries, and the ones that make this a
+// registry rather than a wrapper around diarization.
+//
+// TWO KINDS, ONE FILE. Both write attribution.json, and the ordering rule in
+// lib/attribution.ts is the whole safety of that: the diarized lane may
+// overwrite a text-only record (an UPGRADE — that is what the second kind is
+// for), and the text lane must never overwrite a diarized one (a DOWNGRADE).
+// state() encodes it here and attributeOne re-checks it against disk immediately
+// before writing, because the pool can hold a candidate for minutes after
+// state() ran. It has a test; a comment would not have been enough.
+//
+// THE SPLIT IS ALSO WHAT MAKES THE UPGRADE QUEUE FREE. PLAN.md describes a
+// bespoke "upgrade job" for turning text-only records into diarized ones. It is
+// not needed: `attribution-diarized` reports a text-only record as MISSING work,
+// so the existing lane, sweep and indicators queue the upgrade with no new
+// machinery. What it needs re-acquiring media for is already
+// controller/backfillReacquire.ts.
+//
+// AND IT IS WHY missing-input MATTERS HERE MOST. `attribution-diarized`'s input
+// is diarization.json, of which this corpus has ONE. So its missing-input
+// population is ~73,000 videos on day one — the exact case the reachable /
+// needs-input split exists to stop from poisoning every surface. A single
+// "remaining" number would put every channel at the top of every list forever.
+
+// Shared by both attribution kinds: a video only has speakers worth naming if it
+// has a transcript. Not-applicable rather than missing-input, since re-acquiring
+// media would not help — the video needs transcribing, which is another lane's
+// job entirely.
+function attributionApplies(files: VideoFiles): boolean {
+ return isVideoTranscribed(files) && !files.isUntranscribable;
+}
+
+// The identity, minus the per-video half. The diarized lane's identity also
+// includes the generatedAt of the diarization.json it names clusters from, and
+// that is a disk read — so it is added inside the one state() branch that has
+// already paid for the read. See AttributionProvenance.diarizationGeneratedAt.
+function attributionTargetFor(
+ settings: SiteSettings,
+ method: AttributionMethod,
+): AttributionFreshnessTarget {
+ return resolveAttributionTarget(method, settings.attribution).target;
+}
+
+const attributionText: BackfillKind = {
+ id: "attribution-text",
+ label: "Speaker names (from the transcript)",
+ hint: "Speakers reconstructed from the transcript alone, for videos with no diarization. Cheaper to reach, worse than the diarized lane, and it never overwrites one.",
+ tier: "lane",
+ enabled: (settings) =>
+ settings.attribution.enabled && settings.attribution.textOnlyEnabled,
+ resolveTarget: (settings) => attributionTargetFor(settings, "text-only"),
+ async state({ videoDir, files, target }) {
+ if (!attributionApplies(files)) return "not-applicable";
+ // NEVER missing-input. The input is the cue stream, and a transcribed video
+ // has one by definition — which is exactly why this lane can reach the whole
+ // corpus and why running it over the whole corpus costs ~194,000 model calls.
+ if (!files.entries.includes(ATTRIBUTION_FILENAME)) return "missing";
+ const record = await loadAttribution(videoDir);
+ // A malformed file reads as ABSENT, the same rule the diarization entry
+ // uses. Here it also protects the write path: a half-written sidecar must
+ // not be able to masquerade as a diarized record and block this lane
+ // forever.
+ if (!record) return "missing";
+ // THE DOWNGRADE RULE, in the counter as well as the runner. A diarized
+ // record is not stale for this lane and is not work — there is simply
+ // something better here. Reporting it as work would put this lane in a loop
+ // of "attempt, refuse, still outstanding" across every pass of a sweep.
+ if (isAttributionDowngrade(record, "text-only")) return "present";
+ return isAttributionFresh(record, target as AttributionFreshnessTarget)
+ ? "present"
+ : "stale";
+ },
+ async run(opts) {
+ // Lazy, once, at the point of actually running something. See the import
+ // note at the top of this file.
+ const { attributeOneVideo } = await import("../controller/attributeOne");
+ return toBackfillOutcome(
+ await attributeOneVideo({
+ paths: opts.paths,
+ videoDir: opts.videoDir,
+ videoId: opts.videoId,
+ channelSlug: opts.channelSlug,
+ method: "text-only",
+ force: opts.force,
+ onLog: opts.onLog,
+ signal: opts.signal,
+ }),
+ );
+ },
+};
+
+const attributionDiarized: BackfillKind = {
+ id: "attribution-diarized",
+ label: "Speaker names (from the audio)",
+ hint: "Names put to the speaker clusters in diarization.json — about one model call per video, and better than the text-only lane. Needs diarization to have run first.",
+ tier: "lane",
+ enabled: (settings) =>
+ settings.attribution.enabled && settings.attribution.diarizedEnabled,
+ resolveTarget: (settings) => attributionTargetFor(settings, "diarized"),
+ async state({ videoDir, files, target }) {
+ if (!attributionApplies(files)) return "not-applicable";
+ // The input is diarization.json, NOT audio. That distinction is the whole
+ // reason this kind is cheap: the perishable input was already captured, and
+ // what is left is a naming pass that can be redone at any time.
+ //
+ // Deliberately NOT gated on settings.diarization.enabled — a sidecar
+ // captured during a past run is a perfectly good input after capture is
+ // switched off again, and refusing to name it would strand exactly the work
+ // the capture lane exists to protect.
+ if (!files.hasDiarization) return "missing-input";
+ if (!files.entries.includes(ATTRIBUTION_FILENAME)) return "missing";
+ const record = await loadAttribution(videoDir);
+ if (!record) return "missing";
+ // A text-only record here is THE UPGRADE QUEUE: the diarized record this
+ // kind is responsible for genuinely does not exist yet, so it is `missing`
+ // rather than `stale`. Both are reachable work, but the two words mean
+ // different things to an operator reading a stage card — "stale" says
+ // something changed under a record, "missing" says a better one was never
+ // made.
+ if (record.provenance.method !== "diarized") return "missing";
+ // Only now is the diarization read worth paying for: it is needed solely to
+ // ask whether the clusters these names point at are still the same clusters.
+ const diarization = await loadDiarization(videoDir);
+ if (!diarization) return "missing-input";
+ return isAttributionFresh(record, {
+ ...(target as AttributionFreshnessTarget),
+ diarizationGeneratedAt: diarization.generatedAt,
+ })
+ ? "present"
+ : "stale";
+ },
+ async run(opts) {
+ // Lazy, once, at the point of actually running something. See the import
+ // note at the top of this file.
+ const { attributeOneVideo } = await import("../controller/attributeOne");
+ return toBackfillOutcome(
+ await attributeOneVideo({
+ paths: opts.paths,
+ videoDir: opts.videoDir,
+ videoId: opts.videoId,
+ channelSlug: opts.channelSlug,
+ method: "diarized",
+ force: opts.force,
+ onLog: opts.onLog,
+ signal: opts.signal,
+ }),
+ );
+ },
+};
+
+// One mapping, shared by both kinds, so the two lanes cannot report the same
+// condition differently.
+function toBackfillOutcome(outcome: AttributeOneOutcome): BackfillRunOutcome {
+ switch (outcome) {
+ case "attributed":
+ return "done";
+ case "already-exists":
+ // A better record already exists. Nothing to do here is the SAME answer as
+ // "already current" for the lane's purposes, and reporting it as a failure
+ // would make an untouched corpus look broken.
+ case "outranked":
+ return "already-present";
+ case "no-diarization":
+ return "missing-input";
+ case "disabled":
+ return "disabled";
+ // A transcript that is absent or about to be rewritten. Not a failure of
+ // this lane and not something re-acquiring media fixes — it resolves itself
+ // when the normalize pass catches up, and the next sweep pass will see it.
+ case "no-transcript":
+ return "skipped";
+ default:
+ return "failed";
+ }
+}
+
// One entry per backfill known to the system.
-export const BACKFILL_KINDS: readonly BackfillKind[] = [diarization];
+export const BACKFILL_KINDS: readonly BackfillKind[] = [
+ diarization,
+ attributionDiarized,
+ // Text-only LAST, deliberately. resolveBackfillKinds preserves this order and
+ // backfillBatch walks the kinds in it, so on a video that has diarization the
+ // cheap, better lane gets there first and the text lane then finds a record it
+ // must not overwrite — one wasted classification instead of ~30 model calls.
+ attributionText,
+];
export const BACKFILL_KIND_BY_ID: Record<string, BackfillKind> =
Object.fromEntries(BACKFILL_KINDS.map((k) => [k.id, k]));
diff --git a/common/lib/settings.ts b/common/lib/settings.ts
@@ -25,6 +25,7 @@ import {
sanitizeAutoQueue,
} from "../jobs/autoQueuePolicy";
import { DEFAULT_DIARIZATION_THRESHOLD } from "./diarization";
+import { ATTRIBUTION_PROMPT_VERSION } from "./attribution";
import {
DEFAULT_COOKIE_MODE,
isCookieMode,
@@ -206,6 +207,55 @@ export type SiteSettings = {
// The generic catch-up lane for derived data the existing corpus predates.
// OFF by default, and idle-only when on. See BackfillSettings.
backfill: BackfillSettings;
+ // Naming the speakers diarization found (or reconstructing them from the
+ // transcript when it found none). OFF by default. See AttributionSettings.
+ attribution: AttributionSettings;
+};
+
+// Configuration for speaker attribution — putting names to the speaker turns.
+//
+// OFF by default, and that default is doing real work rather than being
+// cautious. The text-only lane costs roughly one model call per transcript
+// CHUNK, which on this corpus is ~194,000 calls, the same order as the digest
+// sweep — and the digest sweep has completed 0.17% of its own. Arming both at
+// once puts two multi-week lanes on the same 8 GB card with nothing to arbitrate
+// between them (backfillBatch's yield deliberately watches only the
+// transcription lane). Nothing here arms anything; a pilot decides whether the
+// corpus-wide text-only pass is worth 25-55 GPU-days at all.
+export type AttributionSettings = {
+ // Master switch. Off means the backfill registry reports no attribution work
+ // at all — the feature gate every BackfillKind has.
+ enabled: boolean;
+ // Which digest app runs the naming. Attribution IS a digest-app workload —
+ // constrained JSON decoding over transcript text — so it reuses that registry
+ // and that per-app config (settings.digest.apps[appId]) rather than growing a
+ // second copy of the ollama URL, context size and timeout.
+ appId: string;
+ // Model override. Empty = the app's configured model, then its default. It is
+ // separate from the digest's because the two workloads may want different
+ // sizes, and because it is part of the freshness identity: sharing the digest's
+ // model field would make a digest bake-off invalidate every attribution record
+ // on disk as a side effect.
+ model: string;
+ // The lanes, separately. Both default OFF even when `enabled` is on, so
+ // turning the feature on to look at it cannot start a corpus sweep.
+ //
+ // They are not a fallback pair. `diarized` is one call per video and grounded
+ // in acoustic clustering; `textOnly` is ~30 calls and guesses at identity
+ // across chunk seams. An operator may reasonably want the first forever and
+ // the second never.
+ diarizedEnabled: boolean;
+ textOnlyEnabled: boolean;
+ // The prompt generation a record must match to count as fresh.
+ //
+ // Defaults to (and is floored at) ATTRIBUTION_PROMPT_VERSION, the shipped
+ // constant. Raising it forces a corpus-wide regeneration without a code
+ // change, which is the honest way to redo everything after a prompt tweak.
+ // It cannot be set BELOW the shipped constant, and that floor is the lesson
+ // from digestPrompt.ts's version 1 -> 2 note: pinning freshness to an older
+ // generation freezes output from a superseded prompt into the corpus, looking
+ // identical to output from the current one.
+ promptVersion: number;
};
// Configuration for the backfill lane — the generic answer to "a derived-data
@@ -928,6 +978,7 @@ function defaults(): SiteSettings {
digest: defaultDigest(),
diarization: defaultDiarization(),
backfill: defaultBackfill(),
+ attribution: defaultAttribution(),
};
}
@@ -976,6 +1027,43 @@ export function sanitizeBackfill(value: unknown): BackfillSettings {
};
}
+export function defaultAttribution(): AttributionSettings {
+ return {
+ // OFF, and both lanes OFF under it. See AttributionSettings.
+ enabled: false,
+ appId: DEFAULT_DIGEST_APP_ID,
+ model: "",
+ diarizedEnabled: false,
+ textOnlyEnabled: false,
+ promptVersion: ATTRIBUTION_PROMPT_VERSION,
+ };
+}
+
+export function sanitizeAttribution(value: unknown): AttributionSettings {
+ const d = defaultAttribution();
+ if (!value || typeof value !== "object") return d;
+ const r = value as Record<string, unknown>;
+ const str = (v: unknown, fallback: string) =>
+ typeof v === "string" && v.trim() ? v.trim() : fallback;
+ return {
+ enabled: r.enabled === true,
+ appId: str(r.appId, d.appId),
+ // NOT `str(...)` with a fallback: empty is a meaningful value here ("use the
+ // app's own model"), so an empty string must survive rather than reverting
+ // to a default that is also empty by coincidence.
+ model: typeof r.model === "string" ? r.model.trim() : d.model,
+ diarizedEnabled: r.diarizedEnabled === true,
+ textOnlyEnabled: r.textOnlyEnabled === true,
+ // FLOORED at the shipped constant, never merely defaulted. A hand-edited
+ // value below it would pin freshness to a superseded prompt generation and
+ // freeze its output into the corpus — see AttributionSettings.promptVersion.
+ promptVersion:
+ typeof r.promptVersion === "number" && Number.isFinite(r.promptVersion)
+ ? Math.max(ATTRIBUTION_PROMPT_VERSION, Math.floor(r.promptVersion))
+ : d.promptVersion,
+ };
+}
+
export function defaultDiarization(): DiarizationSettings {
return {
// OFF. Capture is opt-in: turning it on makes the cleanup sweep start
@@ -1238,6 +1326,7 @@ export function getSettings(): SiteSettings {
merged.digest = sanitizeDigest(merged.digest);
merged.diarization = sanitizeDiarization(merged.diarization);
merged.backfill = sanitizeBackfill(merged.backfill);
+ merged.attribution = sanitizeAttribution(merged.attribution);
// Workers. When the file predates the worker model (no `workers` key),
// synthesize a default list from the (now-settled) active app + per-app
// configs so existing installs behave identically. Otherwise sanitize the
@@ -1435,6 +1524,7 @@ export async function writeSettings(next: SiteSettings): Promise<void> {
digest: sanitizeDigest(next.digest),
diarization: sanitizeDiarization(next.diarization),
backfill: sanitizeBackfill(next.backfill),
+ attribution: sanitizeAttribution(next.attribution),
};
const tmp = `${file}.tmp-${process.pid}`;
await fs.promises.writeFile(tmp, JSON.stringify(merged, null, 2) + "\n");
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **The archive can now put names to the speakers — two ways, both switched off until you decide they are worth the time.** Diarization already captures *that* the speaker changed; it records anonymous clusters, not people. Naming them is a separate pass, and there are two honest ways to do it that are not interchangeable. **From the audio:** the diarizer has already grouped every voice, so the model is only asked to put a name to each group given samples of what it said — **about one model call for a whole video**, and the speaker boundaries come from the audio rather than from a guess. **From the transcript alone:** for videos with no diarization, which today is nearly all of them, the model has to read the transcript to find the changes in the first place — roughly **one call per chunk**, about 194,000 across this archive, which is the same order as the AI digest sweep and would compete with it for the same card. So both are built and **neither is armed**: they are off, under a master switch that is also off, and the honest next step is a small measured pilot rather than committing 25–55 days of GPU time to the worse of the two lanes. **Quality is not oversold anywhere in the UI.** A name from the audio carries a confidence the model was actually asked for; a name from the transcript carries none, because the model was asked where the speaker changes, not how sure it is who anyone is — and inventing a number there would be exactly the wrong kind of confident. The transcript-only lane misfires on rapid back-and-forth, and on auto-caption channels there are no speaker turns to find at all. **The two lanes share one file, and the rule about that is the important part.** Naming from the audio may replace a transcript-only record — that is an upgrade, and it is queued for you automatically the moment a video gains diarization. The reverse can never happen: a transcript-only pass will not overwrite a record made from audio, even when forced, even when the audio and its diarization have since been deleted and that record is the only thing left that knows who was speaking. **Both lanes plug into the Backfill lane rather than being new machinery**, so they inherit the resource share, the corpus-wide sweep, the restart-survival and every indicator — the channel's Backfill card, `/actionable`, the dashboard and the widget all report them beside diarization, with reachable work and needs-its-media-back still counted separately. That separation matters more here than it did for diarization: only one video in this archive currently *has* the diarization the good lane needs, so a single “remaining” figure would be 73,000 videos of work no button can start. **Staleness is provenance, not age**, as everywhere else: a name is stale when the engine, the model or the prompt generation differs from what would be produced now — and, uniquely here, when the video has been **re-diarized underneath it**, because a cluster number means nothing except relative to the run that produced it. There is a Prompt generation number in Settings you can raise to force a corpus-wide redo; it cannot be set lower, because pinning it to a superseded prompt would freeze that prompt's output into the archive looking current.
- **Backfill is now a first-class thing the system knows about, with a lane of its own and a remainder you can see.** Every derived-data feature that lands runs into the same wall: the corpus that already exists does not have what it needs. Diarization hit it first — the input it wants is audio, which Clean-audio deletes once a video is transcribed — and the only catch-up was a per-channel button written for that one feature, with no way to ask how much of the corpus was missing it and no way to run it alongside new-video work without one starving the other. So a feature now **declares** what it needs and how to tell whether a video has it, and gets the lane, the resource share and the indicator for free. **The number this reports is split in two, and that is the whole design.** On this corpus 836 videos still have their media and about 76,270 do not — 91× more — so a single "remaining" figure would be dominated by work no button can start, and every channel would sit at the top of every list forever. Reachable work and needs-re-acquiring are therefore separate numbers on all four surfaces: the channel's new **Backfill** stage card, a new section on `/actionable` (with the second figure in its own column, exactly as "Est. reclaim" has), a dashboard instrument, and an optional widget strip. **The lane runs on its own queue**, so it is genuinely concurrent with transcription rather than sitting behind it, and how much of the machine it may take is one setting: **0 (the default) means idle-only** — full speed while transcription is quiet, standing aside the instant it isn't — and any value above 0 is a guaranteed share, floored at 1 so a small number is a slow lane and not a stopped one. Catch-up on a corpus that already exists must never slow down new arrivals, and the default enforces that rather than trusting it. There is a corpus-wide **Start Backfill Sweep** on the dashboard that persists its scope along with the flag and resumes itself after a restart. **Staleness now means provenance, not age.** A diarization sidecar records the models and clustering threshold that produced it, and the lane compares that against what the current settings *would* produce — so changing the threshold, which is the single knob most likely to make you want a re-run, finally shows as work instead of leaving the corpus looking finished. A field the record predates compares equal to today's default, so adding one does not invalidate everything on disk. **Re-downloading deleted media is off by default and is the part to read carefully.** With it on, each file is fetched, used, and deleted again immediately in a `finally` — whether the backfill succeeded, failed, or crashed — unless the video is marked "do not clean", and nothing starts at all when free disk is under the configured floor. On a disk at 97% a leak here fills it, which is why those four behaviours have e2e tests of their own.
- **Speaker diarization can now be captured while the audio still exists, and Clean audio will no longer delete audio out from under it.** Audio is the one input in this pipeline that goes away: the Clean-audio sweep removes it as soon as a video has a transcript, and nothing downstream can reconstruct it. Everything built on top of speaker turns — attribution, naming the speakers, badges, quote filtering — can be redone from a saved file at any time, but the turns themselves can only be extracted from audio. So this ships the perishable half only: a new **Diarization** section in Settings, off by default, that captures speaker-labelled time ranges to `diarization.json` beside each transcript, and a **Diarize speakers** button in every channel's Cleanup stage that backfills over whatever audio is still on disk. Turning capture on also arms a guard: the sweep now refuses to delete audio for a transcribed video that has no sidecar yet, reports it as `awaiting diarization`, and drops that video out of the channel's "Est. reclaim" so the estimate matches what the sweep will actually take. That hold is the point — it is what keeps the audio alive long enough to be captured — but it does mean enabling this stops reclaiming disk until diarization catches up, and the setting says so. **Inline-after-transcribe is a separate switch, and it is off by default for a measured reason.** Transcription on this box runs at 221 s per audio-hour on the GPU — 16.3× real time, averaged over 3,602 real videos from their own `transcribe-outcome.json` records — while diarization runs on the CPU at roughly 500–680, because no diarization model has been ported to ggml and neither ONNX nor PyTorch has a Vulkan compute path on Linux. Diarization is therefore 2–3× slower than the transcription it would follow, so running it inline makes the whole pipeline 3–4× slower and idles the GPU. The intended shape for a large batch is the opposite: turn capture on so the audio is held, leave inline off so transcription runs at full speed, then backfill. The clustering threshold defaults to **0.9**, not sherpa-onnx's own 0.5, and that is also measured: on a six-minute excerpt of a known two-person interview, 0.5 produced 22 speaker clusters and 0.9 produced 6, with the top two at 40%/40% of talk time — recognizably the two hosts. It still over-splits, which is why every sidecar records the engine, both model names, the library version and the threshold that produced it: over-splitting is the recoverable direction, and a later pass can re-cluster or re-run selectively without ever needing the audio back. The engine sits behind a path (`DIARIZE_BIN` → `scripts/diarize.mjs`, driving sherpa-onnx through `scripts/diarize-sherpa.py`) exactly as the parakeet wrapper does, so swapping in pyannote later is a settings change rather than a code change. A failing diarizer never fails a transcription, and never counts as "done" — the sweep keeps that video's audio. Full spike results, including the disk arithmetic and what is still unmeasured, are in `plans/diarization-spike-results.md`.
- **The monitor widget's builder is a floorplan you rearrange, not a list of checkboxes.** The widget rendered its nine strips in an order hardcoded in JSX — controls, last sync, scheduler, disk, cleanable total, needs cleaning, workers, active jobs, needs work — and there was no order field in its config at all, so reference data you glance at once (the two cleanable strips) sat permanently above the thing you were actually watching. The builder couldn't help, because it wasn't spatial: 22 checkboxes in three fieldsets, nothing in the control surface corresponding to where anything landed. `/widget/builder` is now a scale model of the widget. Each section is a draggable card carrying a small CSS-drawn likeness of what that strip really renders — four dots for Workers, a part-filled bar for Active jobs, a meter with the low-disk notch for Disk, stacked rows for the two channel lists — laid out in the columns and order it will appear in, with the switched-off ones in a tray underneath. Drag it with the mouse or use the ↑/↓/◀/▶ buttons on every card; nothing here is mouse-only, and a section is one checkbox that changes place rather than a control that vanishes when you turn it off. **Columns are new** — up to three, each header showing the number that actually decides whether columns are viable (the per-column pixel width at the current preview size, flagged when it falls under 150px), alongside two wider presets, **Wide** 640×260 and **Panel** 720×420, because columns need width to be worth having. They are honored at every size by default, since an arrangement you made on purpose shouldn't quietly undo itself; **Stack columns when narrow** opts a widget you intend to resize into collapsing back to a single stack, on a container query measured against the widget's own box rather than the viewport (which, embedded in an iframe, is the host page's). Per-section display options moved onto the section that owns them — the batch bar, heading progress and ETA on the Active jobs card, worker names on Workers, the Sync flags on Controls — leaving only the genuinely global ones (section titles, hide-when-idle, the gear, absolute times) in a row of their own. The in-widget gear gets the same model in a shape that fits its ~320px overlay: one row per section in render order with move buttons and a column picker, driven by the same registry and the same placement functions as the board rather than a second copy of them. **Every link written before this still works, and still renders in the original order.** The arrangement rides in one new `l=` param (`/widget?l=ctl.disk.wk.jobs-cln.act` — `.` between sections, `-` between columns, both characters a URL leaves literal where a comma would come back as `%2C`), it is omitted entirely whenever the arrangement is the one the visibility flags already imply, and an inherited layout is normalized against those flags on arrival: unknown codes and repeats are dropped, a section since switched off is dropped, and a section switched **on** but unlisted is appended rather than silently disappearing — which is what keeps an already-copied link working after you enable something new from the widget's own gear. Underneath, the nine hardcoded conditionals became a section registry that drives both configuration surfaces, so adding a tenth strip is now a table entry (there is one waiting: `/api/widget/sync` already polls a digest field that nothing renders). The builder's remembered form state resets once, since a saved entry from before this has no layout to carry forward. See `editor/app/widget/lib/{sections,placement,config}.ts`, `editor/app/widget/builder/components/{LayoutBoard,SectionCard,SectionMini}.tsx`, and the normalization tests in `editor/app/widget/lib/placement.test.ts`.
diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts
@@ -370,6 +370,29 @@ export async function saveSettingsAction(
}
: dBack;
+ // Speaker attribution. Same hidden-marker discipline, and here it guards the
+ // most expensive switch in the form: `textOnlyEnabled` arms a lane costing
+ // roughly one model call per transcript chunk — ~194,000 across this corpus.
+ // An unrelated save must never be able to flip that on.
+ const dAttr = getSettings().attribution;
+ const attributionFormPresent =
+ formData.get("attributionFormPresent") === "1";
+ const attributionSettings: SiteSettings["attribution"] =
+ attributionFormPresent
+ ? {
+ ...dAttr,
+ enabled: formData.get("attributionEnabled") === "on",
+ appId: String(formData.get("attributionAppId") ?? dAttr.appId),
+ // Trimmed but NOT defaulted: empty is meaningful here ("use the
+ // engine's own model"), so an operator clearing the field must be able
+ // to clear it.
+ model: String(formData.get("attributionModel") ?? dAttr.model).trim(),
+ diarizedEnabled: formData.get("attributionDiarized") === "on",
+ textOnlyEnabled: formData.get("attributionTextOnly") === "on",
+ promptVersion: num("attributionPromptVersion", dAttr.promptVersion),
+ }
+ : dAttr;
+
const dB = defaultBuildPipeline();
const buildModeRaw = String(formData.get("buildMode") ?? "").trim();
const buildPipeline = {
@@ -426,6 +449,7 @@ export async function saveSettingsAction(
digest: digestSettings,
diarization: diarizationSettings,
backfill: backfillSettings,
+ attribution: attributionSettings,
};
try {
await writeSettings(next);
diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx
@@ -865,6 +865,111 @@ export function SettingsForm({ initial, apps, digestApps }: Props) {
</label>
</fieldset>
<fieldset className="flex flex-col gap-3 border border-border rounded p-3">
+ <legend className="px-1 text-sm font-medium">Speaker attribution</legend>
+ {/*
+ The marker, for the third time, and here it guards the most expensive
+ thing in this form: `textOnlyEnabled` arms a lane that costs roughly
+ one model call per transcript CHUNK — ~194,000 across this corpus, the
+ same order as the digest sweep. An unrelated settings save must not be
+ able to switch that on.
+ */}
+ <input type="hidden" name="attributionFormPresent" value="1" readOnly />
+ <p className="text-xs text-muted-foreground">
+ Putting names to the speakers. Nothing here runs on its own — it
+ registers two backfills, and the Backfill lane above decides when they
+ get the machine.
+ </p>
+ <label className="flex items-start gap-2 text-sm">
+ <input
+ type="checkbox"
+ name="attributionEnabled"
+ defaultChecked={initial.attribution.enabled}
+ className="mt-1"
+ />
+ <span className="flex flex-col gap-1">
+ <span className="font-medium">Enable speaker attribution</span>
+ <span className="text-xs text-muted-foreground">
+ Off means neither lane below reports work, whatever their own
+ switches say.
+ </span>
+ </span>
+ </label>
+ <label className="flex items-start gap-2 text-sm">
+ <input
+ type="checkbox"
+ name="attributionDiarized"
+ defaultChecked={initial.attribution.diarizedEnabled}
+ className="mt-1"
+ />
+ <span className="flex flex-col gap-1">
+ <span className="font-medium">
+ Name the speakers found in the audio
+ </span>
+ <span className="text-xs text-muted-foreground">
+ About <strong>one model call per video</strong>: the diarizer
+ already grouped the voices, so the model only has to put a name to
+ each group. Needs <code>diarization.json</code>, which is why
+ videos without it are counted separately.
+ </span>
+ </span>
+ </label>
+ <label className="flex items-start gap-2 text-sm">
+ <input
+ type="checkbox"
+ name="attributionTextOnly"
+ defaultChecked={initial.attribution.textOnlyEnabled}
+ className="mt-1"
+ />
+ <span className="flex flex-col gap-1">
+ <span className="font-medium">
+ Guess the speakers from the transcript
+ </span>
+ <span className="text-xs text-muted-foreground">
+ For videos with no diarization. Roughly{" "}
+ <strong>one call per transcript chunk</strong> — about 194,000
+ across this archive, or 25–55 days of GPU time, which is the same
+ order as the digest sweep and would compete with it. Quality is
+ genuinely uncertain: it misfires on rapid back-and-forth, and
+ auto-caption channels have no speaker turns in the first place. It
+ never overwrites a record made from the audio.
+ </span>
+ </span>
+ </label>
+ <label className="flex flex-col gap-1 text-sm">
+ <span className="font-medium">Engine</span>
+ <select
+ name="attributionAppId"
+ defaultValue={initial.attribution.appId}
+ className="rounded border border-border bg-card px-2 py-1 text-sm"
+ >
+ {digestApps.map((a) => (
+ <option key={a.id} value={a.id}>
+ {a.label}
+ {a.metered ? " (metered)" : ""}
+ </option>
+ ))}
+ </select>
+ <span className="text-xs text-muted-foreground">
+ Uses the digest engines and their configuration above — same
+ constrained-JSON workload, so there is no second copy of the URL and
+ context size to keep in step.
+ </span>
+ </label>
+ <Field
+ label="Model override"
+ name="attributionModel"
+ defaultValue={initial.attribution.model}
+ hint="Empty means the engine's own configured model. Set separately from the digest's so a digest bake-off does not invalidate every attribution record as a side effect — this value is part of what makes a record fresh."
+ />
+ <Field
+ label="Prompt generation"
+ name="attributionPromptVersion"
+ defaultValue={String(initial.attribution.promptVersion)}
+ type="number"
+ hint="A record generated at an older number is stale and will be redone. Raise it to force a corpus-wide regeneration after changing the prompt; it cannot be set below the shipped value, because pinning it lower would freeze superseded output into the archive looking current."
+ />
+ </fieldset>
+ <fieldset className="flex flex-col gap-3 border border-border rounded p-3">
<legend className="px-1 text-sm font-medium">Social links</legend>
<p className="text-xs text-muted-foreground">
Default social links shown in every site's footer. Each site can
diff --git a/editor/e2e/attribution.spec.ts b/editor/e2e/attribution.spec.ts
@@ -0,0 +1,345 @@
+import { rm, writeFile } from "node:fs/promises";
+import { join } from "node:path";
+import { test, expect } from "@playwright/test";
+import type { Page } from "@playwright/test";
+import {
+ pathExists,
+ readJson,
+ resetData,
+ resolvePath,
+ writeChannelConfig,
+ writeDigestVideo,
+ writeSettings,
+ generateReport,
+} from "./helpers";
+import type { AttributionRecord } from "../../common/lib/attribution";
+
+// Speaker attribution: the SECOND and THIRD backfill kinds, and the pair that
+// share one file.
+//
+// What these specs are for, in order of how much they would cost to get wrong:
+//
+// 1. THE TEXT LANE MUST NEVER OVERWRITE A DIARIZED RECORD. Both kinds write
+// attribution.json. The diarized lane costs about one model call per video
+// and names clusters an audio diarizer produced; the text lane costs ~30
+// calls and guesses at identity across chunk seams. A downgrade silently
+// replaces the first with the second, and nothing on disk would say so.
+// 2. THE UPGRADE MUST HAPPEN. A text-only record where diarization exists is
+// the upgrade queue — the whole of PLAN.md's bespoke "upgrade job", falling
+// out of the registry rather than needing new machinery.
+// 3. THE INDICATORS MUST STAY GENERIC. Registering these kinds is supposed to
+// light every surface with no UI change; the last spec is what makes that a
+// claim rather than an assumption, and it also pins that the two numbers
+// are still not summed for a kind whose missing-input population is ~73,000
+// videos on the real corpus.
+//
+// The engine is the ollama HTTP stub the digest specs already use
+// (e2e/fixtures/ollama-stub.mjs, wired in via OLLAMA_URL in dev:test). It
+// answers from the SCHEMA — the diarized lane's cluster enum is read back out of
+// the request — so a spec cannot pass on a cluster the video never had.
+
+const CHANNEL = "attribution-channel";
+// The video that HAS diarization: the cheap, better lane.
+const DIARIZED = "attrvid0001";
+// The video that does not: the text-only lane's territory.
+const TEXTONLY = "attrvid0002";
+
+const SLOW = 120_000;
+
+function dataRel(videoId: string, file: string): string {
+ return join("test-transcripts", "channels", CHANNEL, "data", videoId, file);
+}
+
+const DIARIZED_AT = "2026-08-06T00:00:00.000Z";
+
+// Two speakers, half the video each — so the cluster ranking is a tie broken by
+// index, and the stub's first name lands on cluster 0 deterministically.
+function diarizationRecord(generatedAt = DIARIZED_AT) {
+ return {
+ videoId: DIARIZED,
+ generatedAt,
+ speakers: 2,
+ turns: [
+ { start: 0, end: 300, speaker: 0 },
+ { start: 300, end: 600, speaker: 1 },
+ ],
+ engine: {
+ engine: "fake-diarize",
+ segmentationModel: "null",
+ embeddingModel: "null",
+ threshold: 0.5,
+ },
+ };
+}
+
+function attributionRecord(over: Record<string, unknown> = {}) {
+ return {
+ videoId: DIARIZED,
+ generatedAt: "2020-01-01T00:00:00.000Z",
+ speakers: [{ index: 0, label: "An Earlier Guess" }],
+ segments: [{ start: 0, end: 10, speaker: 0 }],
+ provenance: {
+ method: "text-only",
+ appId: "ollama-direct",
+ model: "qwen2.5:7b",
+ modelRequested: "qwen2.5:7b",
+ promptVersion: 1,
+ generatedAt: "2020-01-01T00:00:00.000Z",
+ ...(over.provenance as Record<string, unknown>),
+ },
+ };
+}
+
+// Settings with both attribution lanes armed and the backfill lane at a real
+// share.
+//
+// weight 1, not the idle-only 0: these specs assert that work HAPPENS, and at
+// weight 0 anything transcribing in the same run would legitimately park the
+// lane. The idle-only default has a pure unit test, which is the right place for
+// it — no pool, no GPU, no timing.
+function attributionSettings(over: {
+ attribution?: Record<string, unknown>;
+} = {}) {
+ return {
+ adminTitle: "Test Admin",
+ maxTranscriptPageBytes: 8388608,
+ sleepBetweenDownloadsSeconds: 0,
+ minFreeDiskGB: 0,
+ verifyAvailabilityBeforeClean: false,
+ syncScheduler: { fullSweepIntervalMinutes: 0 },
+ digest: { localAppId: "ollama-direct", remoteAppId: "claude-code" },
+ // Diarization CAPTURE stays off throughout. Naming clusters that were
+ // already captured must not require the capture lane to still be armed —
+ // otherwise switching capture off would strand exactly the work it exists to
+ // protect.
+ backfill: {
+ enabled: true,
+ weight: 1,
+ concurrency: 1,
+ sweepEnabled: false,
+ sweepKinds: [],
+ sweepChannels: [],
+ allowRedownload: false,
+ },
+ attribution: {
+ enabled: true,
+ appId: "ollama-direct",
+ model: "qwen2.5:7b",
+ diarizedEnabled: true,
+ textOnlyEnabled: true,
+ promptVersion: 1,
+ ...over.attribution,
+ },
+ };
+}
+
+async function seedChannel(opts: { diarization?: boolean } = {}) {
+ await resetData(null);
+ await writeChannelConfig(CHANNEL);
+ await writeDigestVideo({ channelSlug: CHANNEL, videoId: DIARIZED });
+ await writeDigestVideo({ channelSlug: CHANNEL, videoId: TEXTONLY });
+ if (opts.diarization !== false) {
+ await writeFile(
+ resolvePath(dataRel(DIARIZED, "diarization.json")),
+ JSON.stringify(diarizationRecord()) + "\n",
+ );
+ }
+}
+
+async function runBackfill(page: Page): Promise<void> {
+ await page.goto(`/channels/${CHANNEL}`);
+ await page.getByRole("button", { name: "Backfill stage summary" }).click();
+ await page
+ .getByRole("button", { name: "Backfill channel", exact: true })
+ .click();
+ // The batch's closing summary line — emitted after the last write, so it is
+ // the happens-before edge for the sidecar reads below.
+ await expect(page.getByLabel("Backfill channel output")).toContainText(
+ "already current",
+ { timeout: 90_000 },
+ );
+}
+
+function readAttribution(videoId: string) {
+ return readJson<AttributionRecord>(dataRel(videoId, "attribution.json"));
+}
+
+// ---------------------------------------------------------------------------
+// 1. Both lanes produce a sidecar, and each one produces its OWN kind of record
+// ---------------------------------------------------------------------------
+
+test("each lane writes the record it is responsible for", async ({ page }) => {
+ test.setTimeout(SLOW);
+ await seedChannel();
+ await writeSettings(attributionSettings());
+ await generateReport(page, CHANNEL);
+
+ await runBackfill(page);
+
+ // The diarized lane: names read back onto the clusters the diarizer found.
+ const diarized = await readAttribution(DIARIZED);
+ expect(diarized.provenance.method).toBe("diarized");
+ // The identity of the diarization run these names point at. Without it a
+ // re-diarization would leave the names pointing at different clusters, and
+ // nothing on disk could tell.
+ expect(diarized.provenance.diarizationGeneratedAt).toBe(DIARIZED_AT);
+ // The stub answers from the schema's cluster ENUM, so a cluster here proves
+ // the runner sent the clusters the sidecar actually contains.
+ expect(diarized.speakers.map((s) => s.cluster).sort()).toEqual([0, 1]);
+ expect(diarized.speakers[0].label).toBe("Marla Vance");
+ expect(diarized.speakers[0].confidence).toBeGreaterThan(0.5);
+ // Segments come from the diarizer's own acoustic boundaries, not from a model
+ // reading text.
+ expect(diarized.segments.length).toBeGreaterThan(0);
+ expect(diarized.segments[0].start).toBe(0);
+ // ONE call: no chunking happened at all, which is the whole cost argument for
+ // this lane.
+ expect(diarized.provenance.chunks).toBeUndefined();
+
+ // The text-only lane: no clusters to point at, and no confidence claimed.
+ const textOnly = await readAttribution(TEXTONLY);
+ expect(textOnly.provenance.method).toBe("text-only");
+ expect(textOnly.provenance.chunksOk).toBe(textOnly.provenance.chunks);
+ expect(textOnly.speakers.length).toBeGreaterThan(0);
+ expect(textOnly.speakers[0].cluster).toBeUndefined();
+ // Deliberately absent: the model was asked where the speaker changes, not how
+ // sure it is who anyone is. Inventing a number here would be the overclaiming
+ // the plan warns against.
+ expect(textOnly.speakers[0].confidence).toBeUndefined();
+ expect(textOnly.segments.length).toBeGreaterThan(0);
+});
+
+// ---------------------------------------------------------------------------
+// 2. THE UPGRADE — a text-only record where diarization exists is work
+// ---------------------------------------------------------------------------
+
+test("the diarized lane upgrades a text-only record", async ({ page }) => {
+ test.setTimeout(SLOW);
+ await seedChannel();
+ await writeFile(
+ resolvePath(dataRel(DIARIZED, "attribution.json")),
+ JSON.stringify(attributionRecord()) + "\n",
+ );
+ await writeSettings(attributionSettings());
+ await generateReport(page, CHANNEL);
+
+ // The card counts it as work, not as done. `missing` rather than `stale`: the
+ // diarized record genuinely was never made.
+ await expect(
+ page.getByLabel("video needing a backfill attrvid0001"),
+ ).toBeVisible();
+
+ await runBackfill(page);
+
+ const record = await readAttribution(DIARIZED);
+ expect(record.provenance.method).toBe("diarized");
+ expect(record.speakers.map((s) => s.label)).not.toContain("An Earlier Guess");
+ expect(record.provenance.diarizationGeneratedAt).toBe(DIARIZED_AT);
+});
+
+// ---------------------------------------------------------------------------
+// 3. THE DOWNGRADE — and it must not happen, ever
+// ---------------------------------------------------------------------------
+
+test("the text lane will not overwrite a record made from the audio", async ({
+ page,
+}) => {
+ test.setTimeout(SLOW);
+ // No diarization.json at all: the audio was captured, named, and the sidecar
+ // has since been cleaned away. The diarized record is now the ONLY thing that
+ // knows who was speaking, and it is unreproducible.
+ await seedChannel({ diarization: false });
+ const original = attributionRecord({
+ provenance: {
+ method: "diarized",
+ appId: "ollama-direct",
+ model: "qwen2.5:7b",
+ modelRequested: "qwen2.5:7b",
+ promptVersion: 1,
+ generatedAt: "2020-01-01T00:00:00.000Z",
+ diarizationGeneratedAt: DIARIZED_AT,
+ },
+ });
+ await writeFile(
+ resolvePath(dataRel(DIARIZED, "attribution.json")),
+ JSON.stringify(original) + "\n",
+ );
+ // Only the text lane is armed, so nothing else can be what preserved the file.
+ await writeSettings(
+ attributionSettings({ attribution: { diarizedEnabled: false } }),
+ );
+ await generateReport(page, CHANNEL);
+
+ await runBackfill(page);
+
+ const after = await readAttribution(DIARIZED);
+ // Byte-for-byte the record that was there: same method, same names, same
+ // timestamp. Not merely "still diarized" — a rewrite that happened to keep the
+ // method would still have destroyed the labels.
+ expect(after.provenance.method).toBe("diarized");
+ expect(after.generatedAt).toBe("2020-01-01T00:00:00.000Z");
+ expect(after.speakers[0].label).toBe("An Earlier Guess");
+ // The lane was not idle, though — the OTHER video had no record and got one.
+ expect(await pathExists(dataRel(TEXTONLY, "attribution.json"))).toBe(true);
+});
+
+// ---------------------------------------------------------------------------
+// 4. The indicators, which were supposed to need no UI work at all
+// ---------------------------------------------------------------------------
+
+test("the stage card and /actionable show attribution beside diarization, with the two numbers still apart", async ({
+ page,
+}) => {
+ test.setTimeout(SLOW);
+ await seedChannel();
+ await writeSettings(attributionSettings());
+ await rm(resolvePath(`test-transcripts/channels/${CHANNEL}/snapshot.json`), {
+ force: true,
+ });
+ await generateReport(page, CHANNEL);
+
+ await page.goto(`/channels/${CHANNEL}`);
+ await page.getByRole("button", { name: "Backfill stage summary" }).click();
+ const section = page.getByLabel("backfill section");
+
+ // Per-kind lines appear because there is now more than one kind — the same
+ // component, unmodified, driven by a bigger registry.
+ //
+ // The diarized lane can reach ONE video (the one with diarization.json) and
+ // reports the other as needing its input re-acquired. That 1-vs-1 split is the
+ // corpus's handful-vs-73,000 in miniature, and the card must never show "2".
+ await expect(
+ section.getByLabel("backfill kind attribution-diarized"),
+ ).toContainText("1 reachable");
+ await expect(
+ section.getByLabel("backfill kind attribution-diarized"),
+ ).toContainText("1 needing media");
+ // The text lane reaches BOTH: its input is the cue stream, which every
+ // transcribed video has. That is exactly why running it corpus-wide is the
+ // expensive option.
+ await expect(
+ section.getByLabel("backfill kind attribution-text"),
+ ).toContainText("2 reachable");
+ await expect(
+ section.getByLabel("backfill kind attribution-text"),
+ ).toContainText("0 needing media");
+ // Reachable work across kinds: 1 + 2. The needs-re-acquiring figure is on its
+ // own line and is never folded into the heading.
+ await expect(section.getByRole("heading")).toContainText(
+ "Backfill derived data (3)",
+ );
+ await expect(section.getByLabel("backfill needs re-acquiring")).toContainText(
+ "1",
+ );
+
+ // /actionable, which sums across kinds but keeps the two populations in
+ // separate columns.
+ await page.goto("/actionable");
+ const row = page.getByLabel(`backfill row ${CHANNEL}`);
+ await expect(row).toBeVisible();
+ // 3 reachable and 1 needing media, side by side and never added.
+ await expect(row).toContainText("3");
+ await expect(page.getByLabel("backfill", { exact: true })).toContainText(
+ "Needs media",
+ );
+});
diff --git a/editor/e2e/fixtures/ollama-stub.mjs b/editor/e2e/fixtures/ollama-stub.mjs
@@ -89,6 +89,69 @@ function badChapters(range) {
];
}
+// ---------------------------------------------------------------------------
+// Attribution — the two lanes, answered from the schema and the prompt
+// ---------------------------------------------------------------------------
+
+// The names the stub hands out, in cluster-weight order. Real names rather than
+// "Speaker 1", which the parser rejects on purpose (the number is already known).
+const SPEAKERS = ["Marla Vance", "Guest", "Caller", "Producer"];
+
+// Diarized lane: name the clusters the schema ENUMERATES. Reading them back out
+// of the schema is the point — it proves the enum the runner built from the real
+// diarization.json is what reached the engine, so a spec cannot pass on a
+// hardcoded cluster the video never had.
+function namedClusters(speakersSchema, bad) {
+ const enumerated = speakersSchema?.items?.properties?.cluster?.enum ?? [0];
+ const out = enumerated.map((cluster, i) => ({
+ cluster,
+ label: SPEAKERS[i % SPEAKERS.length],
+ confidence: i === 0 ? 0.9 : 0.5,
+ }));
+ if (!bad) return out;
+ // BADOUT: one cluster that does not exist, and one useless label. Both must be
+ // dropped by the parser while the good entries survive — the same "one valid
+ // item so the file is still written" shape the chapter bad-output mode uses.
+ return [
+ ...out,
+ { cluster: 9999, label: "Nobody At All", confidence: 0.9 },
+ { cluster: enumerated[0] ?? 0, label: "Speaker 2", confidence: 0.9 },
+ ];
+}
+
+// Text-only lane: alternate two speakers across the chunk, staying inside the
+// range the prompt states.
+//
+// It REUSES a label it was already given. The prompt carries the roster
+// established by earlier chunks ("These labels were already used earlier in THIS
+// video"), and a model that ignores it produces a fresh cast per chunk — which
+// is precisely the cross-chunk identity failure the lane is judged on. Echoing
+// the roster is what a good model does, so the stub does it too.
+function speakerTurns(range, prompt, bad) {
+ const known = [...(prompt ?? "").matchAll(/^ {2}- (.+)$/gm)].map((m) =>
+ m[1].trim(),
+ );
+ const cast = known.length > 0 ? known : SPEAKERS.slice(0, 2);
+ const span = Math.max(1, range.end - range.start);
+ const out = [];
+ for (let i = 0; i < 2; i++) {
+ out.push({
+ start: hms(range.start + Math.floor((span * i) / 3)),
+ speaker: cast[i % cast.length],
+ });
+ }
+ if (!bad) return out;
+ return [
+ ...out,
+ // GUARD 1 — the malformed shape the naive prompt produced.
+ { start: ":00:27", speaker: "Malformed" },
+ // GUARD 2 — an hour past the end of this chunk's range.
+ { start: hms(range.end + 3600), speaker: "Out Of Range" },
+ // GUARD 3 — the label that carries no information.
+ { start: hms(range.start + 1), speaker: "Speaker 1" },
+ ];
+}
+
function readBody(req) {
return new Promise((resolve, reject) => {
let raw = "";
@@ -123,12 +186,20 @@ const server = createServer(async (req, res) => {
const range = rangeFromPrompt(prompt);
const bad = SENTINEL.test(raw);
- // Tags ask for a `tags` array; chapters ask for `chapters`. Answer whichever
- // the caller's schema names, so the stub covers both sections.
- const wantsTags = Boolean(body.format?.properties?.tags);
- const data = wantsTags
- ? { tags: ["court filings", "podcast", "legal news"] }
- : { chapters: bad ? badChapters(range) : goodChapters(range) };
+ // Answer whichever shape the caller's SCHEMA names, so one stub covers every
+ // workload driven through this engine: digest tags, digest chapters, and the
+ // two attribution lanes.
+ const props = body.format?.properties ?? {};
+ let data;
+ if (props.tags) {
+ data = { tags: ["court filings", "podcast", "legal news"] };
+ } else if (props.speakers) {
+ data = { speakers: namedClusters(props.speakers, bad) };
+ } else if (props.turns) {
+ data = { turns: speakerTurns(range, prompt, bad) };
+ } else {
+ data = { chapters: bad ? badChapters(range) : goodChapters(range) };
+ }
res.writeHead(200, { "content-type": "application/json" });
res.end(
diff --git a/plans/STATE.md b/plans/STATE.md
@@ -3,7 +3,39 @@
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-07-30 — **the cost model was wrong in its UNITS, and fixing
+**Last updated:** 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
@@ -172,7 +204,7 @@ borrowed digest is labelled as borrowed.
| 6 · Ollama `/ask` provider | not started | **No dependencies.** |
| 7 · Tags + chat highlights | not started | Tag generation is now VERIFIED WORKING (it had never run once — 0 of 109 sidecars). 11 videos, 0 warnings, metrics in `digest-validate.ts`, e2e in `digest.spec.ts`. Nothing consumes it yet. Known gap for this phase: only 2.4% cross-video tag reuse, and `antisemitism`/`anti-semitism` are two tags — the vocabulary does not converge, so a consumer needs normalization. |
| 8 · Visibility policy | not started | |
-| 9 · Attribution + quote filtering | **capture half landed** | Only the perishable part: `diarization.json` per video, off by default, plus the cleanup guard that stops the sweep deleting audio for a transcribed-but-undiarized video. Attribution, LLM naming, badges, quote filtering, the text-only lane and snapshot buckets are all still not started — deliberately, because they can be redone from the saved JSON and the audio cannot. Spike results in `plans/diarization-spike-results.md`. |
+| 9 · Attribution + quote filtering | **capture + both attribution lanes landed, nothing armed** | Capture (`diarization.json`) landed 2026-08-07 with the cleanup guard. Attribution followed the same day: `attribution.json`, both lanes (`attribution-diarized` ≈ 1 call/video, `attribution-text` ≈ 1 call/chunk), registered as the second and third **backfill kinds** — so the lane, the share, the sweep and all four indicators came for free. Both lanes are **OFF by default under an off-by-default master switch**; nothing runs until a pilot prices it. Still not started: `attributionMs` through the build, viewer badges, per-channel export counts, and quote filtering (which also depends on Phase 8). PLAN.md's bespoke "upgrade job" is **not needed** — `attribution-diarized` reports a text-only record as `missing`, so the existing lane queues the upgrade, and media re-acquisition is already `controller/backfillReacquire.ts`. Spike results in `plans/diarization-spike-results.md`. |
| 10 · Lead with the derived corpus | not started | |
| 11a · Review queue | **minimum landed** | Total failures now persist warnings (they used to persist none), `buckets.digestWarnings` + a `digest_warnings` filter + an /actionable section, and duplicate-cluster confirm/reject. Approve/dismiss state deferred — needs new persistence. |
| 11b · Viewer feedback | not started | |