Coverage for utilities/rqworker.py: 42%

45 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 18:35 +0000

1import logging 

2 

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 

7 

8from netbox.config import get_config 

9from netbox.constants import RQ_QUEUE_DEFAULT 

10 

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) 

19 

20logger = logging.getLogger('netbox.rqworker') 

21 

22 

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 """ 

31 

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) 

48 

49 

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) 

55 

56 

57def _is_live_worker(worker, queue_name): 

58 """ 

59 Return True if the given Worker is currently servicing queue_name. 

60 

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

71 

72 

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 ) 

82 

83 

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. 

90 

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

97 

98 

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 

114 

115 

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