Files
LLM-Zabbix-HABR/alert-processor/app/queue_repo.py
T

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}"