// Cross-platform duplicate detection — a global pass over every channel's videos // (NOT per-channel, since duplicates are cross-channel and cross-platform by // definition). // // THE PRE-FILTER PROPOSES; THE TRANSCRIPT DISPOSES. // // A blocking strategy (§ "Phase 1" below) only NOMINATES pairs. It is never // evidence on its own. Every nominated pair then goes through the transcript // cascade, which has three outcomes: // // both sides have a transcript, content matches → transcript-exact/near. // Confirmed; may share. // both sides have a transcript, content does NOT → REJECTED, no cluster. // This is the whole value of // testing the suspects. // one side has no transcript → untestable. Kept only when // the pre-filter's own claim // is strong enough to be // worth a human's attention // (title + near-identical // runtime), as a needsReview // `title-duration` suspect // that shares nothing. // // STREAMING, NOT BATCHING. The first implementation built one global candidate // array, then fingerprinted every video it referenced, then evaluated. Both // halves fail corpus-wide: the containment pass alone was an unblocked cartesian // product (498 M pairs — V8 throws RangeError building the array), and holding // 5-word shingle sets for the whole corpus at once (a 3.5 h video is ~40 k // strings) exhausted a 4 GB heap. Neither is inherent to corpus-wide detection. // Blocking fixes the first; this file's block-at-a-time iteration fixes the // second: fingerprint one block, evaluate its pairs, keep only the confirmed // matches, release. Peak memory is O(largest block), not O(corpus), and no // multi-million-element pair array is ever materialised. // // Inputs are the already-aggregated VideoStat index (LMDB statsByPath, from // buildStats) plus the `cues` sub-db (from buildIndex). Output is a single // global transcripts/duplicates.json. Flag-only: nothing is merged or deleted. import path from "node:path"; import { readFile } from "node:fs/promises"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import { createHash } from "node:crypto"; import { open } from "lmdb"; import pLimit from "p-limit"; import type { Paths } from "../lib/paths"; import type { VideoStat } from "../lib/stats"; import { STATS_SCHEMA_VERSION } from "../lib/stats"; import type { Cue } from "../lib/vtt"; import { parseTranscriptJson } from "../lib/whisper"; import { readVideoFiles, pickIndexTranscript } from "../lib/videoStatus"; import { DEFAULT_ALIGNMENT_TOLERANCE_SECONDS, DEFAULT_CONTAINMENT_THRESHOLD, DEFAULT_DURATION_TOLERANCE_SECONDS, DEFAULT_TITLE_DURATION_RATIO, MAX_DURATION_BLOCK_SIZE, MAX_TITLE_GROUP_SIZE, durationsCompatible, measureAlignment, normalizeTitleKey, DEFAULT_NEAR_THRESHOLD, DEFAULT_SHINGLE_SIZE, DEFAULT_SHORT_THRESHOLD_SECONDS, DUPLICATES_FILENAME, DUPLICATE_REPORT_VERSION, DUPLICATE_OVERRIDES_FILENAME, UnionFind, comparisonText, containment, jaccard, pickCanonicalSlug, sanitizeDuplicateOverrides, shingles, strongerMatch, type DuplicateBlocking, type DuplicateOverrides, type DuplicateCluster, type DuplicateMatchKind, type DuplicateReport, type DuplicateRunConfig, type DuplicateVideoRef, } from "../lib/duplicates"; // A normalized title shorter than this carries no blocking information — // pairing on it would rebuild the cartesian product the pass exists to avoid. const MIN_TITLE_KEY_LENGTH = 8; // Hard budget for the shorts-vs-longer containment sweep, which is a genuine // cartesian product and the one pass blocking cannot rescue. Under the budget it // runs (and finds clip-of-longer matches, which nothing else can); over it, it // is SKIPPED AND REPORTED rather than silently truncated or left to die. Scaling // containment properly wants MinHash/LSH signatures — see STATE.md. const MAX_CONTAINMENT_PAIRS = 5_000_000; type PathKey = [string, string]; type IndexKey = [string, string, string]; type StatsRecord = { metaMs: number; stat: VideoStat }; export type DetectDuplicateShortsOptions = { paths: Paths; // Duration cutoff for what counts as a "short". `null` === all durations // (no length filter; enables containment matching against long videos). thresholdSeconds?: number | null; durationToleranceSeconds?: number; nearThreshold?: number; containmentThreshold?: number; shingleSize?: number; // Candidate-nomination strategy. Defaults to "title" corpus-wide (where a // duration sweep is affordable but far slower) and "duration" in shorts mode, // which preserves today's proven behaviour at that scale. blocking?: DuplicateBlocking; titleDurationRatio?: number; // Tolerance for the per-cluster timing measurement written onto each ref. alignmentToleranceSeconds?: number; onLog?: (msg: string) => void; signal?: AbortSignal; }; // Per-participant transcript fingerprint, computed once per BLOCK and released // when the block is done. Holding these for the whole corpus is what exhausted // the heap; the shingle set is the expensive part (~40 k strings for a 3.5 h // video), which is why it never outlives its block. type Fingerprint = { stat: VideoStat; hasTranscript: boolean; hash: string | null; shingleSet: Set; }; type Entry = { stat: VideoStat; videoDir: string }; // One nominated pair. `suspectKind` is what this pair's evidence is worth when // the content CANNOT be compared because a side has no transcript: a shared // title AND a near-identical runtime is a claim worth a human's attention, so it // survives as a `title-duration` suspect; a shared duration alone is not, and // treating it as one is exactly what produced enormous false clusters of // unrelated same-length videos before, so it is null and the pair is dropped. type Nomination = { a: Entry; b: Entry; // Skip the similarity test and go straight to containment — a clip can never // pass Jaccard against the full recording it was cut from. containmentOnly: boolean; suspectKind: DuplicateMatchKind | null; }; // A unit of work: fingerprint these members, evaluate these pairs, release. type Block = { members: Entry[]; nominations: Nomination[] }; type PairResult = { a: string; // slug b: string; // slug kind: DuplicateMatchKind; score: number | null; contained: boolean; }; // Ordered slug-pair key, so "both" can union two nomination streams without // evaluating the overlap twice. function pairKey(a: string, b: string): string { return a < b ? `${a}\n${b}` : `${b}\n${a}`; } function membersOf(nominations: Nomination[]): Entry[] { const m = new Map(); for (const n of nominations) { m.set(n.a.stat.slug, n.a); m.set(n.b.stat.slug, n.b); } return [...m.values()]; } export async function detectDuplicateShorts( opts: DetectDuplicateShortsOptions, ): Promise { const log = opts.onLog ?? ((m: string) => console.log(m)); // Corpus-wide (thresholdSeconds: null) is untested at 74k videos across the // long tail, so the run reports its own wall-clock: the cost of the pass has to // be a measured number before anything is wired to it on a schedule. const startedAt = Date.now(); const thresholdSeconds = opts.thresholdSeconds === undefined ? DEFAULT_SHORT_THRESHOLD_SECONDS : opts.thresholdSeconds; const W = opts.durationToleranceSeconds ?? DEFAULT_DURATION_TOLERANCE_SECONDS; const nearThreshold = opts.nearThreshold ?? DEFAULT_NEAR_THRESHOLD; const containmentThreshold = opts.containmentThreshold ?? DEFAULT_CONTAINMENT_THRESHOLD; const shingleSize = opts.shingleSize ?? DEFAULT_SHINGLE_SIZE; // Corpus-wide defaults to TITLE blocking: it is effectively linear and it is // the signal that finds cross-platform re-uploads, which keep their name. // Duration blocking is viable corpus-wide too since the streaming rewrite (it // is the only strategy that catches a RE-TITLED mirror) but it nominates // millions of pairs where title nominates thousands, so it is opt-in. const blocking = opts.blocking ?? (thresholdSeconds === null ? "title" : "duration"); const titleDurationRatio = opts.titleDurationRatio ?? DEFAULT_TITLE_DURATION_RATIO; const alignmentTolerance = opts.alignmentToleranceSeconds ?? DEFAULT_ALIGNMENT_TOLERANCE_SECONDS; const runConfig: DuplicateRunConfig = { thresholdSeconds, durationToleranceSeconds: W, nearThreshold, containmentThreshold, shingleSize, blocking, }; // ---- Open the index: statsByPath (Phase-1 input) + cues (Phase-2 input) --- const root = open({ path: opts.paths.lmdbPath, maxDbs: 12, compression: true }); const statsByPath = root.openDB({ name: "statsByPath", encoding: "msgpack", }); const cuesDb = root.openDB({ name: "cues", encoding: "msgpack", }); const meta = root.openDB({ name: "statsMeta", encoding: "msgpack", }); const storedSchema = meta.get("schema") as number | undefined; if (storedSchema !== STATS_SCHEMA_VERSION) { log( `Warning: stats cache schema ${storedSchema ?? ""} != ${STATS_SCHEMA_VERSION}; run build:stats first for accurate results.`, ); } const all: Entry[] = []; for (const { key, value } of statsByPath.getRange()) { // statsByPath is keyed [channelSlug, videoDir]; keep videoDir so the // Phase-2 fallback can locate the raw transcript file on disk. const videoDir = (key as PathKey)[1]; all.push({ stat: (value as StatsRecord).stat, videoDir }); } const videosScanned = all.length; const isShort = (d: number) => d > 0 && (thresholdSeconds === null || d <= thresholdSeconds); const includeContainment = thresholdSeconds === null; // In all-durations mode every video participates (long videos are the // containment targets); otherwise only shorts do. const participants = all.filter((e) => thresholdSeconds === null ? e.stat.duration > 0 : isShort(e.stat.duration), ); log( `Scanned ${videosScanned} videos; ${participants.length} in scope ` + `(threshold=${thresholdSeconds === null ? "all" : `${thresholdSeconds}s`}).`, ); // ---- Phase 1 + 2, interleaved: nominate a block, judge it, release it ----- // // The two phases are no longer separate passes. Each block is fingerprinted, // evaluated and dropped before the next is built, so peak memory is O(largest // block) and only the (few, by construction) confirmed matches survive. const stats = { nominated: 0, evaluated: 0, confirmed: 0, rejected: 0, suspects: 0, titleGroups: 0, oversizedGroups: 0, oversizedGroupVideos: 0, oversizedBlocks: 0, oversizedBlockVideos: 0, untitled: 0, fingerprinted: 0, rawFallbacks: 0, rawHits: 0, peakRssMb: 0, }; const noteRss = () => { const mb = Math.round(process.memoryUsage().rss / 1_048_576); if (mb > stats.peakRssMb) stats.peakRssMb = mb; }; // Fingerprint one block's members, re-using anything the previous block // already computed. Primary source is the LMDB `cues` sub-db — but parseVtt // only extracts text from YouTube karaoke-tagged cues, so plain-VTT // transcripts (most non-YouTube captions, whisper-as-vtt, manual subs) store // as zero cues and would look transcript-less. For those, fall back to reading // the raw transcript file off disk and extracting plain text, for comparison // only — nothing is persisted. const fingerprintBlock = async ( members: Entry[], reuse: ReadonlyMap, ): Promise> => { const out = new Map(); const rawFallbacks: Entry[] = []; for (const e of members) { const s = e.stat; if (out.has(s.slug)) continue; const cached = reuse.get(s.slug); if (cached) { out.set(s.slug, cached); continue; } const cues = cuesDb.get([s.uploadDate, s.channelSlug, s.id]); if (cues && cues.length > 0) { stats.fingerprinted++; out.set(s.slug, fingerprintFrom(s, cues, shingleSize)); } else { rawFallbacks.push(e); } } if (rawFallbacks.length > 0) { stats.rawFallbacks += rawFallbacks.length; const limit = pLimit(16); await Promise.all( rawFallbacks.map((e) => limit(async () => { opts.signal?.throwIfAborted(); const cues = await readRawCues(opts.paths.channelsDir, e); if (cues && cues.length > 0) stats.rawHits++; stats.fingerprinted++; out.set(e.stat.slug, fingerprintFrom(e.stat, cues, shingleSize)); }), ), ); } return out; }; // Cues WITH real start times, for the alignment measurement only. Deliberately // LMDB-only: the raw-VTT fallback stamps every cue at start 0, which is fine // for a set-based content fingerprint and worthless — actively misleading — // for measuring an offset. No timed cues → no measurement → `aligned` stays // undefined, which every consumer must read as "not aligned". const timedCues = (e: Entry | undefined): Cue[] | null => { if (!e) return null; const s = e.stat; const cues = cuesDb.get([s.uploadDate, s.channelSlug, s.id]); return cues && cues.length > 0 ? cues : null; }; const streams: Generator[] = []; if (blocking === "title" || blocking === "both") { streams.push(titleBlocks(participants, titleDurationRatio, stats, log)); } if (blocking === "duration" || blocking === "both") { streams.push(durationBlocks(participants, W, stats, log)); } // Only "both" needs cross-stream dedup. It is deliberately one-directional: // the FIRST stream records its nominations and later streams consult them. // Title blocking is pushed first precisely so the remembered set is the small // one — recording the duration stream's millions of pair keys as well would // reintroduce, in the dedup set, exactly the whole-corpus retention the // streaming rewrite exists to remove. const seenPairs = blocking === "both" ? new Set() : null; const matches: PairResult[] = []; const hasTranscriptBySlug = new Map(); // 2-block sliding fingerprint cache: a duration block is bucket[b] ∪ // bucket[b+1], so every member except the last bucket is re-used by the very // next block. Without this, each bucket would be read and shingled twice. let prevFps = new Map(); let blocksDone = 0; for (const [streamIndex, stream] of streams.entries()) { const recordPairs = seenPairs !== null && streamIndex === 0; const skipSeenPairs = seenPairs !== null && streamIndex > 0; for (const block of stream) { opts.signal?.throwIfAborted(); let nominations = block.nominations; if (recordPairs) { for (const n of nominations) { seenPairs.add(pairKey(n.a.stat.slug, n.b.stat.slug)); } } else if (skipSeenPairs) { nominations = nominations.filter( (n) => !seenPairs.has(pairKey(n.a.stat.slug, n.b.stat.slug)), ); } stats.nominated += nominations.length; if (nominations.length === 0) continue; // When dedup narrowed the nominations, re-derive the members from what is // actually left rather than fingerprinting videos no surviving pair needs. const fps = await fingerprintBlock( skipSeenPairs ? membersOf(nominations) : block.members, prevFps, ); for (const [slug, fp] of fps) hasTranscriptBySlug.set(slug, fp.hasTranscript); for (const n of nominations) { const fa = fps.get(n.a.stat.slug); const fb = fps.get(n.b.stat.slug); if (!fa || !fb) continue; stats.evaluated++; const res = n.containmentOnly ? evalContainment(fa, fb, containmentThreshold) : evalBlocked(fa, fb, { nearThreshold, containmentThreshold, allowContainment: includeContainment, suspectKind: n.suspectKind, }); if (!res) { stats.rejected++; continue; } if (res.kind === "title-duration") stats.suspects++; else stats.confirmed++; matches.push({ a: n.a.stat.slug, b: n.b.stat.slug, ...res }); } prevFps = fps; noteRss(); // A corpus-wide duration run is long enough that silence is // indistinguishable from a hang. Heartbeat with the numbers that matter. if (++blocksDone % 100 === 0) { log( ` …${blocksDone} block(s), ${stats.evaluated} pair(s) tested, ` + `${stats.confirmed} confirmed, RSS ${Math.round(process.memoryUsage().rss / 1_048_576)} MB`, ); } } prevFps = new Map(); } // ---- The containment sweep: shorts vs every meaningfully-longer video ----- // // The one pass blocking cannot rescue — a clip and its parent share neither a // duration nor, usually, a title, so nothing nominates them but a cartesian // product. Bounded by an explicit pair budget and skipped-with-a-log when it // does not fit, rather than silently truncated. Memory stays bounded because // the budget implicitly bounds the short side: at 76 k longer videos, fitting // under the budget means only a handful of shorts. if (includeContainment) { const shorts = participants.filter( (e) => e.stat.duration > 0 && e.stat.duration <= DEFAULT_SHORT_THRESHOLD_SECONDS, ); const budget = shorts.length * participants.length; if (shorts.length === 0) { // nothing to sweep } else if (budget > MAX_CONTAINMENT_PAIRS) { log( `Containment sweep SKIPPED: ${shorts.length} short(s) × ${participants.length} video(s) ` + `= ~${budget.toLocaleString("en-US")} pairs, over the ${MAX_CONTAINMENT_PAIRS.toLocaleString("en-US")} budget. ` + `Clip-of-longer duplicates are NOT covered by this run.`, ); } else { const shortFps = await fingerprintBlock(shorts, new Map()); for (const [slug, fp] of shortFps) { hasTranscriptBySlug.set(slug, fp.hasTranscript); } const shortSlugs = new Set(shorts.map((e) => e.stat.slug)); let swept = 0; for (const v of participants) { opts.signal?.throwIfAborted(); // One long video at a time: fingerprint, compare against every short, // release. O(shorts + 1) fingerprints held. const relevant = shorts.filter( (s) => s.stat.slug !== v.stat.slug && v.stat.duration >= s.stat.duration * 1.5, ); if (relevant.length === 0) continue; const fv = shortSlugs.has(v.stat.slug) ? shortFps.get(v.stat.slug) : (await fingerprintBlock([v], new Map())).get(v.stat.slug); if (!fv) continue; hasTranscriptBySlug.set(v.stat.slug, fv.hasTranscript); for (const s of relevant) { const fs = shortFps.get(s.stat.slug); if (!fs) continue; swept++; stats.evaluated++; const res = evalContainment(fs, fv, containmentThreshold); if (!res) { stats.rejected++; continue; } stats.confirmed++; matches.push({ a: s.stat.slug, b: v.stat.slug, ...res }); } noteRss(); } stats.nominated += swept; log(`Containment sweep: ${swept} pair(s) evaluated.`); } } log( `Phase 2: ${stats.evaluated} nominated pair(s) tested → ${stats.confirmed} confirmed, ` + `${stats.rejected} rejected by content, ${stats.suspects} untestable suspect(s). ` + `Fingerprinted ${stats.fingerprinted} video(s) (${stats.rawFallbacks} missing LMDB cues; ` + `${stats.rawHits} recovered from raw transcripts).`, ); prevFps = new Map(); // ---- Cluster matched pairs (union-find) ---------------------------------- const uf = new UnionFind(); for (const m of matches) uf.union(m.a, m.b); const components = uf.groups().filter((g) => g.length >= 2); const componentRoot = new Map(); for (const g of components) for (const slug of g) componentRoot.set(slug, g[0]); type Agg = { matchKind: DuplicateMatchKind; score: number | null; contained: boolean; }; const aggByRoot = new Map(); for (const m of matches) { const root2 = componentRoot.get(m.a); if (!root2) continue; const cur = aggByRoot.get(root2); if (!cur) { aggByRoot.set(root2, { matchKind: m.kind, score: m.score, contained: m.contained, }); } else { cur.matchKind = strongerMatch(cur.matchKind, m.kind); cur.score = maxScore(cur.score, m.score); cur.contained = cur.contained || m.contained; } } // Refs are rebuilt from the participant index, NOT from fingerprints — those // were released with their block. VideoStat is small metadata and the index is // already held for the whole run, so this costs nothing. const bySlug = new Map(); for (const e of participants) bySlug.set(e.stat.slug, e); const clusters: DuplicateCluster[] = components.map((slugs) => { const agg = aggByRoot.get(slugs[0]) as Agg; const refs = slugs .map((slug) => toRef( bySlug.get(slug) as Entry, hasTranscriptBySlug.get(slug) ?? false, ), ) .sort((x, y) => x.slug.localeCompare(y.slug)); const platforms = new Set(refs.map((r) => r.platform)); const channels = new Set(refs.map((r) => r.channelSlug)); const minDuration = Math.min(...refs.map((r) => r.duration)); const cluster: DuplicateCluster = { clusterId: sha1(refs.map((r) => r.slug).join("\n")), matchKind: agg.matchKind, score: agg.score, contained: agg.contained, durationBucket: Math.round(minDuration), crossPlatform: platforms.size > 1, crossChannel: channels.size > 1, // Nothing compared this cluster's CONTENT — every pair in it was // untestable. It is a suspect for human review and shares nothing until // someone records `confirmed`. Written only when true so that a // content-confirmed cluster serialises exactly as it did before. ...(agg.matchKind === "title-duration" ? { needsReview: true } : {}), videoRefs: refs, }; // Record the rule's choice of canonical member: the one derived work (an AI // digest today, attribution later) is generated for and shared FROM. A human // can override it in duplicates.overrides.json — which is a separate file // precisely because this report is rewritten wholesale on every run. return { ...cluster, canonicalSlug: pickCanonicalSlug(cluster) }; }); // ---- Measure timing alignment against each cluster's canonical member ----- // // Content similarity says NOTHING about timing: a mirror with a longer intro // matches on text at shifted times. Anything that seeks into a sibling — a // shared digest's chapters, the search-result "jump to this moment" — is wrong // without this, and wrong in the way that looks right. digestSharing measured // it already but threw the result away; persisting it is what lets the viewer // be honest about when a jump is trustworthy. // // Confirmed clusters only, and matches are few by construction, so re-reading // those cue lists is cheap. Suspects are skipped: they share nothing anyway. let aligned = 0; let alignmentPairs = 0; for (const cluster of clusters) { opts.signal?.throwIfAborted(); if (cluster.needsReview || cluster.contained) continue; const canonicalSlug = cluster.canonicalSlug; if (!canonicalSlug) continue; const canonicalCues = timedCues(bySlug.get(canonicalSlug)); if (!canonicalCues) continue; for (const ref of cluster.videoRefs) { if (ref.slug === canonicalSlug) { ref.offsetSeconds = 0; ref.aligned = true; continue; } const cues = timedCues(bySlug.get(ref.slug)); if (!cues) continue; alignmentPairs++; const a = measureAlignment(canonicalCues, cues, { toleranceSeconds: alignmentTolerance, }); ref.aligned = a.aligned; ref.offsetSeconds = Number.isFinite(a.maxOffsetSeconds) ? Math.round(a.maxOffsetSeconds * 100) / 100 : null; if (a.aligned) aligned++; } } if (alignmentPairs > 0) { log( `Alignment: ${aligned}/${alignmentPairs} mirror(s) within ${alignmentTolerance}s of their canonical member.`, ); } await root.close(); const rank: Record = { "transcript-exact": 2, "transcript-near": 1, // Suspects sort last: they are a review queue, not a result. "title-duration": 0, }; clusters.sort( (a, b) => rank[b.matchKind] - rank[a.matchKind] || b.videoRefs.length - a.videoRefs.length || (b.score ?? 0) - (a.score ?? 0) || a.clusterId.localeCompare(b.clusterId), ); const videosInClusters = clusters.reduce((n, c) => n + c.videoRefs.length, 0); const report: DuplicateReport = { version: DUPLICATE_REPORT_VERSION, generatedAt: new Date().toISOString(), runConfig, totals: { videosScanned, clusters: clusters.length, videosInClusters }, clusters, }; const outPath = path.join(opts.paths.transcriptsDir, DUPLICATES_FILENAME); // Compact, no trailing newline: the report's historical bytes. await writeJsonAtomic(outPath, report, { indent: 0, newline: false, mkdir: true }); noteRss(); const suspectClusters = clusters.filter((c) => c.needsReview).length; log( `Done: ${clusters.length} duplicate cluster(s) over ${videosInClusters} video(s) ` + `(${suspectClusters} awaiting review) → ${DUPLICATES_FILENAME} ` + `in ${Math.round((Date.now() - startedAt) / 1000)}s ` + `(blocking=${blocking}, ${videosScanned} scanned, ${stats.nominated} nominated pair(s), ` + `${stats.confirmed} confirmed, ${stats.rejected} rejected by content, ${stats.suspects} suspect(s), ` + `peak RSS ${stats.peakRssMb} MB).`, ); return report; } export async function readDuplicateReport( paths: Paths, ): Promise { try { const raw = await readFile( path.join(paths.transcriptsDir, DUPLICATES_FILENAME), "utf8", ); const parsed = JSON.parse(raw) as DuplicateReport; if (typeof parsed.version !== "number") return null; return parsed; } catch { return null; } } // --------------------------------------------------------------------------- // Human review decisions (canonical choice / not-a-duplicate) // --------------------------------------------------------------------------- export function duplicateOverridesPath(paths: Paths): string { return path.join(paths.transcriptsDir, DUPLICATE_OVERRIDES_FILENAME); } // Never throws: an unreadable or malformed overrides file reads as "no decisions // recorded", so the duplicates page still renders. export async function readDuplicateOverrides( paths: Paths, ): Promise { try { const raw = await readFile(duplicateOverridesPath(paths), "utf8"); return sanitizeDuplicateOverrides(JSON.parse(raw)); } catch { return sanitizeDuplicateOverrides(null); } } // Record one cluster's decision, preserving every other cluster's. Read-modify- // write with the same atomic tmp+rename the report itself uses. export async function updateDuplicateOverride( paths: Paths, clusterId: string, patch: { canonicalSlug?: string; notDuplicate?: boolean; // "I looked, and these really are the same video." The ONLY thing that lets // a needsReview (title+duration) cluster share derived work. confirmed?: boolean; note?: string; }, ): Promise { const current = await readDuplicateOverrides(paths); const existing = current.clusters[clusterId] ?? {}; const next = { ...existing, ...(patch.canonicalSlug !== undefined ? { canonicalSlug: patch.canonicalSlug } : {}), ...(patch.notDuplicate !== undefined ? { notDuplicate: patch.notDuplicate } : {}), ...(patch.confirmed !== undefined ? { confirmed: patch.confirmed } : {}), ...(patch.note !== undefined ? { note: patch.note } : {}), decidedAt: new Date().toISOString(), }; // An empty patch clears the decision (back to "awaiting review") rather than // leaving a decidedAt-only stub that would read as reviewed. `confirmed` has // to be part of that test: without it, clearing a canonical choice on a // confirmed suspect would DELETE the confirmation and silently un-share the // cluster's derived work. if (!next.canonicalSlug && next.notDuplicate !== true && next.confirmed !== true) { delete current.clusters[clusterId]; } else { current.clusters[clusterId] = next; } const out = sanitizeDuplicateOverrides(current); const file = duplicateOverridesPath(paths); await writeJsonAtomic(file, out, { mkdir: true }); return out; } // --------------------------------------------------------------------------- // Blocking strategies — they NOMINATE, they never decide // --------------------------------------------------------------------------- // Counters the generators fill in as they run, so the caller can report what was // covered AND what was skipped. Skipping without reporting reads as "covered // everything" when it did not. type BlockStats = { titleGroups: number; oversizedGroups: number; oversizedGroupVideos: number; oversizedBlocks: number; oversizedBlockVideos: number; untitled: number; }; // TITLE BLOCKING. One pass to group by exact normalized title, then pair within // each group when the runtimes agree. Quadratic only INSIDE a group, and a group // is a handful of videos, so the whole pass is effectively linear. // // The key is deliberately EXACT rather than fuzzy: a looser key merges "Episode // 12" with "Episode 13", and every such merge is a false cluster a human then // has to reject. The cost is recall on re-titled mirrors, which duration // blocking covers instead — see STATE.md for the deferred middle grounds. function* titleBlocks( participants: Entry[], ratio: number, stats: BlockStats, log: (msg: string) => void, ): Generator { const byTitle = new Map(); for (const e of participants) { const key = normalizeTitleKey(e.stat.title ?? ""); if (key.length < MIN_TITLE_KEY_LENGTH) { stats.untitled++; continue; } const arr = byTitle.get(key); if (arr) arr.push(e); else byTitle.set(key, [e]); } // Decide the whole work list BEFORE yielding any of it, so the summary — and // in particular what was SKIPPED — is reported up front rather than after the // long evaluation it describes. const accepted: Entry[][] = []; for (const [key, group] of byTitle) { if (group.length < 2) continue; if (group.length > MAX_TITLE_GROUP_SIZE) { // A FORMAT, not a title — "live stream", "untitled", a daily show's // date-less name. Reported, never silently truncated. stats.oversizedGroups++; stats.oversizedGroupVideos += group.length; log( ` skipping title group of ${group.length} (over ${MAX_TITLE_GROUP_SIZE}): "${key.slice(0, 60)}"`, ); continue; } accepted.push(group); } log( `Title blocking: ${byTitle.size} distinct title(s), ${accepted.length} group(s) with 2+ members ` + `(${stats.untitled} video(s) with no usable title; ${stats.oversizedGroups} group(s) over ` + `${MAX_TITLE_GROUP_SIZE} skipped, covering ${stats.oversizedGroupVideos} video(s)).`, ); for (const group of accepted) { const nominations: Nomination[] = []; for (let i = 0; i < group.length; i++) { for (let j = i + 1; j < group.length; j++) { if ( durationsCompatible( group[i].stat.duration, group[j].stat.duration, ratio, ) ) { nominations.push({ a: group[i], b: group[j], containmentOnly: false, // Same title AND a near-identical runtime: strong enough to be // worth a human's attention when no transcript can settle it. suspectKind: "title-duration", }); } } } if (nominations.length === 0) continue; stats.titleGroups++; yield { members: membersOf(nominations), nominations }; } } // DURATION BLOCKING. One block per rounded-duration bucket, unioned with the // next so a pair straddling a bucket edge is still nominated. Emitting // here×here and here×next (never next×next) makes every pair appear exactly // once, and yielding bucket b+1's members as part of block b is what lets the // caller's 2-block fingerprint cache halve the transcript reads. // // This is the only strategy that catches a RE-TITLED mirror, and it is viable // corpus-wide only because blocks are evaluated and released one at a time. function* durationBlocks( participants: Entry[], W: number, stats: BlockStats, log: (msg: string) => void, ): Generator { const bucketOf = (d: number) => Math.round(d / W); const buckets = new Map(); for (const e of participants) { const b = bucketOf(e.stat.duration); const arr = buckets.get(b); if (arr) arr.push(e); else buckets.set(b, [e]); } const maxDelta = (d: number) => Math.max(W, Math.ceil(0.02 * d)); const sorted = [...buckets.keys()].sort((x, y) => x - y); // As with title blocking: decide the work list, report it (skips included), // then do it. A summary that only arrives after an hour of evaluation is not a // summary anyone can act on. const accepted: number[] = []; for (const b of sorted) { const size = (buckets.get(b) as Entry[]).length + (buckets.get(b + 1)?.length ?? 0); // A bucket can itself be pathological — round numbers attract videos, and a // corpus can hold thousands that are exactly 60 s. Same cap-and-report // treatment as an oversized title group. if (size > MAX_DURATION_BLOCK_SIZE) { stats.oversizedBlocks++; stats.oversizedBlockVideos += (buckets.get(b) as Entry[]).length; log( ` skipping duration block ~${b * W}s of ${size} video(s) ` + `(over ${MAX_DURATION_BLOCK_SIZE}).`, ); continue; } accepted.push(b); } log( `Duration blocking: ${buckets.size} bucket(s) of ${W}s, ${accepted.length} block(s) to evaluate` + (stats.oversizedBlocks > 0 ? ` (${stats.oversizedBlocks} block(s) over ${MAX_DURATION_BLOCK_SIZE} skipped, ` + `covering ${stats.oversizedBlockVideos} video(s))` : "") + ".", ); for (const b of accepted) { const here = buckets.get(b) as Entry[]; const next = buckets.get(b + 1) ?? []; const nominations: Nomination[] = []; for (let i = 0; i < here.length; i++) { for (let j = i + 1; j < here.length; j++) { if (durationsClose(here[i], here[j], maxDelta)) { // A shared duration ALONE is not evidence — it produced enormous // false clusters of unrelated same-length videos. If the content // cannot be compared, this pair is worth nothing. nominations.push({ a: here[i], b: here[j], containmentOnly: false, suspectKind: null, }); } } } for (const x of here) { for (const y of next) { if (durationsClose(x, y, maxDelta)) { nominations.push({ a: x, b: y, containmentOnly: false, suspectKind: null, }); } } } if (nominations.length === 0) continue; // Members are here ∪ next, so the caller's sliding cache carries bucket b+1 // straight into the next block. yield { members: membersOf(nominations), nominations }; } } // --- helpers --------------------------------------------------------------- function durationsClose( a: Entry, b: Entry, maxDelta: (d: number) => number, ): boolean { const lo = Math.min(a.stat.duration, b.stat.duration); return Math.abs(a.stat.duration - b.stat.duration) <= maxDelta(lo); } function fingerprintFrom( stat: VideoStat, cues: Cue[] | undefined | null, shingleSize: number, ): Fingerprint { if (!cues || cues.length === 0) { return { stat, hasTranscript: false, hash: null, shingleSet: new Set() }; } const text = comparisonText(cues); return { stat, hasTranscript: text.length > 0, hash: text.length > 0 ? sha1(text) : null, shingleSet: shingles(text, shingleSize), }; } // Fallback transcript read for videos whose LMDB cues are empty (plain VTT that // parseVtt skips because it has no karaoke timing tags). Reads the raw // transcript file and extracts plain text for comparison only — nothing is // persisted. Returns null when there's no transcript file. async function readRawCues( channelsDir: string, e: Entry, ): Promise { const dir = path.join(channelsDir, e.stat.channelSlug, "data", e.videoDir); let picked: ReturnType; try { picked = pickIndexTranscript(await readVideoFiles(dir)); } catch { return null; } if (!picked) return null; let raw: string; try { raw = await readFile(path.join(dir, picked.filename), "utf8"); } catch { return null; } return picked.kind === "whisper" ? parseTranscriptJson(raw) : lenientVttCues(raw); } // Lenient VTT text extraction: unlike parseVtt (which only keeps YouTube // karaoke-tagged lines), this collects text from every cue body, strips all // tags/entities, and consecutive-dedups rolling captions. Good enough for a // comparison fingerprint; set-based shingles absorb residual repetition. function lenientVttCues(src: string): Cue[] { const lines = src.replace(/\r\n/g, "\n").split("\n"); const cues: Cue[] = []; let i = 0; while (i < lines.length) { const m = lines[i].match( /(\d{2}:\d{2}:\d{2}\.\d{3})\s+-->\s+(\d{2}:\d{2}:\d{2}\.\d{3})/, ); if (!m) { i++; continue; } i++; const body: string[] = []; while (i < lines.length && lines[i].trim() !== "") { body.push(lines[i]); i++; } const text = stripVttMarkup(body.join(" ")); if (text) cues.push({ start: 0, end: 0, text }); } const out: Cue[] = []; for (const c of cues) { if (out.length > 0 && out[out.length - 1].text === c.text) continue; out.push(c); } return out; } function stripVttMarkup(s: string): string { return s .replace(/ /g, " ") .replace(/&/g, "&") .replace(/</g, "<") .replace(/>/g, ">") .replace(/'/g, "'") .replace(/"/g, '"') .replace(/<[^>]*>/g, "") .replace(/\s+/g, " ") .trim(); } // Similar-duration pair: exact → near. Content agreement is required — a shared // duration alone is NOT evidence of duplication, so when either side lacks a // comparable transcript we return null rather than matching on metadata. (The // old metadata fallback collapsed every same-length video into one cluster.) function evalSimilar( a: Fingerprint, b: Fingerprint, nearThreshold: number, ): Omit | null { if (!a.hasTranscript || !b.hasTranscript) return null; if (a.hash && a.hash === b.hash) { return { kind: "transcript-exact", score: 1, contained: false }; } // Size prune. Exact, not a heuristic: |A∩B| ≤ min and |A∪B| ≥ max, so // J ≤ min/max. A pair whose shingle counts are further apart than the // threshold CANNOT reach it, and skipping it changes no result. Worth having // because it is O(1) where the Jaccard it replaces is O(min set size), and // corpus-wide duration blocking evaluates millions of pairs whose members // share a runtime but not a word count (a music video and a lecture can both // be ten minutes long). const small = Math.min(a.shingleSet.size, b.shingleSet.size); const large = Math.max(a.shingleSet.size, b.shingleSet.size); if (large === 0 || small / large < nearThreshold) return null; const j = jaccard(a.shingleSet, b.shingleSet); if (j >= nearThreshold) { return { kind: "transcript-near", score: round3(j), contained: false }; } return null; } // The verdict on one nominated pair — the point where the pre-filter's proposal // meets the transcript's disposal. Three outcomes, and the middle one is the // reason nominating aggressively is safe: // // both transcripts, content agrees → confirmed (may share) // both transcripts, content differs → null. REJECTED. Two episodes of a daily // show can share a title and a runtime and // be entirely different material; testing // them is what keeps that out. // a transcript is missing → nothing can compare the content, so the // pre-filter's own claim is all there is. // Kept as a needsReview suspect when that // claim is strong (title + runtime), and // dropped when it is not (duration alone). function evalBlocked( a: Fingerprint, b: Fingerprint, opts: { nearThreshold: number; containmentThreshold: number; allowContainment: boolean; suspectKind: DuplicateMatchKind | null; }, ): Omit | null { if (a.hasTranscript && b.hasTranscript) { const similar = evalSimilar(a, b, opts.nearThreshold); if (similar) return similar; // Containment as a FALLBACK inside a pair we already nominated: it costs one // more set intersection over shingle sets that are already in hand, so it is // free relative to the nomination that got us here. if (opts.allowContainment) { const contained = evalContainment(a, b, opts.containmentThreshold); if (contained) return contained; } return null; } if (!opts.suspectKind) return null; return { kind: opts.suspectKind, score: null, contained: false }; } // Containment pair: the shorter transcript's shingles are largely a subset of // the longer one's. Requires both transcripts (no metadata fallback here). function evalContainment( a: Fingerprint, b: Fingerprint, containmentThreshold: number, ): Omit | null { if (!a.hasTranscript || !b.hasTranscript) return null; const c = containment(a.shingleSet, b.shingleSet); if (c >= containmentThreshold) { return { kind: "transcript-near", score: round3(c), contained: true }; } return null; } // Built from the participant index rather than a fingerprint: fingerprints are // released with their block, and VideoStat is the small half of what they held. // `offsetSeconds` / `aligned` are filled in afterwards by the alignment pass. function toRef(e: Entry, hasTranscript: boolean): DuplicateVideoRef { const s = e.stat; return { slug: s.slug, channelSlug: s.channelSlug, channel: s.channel, platform: s.platform, id: s.id, // Only when it differs — see DuplicateVideoRef.videoDir. The detector is the // one place that has both halves in hand, and dropping the directory here is // what silently broke digest cluster-sharing for every Rumble mirror. ...(e.videoDir !== s.id ? { videoDir: e.videoDir } : {}), title: s.title, duration: s.duration, uploadDate: s.uploadDate, hasTranscript, }; } function maxScore(a: number | null, b: number | null): number | null { if (a === null) return b; if (b === null) return a; return Math.max(a, b); } function round3(n: number): number { return Math.round(n * 1000) / 1000; } function sha1(s: string): string { return createHash("sha1").update(s).digest("hex"); }