Files
alembic/app/services/dedup_review_service.py
T
andrew bcd07750bf R7: consolidate duplicated Spotify-token and NDJSON-parse logic
- The client-credentials token fetch was copy-pasted across spotify-retag.py,
  spotify-genre.py, and fix-track-metadata.py (and the app's spotify_client).
  Add pipeline/lib/_spotify_auth.get_token (cached per id/secret); the three
  scripts now source their own credentials but delegate the request to it. The
  scripts run with pipeline/lib on sys.path, so the plain `from _spotify_auth
  import get_token` resolves.
- The identical _parse_json_lines helper in dedup_review_service and
  genre_review_service is now a single app/services/_ndjson.parse_json_lines.

Verified: unit test of the token helper (cache + request), the NDJSON parser
(tests/test_ndjson.py), full suite green (40), and a live spotify-genre dry-run
that fetched a token and queried Spotify through the shared helper.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-10 16:08:15 -06:00

337 lines
14 KiB
Python

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.services._ndjson import parse_json_lines
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 _is_pair_already_pending(db, keep_path: str, delete_path: str) -> bool:
"""True if this exact keep/delete pair is already sitting in the
pending list from an earlier scan. Without this, every re-scan of a
duplicate the user hasn't reviewed yet (e.g. the daily scheduled scan)
added a brand-new row for the same pair, so it piled up multiple
identical entries in the table instead of just staying as one."""
query = select(DedupCandidate.id).where(
DedupCandidate.applied == False, # noqa: E712
DedupCandidate.confirmed == False, # noqa: E712
DedupCandidate.ignored == False, # noqa: E712
DedupCandidate.keep_path == keep_path,
DedupCandidate.delete_path == delete_path,
).limit(1)
return db.execute(query).first() is not None
async def _run_scan(job_key: str, script_name: str, triggered_by: str) -> DedupRun | None:
"""Returns None (persisting nothing) if the job never actually ran --
e.g. skipped_lock because something else was using the pipeline lock
at that moment. Same lesson as genre_review_service.run(): recording a
"run happened" row for a run that didn't is just misleading history,
even though (unlike genres) it isn't user-visible here today since the
pending-candidates list isn't scoped to "latest run only"."""
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
)
if job_run.status != "success":
return None
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
skipped_duplicate = 0
for c in candidates:
if _is_pair_ignored(db, c["keep_path"], c["delete_path"]):
skipped_ignored += 1
continue
if _is_pair_already_pending(db, c["keep_path"], c["delete_path"]):
skipped_duplicate += 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_id=c.get("delete_id"),
delete_size_bytes=c.get("delete_size_bytes"),
)
)
# Flush (not commit) so a duplicate pair emitted twice within
# this same scan's own output -- e.g. two passes agreeing on
# the same file -- is caught by the next iteration's check too,
# not just duplicates from a previous scan's committed rows.
db.flush()
db.commit()
db.refresh(dedup_run)
skipped_total = skipped_ignored + skipped_duplicate
if skipped_total:
dedup_run.kept = (dedup_run.kept or 0) + skipped_total
db.commit()
db.refresh(dedup_run)
return dedup_run
finally:
db.close()
async def scan(triggered_by: str = "manual") -> DedupRun | None:
"""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 | None:
"""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 | None, DedupRun | None]:
"""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.
Either element is None if that particular pass got skipped_lock/failed
(see _run_scan) -- e.g. the fuzzy pass can't start until the tag-based
pass above it has released the lock, so it's not unusual for one to
succeed and the other to be skipped by something else entirely."""
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, "delete_id": c.delete_id})
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()