"""`run` orchestration: poll, then download, then reap — under a lock.""" from __future__ import annotations import contextlib import fcntl import logging import sqlite3 from pathlib import Path from . import config, discovery, download, reap, util, videos from .settings import Settings log = logging.getLogger(__name__) class AlreadyRunning(Exception): pass @contextlib.contextmanager def exclusive_lock(path: Path | None = None): """Non-blocking flock. Raises AlreadyRunning if another run holds it. The cron schedule is hourly and a large backfill can outlast that, so overlapping runs are expected and must be a silent no-op rather than two workers fighting over the same queue. """ path = path or config.LOCK_PATH path.parent.mkdir(parents=True, exist_ok=True) handle = path.open("w") try: try: fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB) except OSError as exc: raise AlreadyRunning(f"another run holds {path}") from exc yield finally: with contextlib.suppress(OSError): fcntl.flock(handle, fcntl.LOCK_UN) handle.close() def recover(conn: sqlite3.Connection) -> dict: """Undo the effects of a killed run before doing anything else.""" requeued = videos.recover_downloading(conn) orphans = download.recover_orphans() if requeued or orphans: log.info( "crash recovery: %d row(s) back to pending, %d orphan file(s) cleared", requeued, orphans, ) return {"requeued": requeued, "orphans": orphans} def run( conn: sqlite3.Connection, settings: Settings, channel_pk: int | None = None ) -> dict: """One full cycle. Assumes the caller holds the lock.""" result = {"recovered": recover(conn)} result["poll"] = discovery.poll_all(conn, settings, channel_pk) try: result["download"] = download.drain(conn, settings) except RuntimeError as exc: # The POT provider being down is loud and fatal for this run, but the # poll results are still worth keeping. log.error("%s", exc) result["download"] = {"error": str(exc)} return result result["reap"] = reap.run(conn, settings) settings.set("last_run_at", util.utcnow_iso()) return result def summarise(result: dict) -> str: poll = result.get("poll", {}) down = result.get("download", {}) reaped = result.get("reap", {}) parts = [ f"discovered={poll.get('queued', 0)}", f"repaired={poll.get('repaired', 0)}", f"poll_failures={poll.get('failed', 0)}", f"downloaded={down.get(videos.DOWNLOADED, 0)}", f"failed={down.get(videos.FAILED, 0)}", f"skipped={down.get(videos.SKIPPED_SHORT, 0) + down.get(videos.SKIPPED_LIVE, 0)}", f"deferred={down.get(videos.DEFERRED, 0)}", f"reaped={reaped.get('deleted', 0)}", f"evicted={reaped.get('evicted', 0)}", ] if "error" in down: parts.append(f"ERROR={down['error']}") return " ".join(parts)