Files
2026-08-09 02:31:54 +00:00

72 lines
2.2 KiB
Python

"""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)