commit 4826443681a9fd288e55c2e459b3f089ecd0daab
parent 4737bb2a3cc663585d56e4e920a9ebc7a33a1f07
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 29 Jul 2026 11:20:03 -0400
Price the digest backfill in audio-hours, and find out the levers are small
`digest-plan.ts` builds the sweep's work-list from the duplicate report and
build:stats' cache, and reports what remains in AUDIO-HOURS split by role —
canonical (must generate), mirror (shared free), unclustered — plus the day
projection at the measured 90 s/audio-hour. Audio-hours because a saving
counted in videos says nothing: the sweep's wall clock is audio-hours times
seconds-per-audio-hour, and the two are not proportional per channel.
It reproduces the roadmap's headline exactly — 77,298 digestable audio-hours,
80.0 sweep days — which is what makes the rest of its output trustworthy.
The mirror split is deliberately pessimistic. A mirror is only free if its cues
ALIGN with its canonical member's; detection measures that per ref and the
sharing pass refuses to place a digest without it, so `mirror-unaligned` counts
as work. On the current corpus that is nearly half of all mirrors, and counting
them as free would overstate the saving by that much.
What it immediately showed, and the reason it exists: **cluster sharing is not a
meaningful cost lever.** Cluster members are 19% of the corpus by video count
but 4.3% of its audio-hours. The plan assumed mirrors skew long; they skew
short. Sharing is worth ~1.7 sweep days of 80, not the "~11%" digestSharing.ts
claims — the sweep is dominated by unclustered long-form VODs, one of which
(HasanAbiVODs3, 8,331 h) outweighs every mirror in the corpus combined.
The plan object is also the sweep orchestrator's work queue and ETA
denominator, which is why the logic is a library and this is a thin CLI.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Diffstat:
2 files changed, 455 insertions(+), 0 deletions(-)
diff --git a/common/bin/digest-plan.ts b/common/bin/digest-plan.ts
@@ -0,0 +1,136 @@
+#!/usr/bin/env tsx
+// Price the digest backfill, in audio-hours, before spending GPU-weeks on it.
+//
+// Reads the last duplicate-detection run and build:stats' cache; writes nothing.
+// Run it before and after a cheapness lever (a --near threshold change, a batch
+// of confirmed clusters) and diff `generate`: that number, not a video count, is
+// what the sweep's wall clock is made of.
+//
+// Flags:
+// --lane local|remote which engine identity to test freshness against
+// --rate N seconds per audio-hour (default 90, the measured local
+// rate from the 102-video validation run)
+// --no-freshness skip the per-video sidecar read; reports the
+// from-scratch cost instead of the remaining one
+// --channels a,b restrict to these channel slugs
+// --top N how many channels to list (default 15, 0 = all)
+// --json machine-readable, for diffing two runs
+import { getPaths } from "../lib/paths";
+import {
+ buildDigestSweepPlan,
+ audioHours,
+ sweepDays,
+ DIGEST_PLAN_ROLES,
+ MEASURED_SECONDS_PER_AUDIO_HOUR,
+ roleMustGenerate,
+} from "../controller/digestPlan";
+import { parseFlags } from "./_parseFlags";
+
+const flags = parseFlags(process.argv.slice(2));
+
+const lane = flags.lane === "remote" ? "remote" : "local";
+const rate =
+ flags.rate !== undefined
+ ? Number(flags.rate)
+ : MEASURED_SECONDS_PER_AUDIO_HOUR;
+const top = flags.top !== undefined ? Number(flags.top) : 15;
+const asJson = flags.json === "true";
+
+function hours(seconds: number): string {
+ return audioHours(seconds).toLocaleString("en-US", {
+ maximumFractionDigits: 0,
+ });
+}
+
+function pct(part: number, whole: number): string {
+ return whole > 0 ? `${((part / whole) * 100).toFixed(1)}%` : "—";
+}
+
+async function main(): Promise<void> {
+ const plan = await buildDigestSweepPlan({
+ paths: getPaths(),
+ lane,
+ checkFreshness: flags["no-freshness"] !== "true",
+ channelSlugs: flags.channels ? flags.channels.split(",") : undefined,
+ onLog: (m) => {
+ if (!asJson) console.log(m);
+ },
+ });
+
+ if (asJson) {
+ console.log(
+ JSON.stringify(
+ {
+ ...plan,
+ secondsPerAudioHour: rate,
+ generateAudioHours: audioHours(plan.generateSeconds),
+ sharedAudioHours: audioHours(plan.sharedSeconds),
+ sweepDays: sweepDays(plan.generateSeconds, rate),
+ },
+ null,
+ 2,
+ ),
+ );
+ return;
+ }
+
+ const eligible =
+ plan.generateSeconds + plan.sharedSeconds + plan.fresh.audioSeconds;
+
+ console.log("");
+ console.log(
+ `Digest sweep plan — ${lane} lane, ${rate}s per audio-hour` +
+ (plan.freshnessChecked ? "" : ", FROM SCRATCH (freshness not checked)"),
+ );
+ console.log(
+ `Duplicate report: ${plan.clusters.toLocaleString()} cluster(s), ${plan.clusterMembersMapped.toLocaleString()} member(s) mapped.`,
+ );
+ if (plan.statsSchemaStale)
+ console.log("!! stats cache is stale — run build:stats.");
+ console.log("");
+
+ console.log("Remaining work by role (audio-hours):");
+ for (const role of DIGEST_PLAN_ROLES) {
+ const t = plan.remaining[role];
+ console.log(
+ ` ${role.padEnd(17)} ${hours(t.audioSeconds).padStart(8)} h ` +
+ `${t.videos.toLocaleString().padStart(7)} videos ` +
+ `${roleMustGenerate(role) ? "GENERATE" : "shared free"}`,
+ );
+ }
+ console.log(
+ ` ${"already fresh".padEnd(17)} ${hours(plan.fresh.audioSeconds).padStart(8)} h ` +
+ `${plan.fresh.videos.toLocaleString().padStart(7)} videos skipped`,
+ );
+ console.log(
+ ` ${"ineligible".padEnd(17)} ${"—".padStart(8)} ${plan.ineligible.toLocaleString().padStart(7)} videos no transcript`,
+ );
+ console.log("");
+
+ console.log(
+ `TO GENERATE : ${hours(plan.generateSeconds)} audio-hours → ${sweepDays(plan.generateSeconds, rate).toFixed(1)} sweep days`,
+ );
+ console.log(
+ `SAVED by sharing: ${hours(plan.sharedSeconds)} audio-hours (${pct(plan.sharedSeconds, eligible)} of eligible) → ${sweepDays(plan.sharedSeconds, rate).toFixed(1)} days avoided`,
+ );
+ console.log("");
+
+ const listed = top > 0 ? plan.channels.slice(0, top) : plan.channels;
+ console.log(`Heaviest channels (the sweep's work queue order):`);
+ for (const c of listed) {
+ if (c.generateSeconds <= 0) continue;
+ console.log(
+ ` ${c.channelSlug.padEnd(28)} ${hours(c.generateSeconds).padStart(7)} h generate ` +
+ `${hours(c.sharedSeconds).padStart(6)} h shared ` +
+ `${sweepDays(c.generateSeconds, rate).toFixed(1).padStart(5)} d`,
+ );
+ }
+ const rest = plan.channels.length - listed.length;
+ if (rest > 0) console.log(` … and ${rest} more channel(s).`);
+ console.log("");
+}
+
+main().catch((err) => {
+ console.error(err);
+ process.exit(1);
+});
diff --git a/common/controller/digestPlan.ts b/common/controller/digestPlan.ts
@@ -0,0 +1,319 @@
+// What the digest backfill actually costs, measured in AUDIO-HOURS.
+//
+// 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.
+// The sweep's wall clock is (audio-hours × seconds-per-audio-hour), so a plan
+// that moves 4,000 videos and 200 audio-hours has moved nothing. Everything
+// here is therefore denominated in seconds of audio.
+//
+// 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 in one scan. It can lag the corpus; the
+// schema version is checked and a mismatch is reported rather than swallowed.
+
+import { open } from "lmdb";
+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 {
+ buildDigestClusterPlan,
+ type DigestClusterPlan,
+} from "./digestSharing";
+import { resolveDigestTarget, type DigestLaneChoice } from "./digestTarget";
+
+// The measured cost of the local lane, from the 102-video validation run on the
+// real corpus (NOT the bake-off's 27 s, which was projected from long videos on
+// an idle box — see plans/FACTS.md). Overridable, because it is the one number
+// here that is an estimate rather than a measurement of this corpus.
+export const MEASURED_SECONDS_PER_AUDIO_HOUR = 90;
+
+// 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;
+};
+
+function emptyTotals(): DigestPlanTotals {
+ return { videos: 0, audioSeconds: 0 };
+}
+
+function addTo(t: DigestPlanTotals, seconds: number): void {
+ t.videos++;
+ t.audioSeconds += seconds;
+}
+
+export type DigestChannelPlan = {
+ channelSlug: string;
+ // Videos with no fresh digest at the current identity, split by role.
+ remaining: Record<DigestPlanRole, DigestPlanTotals>;
+ // 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; // what this channel costs the GPU
+ sharedSeconds: number; // what cluster sharing takes off the bill
+};
+
+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<DigestPlanRole, DigestPlanTotals>;
+ fresh: DigestPlanTotals;
+ ineligible: number;
+ generateSeconds: number;
+ sharedSeconds: 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;
+ onLog?: (msg: string) => void;
+};
+
+function emptyRoleTotals(): Record<DigestPlanRole, DigestPlanTotals> {
+ return {
+ canonical: emptyTotals(),
+ "mirror-aligned": emptyTotals(),
+ "mirror-unaligned": emptyTotals(),
+ unclustered: emptyTotals(),
+ };
+}
+
+function mergeInto(
+ into: Record<DigestPlanRole, DigestPlanTotals>,
+ from: Record<DigestPlanRole, DigestPlanTotals>,
+): void {
+ for (const role of DIGEST_PLAN_ROLES) {
+ into[role].videos += from[role].videos;
+ into[role].audioSeconds += from[role].audioSeconds;
+ }
+}
+
+// Which mirrors the alignment gate will actually accept. Read off the report's
+// per-ref `aligned`, measured at detection time. Absent means "not measured",
+// which sharing treats as NOT aligned — so absence lands in mirror-unaligned
+// and the plan under-promises rather than over-promises.
+function mirrorIsFree(
+ alignedBySlug: ReadonlyMap<string, boolean>,
+ slug: string,
+): boolean {
+ return alignedBySlug.get(slug) === true;
+}
+
+export async function buildDigestSweepPlan(
+ opts: BuildDigestSweepPlanOptions,
+): Promise<DigestSweepPlan> {
+ 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<unknown, string>({
+ 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") ?? "<none>"} != ${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<string, Row[]>();
+ 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 });
+ }
+
+ const channels: DigestChannelPlan[] = [];
+ for (const [channelSlug, rows] of byChannel) {
+ const entry: DigestChannelPlan = {
+ channelSlug,
+ remaining: emptyRoleTotals(),
+ fresh: emptyTotals(),
+ ineligible: 0,
+ generateSeconds: 0,
+ sharedSeconds: 0,
+ };
+
+ 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;
+ }
+
+ 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);
+ continue;
+ }
+ }
+
+ const role = clusterPlan.bySlug.get(stat.slug);
+ const planRole: DigestPlanRole =
+ role?.kind === "canonical"
+ ? "canonical"
+ : role?.kind === "mirror"
+ ? mirrorIsFree(alignedBySlug, stat.slug)
+ ? "mirror-aligned"
+ : "mirror-unaligned"
+ : "unclustered";
+ addTo(entry.remaining[planRole], stat.duration);
+ }
+
+ for (const role of DIGEST_PLAN_ROLES) {
+ if (roleMustGenerate(role))
+ entry.generateSeconds += entry.remaining[role].audioSeconds;
+ else entry.sharedSeconds += entry.remaining[role].audioSeconds;
+ }
+ channels.push(entry);
+ }
+
+ channels.sort(
+ (a, b) =>
+ b.generateSeconds - a.generateSeconds ||
+ a.channelSlug.localeCompare(b.channelSlug),
+ );
+
+ const totals: DigestSweepPlan = {
+ channels,
+ remaining: emptyRoleTotals(),
+ fresh: emptyTotals(),
+ ineligible: 0,
+ generateSeconds: 0,
+ sharedSeconds: 0,
+ 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.ineligible += c.ineligible;
+ totals.generateSeconds += c.generateSeconds;
+ totals.sharedSeconds += c.sharedSeconds;
+ }
+ 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<Map<string, boolean>> {
+ const { readDuplicateReport } = await import("./duplicateShorts");
+ const report = await readDuplicateReport(paths);
+ const map = new Map<string, boolean>();
+ 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.
+export function sweepDays(
+ audioSeconds: number,
+ secondsPerAudioHour: number = MEASURED_SECONDS_PER_AUDIO_HOUR,
+): number {
+ return (audioHours(audioSeconds) * secondsPerAudioHour) / 86400;
+}