# Trouve-KA — coordination Redis (politesse, pause, enrichissement) # Author: Simon-Pierre Boucher # Contact: contact@spboucher.ai """Coordination inter-workers via Redis. - Politesse par hôte : lock SET NX PX — un seul fetch par hôte par fenêtre, quel que soit le nombre de workers. - Pause globale du crawler : simple clé drapeau. - Enrichissement asynchrone : Redis Stream (jamais bloquant pour l'indexation). """ import json from typing import Any import redis.asyncio as aioredis PAUSE_KEY = "trouveka:crawler:paused" HOST_LOCK_PREFIX = "trouveka:host-lock:" ENRICH_STREAM = "trouveka:enrich" class Coordination: def __init__(self, redis_url: str): self.redis: aioredis.Redis = aioredis.from_url(redis_url, decode_responses=True) async def close(self) -> None: await self.redis.aclose() # ------------------------------------------------------------- politesse async def acquire_host_slot(self, host: str, delay_seconds: float) -> bool: """Réserve le droit de fetcher cet hôte. False = trop tôt, repasser plus tard.""" px = max(int(delay_seconds * 1000), 100) return bool(await self.redis.set(HOST_LOCK_PREFIX + host, "1", nx=True, px=px)) # ------------------------------------------------------------- pause async def pause_crawler(self) -> None: await self.redis.set(PAUSE_KEY, "1") async def resume_crawler(self) -> None: await self.redis.delete(PAUSE_KEY) async def is_paused(self) -> bool: return await self.redis.exists(PAUSE_KEY) == 1 # ------------------------------------------------------------- enrichissement async def enqueue_enrichment(self, payload: dict[str, Any]) -> None: await self.redis.xadd(ENRICH_STREAM, {"data": json.dumps(payload, default=str)}, maxlen=100_000) async def read_enrichment( self, group: str, consumer: str, count: int = 10, block_ms: int = 5000 ) -> list[tuple[str, dict[str, Any]]]: try: await self.redis.xgroup_create(ENRICH_STREAM, group, id="0", mkstream=True) except aioredis.ResponseError as exc: if "BUSYGROUP" not in str(exc): raise entries = await self.redis.xreadgroup( group, consumer, {ENRICH_STREAM: ">"}, count=count, block=block_ms ) out: list[tuple[str, dict[str, Any]]] = [] for _stream, items in entries or []: for msg_id, fields in items: out.append((msg_id, json.loads(fields["data"]))) return out async def ack_enrichment(self, group: str, msg_id: str) -> None: await self.redis.xack(ENRICH_STREAM, group, msg_id) async def enrich_backlog(self) -> int: return await self.redis.xlen(ENRICH_STREAM)