from __future__ import annotations import json from dataclasses import dataclass from datetime import datetime, timezone from redis.asyncio import Redis from app.models import ProcessorForwardEnvelope @dataclass class QueueMessage: job_id: str attempt: int envelope: ProcessorForwardEnvelope raw_message: str class RedisQueueRepository: def __init__( self, client: Redis, queue_name: str, processing_name: str, deadletter_name: str, dedup_ttl_seconds: int, ) -> None: self.client = client self.queue_name = queue_name self.processing_name = processing_name self.deadletter_name = deadletter_name self.dedup_ttl_seconds = dedup_ttl_seconds async def ping(self) -> bool: result = await self.client.ping() return bool(result) async def enqueue_event(self, envelope: ProcessorForwardEnvelope) -> str: job_id = envelope.event.correlation_id payload = { "job_id": job_id, "attempt": 0, "enqueued_at": datetime.now(timezone.utc).isoformat(), "envelope": envelope.model_dump(mode="json"), } raw = json.dumps(payload, ensure_ascii=False) await self.client.lpush(self.queue_name, raw) return job_id async def claim_event(self, timeout_seconds: int) -> QueueMessage | None: raw = await self.client.brpoplpush( self.queue_name, self.processing_name, timeout=timeout_seconds, ) if raw is None: return None raw_str = raw.decode() if isinstance(raw, bytes) else str(raw) payload = json.loads(raw_str) return QueueMessage( job_id=str(payload["job_id"]), attempt=int(payload.get("attempt", 0)), envelope=ProcessorForwardEnvelope.model_validate(payload["envelope"]), raw_message=raw_str, ) async def ack_message(self, raw_message: str) -> None: await self.client.lrem(self.processing_name, 1, raw_message) async def requeue_message( self, message: QueueMessage, error: str | None = None, ) -> None: payload = json.loads(message.raw_message) payload["attempt"] = int(payload.get("attempt", 0)) + 1 payload["last_error"] = error or "" payload["requeued_at"] = datetime.now(timezone.utc).isoformat() new_raw = json.dumps(payload, ensure_ascii=False) await self.client.lpush(self.queue_name, new_raw) await self.ack_message(message.raw_message) async def deadletter_message( self, message: QueueMessage, error: str | None = None, ) -> None: payload = json.loads(message.raw_message) payload["deadlettered_at"] = datetime.now(timezone.utc).isoformat() payload["last_error"] = error or "" raw = json.dumps(payload, ensure_ascii=False) await self.client.lpush(self.deadletter_name, raw) await self.ack_message(message.raw_message) async def requeue_processing_messages(self) -> int: moved = 0 while True: raw = await self.client.rpoplpush(self.processing_name, self.queue_name) if raw is None: break moved += 1 return moved def build_processing_identity(self, envelope: ProcessorForwardEnvelope) -> str: event = envelope.event phase = (event.event_type or "problem").strip().lower() stable_id = event.event_id or event.correlation_id return f"{stable_id}:{phase}" async def is_processed(self, identity: str) -> bool: key = self._processed_key(identity) exists = await self.client.exists(key) return bool(exists) async def mark_processed(self, identity: str) -> None: key = self._processed_key(identity) now = datetime.now(timezone.utc).isoformat() await self.client.set(key, now, ex=self.dedup_ttl_seconds) def _processed_key(self, identity: str) -> str: return f"alert:queue:processed:{identity}"