Coverage for utilities/rqworker.py: 42%
45 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 18:35 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 18:35 +0000
1import logging
3from django_rq.queues import get_connection
4from rq import Retry, Worker
5from rq.worker_registration import REDIS_WORKER_KEYS
6from rq.worker_registration import register as register_worker
8from netbox.config import get_config
9from netbox.constants import RQ_QUEUE_DEFAULT
11__all__ = (
12 'NetBoxRQWorker',
13 'any_workers_for_queue',
14 'get_all_workers',
15 'get_queue_for_model',
16 'get_rq_retry',
17 'get_workers_for_queue',
18)
20logger = logging.getLogger('netbox.rqworker')
23class NetBoxRQWorker(Worker):
24 """
25 RQ worker subclass which self-heals its registration. If the worker's
26 registration is missing from Redis (e.g. because the tasks Redis database
27 was lost and rebuilt while the worker was running), the next heartbeat
28 will re-register the worker so that Worker.all() / Worker.find_by_key()
29 can locate it again.
30 """
32 def heartbeat(self, *args, **kwargs):
33 try:
34 if not self.connection.sismember(REDIS_WORKER_KEYS, self.key):
35 logger.warning(f"Worker {self.name} not found in registry; re-registering.")
36 # If the worker hash still exists (partial Redis data loss),
37 # register_birth() would raise because rq treats an existing,
38 # non-dead hash as an active worker. Re-add to the registry
39 # sets directly in that case; the heartbeat below will refresh
40 # the hash TTL.
41 if self.connection.exists(self.key) and not self.connection.hexists(self.key, 'death'):
42 register_worker(self, self.connection)
43 else:
44 self.register_birth()
45 except Exception:
46 logger.exception("Failed to verify worker registration.")
47 super().heartbeat(*args, **kwargs)
50def get_queue_for_model(model):
51 """
52 Return the configured queue name for jobs associated with the given model.
53 """
54 return get_config().QUEUE_MAPPINGS.get(model, RQ_QUEUE_DEFAULT)
57def _is_live_worker(worker, queue_name):
58 """
59 Return True if the given Worker is currently servicing queue_name.
61 Liveness itself is enforced by RQ: Worker.all() / Worker.find_by_key()
62 only return workers whose Redis hash still exists, and RQ resets that
63 hash's expiry to (worker_ttl + 60s) on every heartbeat. So any worker
64 returned by RQ has heartbeat'd within its configured TTL -- we only need
65 to confirm it's listening on the requested queue. (Reconstructing
66 worker_ttl ourselves would be unsafe: RQ does not persist worker_ttl in
67 the hash, so a worker started with a non-default --worker-ttl is
68 reconstructed with DEFAULT_WORKER_TTL regardless of its real TTL.)
69 """
70 return queue_name in worker.queue_names()
73def get_workers_for_queue(queue_name):
74 """
75 Return the number of live workers currently servicing the given queue.
76 """
77 connection = get_connection(queue_name)
78 return sum(
79 1 for worker in Worker.all(connection=connection)
80 if _is_live_worker(worker, queue_name)
81 )
84def get_all_workers():
85 """
86 Return the set of worker names currently registered on the tasks Redis
87 connection, regardless of which queue(s) each worker is servicing. Stale
88 registrations (workers whose Redis hash has expired) are filtered out by
89 RQ via Worker.all() -- see _is_live_worker() for details.
91 Used for system-wide worker counts (dashboard, status API), where the
92 intent is "are any RQ workers running" rather than "are workers handling
93 a specific queue."
94 """
95 connection = get_connection(RQ_QUEUE_DEFAULT)
96 return {worker.name for worker in Worker.all(connection=connection)}
99def any_workers_for_queue(queue_name):
100 """
101 Return True if at least one live worker is currently servicing the given
102 queue. Cheaper than get_workers_for_queue() when only a liveness check is
103 needed: workers are fetched one at a time and iteration stops at the first
104 live match.
105 """
106 connection = get_connection(queue_name)
107 for key in Worker.all_keys(connection=connection): 107 ↛ 108line 107 didn't jump to line 108 because the loop on line 107 never started
108 worker = Worker.find_by_key(key, connection=connection)
109 if worker is None:
110 continue
111 if _is_live_worker(worker, queue_name):
112 return True
113 return False
116def get_rq_retry():
117 """
118 If RQ_RETRY_MAX is defined and greater than zero, instantiate and return a Retry object to be
119 used when queuing a job. Otherwise, return None.
120 """
121 retry_max = get_config().RQ_RETRY_MAX
122 retry_interval = get_config().RQ_RETRY_INTERVAL
123 if retry_max:
124 return Retry(max=retry_max, interval=retry_interval)
125 return None