Recently Written · git

subplz-web

git clone https://github.com/equwal/subplz-web

Log | Files | Refs


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