# Trouve-KA — worker d'enrichissement # Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai """Worker d'enrichissement : consomme trouveka:enrich (Redis Streams). Étape 2 (implémentée) : propage domain_quebec_score et authority_score à jour dans le document indexé (le domaine apprend au fil du crawl, les documents déjà indexés en profitent rétroactivement). Étapes futures : embeddings, entités, classification thématique — même canal, même contrat : mise à jour partielle du document, jamais bloquante. """ import asyncio import os import signal import uuid from trouveka.config import get_settings from trouveka.database import Database from trouveka.logging import get_logger from trouveka.queue import Coordination from trouveka.search_core import SearchCore log = get_logger("enrichment") GROUP = "enrichers" class EnrichmentWorker: def __init__(self) -> None: self.s = get_settings() self.consumer = f"enrich-{os.getpid()}-{uuid.uuid4().hex[:6]}" self.db = Database(self.s.database_url, pool_min=1, pool_max=3) self.coord = Coordination(self.s.redis_url) self.search = SearchCore(self.s.search_url, self.s.search_index) self.stop_event = asyncio.Event() async def enrich(self, payload: dict) -> None: domain_row = await self.db.get_domain_by_id(int(payload["domain_id"])) if not domain_row: return await self.search.update_document( payload["url"], { "domain_quebec_score": round(float(domain_row["quebec_score"]), 4), "authority_score": round(float(domain_row["authority_score"]), 4), }, ) 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("worker d'enrichissement démarré", extra={"ctx": {"consumer": self.consumer}}) while not self.stop_event.is_set(): try: messages = await self.coord.read_enrichment(GROUP, self.consumer, count=20, block_ms=5000) for msg_id, payload in messages: try: await self.enrich(payload) except Exception: log.exception("échec enrichissement", extra={"ctx": payload}) finally: await self.coord.ack_enrichment(GROUP, msg_id) except Exception: log.exception("erreur boucle enrichissement (on continue)") await asyncio.sleep(2) await self.search.close() await self.coord.close() await self.db.close() def main() -> None: asyncio.run(EnrichmentWorker().start()) if __name__ == "__main__": main()