backend/queue.py (2888 bytes)
1 """Job dispatch. 2 3 Localhost runs jobs on a small thread pool inside the API process. The public 4 deployment sets SUBPLZ_WEB_QUEUE_BACKEND=redis and runs `worker.py` on separate 5 machines; the API side then only enqueues. Same `enqueue(job_id)` call either way. 6 """ 7 8 from __future__ import annotations 9 10 import logging 11 import threading 12 from abc import ABC, abstractmethod 13 from concurrent.futures import ThreadPoolExecutor 14 15 from .settings import settings 16 17 log = logging.getLogger(__name__) 18 19 20 class JobQueue(ABC): 21 @abstractmethod 22 def enqueue(self, job_id: str) -> None: 23 ... 24 25 @abstractmethod 26 def depth(self) -> int: 27 """Jobs waiting or running. Drives the 'N ahead of you' hint in the UI.""" 28 29 def shutdown(self) -> None: 30 ... 31 32 33 class InProcessQueue(JobQueue): 34 """Thread pool in the API process. Fine for one person on one machine. 35 36 Alignment is CPU-bound and subplz is a subprocess, so the GIL is not the 37 limit here - max_concurrent_jobs is. 38 """ 39 40 def __init__(self, workers: int): 41 self._pool = ThreadPoolExecutor( 42 max_workers=workers, thread_name_prefix="subplz-job" 43 ) 44 self._lock = threading.Lock() 45 self._pending: set[str] = set() 46 47 def enqueue(self, job_id: str) -> None: 48 with self._lock: 49 if job_id in self._pending: 50 return 51 self._pending.add(job_id) 52 self._pool.submit(self._run, job_id) 53 54 def _run(self, job_id: str) -> None: 55 # Imported here to avoid a circular import at module load. 56 from .runner import run_job 57 58 try: 59 run_job(job_id) 60 except Exception: 61 log.exception("job %s crashed outside the runner", job_id) 62 finally: 63 with self._lock: 64 self._pending.discard(job_id) 65 66 def depth(self) -> int: 67 with self._lock: 68 return len(self._pending) 69 70 def shutdown(self) -> None: 71 self._pool.shutdown(wait=False, cancel_futures=True) 72 73 74 class RedisQueue(JobQueue): 75 """Public-release backend. Needs `rq` and a reachable Redis.""" 76 77 def __init__(self, url: str): 78 from redis import Redis # lazy: localhost needs neither package 79 from rq import Queue as RQQueue 80 81 self._conn = Redis.from_url(url) 82 self._q = RQQueue("subplz-jobs", connection=self._conn) 83 84 def enqueue(self, job_id: str) -> None: 85 self._q.enqueue( 86 "backend.runner.run_job", 87 job_id, 88 job_timeout=settings.job_timeout_seconds, 89 result_ttl=86400, 90 ) 91 92 def depth(self) -> int: 93 return self._q.count + self._q.started_job_registry.count 94 95 96 def build_queue() -> JobQueue: 97 if settings.queue_backend == "redis": 98 return RedisQueue(settings.redis_url) 99 return InProcessQueue(settings.max_concurrent_jobs) 100 101 102 queue: JobQueue = build_queue()