// What the digest backfill actually costs, measured in CHUNKS — and reported in // audio-hours too, because that is the unit a human thinks in. // // A saving counted in videos is close to meaningless here: the corpus is 77k // videos but 77k audio-HOURS, and the two are not proportional per channel — // mirrors skew long (Hasan VODs, Quartering re-uploads), shorts skew numerous. // So audio-seconds are tracked per video and reported everywhere. // // But audio-hours are not what the GPU is billed in. ONE CHUNK IS ONE MODEL // CALL, and chunk density varies FOURFOLD across the corpus — 6.8 chunks per // audio-hour under 15 minutes against 1.7 over 8 hours. A single // seconds-per-audio-hour figure therefore is not a stable unit, and that unit // error is the whole reason three separate cost measurements looked // irreconcilable (27 vs 90 vs 151 s/audio-hour) when re-pricing them per chunk // reproduces all three to within ~2%. See MEASURED_SECONDS_PER_CHUNK. // // Two consumers, deliberately one implementation: // - `bin/digest-plan.ts`, to price the sweep before committing GPU-weeks to it // and to prove that a cheapness lever (a duplicate threshold change, a batch // of confirmed clusters) actually moved the number; // - the sweep orchestrator, which needs the same remaining-audio-hours figure // as an ETA denominator and the same per-channel ordering as a work queue. // // Source of truth is the `statsByPath` LMDB sub-DB — build:stats' output, keyed // [channelSlug, videoDir], which is the batch runner's own enumeration unit and // carries duration + hasTranscript + cueCount in one scan. It can lag the corpus; // the schema version is checked and a mismatch is reported rather than swallowed. // `cueCount` being on that row is what makes the chunk census free: it needed no // GPU time and no transcript reads to compute. import { open } from "lmdb"; import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; import { orderPlanByRecency } from "./planOrder"; import path from "node:path"; import type { Paths } from "../lib/paths"; import type { VideoStat } from "../lib/stats"; import { STATS_SCHEMA_VERSION } from "../lib/stats"; import { isSectionFresh } from "../lib/digest"; import { loadDigest } from "../lib/digest-server"; import { countCueChunks } from "../lib/transcriptWindow"; import { DIGEST_OVERLAP_CUES } from "../lib/digestPrompt"; import { buildDigestClusterPlan, type DigestClusterPlan, type DigestClusterRole, } from "./digestSharing"; import { resolveDigestChunking, resolveDigestTarget, type DigestLaneChoice, } from "./digestTarget"; // The measured cost of the local lane, PER CHUNK — one model call. // // Re-pricing every measurement taken so far in chunks reconciles all of them, // which the audio-hour model does for none (× 191,116 corpus chunks): // // 11.2 s idle box, ENGINE time (bake-off round 2) → 24.8 days [24.2 ✓] // 24.7 s contended box, WALL time (102-video run) → 54.6 days [81 ✗] // 60.6 s idle box, engine, gemma2 (bake-off round 2) → 134 days [132 ✓] // // The 24.7 is recomputed from the validation run's OWN artifacts, not from its // report: 62.9 min of wall over the 153 chunks its 102 sidecars actually record in // `provenance.chunks`. The old 27 → 90 s/audio-hour "contention penalty" then // factors almost exactly — 1.49× sample length mix (3.63 chunks/audio-hour against // round 2's 2.44) × 2.20× per-chunk cost = 3.27×, against the 3.33× observed. So // the honest range is ~25 days idle to ~55 contended — not 81. // // The DEFAULT is the contended wall figure, deliberately: it is the only one // measured in production, on this corpus, with all the sidecar work a bake-off // skips. It is pessimistic on an idle box by design. The 2.20× residual is NOT // yet attributable to contention: it compares wall against engine time, so it also // contains model-load, prefill, and the yield gate idling at its 3 s poll. Only a // production run can separate those, which is why the optional load/prefill/decode // fields on DigestRunResult exist. export const MEASURED_SECONDS_PER_CHUNK = 24.7; // Corpus-wide chunk density, from the free census over `statsByPath.cueCount` — // all 73,367 transcribed videos, none missing a cue count, at the shipped maxCues // 600 / overlap 40. It cost zero GPU time and no transcript reads: // // band videos audio-h % audio chunks chunks/audio-h % chunks // < 15 min 41,960 6,203 8.0% 41,960 6.76 22.0% // 15–60 min 14,036 6,631 8.6% 22,446 3.39 11.7% // 1–2 h 5,847 8,480 11.0% 22,353 2.64 11.7% // 2–4 h 5,486 15,832 20.5% 33,106 2.09 17.3% // 4–8 h 4,504 25,304 32.7% 46,363 1.83 24.3% // > 8 h 1,534 14,850 19.2% 24,888 1.68 13.0% // total 73,367 77,298 191,116 2.47 // // The counter-intuitive result worth keeping in view: 4 h+ videos are 52% of the // AUDIO but only 37% of the WORK, while sub-hour videos are 16% of audio and 34% // of chunks. Any reasoning that prices the sweep in audio-hours gets this // backwards. // // THE TABLE ABOVE IS ONE RUN, not a live number, and the corpus grows underneath // it: re-measured 2026-08-08 it reads 74,321 videos / 194,053 chunks / density // 2.46. That drift is why the census is `bin/digest-plan.ts --census` and not a // unit test — an absolute count of a live corpus can only ever go red. Re-run it // after a large ingest or ANY change to a DEFAULT_DIGEST_*; it prints the drift // against this constant and says when the constant needs updating. The density // is what the cost model consumes, and it has held to within 0.01 across ~3,000 // new videos, which is the property worth having. export const CORPUS_CHUNKS_PER_AUDIO_HOUR = 2.47; // DERIVED, not measured — kept only so `bin/digest-plan.ts --rate` and anything // else thinking in audio-hours still has a number to start from. It is the // corpus-average conversion and is wrong for any individual channel by up to 2.7× // in either direction. export const MEASURED_SECONDS_PER_AUDIO_HOUR = MEASURED_SECONDS_PER_CHUNK * CORPUS_CHUNKS_PER_AUDIO_HOUR; // Why a video is or is not this sweep's work. The mirror split is the point: a // mirror is only free if its cues actually ALIGN with its canonical member's. // Detection measures that and records it per ref, and the sharing pass refuses // to place a digest without it — so counting all mirrors as free overstates the // saving by however many the alignment gate will later reject. On the current // corpus that is nearly half of them. export type DigestPlanRole = | "canonical" // owns its cluster's digest: must generate | "mirror-aligned" // receives a share: free | "mirror-unaligned" // share will be refused: must generate after all | "unclustered"; // in no cluster: must generate export const DIGEST_PLAN_ROLES: readonly DigestPlanRole[] = [ "canonical", "mirror-aligned", "mirror-unaligned", "unclustered", ]; // A role must generate unless the digest arrives by sharing. export function roleMustGenerate(role: DigestPlanRole): boolean { return role !== "mirror-aligned"; } export type DigestPlanTotals = { videos: number; audioSeconds: number; // Model calls. THIS is what the sweep's wall clock is made of; audioSeconds is // for human comprehension and for comparing against the corpus census. chunks: number; }; function emptyTotals(): DigestPlanTotals { return { videos: 0, audioSeconds: 0, chunks: 0 }; } function addTo(t: DigestPlanTotals, seconds: number, chunks: number): void { t.videos++; t.audioSeconds += seconds; t.chunks += chunks; } export type DigestChannelPlan = { channelSlug: string; // Videos with no fresh digest at the current identity, split by role. remaining: Record; // Already digested at the current identity — the sweep skips these. fresh: DigestPlanTotals; // Indexed but not digestable: no transcript, or no usable duration. ineligible: number; // Convenience rollups over `remaining`. generateSeconds: number; // audio this channel still has to be digested sharedSeconds: number; // audio cluster sharing takes off the bill generateChunks: number; // what this channel actually costs the GPU sharedChunks: number; // model calls cluster sharing avoids // Videos counted as remaining work whose cueCount is missing from the stats // cache, i.e. whose chunk count had to be estimated from duration. Surfaced // because it is the one way the chunk total can be wrong, and it must not be // invisible. chunksEstimated: number; // The newest and oldest upload dates (YYYYMMDD) among the videos this channel // will actually GENERATE — not among everything remaining, because a video a // duplicate cluster will share rather than generate costs the queue nothing // and must not pull its channel up the plan. // // "" means no dated pending video: either the channel has nothing to generate // or build:stats has no date for what it has. It sorts LAST under "newest" // and FIRST under "oldest" — the same rule makeRecencyComparator applies to // an individual video, so the two levels cannot disagree, and a corpus with // no dates at all reproduces the historical plan exactly. newestPending: string; oldestPending: string; }; export type DigestSweepPlan = { // Ordered by generateSeconds descending — the sweep's work queue. Channels // with nothing to do are still present, with zeroed totals, so a caller can // tell "done" apart from "not in the corpus". channels: DigestChannelPlan[]; remaining: Record; fresh: DigestPlanTotals; ineligible: number; generateSeconds: number; sharedSeconds: number; generateChunks: number; sharedChunks: number; chunksEstimated: number; // Cues per chunk the census was computed at. Recorded because the chunk total // is only meaningful against a context size: halving numCtx roughly doubles it. maxCues: number; clusters: number; clusterMembersMapped: number; // True when build:stats' cache is not at the version this code expects, i.e. // every number below may be stale. Reported, never silently tolerated. statsSchemaStale: boolean; // Set when freshness was NOT consulted, so `fresh` is 0 by construction and // `remaining` is the cost of a sweep from scratch. freshnessChecked: boolean; }; export type BuildDigestSweepPlanOptions = { paths: Paths; lane?: DigestLaneChoice; // Off makes the scan pure-LMDB and near-instant, at the cost of reporting the // from-scratch cost rather than the remaining one. At 0.13% coverage the two // are nearly the same number; mid-sweep they are not. checkFreshness?: boolean; // Restrict to these channels (the orchestrator uses it to re-price one). channelSlugs?: string[]; clusterPlan?: DigestClusterPlan; // Order the CHANNEL list by each channel's freshest (or oldest) pending // video instead of heaviest-first. Absent/"listed" leaves today's // chunks-descending order untouched, which is what the census wants — this is // a dispatch decision, not a costing one. order?: AutoQueueOrder; onLog?: (msg: string) => void; }; function emptyRoleTotals(): Record { return { canonical: emptyTotals(), "mirror-aligned": emptyTotals(), "mirror-unaligned": emptyTotals(), unclustered: emptyTotals(), }; } function mergeInto( into: Record, from: Record, ): void { for (const role of DIGEST_PLAN_ROLES) { into[role].videos += from[role].videos; into[role].audioSeconds += from[role].audioSeconds; into[role].chunks += from[role].chunks; } } // Chunks for one video, from the cue count build:stats already recorded. // // `cueCount` is null only for a row written before stats carried it. Rather than // drop such a video from the bill (which would understate the sweep), fall back to // the census density — and the caller counts how often that happened, so an // estimate can never quietly become the headline. export function chunksForVideo( stat: Pick, maxCues: number, ): { chunks: number; estimated: boolean } { if (typeof stat.cueCount === "number" && stat.cueCount >= 0) { return { chunks: countCueChunks(stat.cueCount, { maxCues, overlapCues: DIGEST_OVERLAP_CUES, }), estimated: false, }; } const hours = Math.max(0, stat.duration) / 3600; return { chunks: Math.max(1, Math.round(hours * CORPUS_CHUNKS_PER_AUDIO_HOUR)), estimated: true, }; } // --------------------------------------------------------------------------- // The chunk census // --------------------------------------------------------------------------- // The duration bands the census is reported in. They exist because the headline // result is a SHAPE, not a total: chunk density falls monotonically from 6.8 to // 1.7 chunks per audio-hour as videos get longer, so 4 h+ videos are 52% of the // audio but only 37% of the work. A single total hides exactly the thing that // makes seconds-per-audio-hour an invalid unit. export const CENSUS_BANDS: readonly { label: string; maxSeconds: number }[] = [ { label: "< 15 min", maxSeconds: 15 * 60 }, { label: "15–60 min", maxSeconds: 3600 }, { label: "1–2 h", maxSeconds: 2 * 3600 }, { label: "2–4 h", maxSeconds: 4 * 3600 }, { label: "4–8 h", maxSeconds: 8 * 3600 }, { label: "> 8 h", maxSeconds: Infinity }, ]; export type CensusBand = { label: string; videos: number; audioSeconds: number; chunks: number; }; export type ChunkCensus = { bands: CensusBand[]; videos: number; audioSeconds: number; chunks: number; // Rows with no recorded cueCount, whose chunks were ESTIMATED from duration. // Reported rather than silently absorbed: the census's whole claim to being // free is that cueCount was already on every row, and a non-zero here means // part of the headline is a guess. estimated: number; chunksPerAudioHour: number; }; // Bucket transcribed videos into the bands above and total their chunks. // // PURE, over an iterable of stats, and that is the point. This used to be // inlined in a TEST that opened the operator's real LMDB and asserted the corpus // totalled exactly 191,116 chunks across 73,367 videos. That test could only run // on one machine, and it failed there the moment a video was downloaded — which // is not a regression, it is the corpus growing. An absolute count of a live, // changing corpus is a MEASUREMENT, not an invariant, and pinning it as a test // meant the suite reported a permanent red for doing its job correctly. // // The two things that test was actually for are both kept, split apart: // - "the plan and the chunker have drifted" is a pure property, asserted // hermetically against this function and against countCueChunks. // - "a DEFAULT_DIGEST_* changed the shipped chunk size and nobody re-priced // the sweep" is pinned on the DENSITY constant, which is what the cost model // consumes, rather than on a video count that moves on its own. // The live numbers are printed by `bin/digest-plan.ts --census`, where a // measurement belongs. export function buildChunkCensus( rows: Iterable>, maxCues: number, ): ChunkCensus { const bands: CensusBand[] = CENSUS_BANDS.map((b) => ({ label: b.label, videos: 0, audioSeconds: 0, chunks: 0, })); let videos = 0; let audioSeconds = 0; let chunks = 0; let estimated = 0; for (const stat of rows) { // The same eligibility the plan uses: a video with no transcript is not work // and a zero duration is a row we cannot band. if (!stat.hasTranscript || !(stat.duration > 0)) continue; const index = CENSUS_BANDS.findIndex((b) => stat.duration < b.maxSeconds); const band = bands[index === -1 ? bands.length - 1 : index]; const forVideo = chunksForVideo(stat, maxCues); videos++; audioSeconds += stat.duration; chunks += forVideo.chunks; if (forVideo.estimated) estimated++; band.videos++; band.audioSeconds += stat.duration; band.chunks += forVideo.chunks; } return { bands, videos, audioSeconds, chunks, estimated, chunksPerAudioHour: chunksPerAudioHour(chunks, audioSeconds), }; } // Run the census over the real corpus. The only impure part, and it does nothing // but feed rows to the function above. // // Returns null when there is no stats cache on this machine, so the caller can // say so rather than reporting a corpus of zero. export async function readChunkCensus( paths: Paths, maxCues: number, ): Promise<(ChunkCensus & { statsSchemaStale: boolean }) | null> { const { existsSync } = await import("node:fs"); if (!existsSync(paths.lmdbPath)) return null; const root = open({ path: paths.lmdbPath, maxDbs: 12, compression: true }); const statsByPath = root.openDB< { metaMs: number; stat: VideoStat }, [string, string] >({ name: "statsByPath", encoding: "msgpack" }); const meta = root.openDB({ name: "statsMeta", encoding: "msgpack", }); const statsSchemaStale = (meta.get("schema") as number | undefined) !== STATS_SCHEMA_VERSION; function* rows(): Generator { for (const { value } of statsByPath.getRange()) yield value.stat; } return { ...buildChunkCensus(rows(), maxCues), statsSchemaStale }; } // The whole role decision, as one pure function so it can be asserted on // without an LMDB corpus behind it. // // `aligned` comes from the report's per-ref measurement and is deliberately // tri-state at the call site: absent means NOT MEASURED, which the sharing pass // treats as not aligned. So anything other than an explicit `true` lands in // mirror-unaligned and the plan under-promises rather than over-promises. export function classifyDigestRole( clusterRole: DigestClusterRole | undefined, aligned: boolean | undefined, ): DigestPlanRole { if (clusterRole?.kind === "canonical") return "canonical"; if (clusterRole?.kind === "mirror") { return aligned === true ? "mirror-aligned" : "mirror-unaligned"; } return "unclustered"; } export async function buildDigestSweepPlan( opts: BuildDigestSweepPlanOptions, ): Promise { const log = opts.onLog ?? (() => {}); const checkFreshness = opts.checkFreshness !== false; const clusterPlan = opts.clusterPlan ?? (await buildDigestClusterPlan(opts.paths)); // `aligned` lives on the report's refs, not on the derived plan, so read it // straight from the report rather than widening DigestClusterRole for one // consumer. const alignedBySlug = await readAlignmentBySlug(opts.paths); const root = open({ path: opts.paths.lmdbPath, maxDbs: 12, compression: true, }); const statsByPath = root.openDB< { metaMs: number; stat: VideoStat }, [string, string] >({ name: "statsByPath", encoding: "msgpack" }); const meta = root.openDB({ name: "statsMeta", encoding: "msgpack", }); const statsSchemaStale = (meta.get("schema") as number | undefined) !== STATS_SCHEMA_VERSION; if (statsSchemaStale) { log( `Warning: stats cache schema ${meta.get("schema") ?? ""} != ${STATS_SCHEMA_VERSION}; run build:stats for accurate audio-hours.`, ); } const wanted = opts.channelSlugs ? new Set(opts.channelSlugs) : null; // Group by channel first: the freshness target is resolved ONCE per channel // (it reads settings + the channel-context note), and resolving it per video // would dominate the scan. type Row = { videoDir: string; stat: VideoStat }; const byChannel = new Map(); for (const { key, value } of statsByPath.getRange()) { const [channelSlug, videoDir] = key as [string, string]; if (wanted && !wanted.has(channelSlug)) continue; let rows = byChannel.get(channelSlug); if (!rows) byChannel.set(channelSlug, (rows = [])); rows.push({ videoDir, stat: value.stat }); } // Chunk size for the lane being priced. Resolved ONCE: it depends only on the // app's numCtx, so re-deriving it per channel would just be slower. const { maxCues } = resolveDigestChunking({ lane: opts.lane }); const channels: DigestChannelPlan[] = []; for (const [channelSlug, rows] of byChannel) { const entry: DigestChannelPlan = { channelSlug, remaining: emptyRoleTotals(), fresh: emptyTotals(), ineligible: 0, generateSeconds: 0, sharedSeconds: 0, generateChunks: 0, sharedChunks: 0, chunksEstimated: 0, newestPending: "", oldestPending: "", }; const resolved = checkFreshness ? await resolveDigestTarget({ paths: opts.paths, channelSlug, lane: opts.lane, }) : null; const dataDir = path.join(opts.paths.channelsDir, channelSlug, "data"); for (const { videoDir, stat } of rows) { // Exactly the batch runner's eligibility: a video with no transcript is a // transcription problem, not a digest one, and it is not this sweep's work. if (!stat.hasTranscript || !(stat.duration > 0)) { entry.ineligible++; continue; } const { chunks, estimated } = chunksForVideo(stat, maxCues); if (resolved) { const record = await loadDigest(path.join(dataDir, videoDir)); const allFresh = resolved.sections.every((section) => isSectionFresh(record, section, resolved.target), ); if (allFresh) { addTo(entry.fresh, stat.duration, chunks); continue; } } const planRole = classifyDigestRole( clusterPlan.bySlug.get(stat.slug), alignedBySlug.get(stat.slug), ); addTo(entry.remaining[planRole], stat.duration, chunks); if (estimated) entry.chunksEstimated++; // Free: `stat` is already in hand and `uploadDate` is already on it, so // cross-channel ordering costs this comparison and no extra I/O. Only // rows that will actually be generated count — see newestPending. if (roleMustGenerate(planRole) && stat.uploadDate) { if (stat.uploadDate > entry.newestPending) { entry.newestPending = stat.uploadDate; } if (!entry.oldestPending || stat.uploadDate < entry.oldestPending) { entry.oldestPending = stat.uploadDate; } } } for (const role of DIGEST_PLAN_ROLES) { if (roleMustGenerate(role)) { entry.generateSeconds += entry.remaining[role].audioSeconds; entry.generateChunks += entry.remaining[role].chunks; } else { entry.sharedSeconds += entry.remaining[role].audioSeconds; entry.sharedChunks += entry.remaining[role].chunks; } } channels.push(entry); } // Ordered by CHUNKS, not audio-hours: chunks are what the queue actually spends. // The two orders differ — long-VOD channels have the lowest chunk density — and // ordering by audio would put the cheapest-per-hour work first while claiming to // be heaviest-first. (It is still longest-first in practice, so early throughput // will look better than the corpus average; see the census table above.) // // With `order` set, recency leads and the chunk order becomes the tiebreak — // so two channels whose freshest pending video landed the same day are still // visited heaviest-first, and a corpus where nothing is dated (every key "") // sorts exactly as it does today. orderPlanByRecency( channels, opts.order ?? "listed", (a, b) => b.generateChunks - a.generateChunks || b.generateSeconds - a.generateSeconds || a.channelSlug.localeCompare(b.channelSlug), ); const totals: DigestSweepPlan = { channels, remaining: emptyRoleTotals(), fresh: emptyTotals(), ineligible: 0, generateSeconds: 0, sharedSeconds: 0, generateChunks: 0, sharedChunks: 0, chunksEstimated: 0, maxCues, clusters: clusterPlan.clusters, clusterMembersMapped: clusterPlan.bySlug.size, statsSchemaStale, freshnessChecked: checkFreshness, }; for (const c of channels) { mergeInto(totals.remaining, c.remaining); totals.fresh.videos += c.fresh.videos; totals.fresh.audioSeconds += c.fresh.audioSeconds; totals.fresh.chunks += c.fresh.chunks; totals.ineligible += c.ineligible; totals.generateSeconds += c.generateSeconds; totals.sharedSeconds += c.sharedSeconds; totals.generateChunks += c.generateChunks; totals.sharedChunks += c.sharedChunks; totals.chunksEstimated += c.chunksEstimated; } return totals; } // Per-ref `aligned` from the last detection run, keyed by slug. Kept separate // from DigestClusterPlan because only the pricing path needs it. async function readAlignmentBySlug( paths: Paths, ): Promise> { const { readDuplicateReport } = await import("./duplicateShorts"); const report = await readDuplicateReport(paths); const map = new Map(); if (!report) return map; for (const cluster of report.clusters) { for (const ref of cluster.videoRefs) { if (ref.aligned === true) map.set(ref.slug, true); } } return map; } export function audioHours(seconds: number): number { return seconds / 3600; } // The projection the whole plan exists to produce, priced in the unit the GPU is // actually billed in. export function sweepDays( chunks: number, secondsPerChunk: number = MEASURED_SECONDS_PER_CHUNK, ): number { return (chunks * secondsPerChunk) / 86400; } // The same projection from audio-hours, for the one caller that thinks in them // (`bin/digest-plan.ts --rate`). Kept because a rate in s/audio-hour is what every // past measurement was quoted in and a reader will want to reproduce them — but it // is corpus-average only. Prefer sweepDays(). export function sweepDaysFromAudio( audioSeconds: number, secondsPerAudioHour: number = MEASURED_SECONDS_PER_AUDIO_HOUR, ): number { return (audioHours(audioSeconds) * secondsPerAudioHour) / 86400; } // Chunks per audio-hour for an arbitrary slice of the plan — the density figure // that makes a per-chunk cost comparable to the census. Report THIS beside a // throughput number, never a bare s/audio-hour: mid-sweep the queue is // longest-first and its density drifts from 1.7 toward 6.8 as it drains, so a // running s/audio-hour will look like the projection is falling apart when it is // only reaching shorter videos. export function chunksPerAudioHour( chunks: number, audioSeconds: number, ): number { const h = audioHours(audioSeconds); return h > 0 ? chunks / h : 0; }