// 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, transcriptSourceOf, 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 { assembleTurns, createSpeakerRoster, type SpeakerMark, } from "../lib/attributionTurns"; import { DIGEST_OVERLAP_CUES, hmsToSeconds, maxCuesForContext, toHms, } from "../lib/digestPrompt"; import { chunkCuesForContext } from "../lib/transcriptWindow"; import { transcriptToMarkdown } from "../lib/transcriptToMarkdown"; import { readDigestContext, type DigestContext, } from "../lib/digestContext-server"; import type { Cue } from "../lib/vtt"; import { isCuesJsonFresh, readNormalizedTranscript, } from "./normalizeTranscript"; import { isModelResolutionSuspicious } from "../lib/digest"; 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; // Overrides the digest app's config block (settings.digest.apps[appId]). // Two callers: the LLM fan-out, which passes the primary's resolved config // with only `baseUrl` swapped to the leased endpoint (baseUrl is not part of // the freshness identity, so endpoint choice causes zero churn); and the // unit executor, which passes the primary's injected identity config so a // bare box's default settings can never leak into provenance. appConfig?: DigestAppConfig; // Pre-read channel context (note + hash). A unit executor injects the // primary's, since its scratch corpus has no digest-context.md to read. context?: DigestContext; // 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 { 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, opts.appConfig, ); // The per-video half of the diarized identity. See // AttributionProvenance.diarizationGeneratedAt: a cluster index means nothing // except relative to the run that produced it. // // The transcript half of the identity applies to BOTH lanes: the names were // found in this text, and if it is replaced they are unverified. Free — the // normalized transcript is already in hand. const src = transcriptSourceOf(transcript.source); const fullTarget: AttributionFreshnessTarget = { ...target, ...(diarization?.generatedAt ? { diarizationGeneratedAt: diarization.generatedAt } : {}), ...(src ? { transcriptSource: src } : {}), }; if (!opts.force && isAttributionFresh(existing, fullTarget)) { return "already-exists"; } const context = opts.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"; } // Same tripwire as digestVideo's: freshness compares the REQUESTED model // string, so an endpoint serving the wrong weights would never invalidate // anything — the log line is the only thing that can catch it. if (isModelResolutionSuspicious(modelRequested, reportedModel)) { log( `Attribute ${opts.videoId}: engine reported model "${reportedModel}" for requested "${modelRequested}" — check the serving endpoint has the right weights.`, ); } 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 } : {}), ...(src ? { transcriptSource: src } : {}), ...(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 { 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(); 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(); 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. Lives in lib/attributionTurns.ts so a bake-off harness assembles // a record the same way this does rather than re-implementing it. const roster = createSpeakerRoster(); const labels = roster.labels; const marks: SpeakerMark[] = []; let reportedModel = ""; let costUsd = 0; let chunksOk = 0; // Turns thrown away for landing outside the chunk that produced them, and the // total emitted, so the summary can report a RATE rather than a bare count — // 520 discards means nothing without the 1,705 it is out of. let outOfRangeTotal = 0; let emittedTotal = 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; let outOfRange = 0; emittedTotal += raw.length; 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. COUNTED, because the // drop rate is a first-class quality signal: see AttributionWarning's // "out-of-range". if (at < startSeconds || at > endSeconds) { outOfRange++; continue; } if (isUselessSpeakerLabel(e.speaker)) continue; marks.push({ at, speaker: roster.intern(e.speaker) }); } // One warning per chunk carrying the count, not one per discarded turn: a // 30% drop rate would otherwise bury every other warning in the record. if (outOfRange > 0) { outOfRangeTotal += outOfRange; args.warnings.push({ code: "out-of-range", chunk: i, detail: `${outOfRange} of ${raw.length} turn(s) outside ${startSeconds}-${endSeconds}s`, }); } } 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 (outOfRangeTotal > 0) { const pct = emittedTotal > 0 ? (outOfRangeTotal / emittedTotal) * 100 : 0; args.log( `Attribute ${args.videoId}: discarded ${outOfRangeTotal} of ${emittedTotal} emitted turn(s) (${pct.toFixed(0)}%) for landing outside their own chunk.`, ); } if (labels.length === 0) { return { speakers: [], segments: [], model: reportedModel, costUsd, chunks: chunks.length, chunksOk }; } // Marks -> segments -> speakers, in lib/attributionTurns.ts. See its header: // the assembly is shared so the shipped lane and any harness scoring a prompt // variant cannot produce differently-shaped records from the same marks. const lastCue = args.cues[args.cues.length - 1]; const { speakers, segments } = assembleTurns({ labels, marks, transcriptEnd: lastCue.end || lastCue.start, }); return { speakers, segments, model: reportedModel, costUsd, chunks: chunks.length, chunksOk, }; } // --------------------------------------------------------------------------- function clamp01(n: number): number { if (!Number.isFinite(n)) return 0; return Math.min(1, Math.max(0, n)); } 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)), }, ); }