"""Internal recurring-task scheduler.""" from __future__ import annotations import threading import time from .discovery import BackgroundTask, TaskCatalog from .runner import RunManager class BackgroundScheduler: """Starts each declared background task at most once per configured interval.""" def __init__(self, catalog: TaskCatalog, manager: RunManager) -> None: self.catalog = catalog self.manager = manager self.tasks = catalog.background_tasks() self._stopped = threading.Event() self._lock = threading.Lock() self._next_due: dict[str, float] = {} self._next_cleanup = 0.0 self._running: set[str] = set() self._thread: threading.Thread | None = None def start(self) -> None: if self._thread is not None: return self._thread = threading.Thread( target=self._loop, name="background-task-scheduler", daemon=True, ) self._thread.start() def stop(self) -> None: self._stopped.set() if self._thread is not None: self._thread.join(timeout=1) def _loop(self) -> None: while not self._stopped.is_set(): current = time.monotonic() if current >= self._next_cleanup: self._next_cleanup = current + 60 self.manager.cleanup_transient_files() for task in self.tasks: due = self._next_due.setdefault(task.id, current) if current >= due: self._next_due[task.id] = current + task.interval_seconds self._start_task(task) self._stopped.wait(0.25) def _start_task(self, task: BackgroundTask) -> None: with self._lock: if task.id in self._running: return self._running.add(task.id) threading.Thread( target=self._run_task, args=(task,), name=f"background-task-{task.id}", daemon=True, ).start() def _run_task(self, task: BackgroundTask) -> None: try: self.manager.run_background_task(task) finally: with self._lock: self._running.discard(task.id)