import json import time from pathlib import Path from sqlalchemy import or_, select from app.db import SessionLocal from app.models import DedupCandidate, DedupRun from app.services import pipeline_runner from app.settings import settings _SCRIPT = "dedup-library.sh" _FUZZY_SCRIPT = "find-fuzzy-dupes.py" _FUZZY_PASS = "fuzzy_audio" # Generous ceiling for any pipeline script invoked here. Without one, a # stuck subprocess (e.g. a stalled stat() on the NAS mount) holds the single # global pipeline lock forever -- 2026-07-08 incident: one stuck dedup:apply # call held the lock for 16+ hours, silently skipping every later confirm # click, each of which still got marked confirmed=True (see confirm_and_apply) # and vanished from the pending list without anything actually being deleted. _JOB_TIMEOUT_SECONDS = 1800 def _script_for_pass(pass_name: str) -> str: return _FUZZY_SCRIPT if pass_name == _FUZZY_PASS else _SCRIPT def _is_pair_ignored(db, path_a: str, path_b: str) -> bool: """True if this path pair was ever marked 'keep both', regardless of which path was on the keep/delete side that time -- a later scan can flip the ranking (e.g. file sizes changed) without changing the fact that the user already decided this pair is fine as two copies.""" query = select(DedupCandidate.id).where( DedupCandidate.ignored == True, # noqa: E712 or_( (DedupCandidate.keep_path == path_a) & (DedupCandidate.delete_path == path_b), (DedupCandidate.keep_path == path_b) & (DedupCandidate.delete_path == path_a), ), ).limit(1) return db.execute(query).first() is not None def _parse_json_lines(output: str) -> list[dict]: candidates = [] for line in output.splitlines(): line = line.strip() if not line.startswith("{"): continue try: candidates.append(json.loads(line)) except json.JSONDecodeError: continue return candidates async def _run_scan(job_key: str, script_name: str, triggered_by: str) -> DedupRun: script = str(settings.pipeline_dir / "lib" / script_name) job_run, output = await pipeline_runner.run_job_capture( job_key, [script, "--json"], triggered_by=triggered_by, timeout=_JOB_TIMEOUT_SECONDS ) candidates = _parse_json_lines(output) db = SessionLocal() try: dedup_run = DedupRun( started_at=job_run.started_at, finished_at=job_run.finished_at, mode="dry_run", groups_found=len({(c["pass"], c["keep_path"]) for c in candidates}), kept=len({(c["pass"], c["keep_path"]) for c in candidates}), deleted=0, log_path=job_run.log_path, ) db.add(dedup_run) db.commit() db.refresh(dedup_run) skipped_ignored = 0 for c in candidates: if _is_pair_ignored(db, c["keep_path"], c["delete_path"]): skipped_ignored += 1 continue db.add( DedupCandidate( dedup_run_id=dedup_run.id, pass_name=c.get("pass", "unknown"), keep_path=c["keep_path"], delete_path=c["delete_path"], delete_size_bytes=c.get("delete_size_bytes"), ) ) db.commit() db.refresh(dedup_run) if skipped_ignored: dedup_run.kept = (dedup_run.kept or 0) + skipped_ignored db.commit() db.refresh(dedup_run) return dedup_run finally: db.close() async def scan(triggered_by: str = "manual") -> DedupRun: """Dry-run dedup-library.sh --json, persist every candidate deletion into a fresh dedup_runs/dedup_candidates pair. Never deletes anything -- the scheduled maintenance:dedup job also only ever calls this (no --apply), matching the false-negative-biased dedup preference; actual deletion only ever happens through confirm_and_apply() below.""" return await _run_scan("dedup:scan", _SCRIPT, triggered_by) async def scan_fuzzy(triggered_by: str = "manual") -> DedupRun: """Dry-run find-fuzzy-dupes.py --json. Tag-based passes (scan() above) only catch duplicates whose artist/title tags overlap; this compares Chromaprint audio fingerprints instead, so it also catches the same recording filed under different tags (e.g. a remix credited to different artists between two copies).""" return await _run_scan("dedup:scan_fuzzy", _FUZZY_SCRIPT, triggered_by) async def scan_all(triggered_by: str = "manual") -> tuple[DedupRun, DedupRun]: """Run both passes back to back for a single "Scan now" click -- each call acquires and releases the shared pipeline lock on its own, so running them sequentially here is safe, just one after the other.""" tag_run = await scan(triggered_by=triggered_by) fuzzy_run = await scan_fuzzy(triggered_by=triggered_by) return tag_run, fuzzy_run def _write_apply_args(script_name: str, script: str, group: list[DedupCandidate], stamp: int) -> list[str]: """Build the --apply invocation for one pass's confirmed group. find-fuzzy-dupes.py gets --apply-pairs: it applies the exact keep/delete pairs the user already reviewed without re-scanning the whole library for acoustic matches (that rescan is what hung for 16+ hours on 2026-07-08, holding the global pipeline lock). dedup-library.sh has no equivalent fast path -- its tag-based passes are cheap to rescan, so --apply --only-paths (full rescan, filtered to the confirmed list) is fine there. """ if script_name == _FUZZY_SCRIPT: pairs_file = settings.logs_dir / f"dedup-confirm-{stamp}-{script_name}.jsonl" pairs_file.parent.mkdir(parents=True, exist_ok=True) pairs_file.write_text( "\n".join( json.dumps({"keep_path": c.keep_path, "delete_path": c.delete_path}) for c in group ) + "\n" ) return [script, "--apply-pairs", str(pairs_file), "--json"] confirm_file = settings.logs_dir / f"dedup-confirm-{stamp}-{script_name}.txt" confirm_file.parent.mkdir(parents=True, exist_ok=True) confirm_file.write_text("\n".join(c.delete_path for c in group) + "\n") return [script, "--apply", "--only-paths", str(confirm_file), "--json"] async def confirm_and_apply(candidate_ids: list[int], confirmed_by: str) -> DedupRun | None: """Apply only the confirmed candidate deletions. Candidates are only marked confirmed=True once we know their apply job actually ran (status == "success") -- if the shared pipeline lock is held by something else, or the job times out/crashes, those candidates are left untouched so they stay visible in the pending list instead of silently vanishing without being deleted (see 2026-07-08 incident). Cheap pre-check here: skip anything where delete_path or keep_path no longer exists (something already changed it since the scan). Further re-verification of ranking happens for free inside the underlying script itself -- see _write_apply_args(). dedup-library.sh (tag-based passes) and find-fuzzy-dupes.py (acoustic fingerprint pass) are separate scripts -- candidates are grouped by pass_name and each group is applied through the script that actually produced it. """ db = SessionLocal() try: candidates = [db.get(DedupCandidate, cid) for cid in candidate_ids] candidates = [c for c in candidates if c is not None and not c.applied] candidates = [c for c in candidates if Path(c.delete_path).exists() and Path(c.keep_path).exists()] if not candidates: return None by_script: dict[str, list[DedupCandidate]] = {} for c in candidates: by_script.setdefault(_script_for_pass(c.pass_name), []).append(c) now = time.time() finished_at = now log_paths = [] ran_candidates: list[DedupCandidate] = [] for script_name, group in by_script.items(): script = str(settings.pipeline_dir / "lib" / script_name) argv = _write_apply_args(script_name, script, group, int(now * 1000)) job_run, _output = await pipeline_runner.run_job_capture( "dedup:apply", argv, triggered_by=f"manual:{confirmed_by}", timeout=_JOB_TIMEOUT_SECONDS, ) finished_at = job_run.finished_at or finished_at if job_run.log_path: log_paths.append(job_run.log_path) if job_run.status != "success": # Lock contention, crash, or timeout -- nothing in this # group was actually run. Leave confirmed=False so these # stay in the pending list for a retry. continue for c in group: c.confirmed = True c.confirmed_by = confirmed_by c.confirmed_at = now ran_candidates.extend(group) db.commit() if not ran_candidates: return None still_there = {c.delete_path for c in ran_candidates if Path(c.delete_path).exists()} actually_deleted = 0 for c in ran_candidates: if c.delete_path not in still_there: c.applied = True actually_deleted += 1 db.commit() apply_run = DedupRun( started_at=now, finished_at=finished_at, mode="apply", deleted=actually_deleted, kept=len(ran_candidates) - actually_deleted, log_path=";".join(log_paths) or None, ) db.add(apply_run) db.commit() db.refresh(apply_run) return apply_run finally: db.close() def list_pending_candidates(dedup_run_id: int | None = None) -> list[DedupCandidate]: db = SessionLocal() try: query = select(DedupCandidate).where( DedupCandidate.applied == False, # noqa: E712 DedupCandidate.confirmed == False, # noqa: E712 DedupCandidate.ignored == False, # noqa: E712 ) if dedup_run_id is not None: query = query.where(DedupCandidate.dedup_run_id == dedup_run_id) return list(db.execute(query).scalars()) finally: db.close() def list_ignored_candidates() -> list[DedupCandidate]: db = SessionLocal() try: query = select(DedupCandidate).where(DedupCandidate.ignored == True).order_by( # noqa: E712 DedupCandidate.ignored_at.desc() ) return list(db.execute(query).scalars()) finally: db.close() def ignore_candidate(candidate_id: int, ignored_by: str) -> DedupCandidate | None: """Mark a pending candidate 'keep both' -- it drops off the pending list immediately, and future scans skip re-creating a candidate for the same path pair (see _is_pair_ignored).""" db = SessionLocal() try: c = db.get(DedupCandidate, candidate_id) if c is None or c.applied: return None c.ignored = True c.ignored_by = ignored_by c.ignored_at = time.time() db.commit() db.refresh(c) return c finally: db.close() def unignore_candidate(candidate_id: int) -> DedupCandidate | None: """Undo a 'keep both' -- the candidate returns to the pending list and future scans are free to re-flag the same pair again.""" db = SessionLocal() try: c = db.get(DedupCandidate, candidate_id) if c is None: return None c.ignored = False c.ignored_by = None c.ignored_at = None db.commit() db.refresh(c) return c finally: db.close()