// Where duplicate detection meets the digest layer: generate ONCE per cluster, // share to the mirrors that are safe to share to. // // This is worth ~11% of the sweep on its own — 8,204 exact-title redundancies // across the corpus, 5,863 of them Quartering YouTube↔Rumble mirror pairs — and // at weeks of wall-clock per pass, not re-generating those is real time. // // Two rules keep the sharing honest, and both are the difference between a // correct optimisation and a plausible-looking wrong one: // // 1. A `contained` cluster NEVER shares. Containment means one member is a CLIP // of a longer video; the longer video's chapters describe material the clip // does not contain. // 2. A mirror only receives a digest when its cue TIMINGS align with the // canonical member's. Content similarity says nothing about timing — a // mirror with a longer intro matches on text at shifted times — so a shared // digest would place every chapter wrong while looking perfectly fine. The // gate measures the offset at several anchors and requires near-zero. import path from "node:path"; import type { Paths } from "../lib/paths"; import { DEFAULT_ALIGNMENT_TOLERANCE_SECONDS, clusterMaySharePartial, measureAlignment, resolveCanonicalSlug, type AlignmentResult, type DuplicateCluster, type DuplicateOverrides, } from "../lib/duplicates"; import { CUES_JSON_FILENAME } from "../lib/videoStatus"; import type { Cue } from "../lib/vtt"; import { loadDigest, writeSharedDigest } from "../lib/digest-server"; import type { DigestRecord } from "../lib/digest"; import { readNormalizedTranscript } from "./normalizeTranscript"; import { readDuplicateOverrides, readDuplicateReport } from "./duplicateShorts"; // What a video's cluster membership means for the batch. export type DigestClusterRole = // Generate for this video: it owns its cluster's derived work. | { kind: "canonical"; clusterId: string; mirrors: string[] } // Do NOT generate: the canonical member owns it and will share it here. | { kind: "mirror"; clusterId: string; canonicalSlug: string }; export type DigestClusterPlan = { // `${channelSlug}/${id}` → role. Videos absent from the map are in no cluster // and are generated normally, which is the overwhelming majority. bySlug: Map; // The two directions of the id ↔ on-disk-directory mapping, populated ONLY for // cluster members whose directory name is not their metadata id. // // These exist because a slug is `${channelSlug}/${id}` and is not a path, // while every caller that reaches this plan starts from one side or the other: // the batch runner enumerates DIRECTORIES off disk, and the sharing pass has // SLUGS from the report and has to open their files. 14.5% of the corpus has // dir !== id, and it is not spread evenly — it is 100% of `the-quartering- // rumble` (7,870 videos), i.e. exactly the mirror set cluster-sharing was // written to exploit. Keying either lookup with the wrong half misses silently // and simply regenerates the mirror, which is why this went unnoticed. // // Empty for a report written before `DuplicateVideoRef.videoDir` existed; // lookups fall back to treating the id as the directory, the old behaviour. slugByDir: Map; // `${channelSlug}/${videoDir}` → slug dirBySlug: Map; // slug → videoDir (bare name, not a path) // Clusters that exist but share nothing (contained, or marked not-a-duplicate). // Their members are absent from bySlug and are each generated independently — // a clip is a different artifact and deserves its own digest. independentClusters: number; clusters: number; }; export function emptyDigestClusterPlan(): DigestClusterPlan { return { bySlug: new Map(), slugByDir: new Map(), dirBySlug: new Map(), independentClusters: 0, clusters: 0, }; } // The batch runner enumerates on-disk directories; the plan is keyed by slug. // Falls back to the directory name as the id, which is correct for the ~85.5% // where they agree and is what the code did before the mapping existed. export function planSlugForDir( plan: DigestClusterPlan | null | undefined, channelSlug: string, videoDir: string, ): string { const key = `${channelSlug}/${videoDir}`; return plan?.slugByDir.get(key) ?? key; } // Build the plan from the last detection run. Absent report → an empty plan, so // the digest batch works fine before duplicates have ever been detected (it just // generates for mirrors twice, which is correct, only slower). export async function buildDigestClusterPlan( paths: Paths, opts: { overrides?: DuplicateOverrides | null } = {}, ): Promise { const report = await readDuplicateReport(paths); if (!report) return emptyDigestClusterPlan(); const overrides = opts.overrides ?? (await readDuplicateOverrides(paths)); const plan = emptyDigestClusterPlan(); plan.clusters = report.clusters.length; for (const cluster of report.clusters) { // Recorded for EVERY cluster, including the ones that share nothing: an // independent cluster's members are still generated, and a caller that // reaches the plan at all deserves a correct path for them. for (const ref of cluster.videoRefs) { if (!ref.videoDir || ref.videoDir === ref.id) continue; plan.dirBySlug.set(ref.slug, ref.videoDir); plan.slugByDir.set(`${ref.channelSlug}/${ref.videoDir}`, ref.slug); } const canonicalSlug = resolveCanonicalSlug(cluster, overrides); // null → a human said "not a duplicate": every member stands alone. // The overrides also carry the `confirmed` flag that is the ONLY thing // letting a needsReview cluster share, so they must be passed here — without // them every confirmed suspect would silently keep generating twice. if (!canonicalSlug || !clusterMaySharePartial(cluster, overrides)) { plan.independentClusters++; continue; } const mirrors = cluster.videoRefs .map((r) => r.slug) .filter((slug) => slug !== canonicalSlug); if (mirrors.length === 0) { plan.independentClusters++; continue; } plan.bySlug.set(canonicalSlug, { kind: "canonical", clusterId: cluster.clusterId, mirrors, }); for (const slug of mirrors) { plan.bySlug.set(slug, { kind: "mirror", clusterId: cluster.clusterId, canonicalSlug, }); } } return plan; } // A slug is `${channelSlug}/${id}`, and an id is NOT always the directory name — // see DigestClusterPlan.slugByDir. `dirBySlug` carries the exceptions; without // it this resolves the 14.5% of the corpus with dir !== id to a path that does // not exist, which reads as "no transcript" and silently declines to share. function videoDirForSlug( paths: Paths, slug: string, dirBySlug?: ReadonlyMap, ): string | null { const at = slug.indexOf("/"); if (at <= 0) return null; return path.join( paths.channelsDir, slug.slice(0, at), "data", dirBySlug?.get(slug) ?? slug.slice(at + 1), ); } async function readCues( paths: Paths, slug: string, dirBySlug?: ReadonlyMap, ): Promise { const dir = videoDirForSlug(paths, slug, dirBySlug); if (!dir) return null; const t = await readNormalizedTranscript(path.join(dir, CUES_JSON_FILENAME)); return t?.cues ?? null; } export type ShareOutcome = { slug: string; status: "shared" | "misaligned" | "no-transcript" | "already-shared" | "failed"; alignment?: AlignmentResult; error?: string; }; export type ShareDigestOptions = { paths: Paths; clusterId: string; canonicalSlug: string; mirrors: string[]; // slug → on-disk directory, for members where the two differ. Take it from // DigestClusterPlan.dirBySlug (or build it from the cluster's own refs). // Omitting it is not an error, but every member with dir !== id will resolve // to a nonexistent path and be reported "no-transcript". dirBySlug?: ReadonlyMap; toleranceSeconds?: number; onLog?: (msg: string) => void; }; // Copy the canonical member's digest onto each mirror that passes the alignment // gate. Returns one outcome per mirror — a misaligned mirror is a normal, // expected result, not an error, and it simply stays on the generation worklist. export async function shareDigestToCluster( opts: ShareDigestOptions, ): Promise { const log = opts.onLog ?? (() => {}); const tolerance = opts.toleranceSeconds ?? DEFAULT_ALIGNMENT_TOLERANCE_SECONDS; const canonicalDir = videoDirForSlug( opts.paths, opts.canonicalSlug, opts.dirBySlug, ); if (!canonicalDir) return []; const source: DigestRecord | null = await loadDigest(canonicalDir); if (!source) return []; const canonicalCues = await readCues( opts.paths, opts.canonicalSlug, opts.dirBySlug, ); if (!canonicalCues || canonicalCues.length === 0) return []; const sharedAt = new Date().toISOString(); const outcomes: ShareOutcome[] = []; for (const slug of opts.mirrors) { const dir = videoDirForSlug(opts.paths, slug, opts.dirBySlug); if (!dir) { outcomes.push({ slug, status: "failed", error: "unparseable slug" }); continue; } try { const existing = await loadDigest(dir); // Already carrying this canonical member's digest at the same provenance: // nothing to do. Keeps a re-run a genuine no-op. if ( existing?.derivedFrom?.slug === opts.canonicalSlug && existing.promptVersion === source.promptVersion && existing.contextHash === source.contextHash ) { outcomes.push({ slug, status: "already-shared" }); continue; } const cues = await readCues(opts.paths, slug, opts.dirBySlug); if (!cues || cues.length === 0) { outcomes.push({ slug, status: "no-transcript" }); continue; } const alignment = measureAlignment(canonicalCues, cues, { toleranceSeconds: tolerance, }); if (!alignment.aligned) { log( `Not sharing to ${slug}: ${alignment.reason} (max offset ${ Number.isFinite(alignment.maxOffsetSeconds) ? `${alignment.maxOffsetSeconds.toFixed(1)}s` : "n/a" }, ${alignment.matchedAnchors}/${alignment.totalAnchors} anchors matched).`, ); outcomes.push({ slug, status: "misaligned", alignment }); continue; } await writeSharedDigest(dir, source, { slug: opts.canonicalSlug, clusterId: opts.clusterId, sharedAt, offsetSeconds: Math.round(alignment.maxOffsetSeconds * 100) / 100, }); log( `Shared digest ${opts.canonicalSlug} → ${slug} (max offset ${alignment.maxOffsetSeconds.toFixed(2)}s).`, ); outcomes.push({ slug, status: "shared", alignment }); } catch (err) { outcomes.push({ slug, status: "failed", error: (err as Error)?.message ?? String(err), }); } } return outcomes; } // Convenience for /review's per-cluster action: share from whatever the // effective canonical member currently is. export async function shareClusterFromCanonical( paths: Paths, cluster: DuplicateCluster, opts: { overrides?: DuplicateOverrides | null; onLog?: (msg: string) => void; } = {}, ): Promise { // Load the overrides BEFORE the gate, not after: `confirmed` lives in them and // is what unblocks a needsReview cluster, so testing the gate first would // refuse to share from every cluster a human had just approved. const overrides = opts.overrides ?? (await readDuplicateOverrides(paths)); if (!clusterMaySharePartial(cluster, overrides)) return []; const canonicalSlug = resolveCanonicalSlug(cluster, overrides); if (!canonicalSlug) return []; return shareDigestToCluster({ paths, clusterId: cluster.clusterId, canonicalSlug, mirrors: cluster.videoRefs .map((r) => r.slug) .filter((slug) => slug !== canonicalSlug), dirBySlug: dirBySlugForCluster(cluster), onLog: opts.onLog, }); } // The single-cluster equivalent of DigestClusterPlan.dirBySlug, for callers that // hold one cluster rather than a whole plan. export function dirBySlugForCluster( cluster: DuplicateCluster, ): Map { const map = new Map(); for (const ref of cluster.videoRefs) { if (ref.videoDir && ref.videoDir !== ref.id) map.set(ref.slug, ref.videoDir); } return map; }