// Generate the AI digest for ONE video: chapters and/or topic tags, written to // the ai-digest.json sidecar next to the transcript. // // Pipeline (each step reusing the declared single source of truth for its job): // isCuesJsonFresh — don't digest a transcript that's about to change // readNormalizedTranscript // chunkCuesForContext — sequential overlapping slices // transcriptToMarkdown — cues -> text for an AI, with hms() stamps // digest app — ollama (local) or claude CLI (metered, opt-in) // parseChapters/parseTags — the guards; every rejection recorded // writeDigestSection — read-modify-write ONE section, atomic rename // // THE FRESHNESS SKIP IS THE POINT. A full sweep of the 74k-transcript corpus is // weeks of wall-clock on one local lane, so a redo is unaffordable. A section is // regenerated only when its recorded (schemaVersion, appId, model, promptVersion, // contextHash) differs from what we would produce now — which makes a re-run // after a prompt change minutes long over a sample instead of weeks over the // corpus. import path from "node:path"; import type { Paths } from "../lib/paths"; import { DIGEST_OVERLAP_CUES, PROMPT_VERSION, CHAPTER_SYSTEM_PROMPT, TAG_SYSTEM_PROMPT, buildChapterPrompt, buildTagPrompt, chapterSchema, maxCuesForContext, tagSchema, toHms, } from "../lib/digestPrompt"; import { getDigestApp, type DigestAppConfig } from "../lib/digestApps"; import { digestPromptVariant, isModelResolutionSuspicious, isSectionFresh, type DigestItem, type DigestTimestampMode, type DigestProvenance, type DigestRecord, type DigestSectionKind, type DigestWarning, } from "../lib/digest"; import { loadDigest, writeDigestFailure, writeDigestSection, } from "../lib/digest-server"; import { readDigestContext, type DigestContext } from "../lib/digestContext-server"; import { parseChapters, parseTags, type DigestChunkOutput } from "../lib/digestParse"; import { chunkCuesForContext } from "../lib/transcriptWindow"; import { transcriptToMarkdown } from "../lib/transcriptToMarkdown"; import type { Cue } from "../lib/vtt"; import { isCuesJsonFresh, readNormalizedTranscript, } from "./normalizeTranscript"; export type DigestVideoOptions = { paths: Paths; channelSlug: string; videoId: string; // Which sections to generate. Defaults to chapters only — tags are cheap but // double the call count, so the sweep opts into them explicitly. sections?: DigestSectionKind[]; appId?: string; config?: DigestAppConfig; // Prompt SHAPE. Both default, so an existing caller keeps producing records // with no promptVariant and every digest already on disk stays fresh. timestampMode?: DigestTimestampMode; promptVariant?: string; // Pre-read channel context, so a batch reads it once per channel instead of // once per video. Omitted → read here. context?: DigestContext; // Regenerate even when the recorded provenance matches. For Stage B iteration // on a fixed sample; never set for a sweep. force?: boolean; onLog?: (msg: string) => void; signal?: AbortSignal; }; export type DigestVideoOutcome = | { status: "wrote"; record: DigestRecord; sections: DigestSectionKind[]; itemCount: number; warningCount: number; costUsd: number; engineCalls: number; } | { status: "fresh" } | { status: "skipped"; reason: "no-transcript" | "stale-cues" | "no-cues" | "no-metadata"; }; // A per-chunk engine failure must not lose the chunks that DID work: a 12-hour // video is 30+ calls and one 500 from ollama should cost that chunk, not the // video. The failure is recorded as a warning and the section is written from // what survived, with chunks/chunksOk in provenance so a partially-generated // section is identifiable later. export async function digestVideo( opts: DigestVideoOptions, ): Promise { const log = opts.onLog ?? (() => {}); const videoDir = path.join( opts.paths.channelsDir, opts.channelSlug, "data", opts.videoId, ); const sections = opts.sections ?? ["chapters"]; // Don't digest a transcript that is about to be rewritten: a stale cues.json // means the raw transcript changed under it, so the digest would describe // superseded text and then look "fresh" forever. const { fresh: cuesFresh, cuesPath } = await isCuesJsonFresh(videoDir); if (!cuesFresh) { const existing = await readNormalizedTranscript(cuesPath); return { status: "skipped", reason: existing ? "stale-cues" : "no-transcript", }; } const transcript = await readNormalizedTranscript(cuesPath); if (!transcript) return { status: "skipped", reason: "no-transcript" }; const cues = transcript.cues ?? []; if (cues.length === 0) return { status: "skipped", reason: "no-cues" }; const app = getDigestApp(opts.appId); const config = opts.config ?? {}; const modelRequested = config.model?.trim() || app.defaultModel(); const context = opts.context ?? (await readDigestContext(opts.paths, opts.channelSlug)); const timestampMode = opts.timestampMode ?? "absolute"; // Sized to the CONFIGURED context, not to a constant. A 8192-token window with // 1200-cue chunks overflows and ollama truncates without saying so. const maxCues = maxCuesForContext(config.numCtx); // One derivation, shared with the batch and the bake-off, so the same config // never produces two different identities. const promptVariant = digestPromptVariant({ promptVariant: opts.promptVariant, timestampMode, maxCues, }); const existing = await loadDigest(videoDir); const target = { appId: app.id, model: modelRequested, promptVersion: PROMPT_VERSION, contextHash: context.hash, ...(promptVariant ? { promptVariant } : {}), }; const stale = sections.filter( (section) => opts.force || !isSectionFresh(existing, section, target), ); if (stale.length === 0) return { status: "fresh" }; const chunks = chunkCuesForContext(cues, { maxCues, overlapCues: DIGEST_OVERLAP_CUES, }); log( `${opts.channelSlug}/${opts.videoId}: ${cues.length} cues → ${chunks.length} chunk(s), sections: ${stale.join(", ")}`, ); let record: DigestRecord | null = existing; let itemCount = 0; let warningCount = 0; let costUsd = 0; let engineCalls = 0; for (const section of stale) { const outputs: DigestChunkOutput[] = []; const runWarnings: DigestWarning[] = []; let reportedModel = modelRequested; let sectionCost = 0; for (let i = 0; i < chunks.length; i++) { opts.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), ); const promptInput = { title: transcript.title || opts.videoId, channel: transcript.channel || opts.channelSlug, startSeconds, endSeconds, // The renderer and the prompt MUST agree on the numbering, so both are // driven from the same value rather than each deciding for itself. transcript: renderChunk( transcript, chunk, timestampMode === "chunk-local" ? startSeconds : 0, ), timestampMode, ...(context.note ? { contextNote: context.note } : {}), }; const span = endSeconds - startSeconds; const request = section === "chapters" ? { system: CHAPTER_SYSTEM_PROMPT, prompt: buildChapterPrompt(promptInput), schema: chapterSchema(span), } : { system: TAG_SYSTEM_PROMPT, prompt: buildTagPrompt(promptInput), schema: tagSchema(), }; try { const result = await app.run({ ...request, config, signal: opts.signal, onLog: opts.onLog, }); engineCalls++; reportedModel = result.model || reportedModel; if (typeof result.costUsd === "number") sectionCost += result.costUsd; outputs.push({ index: i, startSeconds, endSeconds, data: result.data, timestampMode, }); } catch (err) { // A cancel is not a chunk failure — let it propagate so the batch stops. if (opts.signal?.aborted) throw err; const message = (err as Error)?.message ?? String(err); runWarnings.push({ code: "chunk-failed", section, chunk: i, detail: message.slice(0, 300), }); log( `${opts.channelSlug}/${opts.videoId}: chunk ${i + 1}/${chunks.length} failed: ${message}`, ); } } // A tag completion ("qwen2.5" → "qwen2.5:7b") is normal; a different model // is an endpoint serving the wrong weights. Freshness compares the // REQUESTED string, so this would never invalidate anything — the loud log // line is the only tripwire, which is why it is here and not optional. if (isModelResolutionSuspicious(modelRequested, reportedModel)) { log( `${opts.channelSlug}/${opts.videoId}: engine reported model "${reportedModel}" for requested "${modelRequested}" — check the serving endpoint has the right weights.`, ); } if (outputs.length === 0) { // Nothing usable. Do NOT write a section — an empty section with current // provenance would read as "fresh" and the video would never be retried. // But DO persist the warnings, in the sibling `failures` field that // freshness never reads: this is the worst outcome the generator has, and // it used to leave no trace anywhere except a job log that rotates. log( `${opts.channelSlug}/${opts.videoId}: ${section} produced no usable output (${runWarnings.length} warning(s)); leaving the sidecar untouched so it retries.`, ); record = await writeDigestFailure(videoDir, { section, at: new Date().toISOString(), appId: app.id, model: reportedModel, promptVersion: PROMPT_VERSION, reason: "no-output", chunks: chunks.length, chunksOk: outputs.length, warnings: runWarnings, }); warningCount += runWarnings.length; continue; } const parsed = section === "chapters" ? parseChapters(outputs, cues) : parseTags(outputs); const items: DigestItem[] = "chapters" in parsed ? parsed.chapters : parsed.tags; const warnings = [...runWarnings, ...parsed.warnings]; if (items.length === 0) { // The model DID propose content and every item failed a guard — a // different failure from "produced nothing", and one the warnings can // actually explain (the validation run's worst video had 13 chapters // clamped away as out-of-range). Recording which of the two happened is // the difference between a review queue that can act and one that can only // report a blank. log( `${opts.channelSlug}/${opts.videoId}: every ${section} entry was rejected (${warnings.length} warning(s)); leaving the sidecar untouched so it retries.`, ); record = await writeDigestFailure(videoDir, { section, at: new Date().toISOString(), appId: app.id, model: reportedModel, promptVersion: PROMPT_VERSION, reason: "all-rejected", chunks: chunks.length, chunksOk: outputs.length, warnings, }); warningCount += warnings.length; continue; } const provenance: DigestProvenance = { appId: app.id, model: reportedModel, modelRequested, lane: app.lane, generatedAt: new Date().toISOString(), promptVersion: PROMPT_VERSION, ...(promptVariant ? { promptVariant } : {}), contextHash: context.hash, chunks: chunks.length, chunksOk: outputs.length, ...(app.metered && sectionCost > 0 ? { costUsd: sectionCost } : {}), }; record = await writeDigestSection(videoDir, { section, items, provenance, warnings, }); itemCount += items.length; warningCount += warnings.length; costUsd += sectionCost; log( `${opts.channelSlug}/${opts.videoId}: ${section} → ${items.length} item(s), ${warnings.length} warning(s)` + (sectionCost > 0 ? `, $${sectionCost.toFixed(4)}` : ""), ); } if (!record || itemCount === 0) { // Every requested section failed. Report it as a skip rather than a write so // the batch's counters stay honest. return { status: "skipped", reason: "no-cues" }; } return { status: "wrote", record, sections: stale, itemCount, warningCount, costUsd, engineCalls, }; } // Render one chunk as the text the engine sees. transcriptToMarkdown is the // declared single source of truth for "transcript -> text for an AI"; the // stampForCue override makes every line carry a HH:MM:SS marker, which is what // the prompt tells the model to copy its `start` values from. // // toHms — NOT aiHandoff's hms() — because hms drops the hour field below an hour // ("2:36") and abbreviates it above one ("1:00:00"), while the schema pattern // requires two digits in all three fields. Feeding the model markers it cannot // legally echo would reintroduce the malformed-timestamp failure the pin fixed. // // The description is omitted: it is the uploader's own promotional copy and // biases titles toward it. // // `offsetSeconds` is subtracted from every marker. It is 0 in absolute mode and // the chunk's own start in chunk-local mode, which is the whole of what // "re-basing" means — the cue list itself is never modified, only how it is // stamped, so the parser's cue snap still works against real cue times. function renderChunk( transcript: { id: string; title: string; channel?: string; duration?: number }, cues: Cue[], offsetSeconds = 0, ): 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 - offsetSeconds)), }, ); }