from __future__ import annotations import asyncio import logging from apscheduler.schedulers.asyncio import AsyncIOScheduler from app.collectors import import_all_sources from app.config import Settings, get_settings from app.db import get_session_factory from app.db.queries import prune_probe_results from app.probes import probe_due_batch logger = logging.getLogger(__name__) _scheduler: AsyncIOScheduler | None = None _locks = { "import": asyncio.Lock(), "probe": asyncio.Lock(), "prune": asyncio.Lock(), } async def job_import_sources() -> None: if _locks["import"].locked(): logger.warning("import job already running; skipping") return async with _locks["import"]: settings = get_settings() factory = get_session_factory() async with factory() as session: await import_all_sources(session, settings) async def job_probe_due() -> None: if _locks["probe"].locked(): logger.warning("probe job already running; skipping") return async with _locks["probe"]: settings = get_settings() factory = get_session_factory() async with factory() as session: count = await probe_due_batch(session, settings) logger.info("probed %s trackers", count) async def job_prune() -> None: if _locks["prune"].locked(): return async with _locks["prune"]: settings = get_settings() factory = get_session_factory() async with factory() as session: removed = await prune_probe_results(session, settings.raw_result_retention_days) logger.info("pruned %s probe results", removed) def start_scheduler(settings: Settings | None = None) -> AsyncIOScheduler: global _scheduler settings = settings or get_settings() if _scheduler is not None: return _scheduler scheduler = AsyncIOScheduler(timezone="UTC") scheduler.add_job( job_import_sources, "cron", day_of_week=settings.import_day_of_week, hour=settings.import_hour_utc, minute=0, id="import_sources", max_instances=1, coalesce=True, ) scheduler.add_job( job_probe_due, "interval", seconds=60, id="probe_due", max_instances=1, coalesce=True, ) scheduler.add_job( job_prune, "interval", hours=24, id="prune_results", max_instances=1, coalesce=True, ) scheduler.start() _scheduler = scheduler logger.info("scheduler started (single process / single worker required)") return scheduler def stop_scheduler() -> None: global _scheduler if _scheduler is not None: _scheduler.shutdown(wait=False) _scheduler = None