diff --git a/src/app.rs b/src/app.rs index b86dd19..66accb5 100644 --- a/src/app.rs +++ b/src/app.rs @@ -43,6 +43,23 @@ enum Screen { Maintenance, } +/// One video in a perceptual-similarity group (desktop dedup review). +#[derive(Clone)] +struct SimVideo { + video_id: String, + title: String, + channel: String, + file_size: Option, + files: Vec, + recommended_keep: bool, +} + +/// A cluster of videos that share visual content across different IDs. +#[derive(Clone)] +struct SimGroup { + videos: Vec, +} + #[derive(Clone, PartialEq)] enum SortMode { Title, @@ -168,6 +185,13 @@ pub struct App { transcript_video: Option, transcript_cues: Vec, transcript_query: String, + // Perceptual dedup (background fingerprint job + results) + dedup_running: bool, + dedup_started: bool, + dedup_progress: std::sync::Arc<(std::sync::atomic::AtomicUsize, std::sync::atomic::AtomicUsize)>, + dedup_rx: Option, String>>>, + dedup_groups: Vec, + dedup_error: Option, // Scheduler last_scheduled_check: Option, // Cards cache — recomputed only when inputs change @@ -496,6 +520,15 @@ impl App { transcript_video: None, transcript_cues: Vec::new(), transcript_query: String::new(), + dedup_running: false, + dedup_started: false, + dedup_progress: std::sync::Arc::new(( + std::sync::atomic::AtomicUsize::new(0), + std::sync::atomic::AtomicUsize::new(0), + )), + dedup_rx: None, + dedup_groups: Vec::new(), + dedup_error: None, last_scheduled_check: None, cards_cache: Vec::new(), cards_cache_key: None, @@ -871,6 +904,69 @@ impl App { } } + /// Kick off the background perceptual-dedup job: fingerprint new/changed + /// videos (cached by mtime), group by visual similarity, deliver the + /// result over a channel. No-op if a job is already running. + fn start_dedup(&mut self) { + if self.dedup_running { return; } + let mut inputs = Vec::new(); + let mut by_path: HashMap = HashMap::new(); + let mut valid_paths: HashSet = HashSet::new(); + for ch in &self.library { + let channel = ch.name.clone(); + for v in ch.videos.iter().chain(ch.playlists.iter().flat_map(|p| p.videos.iter())) { + let Some(vp) = &v.video_path else { continue }; + let path_str = vp.display().to_string(); + valid_paths.insert(path_str.clone()); + inputs.push(crate::fingerprint::FpInput { + path: vp.clone(), + mtime_unix: v.mtime_unix.map(|m| m as i64).unwrap_or(0), + video_id: v.id.clone(), + duration_secs: v.duration_secs.unwrap_or(0.0), + }); + let mut files = vec![vp.clone()]; + for p in [v.thumb_path.as_ref(), v.info_path.as_ref(), v.description_path.as_ref()] + .into_iter().flatten() { files.push(p.clone()); } + for s in &v.subtitles { files.push(s.path.clone()); } + by_path.insert(path_str, SimVideo { + video_id: v.id.clone(), title: v.title.clone(), channel: channel.clone(), + file_size: v.file_size, files, recommended_keep: false, + }); + } + } + use std::sync::atomic::Ordering; + self.dedup_progress.0.store(0, Ordering::Relaxed); + self.dedup_progress.1.store(0, Ordering::Relaxed); + self.dedup_running = true; + self.dedup_started = true; + self.dedup_error = None; + self.dedup_groups.clear(); + + let db = self.db.clone(); + let progress = self.dedup_progress.clone(); + let (tx, rx) = std::sync::mpsc::channel(); + self.dedup_rx = Some(rx); + let workers = std::thread::available_parallelism().map(|n| n.get()).unwrap_or(4); + std::thread::spawn(move || { + let res = crate::fingerprint::rebuild_and_group( + &db, inputs, &valid_paths, workers, &progress.0, &progress.1, + ).map(|path_groups| { + let mut groups = Vec::new(); + for paths in path_groups { + let mut vids: Vec = + paths.iter().filter_map(|p| by_path.get(p).cloned()).collect(); + if vids.len() < 2 { continue; } + let keep = vids.iter().enumerate() + .max_by_key(|(_, v)| v.file_size.unwrap_or(0)).map(|(i, _)| i); + for (i, v) in vids.iter_mut().enumerate() { v.recommended_keep = Some(i) == keep; } + groups.push(SimGroup { videos: vids }); + } + groups + }); + let _ = tx.send(res); + }); + } + /// Find a video by id anywhere in the library and play it. Used by the /// full-text search results, where we only have the id. fn play_by_id(&mut self, id: &str) { @@ -1942,11 +2038,35 @@ impl App { fn maintenance_screen(&mut self, ctx: &egui::Context) { let report = self.health_report.clone().unwrap_or_default(); + // Drain the background dedup result if it's ready. + let mut dedup_done = None; + if let Some(rx) = &self.dedup_rx { + if let Ok(res) = rx.try_recv() { dedup_done = Some(res); } + } + if let Some(res) = dedup_done { + match res { + Ok(groups) => self.dedup_groups = groups, + Err(e) => self.dedup_error = Some(e), + } + self.dedup_running = false; + self.dedup_rx = None; + } + // Actions are collected during rendering and applied after the closure // to avoid borrowing `self` while the report is borrowed immutably. let mut to_remove: Vec = Vec::new(); let mut to_repair: Vec = Vec::new(); let mut rescan_health = false; + let mut start_dedup = false; + let mut dedup_remove: Vec = Vec::new(); + let dedup_groups = self.dedup_groups.clone(); + let dedup_running = self.dedup_running; + let dedup_started = self.dedup_started; + let dedup_error = self.dedup_error.clone(); + let (dedup_done_n, dedup_total_n) = { + use std::sync::atomic::Ordering; + (self.dedup_progress.0.load(Ordering::Relaxed), self.dedup_progress.1.load(Ordering::Relaxed)) + }; egui::CentralPanel::default().show(ctx, |ui| { ui.horizontal(|ui| { @@ -2026,6 +2146,55 @@ impl App { } }); } + + // ── Similar content (perceptual dedup) ────────────────── + ui.add_space(12.0); + ui.heading("Similar content (perceptual)"); + ui.label(egui::RichText::new( + "Finds the same video re-uploaded under a different ID — mirrors, \ + re-encodes, resolution changes — by comparing sampled frames. The \ + first scan fingerprints your library (a few minutes); it's cached \ + after, so re-scans are quick.").weak().small()); + if dedup_running { + ui.horizontal(|ui| { + ui.spinner(); + let t = if dedup_total_n > 0 { dedup_total_n.to_string() } else { "…".into() }; + ui.label(format!("Fingerprinting {dedup_done_n} / {t} new videos…")); + }); + ctx.request_repaint(); // keep progress live + } else { + if ui.button("🔍 Scan for similar content").clicked() { + start_dedup = true; + } + if let Some(err) = &dedup_error { + ui.colored_label(egui::Color32::from_rgb(0xf8, 0x71, 0x71), + format!("Last scan error: {err}")); + } + if dedup_started { + ui.label(egui::RichText::new(format!( + "{} similar group(s)", dedup_groups.len())).weak()); + } + for (gi, g) in dedup_groups.iter().enumerate() { + ui.group(|ui| { + ui.label(egui::RichText::new( + format!("Group {} — {} copies", gi + 1, g.videos.len())).strong()); + for v in &g.videos { + let size = v.file_size.map(format_size) + .unwrap_or_else(|| "no video".to_string()); + let tag = if v.recommended_keep { "✓ keep" } else { "✗ remove" }; + ui.label(format!(" {} · {} · {} — {}", + v.title, v.channel, size, tag)); + } + if ui.button("🗑 Remove non-recommended copies").clicked() { + for v in &g.videos { + if !v.recommended_keep { + dedup_remove.extend(v.files.iter().cloned()); + } + } + } + }); + } + } }); }); @@ -2050,11 +2219,25 @@ impl App { self.status = format!("Queued {} repair(s) — see Downloads", to_repair.len()); self.show_downloads = true; } + if !dedup_remove.is_empty() { + let (removed, errors) = + crate::maintenance::remove_files(&self.library_root, &dedup_remove); + self.status = if errors.is_empty() { + format!("Removed {removed} file(s)") + } else { + format!("Removed {removed} file(s), {} error(s)", errors.len()) + }; + changed = true; + start_dedup = true; // re-scan so the deleted copies drop out + } if changed || rescan_health { self.rescan(); self.health_report = Some(crate::maintenance::scan(&self.library_root, &self.library)); } + if start_dedup { + self.start_dedup(); + } } fn stats_screen(&mut self, ctx: &egui::Context) { diff --git a/src/fingerprint.rs b/src/fingerprint.rs index fa333fc..b81efac 100644 --- a/src/fingerprint.rs +++ b/src/fingerprint.rs @@ -265,6 +265,49 @@ pub const DEFAULT_DUR_TOL: f64 = 3.0; pub const DEFAULT_MAX_HAM: u32 = 8; pub const DEFAULT_MIN_MATCH: usize = 3; +/// The whole dedup pipeline, shared by both front-ends: mtime-gate `inputs` +/// against what's stored, fingerprint the new/changed ones in parallel +/// (bumping `progress`, after first storing the count in `total`), upsert, +/// prune anything not in `valid_paths`, then group by visual similarity. +/// Returns groups (≥2) of **file paths**; each UI maps those to its own +/// display rows. Errors are DB errors stringified. +pub fn rebuild_and_group( + db: &crate::database::Database, + inputs: Vec, + valid_paths: &std::collections::HashSet, + workers: usize, + progress: &std::sync::atomic::AtomicUsize, + total: &std::sync::atomic::AtomicUsize, +) -> Result>, String> { + use std::sync::atomic::Ordering; + let known = db.fingerprint_mtimes().map_err(|e| e.to_string())?; + let todo: Vec = inputs + .into_iter() + .filter(|i| known.get(&i.path.display().to_string()) != Some(&i.mtime_unix)) + .collect(); + total.store(todo.len(), Ordering::Relaxed); + + let computed = compute_batch(todo, workers, progress); + for c in &computed { + let _ = db.upsert_fingerprint( + &c.input.path.display().to_string(), c.input.mtime_unix, + &c.input.video_id, c.input.duration_secs, &c.hashes, + ); + } + let _ = db.prune_fingerprints(valid_paths); + + let stored = db.load_fingerprints().map_err(|e| e.to_string())?; + let records: Vec = stored.iter().map(|s| FpRecord { + video_id: s.video_id.clone(), duration_secs: s.duration_secs, hashes: s.hashes.clone(), + }).collect(); + let groups_idx = + group_similar(&records, DEFAULT_DUR_TOL, DEFAULT_MAX_HAM, DEFAULT_MIN_MATCH); + Ok(groups_idx + .into_iter() + .map(|g| g.into_iter().map(|i| stored[i].path.clone()).collect()) + .collect()) +} + #[cfg(test)] mod tests { use super::*; diff --git a/src/web.rs b/src/web.rs index 6b75594..b3c1237 100644 --- a/src/web.rs +++ b/src/web.rs @@ -2090,47 +2090,23 @@ fn run_dedup( by_path: HashMap, valid_paths: HashSet, ) { - use crate::fingerprint; - let outcome: Result, String> = (|| { - let known = db.fingerprint_mtimes().map_err(|e| e.to_string())?; - let todo: Vec = inputs - .into_iter() - .filter(|i| known.get(&i.path.display().to_string()) != Some(&i.mtime_unix)) - .collect(); - dedup.total.store(todo.len(), Ordering::Relaxed); - - let workers = std::thread::available_parallelism().map(|n| n.get()).unwrap_or(4); - let computed = fingerprint::compute_batch(todo, workers, &dedup.done); - for c in &computed { - let _ = db.upsert_fingerprint( - &c.input.path.display().to_string(), c.input.mtime_unix, - &c.input.video_id, c.input.duration_secs, &c.hashes, - ); - } - let _ = db.prune_fingerprints(&valid_paths); - - let stored = db.load_fingerprints().map_err(|e| e.to_string())?; - let records: Vec = stored.iter().map(|s| fingerprint::FpRecord { - video_id: s.video_id.clone(), duration_secs: s.duration_secs, hashes: s.hashes.clone(), - }).collect(); - let groups_idx = fingerprint::group_similar( - &records, fingerprint::DEFAULT_DUR_TOL, - fingerprint::DEFAULT_MAX_HAM, fingerprint::DEFAULT_MIN_MATCH, - ); - - let mut groups = Vec::new(); - for g in groups_idx { - let mut vids: Vec = - g.iter().filter_map(|&i| by_path.get(&stored[i].path).cloned()).collect(); - if vids.len() < 2 { continue; } // stale paths dropped out - // Recommend keeping the largest file (best quality) in each group. - let keep = vids.iter().enumerate() - .max_by_key(|(_, v)| v.file_size.unwrap_or(0)).map(|(i, _)| i); - for (i, v) in vids.iter_mut().enumerate() { v.recommended_keep = Some(i) == keep; } - groups.push(SimilarGroup { videos: vids }); - } - Ok(groups) - })(); + let workers = std::thread::available_parallelism().map(|n| n.get()).unwrap_or(4); + let outcome: Result, String> = + crate::fingerprint::rebuild_and_group(&db, inputs, &valid_paths, workers, &dedup.done, &dedup.total) + .map(|path_groups| { + let mut groups = Vec::new(); + for paths in path_groups { + let mut vids: Vec = + paths.iter().filter_map(|p| by_path.get(p).cloned()).collect(); + if vids.len() < 2 { continue; } // stale paths dropped out + // Recommend keeping the largest file (best quality). + let keep = vids.iter().enumerate() + .max_by_key(|(_, v)| v.file_size.unwrap_or(0)).map(|(i, _)| i); + for (i, v) in vids.iter_mut().enumerate() { v.recommended_keep = Some(i) == keep; } + groups.push(SimilarGroup { videos: vids }); + } + groups + }); match outcome { Ok(groups) => *dedup.result.lock_recover() = Some(groups),