2c86d7973f
- Navidrome password is now passed to curl via stdin (--data-urlencode "p@-") in navidrome-scan.sh and pipeline-status.sh, so it no longer appears in ps/proc. Verified the query sent is identical and a live scan still triggers. - MIN_ARTIST_DIRS (the share-health gate) is now a setting, threaded through to the pipeline env, so a user with a small library can lower it instead of the scan/sync being permanently blocked by the hardcoded 500. - /auth/logout is now POST-only (with a nav form + aria-label), so a drive-by GET can't log the user out; enforced allowed_email already landed separately. - view_log now confirms the run's log_path resolves under the logs dir before serving it (defense in depth). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
282 lines
9.8 KiB
Python
282 lines
9.8 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)
|
|
env["MIN_ARTIST_DIRS"] = str(settings.min_artist_dirs)
|
|
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,
|
|
use_lock: bool = True,
|
|
) -> 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.
|
|
|
|
use_lock=False runs the command WITHOUT taking the pipeline lock, for
|
|
read-only jobs that never touch the library (e.g. the status report).
|
|
Those must never be starved by a long-running write job, and running
|
|
them concurrently is safe."""
|
|
started_at = time.time()
|
|
|
|
acquired = False
|
|
if use_lock:
|
|
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
|
|
acquired = True
|
|
|
|
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:
|
|
if acquired:
|
|
_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",
|
|
timeout: float | None = None,
|
|
use_lock: bool = True,
|
|
) -> JobRun:
|
|
script = settings.pipeline_dir / "lib" / script_name
|
|
return await run_job(
|
|
job_key,
|
|
[str(script)] + (args or []),
|
|
triggered_by=triggered_by,
|
|
timeout=timeout,
|
|
use_lock=use_lock,
|
|
)
|
|
|
|
|
|
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)
|