1aaef71f52
Document-driven MVP with FastAPI backend, Vue H5, WeChat mini shell, product demo, and Docker dev stack. Co-authored-by: Cursor <cursoragent@cursor.com>
58 lines
1.4 KiB
Python
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()
|