# Trouve-KA — boucle de maintenance # Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai """Scheduler : tâches périodiques légères. - Relance les items in_progress abandonnés (worker mort) — dégradation gracieuse §13. - Recalcule l'autorité de domaine à partir du graphe de liens (inlinks pondérés par le score Québec des domaines source) — version simple, PageRank-like plus tard. """ import asyncio import signal from trouveka.config import get_settings from trouveka.database import Database from trouveka.logging import get_logger log = get_logger("scheduler") STALE_RESET_INTERVAL = 60 # secondes AUTHORITY_INTERVAL = 15 * 60 # secondes class SchedulerLoop: def __init__(self) -> None: self.s = get_settings() self.db = Database(self.s.database_url, pool_min=1, pool_max=3) self.stop_event = asyncio.Event() async def recompute_authority(self) -> None: """Autorité ∈ [0,1] : log-saturation des inlinks pondérés par le Québec-score des sources.""" await self.db.pool.execute( """ WITH weighted AS ( SELECT dl.to_domain_id AS id, sum(least(dl.link_count, 50) * greatest(d.quebec_score, 0.1)) AS w, count(DISTINCT dl.from_domain_id) AS in_domains FROM domain_links dl JOIN domains d ON d.id = dl.from_domain_id GROUP BY dl.to_domain_id ) UPDATE domains SET authority_score = least(1.0, ln(1 + w.w) / ln(1 + 5000)), inlink_domains = w.in_domains FROM weighted w WHERE domains.id = w.id """ ) log.info("autorité de domaine recalculée") async def start(self) -> None: await self.db.connect() loop = asyncio.get_running_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, self.stop_event.set) log.info("scheduler démarré") elapsed_authority = AUTHORITY_INTERVAL # premier calcul immédiat while not self.stop_event.is_set(): try: reset = await self.db.reset_stale_items(older_than_minutes=30) if reset: log.info("items abandonnés relancés", extra={"ctx": {"count": reset}}) if elapsed_authority >= AUTHORITY_INTERVAL: await self.recompute_authority() elapsed_authority = 0 except Exception: log.exception("erreur scheduler (on continue)") try: await asyncio.wait_for(self.stop_event.wait(), timeout=STALE_RESET_INTERVAL) except TimeoutError: pass elapsed_authority += STALE_RESET_INTERVAL await self.db.close() def main() -> None: asyncio.run(SchedulerLoop().start()) if __name__ == "__main__": main()