import os import asyncio import logging from typing import Optional from core.queue import RedisQueue from core.metrics import QUEUE_POSTS_GAUGE, REDIS_QUEUE_SIZE_GAUGE logger = logging.getLogger(__name__) FETCH_INTERVAL_SECONDS = int(os.getenv("AI_PROCESSING_INTERVAL_SECONDS", "120")) class QueueConsumerService: def __init__(self, queue: RedisQueue, on_post_popped = None, fetch_interval: int = FETCH_INTERVAL_SECONDS): self.queue = queue self.on_post_popped = on_post_popped self.fetch_interval = fetch_interval self._running = False self._task: Optional[asyncio.Task] = None async def start(self): logger.info(f"Starting Redis Queue Consumer Service (interval: {self.fetch_interval}s / {self.fetch_interval // 60}m)...") await self.queue.connect() self._running = True self._task = asyncio.create_task(self._consumer_loop()) async def _consumer_loop(self): while self._running: try: qsize = await self.queue.qsize() QUEUE_POSTS_GAUGE.labels(status="redis_incoming").set(qsize) REDIS_QUEUE_SIZE_GAUGE.set(qsize) if qsize > 0: post_id = await self.queue.pop() if post_id and self.on_post_popped: logger.info(f"Dispatching raw post ID {post_id} to Admin Review channel (remaining in Redis: {qsize - 1})") await self.on_post_popped(post_id) new_qsize = await self.queue.qsize() QUEUE_POSTS_GAUGE.labels(status="redis_incoming").set(new_qsize) REDIS_QUEUE_SIZE_GAUGE.set(new_qsize) except Exception as e: logger.error(f"Error in queue consumer loop: {e}", exc_info=True) await asyncio.sleep(self.fetch_interval) async def stop(self): self._running = False if self._task: self._task.cancel() logger.info("Redis Queue Consumer Service stopped.")