Spaces:
Paused
Paused
| """ | |
| Message Queue — Redis | |
| ───────────────────── | |
| Push sentiment results → Dashboard (Lovable AI) รับแบบ real-time | |
| ใช้ Redis List + Pub/Sub: | |
| - LPUSH sentiment:results → queue สำหรับ polling | |
| - PUBLISH sentiment:live → pub/sub สำหรับ WebSocket / SSE | |
| """ | |
| import json | |
| import logging | |
| import redis.asyncio as aioredis | |
| from app.config import settings | |
| logger = logging.getLogger(__name__) | |
| QUEUE_KEY = "sentiment:results" | |
| ALERT_KEY = "sentiment:alerts" | |
| PUBSUB_CH = "sentiment:live" | |
| ALERT_CH = "sentiment:alert_live" | |
| class MessageQueue: | |
| def __init__(self): | |
| self._redis: aioredis.Redis | None = None | |
| async def _get_redis(self) -> aioredis.Redis: | |
| if self._redis is None: | |
| self._redis = await aioredis.from_url( | |
| settings.REDIS_URL, | |
| encoding="utf-8", | |
| decode_responses=True, | |
| ) | |
| return self._redis | |
| async def push(self, data: dict) -> None: | |
| """ | |
| Push sentiment result เข้า queue และ publish ให้ subscribers | |
| """ | |
| try: | |
| r = await self._get_redis() | |
| serialized = json.dumps(data, ensure_ascii=False) | |
| # List (Lovable polling หรือ worker ดึงไปประมวลผลต่อ) | |
| await r.lpush(QUEUE_KEY, serialized) | |
| await r.ltrim(QUEUE_KEY, 0, 999) # เก็บแค่ 1,000 รายการล่าสุด | |
| # Pub/Sub (real-time Dashboard WebSocket) | |
| await r.publish(PUBSUB_CH, serialized) | |
| logger.debug(f"✅ Queued: conv_id={data.get('conv_id')}") | |
| except Exception as e: | |
| logger.error(f"❌ Queue push failed: {e}") | |
| async def push_alert(self, data: dict) -> None: | |
| """Push high-negative alert แยก channel""" | |
| try: | |
| r = await self._get_redis() | |
| serialized = json.dumps(data, ensure_ascii=False) | |
| await r.lpush(ALERT_KEY, serialized) | |
| await r.publish(ALERT_CH, serialized) | |
| logger.warning(f"🚨 Alert pushed: {data}") | |
| except Exception as e: | |
| logger.error(f"❌ Alert push failed: {e}") | |
| async def get_recent(self, count: int = 50) -> list[dict]: | |
| """ดึง result ล่าสุดจาก queue สำหรับ initial load""" | |
| try: | |
| r = await self._get_redis() | |
| items = await r.lrange(QUEUE_KEY, 0, count - 1) | |
| return [json.loads(i) for i in items] | |
| except Exception as e: | |
| logger.error(f"❌ Queue read failed: {e}") | |
| return [] | |
| # Singleton instance | |
| message_queue = MessageQueue() | |