Files
2026-06-27 13:50:11 +03:00

45 lines
1.6 KiB
Python

import json
import logging
from datetime import datetime
from typing import Any
import redis
from app.config import get_settings
from app.domain.schemas import EventType
logger = logging.getLogger(__name__)
class EventPublisher:
"""Publishes orchestration events to Redis Streams with in-memory fallback."""
def __init__(self, redis_url: str | None = None) -> None:
self.stream = "orchestrator.events"
self._fallback: list[dict[str, Any]] = []
self._redis = None
try:
self._redis = redis.Redis.from_url(redis_url or get_settings().redis_url, decode_responses=True)
self._redis.ping()
except Exception as exc: # pragma: no cover - depends on local infra
self._redis = None
logger.warning("Redis unavailable, using in-memory event log: %s", exc)
def publish(self, event_type: EventType, request_id: str, payload: dict[str, Any]) -> dict[str, Any]:
event = {
"type": event_type.value,
"request_id": request_id,
"timestamp": datetime.utcnow().isoformat(),
"payload": payload,
}
if self._redis:
self._redis.xadd(self.stream, {"event": json.dumps(event, default=str)}, maxlen=500, approximate=True)
self._fallback.append(event)
return event
def recent(self, limit: int = 50) -> list[dict[str, Any]]:
if self._redis:
rows = self._redis.xrevrange(self.stream, count=limit)
return [json.loads(fields["event"]) for _, fields in reversed(rows)]
return self._fallback[-limit:]