51 lines
2.0 KiB
Python
51 lines
2.0 KiB
Python
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.")
|