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:]