Files
john 1aaef71f52 Initial commit: Happy Up monorepo through Sprint 5.
Document-driven MVP with FastAPI backend, Vue H5, WeChat mini shell, product demo, and Docker dev stack.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-23 11:42:40 +08:00

58 lines
1.4 KiB
Python

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()