"""One APScheduler job bridged into the asyncio loop — the shape the four *_scheduler modules (recurrence spawn, auto-pin scan, trash purge, DB maintenance) each used to carry a private copy of. APScheduler's BackgroundScheduler fires from a worker thread; the work is async and must run on the app's loop, so the fire is bridged with ``run_coroutine_threadsafe``. Each job is a module-level singleton: start is idempotent, stop shuts the scheduler down, and a job whose trigger the operator can change (the maintenance hour) reschedules the live job instead of restarting. """ from __future__ import annotations import asyncio import logging from collections.abc import Awaitable, Callable from apscheduler.schedulers.background import BackgroundScheduler logger = logging.getLogger(__name__) class ScheduledJob: """A named APScheduler job that awaits ``work()`` on the asyncio loop. ``work`` is an async callable; exceptions it raises are logged under ``label`` and never propagate into APScheduler's thread. """ def __init__(self, job_id: str, work: Callable[[], Awaitable[None]], *, label: str) -> None: self.job_id = job_id self._work = work self.label = label self._scheduler: BackgroundScheduler | None = None self._loop: asyncio.AbstractEventLoop | None = None @property def running(self) -> bool: return self._scheduler is not None def _fire(self) -> None: """APScheduler invokes this from its worker thread; bridge into the loop.""" if self._loop is None: logger.warning("%s scheduler: no loop registered", self.label) return async def _runner() -> None: try: await self._work() except Exception: logger.exception("%s run failed", self.label) asyncio.run_coroutine_threadsafe(_runner(), self._loop) def start(self, loop: asyncio.AbstractEventLoop, trigger, *, describe: str = "") -> None: """Start the job on ``trigger``. Idempotent — a second start is a no-op.""" if self._scheduler is not None: return self._loop = loop self._scheduler = BackgroundScheduler() self._scheduler.add_job( self._fire, trigger=trigger, id=self.job_id, replace_existing=True, ) self._scheduler.start() logger.info("%s scheduler started%s", self.label, f" ({describe})" if describe else "") def reschedule(self, trigger, *, describe: str = "") -> None: """Move the live job to a new trigger; a no-op when not running.""" if self._scheduler is None: return self._scheduler.reschedule_job(self.job_id, trigger=trigger) logger.info("%s scheduler rescheduled%s", self.label, f" to {describe}" if describe else "") def stop(self) -> None: if self._scheduler is not None: self._scheduler.shutdown(wait=False) self._scheduler = None logger.info("%s scheduler stopped", self.label)