126 lines
4.0 KiB
Python
126 lines
4.0 KiB
Python
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}" |