0.6.4: Truthful maintenance digest + scheduled jobs queue for the lock instead of skipping
The Telegram digest reported healthy weekly jobs as "never run yet": its maintenance section still globbed for cron-era per-script log filenames (clean-index-*.log etc.) that jobs run through pipeline_runner no longer write. It now reads job_runs, the scheduler's own record, so it reports the last success (age + summary snippet from that run's log), flags a newer failed or lock-skipped attempt on the same line, and distinguishes "never succeeded (last attempt: skipped, lock busy)" from genuinely never scheduled. Also matched the clear-bad-genres snippet grep to the script's current summary wording, and dropped the now-unused newest_log helper. The DB also showed WHY two jobs had never run: on Sunday the 08:30 strip-mb-tags run took 10m19s while holding the shared pipeline lock, so strip-watermark-art (08:35) and scrub-watermark-text (08:40) hit flock-style instant skip and lost their only slot of the week -- the 08:40 job missed by 19 seconds. Scheduled runs now wait up to 30 minutes for the lock (SCHEDULED_LOCK_WAIT_SECONDS) and only then record skipped_lock, so fixed-time blocks queue instead of starving; manual "Run now" keeps the instant skip since a person expects an immediate answer. Tests: scheduled-run queueing, lock-wait timeout, and manual instant-skip. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -12,6 +12,15 @@ from app.settings import settings
|
||||
# all-jobs-share-one-lock behavior exactly rather than over-engineering it.
|
||||
_lock = asyncio.Lock()
|
||||
|
||||
# How long a SCHEDULED run waits for the lock before giving up (recorded as
|
||||
# skipped_lock). The old flock -n instant-skip starves fixed-time schedules:
|
||||
# the Sunday 08:30 strip-mb-tags run took 10m19s on 2026-07-12, so the 08:35
|
||||
# and 08:40 jobs both hit the held lock and silently lost their only slot of
|
||||
# the week. Scheduled jobs have nobody watching, so queueing (bounded) beats
|
||||
# skipping; manual runs keep the instant skip because a person clicking "Run
|
||||
# now" should get an immediate answer, not a silent 30-minute wait.
|
||||
SCHEDULED_LOCK_WAIT_SECONDS = 1800.0
|
||||
|
||||
_ENV_PASSTHROUGH_KEYS = ("PATH", "HOME", "LANG", "LC_ALL", "TZ")
|
||||
|
||||
|
||||
@@ -94,18 +103,30 @@ async def _execute(
|
||||
timeout: float | None,
|
||||
use_lock: bool,
|
||||
capture: bool,
|
||||
lock_wait: float | None,
|
||||
) -> tuple[JobRun, str]:
|
||||
"""Shared core for run_job/run_job_capture: acquire the lock (unless
|
||||
use_lock=False), record the run, exec the subprocess, and finalize. When
|
||||
capture=True, stdout+stderr is captured as text and returned; otherwise it
|
||||
streams straight to the log file. Returns (JobRun, output_text) -- the text
|
||||
is "" in the non-capture and skipped-lock cases."""
|
||||
is "" in the non-capture and skipped-lock cases.
|
||||
|
||||
lock_wait: how long to wait for a held lock before recording
|
||||
skipped_lock. None picks the policy default: SCHEDULED_LOCK_WAIT_SECONDS
|
||||
for triggered_by="schedule", instant skip for everything else."""
|
||||
started_at = time.time()
|
||||
if lock_wait is None:
|
||||
lock_wait = SCHEDULED_LOCK_WAIT_SECONDS if triggered_by == "schedule" else 0.0
|
||||
|
||||
acquired = False
|
||||
if use_lock:
|
||||
if not await _try_acquire_nowait():
|
||||
return _record_run(job_key, started_at, triggered_by, finished_at=started_at, status="skipped_lock"), ""
|
||||
if lock_wait <= 0:
|
||||
return _record_run(job_key, started_at, triggered_by, finished_at=started_at, status="skipped_lock"), ""
|
||||
try:
|
||||
await asyncio.wait_for(_lock.acquire(), timeout=lock_wait)
|
||||
except asyncio.TimeoutError:
|
||||
return _record_run(job_key, started_at, triggered_by, finished_at=time.time(), status="skipped_lock"), ""
|
||||
acquired = True
|
||||
|
||||
try:
|
||||
@@ -170,18 +191,21 @@ async def run_job(
|
||||
triggered_by: str = "schedule",
|
||||
timeout: float | None = None,
|
||||
use_lock: bool = True,
|
||||
lock_wait: 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.
|
||||
recording one job_runs row start-to-finish. If the lock is already held,
|
||||
scheduled runs (triggered_by="schedule") wait up to
|
||||
SCHEDULED_LOCK_WAIT_SECONDS for it before recording skipped_lock -- see
|
||||
that constant for why -- while manual/other runs record skipped_lock
|
||||
immediately (the old flock -n behavior, visible in the UI instead of
|
||||
silently skipping). Pass lock_wait to override either way.
|
||||
|
||||
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."""
|
||||
run, _ = await _execute(job_key, argv, triggered_by, timeout, use_lock, capture=False)
|
||||
run, _ = await _execute(job_key, argv, triggered_by, timeout, use_lock, capture=False, lock_wait=lock_wait)
|
||||
return run
|
||||
|
||||
|
||||
@@ -191,6 +215,7 @@ async def run_job_capture(
|
||||
triggered_by: str = "manual",
|
||||
timeout: float | None = None,
|
||||
use_lock: bool = True,
|
||||
lock_wait: 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 --
|
||||
@@ -198,7 +223,7 @@ async def run_job_capture(
|
||||
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."""
|
||||
return await _execute(job_key, argv, triggered_by, timeout, use_lock, capture=True)
|
||||
return await _execute(job_key, argv, triggered_by, timeout, use_lock, capture=True, lock_wait=lock_wait)
|
||||
|
||||
|
||||
async def run_playlist(playlist_name: str, no_m3u: bool = False, triggered_by: str = "schedule") -> JobRun:
|
||||
|
||||
Reference in New Issue
Block a user