2075d6cf66
dedup-library.sh: additive --json (one NDJSON line per candidate deletion to stdout, alongside the unchanged human log) and --only-paths FILE (in --apply mode, only actually delete entries whose path is in FILE; without it, --apply deletes everything as before -- existing direct callers are unaffected). Without --only-paths the safety re-verification is free: --apply --only-paths re-runs all 4 passes from scratch on every invocation, so if a group's ranking changed since a scan (e.g. the old keep_path is gone), the fresh pass assigns the previously-"delete" path the KEEP role instead and the only-paths allowlist naming it is simply never consulted -- no duplicate ranking logic needed in the review service. spotify-genre.py: additive --json emitting one JSON line per genre change (dry-run or --apply) for genre_review_service to persist. pipeline_runner.run_job_capture(): like run_job() but captures stdout as text (still under the same shared lock, still writes a job_runs row) for callers that need to parse structured output rather than just log it. dedup_review_service.scan() persists dry-run candidates into dedup_runs/dedup_candidates. confirm_and_apply() re-checks confirmed candidates still exist before invoking --apply --only-paths, so nothing is ever deleted without an explicit confirm -- matches the false-negative- biased dedup preference. scheduler_service's maintenance:dedup job now goes through this (still dry-run only, every day). genre_review_service.run() wraps spotify-genre.py for both dry-run preview and the real scheduled --apply run, persisting every run's diff into genre_runs/genre_candidates either way -- genre writes keep their current auto-apply behavior (low-risk, reversible, GENRE_LOCK-protected) but are now reviewable after the fact. lock_artist_genre() gives a one-click revert path when a run gets something wrong. Added minimal routers+templates for /dedup (scan, review, confirm-and- delete) and /genres (preview, review, lock-old-genre). Verified end-to-end against REAL duplicate files (not mocked): built an actual FLAC+MP3 duplicate pair in a real beets library, ran dedup-library.sh --json and confirmed correct JSON output, verified --apply --only-paths with an empty confirm list deletes nothing and with the real confirmed path deletes exactly that file (DB + disk) while preserving the FLAC, and ran the full dedup_review_service scan->confirm->apply flow through the same fixture. genre_review_service and spotify-genre.py --json verified against mocked/direct output (spotify-genre.py's own artist-genre lookup needs a live Spotify API call, out of reach in this sandbox). Confirmed the full app boots with all five routers registered.
258 lines
9.2 KiB
Python
258 lines
9.2 KiB
Python
import asyncio
|
|
import time
|
|
from pathlib import Path
|
|
|
|
from app.db import SessionLocal
|
|
from app.models import JobRun
|
|
from app.settings import settings
|
|
|
|
# Single global lock replacing /var/lock/sldl-pipeline.lock + flock -n. One
|
|
# process (this one) now owns every pipeline invocation, so one asyncio.Lock
|
|
# is enough -- deliberately not per-resource, matching the original flock's
|
|
# all-jobs-share-one-lock behavior exactly rather than over-engineering it.
|
|
_lock = asyncio.Lock()
|
|
|
|
_ENV_PASSTHROUGH_KEYS = ("PATH", "HOME", "LANG", "LC_ALL", "TZ")
|
|
|
|
|
|
async def _try_acquire_nowait() -> bool:
|
|
"""asyncio.Lock has no acquire_nowait(). asyncio.wait_for(lock.acquire(),
|
|
timeout=0) looks like the obvious idiom but is broken in practice: it
|
|
wraps acquire() in a Task, and the Task's first iteration and the
|
|
timeout-0 callback are both scheduled via the event loop with no
|
|
guaranteed ordering, so it times out almost every time even when the
|
|
lock is completely uncontended (confirmed empirically -- the very first,
|
|
uncontended call failed every time). The correct pattern relies on
|
|
Lock.acquire()'s fast path never suspending when the lock is free: the
|
|
.locked() check and the subsequent acquire() happen with no `await`
|
|
point in between, so nothing else can interleave in this single-
|
|
threaded event loop."""
|
|
if _lock.locked():
|
|
return False
|
|
await _lock.acquire()
|
|
return True
|
|
|
|
|
|
def _subprocess_env() -> dict:
|
|
import os
|
|
|
|
env = {k: v for k, v in os.environ.items() if k in _ENV_PASSTHROUGH_KEYS}
|
|
env.setdefault("PATH", "/opt/venv/bin:/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin")
|
|
env["ALEMBIC_CONFIG_DIR"] = str(settings.alembic_config_dir)
|
|
env["MUSIC_DATA_DIR"] = str(settings.music_data_dir)
|
|
env["PIPELINE_DIR"] = str(settings.pipeline_dir)
|
|
env["BEETSDIR"] = str(settings.beets_dir)
|
|
return env
|
|
|
|
|
|
def _summarize_log(log_path: Path, max_len: int = 200) -> str:
|
|
"""Best-effort one-liner for the jobs list: the log's last non-empty
|
|
line. Scripts here already end their own runs with a human-readable
|
|
'=== ... done ===' / summary line (see pipeline-status.sh's convention),
|
|
so this is usually meaningful without any per-script special-casing."""
|
|
try:
|
|
lines = [line for line in log_path.read_text(errors="replace").splitlines() if line.strip()]
|
|
return lines[-1][:max_len] if lines else ""
|
|
except OSError:
|
|
return ""
|
|
|
|
|
|
async def run_job(
|
|
job_key: str,
|
|
argv: list[str],
|
|
triggered_by: str = "schedule",
|
|
timeout: float | None = None,
|
|
) -> JobRun:
|
|
"""Run a pipeline command under the shared mutual-exclusion lock,
|
|
recording one job_runs row start-to-finish. If the lock is already
|
|
held, records status='skipped_lock' immediately and returns without
|
|
running anything -- the old flock -n behavior, now visible in the UI
|
|
instead of silently skipping."""
|
|
started_at = time.time()
|
|
|
|
if not await _try_acquire_nowait():
|
|
db = SessionLocal()
|
|
run = JobRun(
|
|
job_key=job_key,
|
|
started_at=started_at,
|
|
finished_at=started_at,
|
|
status="skipped_lock",
|
|
triggered_by=triggered_by,
|
|
)
|
|
db.add(run)
|
|
db.commit()
|
|
db.refresh(run)
|
|
db.close()
|
|
return run
|
|
|
|
try:
|
|
log_dir = settings.logs_dir / job_key.replace(":", "_")
|
|
log_dir.mkdir(parents=True, exist_ok=True)
|
|
log_path = log_dir / f"{time.strftime('%Y%m%d-%H%M%S')}.log"
|
|
|
|
db = SessionLocal()
|
|
run = JobRun(
|
|
job_key=job_key,
|
|
started_at=started_at,
|
|
status="running",
|
|
triggered_by=triggered_by,
|
|
log_path=str(log_path),
|
|
)
|
|
db.add(run)
|
|
db.commit()
|
|
db.refresh(run)
|
|
run_id = run.id
|
|
db.close()
|
|
|
|
exit_code: int | None
|
|
try:
|
|
with open(log_path, "wb") as log_file:
|
|
proc = await asyncio.create_subprocess_exec(
|
|
*argv,
|
|
stdout=log_file,
|
|
stderr=asyncio.subprocess.STDOUT,
|
|
env=_subprocess_env(),
|
|
)
|
|
try:
|
|
exit_code = await asyncio.wait_for(proc.wait(), timeout=timeout)
|
|
except asyncio.TimeoutError:
|
|
proc.kill()
|
|
await proc.wait()
|
|
exit_code = -1
|
|
with open(log_path, "a") as f:
|
|
f.write(f"\n[pipeline_runner] TIMEOUT after {timeout}s -- process killed\n")
|
|
status = "success" if exit_code == 0 else "failed"
|
|
except Exception as exc:
|
|
exit_code = None
|
|
status = "failed"
|
|
with open(log_path, "a") as f:
|
|
f.write(f"\n[pipeline_runner] exception before/while running: {exc!r}\n")
|
|
|
|
finished_at = time.time()
|
|
db = SessionLocal()
|
|
run = db.get(JobRun, run_id)
|
|
run.finished_at = finished_at
|
|
run.status = status
|
|
run.exit_code = exit_code
|
|
run.summary = _summarize_log(log_path)
|
|
db.commit()
|
|
db.refresh(run)
|
|
db.close()
|
|
return run
|
|
finally:
|
|
_lock.release()
|
|
|
|
|
|
async def run_job_capture(
|
|
job_key: str,
|
|
argv: list[str],
|
|
triggered_by: str = "manual",
|
|
timeout: float | None = None,
|
|
) -> tuple[JobRun, str]:
|
|
"""Like run_job(), but captures stdout+stderr as text and returns it
|
|
alongside the JobRun, instead of only writing it to the log file --
|
|
for callers that need to parse structured (JSON-lines) output, e.g.
|
|
dedup_review_service's dry-run scan. The full output is still written
|
|
to a log file afterward so job_runs.log_path works the same as any
|
|
other job. Shares the same lock as run_job()."""
|
|
started_at = time.time()
|
|
|
|
if not await _try_acquire_nowait():
|
|
db = SessionLocal()
|
|
run = JobRun(
|
|
job_key=job_key,
|
|
started_at=started_at,
|
|
finished_at=started_at,
|
|
status="skipped_lock",
|
|
triggered_by=triggered_by,
|
|
)
|
|
db.add(run)
|
|
db.commit()
|
|
db.refresh(run)
|
|
db.close()
|
|
return run, ""
|
|
|
|
try:
|
|
log_dir = settings.logs_dir / job_key.replace(":", "_")
|
|
log_dir.mkdir(parents=True, exist_ok=True)
|
|
log_path = log_dir / f"{time.strftime('%Y%m%d-%H%M%S')}.log"
|
|
|
|
db = SessionLocal()
|
|
run = JobRun(
|
|
job_key=job_key,
|
|
started_at=started_at,
|
|
status="running",
|
|
triggered_by=triggered_by,
|
|
log_path=str(log_path),
|
|
)
|
|
db.add(run)
|
|
db.commit()
|
|
db.refresh(run)
|
|
run_id = run.id
|
|
db.close()
|
|
|
|
output_text = ""
|
|
exit_code: int | None
|
|
try:
|
|
proc = await asyncio.create_subprocess_exec(
|
|
*argv,
|
|
stdout=asyncio.subprocess.PIPE,
|
|
stderr=asyncio.subprocess.STDOUT,
|
|
env=_subprocess_env(),
|
|
)
|
|
try:
|
|
stdout_bytes, _ = await asyncio.wait_for(proc.communicate(), timeout=timeout)
|
|
exit_code = proc.returncode
|
|
output_text = stdout_bytes.decode(errors="replace")
|
|
except asyncio.TimeoutError:
|
|
proc.kill()
|
|
await proc.wait()
|
|
exit_code = -1
|
|
output_text += f"\n[pipeline_runner] TIMEOUT after {timeout}s -- process killed\n"
|
|
status = "success" if exit_code == 0 else "failed"
|
|
except Exception as exc:
|
|
exit_code = None
|
|
status = "failed"
|
|
output_text += f"\n[pipeline_runner] exception before/while running: {exc!r}\n"
|
|
|
|
log_path.write_text(output_text)
|
|
|
|
finished_at = time.time()
|
|
db = SessionLocal()
|
|
run = db.get(JobRun, run_id)
|
|
run.finished_at = finished_at
|
|
run.status = status
|
|
run.exit_code = exit_code
|
|
run.summary = _summarize_log(log_path)
|
|
db.commit()
|
|
db.refresh(run)
|
|
db.close()
|
|
return run, output_text
|
|
finally:
|
|
_lock.release()
|
|
|
|
|
|
async def run_playlist(playlist_name: str, no_m3u: bool = False, triggered_by: str = "schedule") -> JobRun:
|
|
script = settings.pipeline_dir / "bin" / "run-playlist.sh"
|
|
argv = [str(script)]
|
|
if no_m3u:
|
|
argv.append("--no-m3u")
|
|
argv.append(playlist_name)
|
|
# 45 min SLDL_TIMEOUT (run-playlist.sh) + generous slack for the
|
|
# tagging/beets-import/M3U steps that follow it in the same script.
|
|
return await run_job(f"playlist:{playlist_name}", argv, triggered_by=triggered_by, timeout=3300)
|
|
|
|
|
|
async def run_import_track(args: list[str], triggered_by: str = "manual") -> JobRun:
|
|
script = settings.pipeline_dir / "bin" / "import-track.sh"
|
|
return await run_job("manual:import", [str(script)] + args, triggered_by=triggered_by)
|
|
|
|
|
|
async def run_lib_script(job_key: str, script_name: str, args: list[str] | None = None, triggered_by: str = "schedule") -> JobRun:
|
|
script = settings.pipeline_dir / "lib" / script_name
|
|
return await run_job(job_key, [str(script)] + (args or []), triggered_by=triggered_by)
|
|
|
|
|
|
async def run_beet(job_key: str, args: list[str], triggered_by: str = "schedule") -> JobRun:
|
|
return await run_job(job_key, ["beet"] + args, triggered_by=triggered_by)
|