from app.config import settings class AnalysisQueue: def push(self, task_id: int) -> None: raise NotImplementedError def pop(self, timeout: int = 5) -> int | None: raise NotImplementedError _memory_queue: list[int] = [] class MemoryAnalysisQueue(AnalysisQueue): def push(self, task_id: int) -> None: _memory_queue.append(task_id) def pop(self, timeout: int = 5) -> int | None: return _memory_queue.pop(0) if _memory_queue else None class RedisAnalysisQueue(AnalysisQueue): def __init__(self) -> None: import redis self._client = redis.from_url(settings.redis_url, decode_responses=True) def push(self, task_id: int) -> None: self._client.lpush(settings.analysis_queue_key, str(task_id)) def pop(self, timeout: int = 5) -> int | None: item = self._client.brpop(settings.analysis_queue_key, timeout=timeout) if not item: return None return int(item[1]) _queue: AnalysisQueue | None = None def get_analysis_queue() -> AnalysisQueue: global _queue if _queue is not None: return _queue try: _queue = RedisAnalysisQueue() _queue._client.ping() except Exception: _queue = MemoryAnalysisQueue() return _queue def reset_analysis_queue_for_tests() -> None: global _queue _queue = MemoryAnalysisQueue() _memory_queue.clear()