diff --git a/django_q/cluster.py b/django_q/cluster.py index 9e4a065..e97d154 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -31,6 +31,7 @@ from django_q.conf import ( psutil, setproctitle, ) +from django_q.enums import TimerStatus from django_q.humanhash import humanize from django_q.monitor import monitor from django_q.pusher import pusher @@ -234,7 +235,11 @@ class Sentinel: def spawn_worker(self): self.spawn_process( - worker, self.task_queue, self.result_queue, Value("f", -1), self.timeout + worker, + self.task_queue, + self.result_queue, + Value("f", TimerStatus.IDLE), + self.timeout, ) def spawn_monitor(self) -> Process: @@ -271,7 +276,7 @@ class Sentinel: self.pool.remove(process) self.spawn_worker() - if process.timer.value == 0: + if process.timer.value == TimerStatus.TIMEOUT: # only need to terminate on timeout, otherwise we risk destabilizing # the queues task_name = "" @@ -296,7 +301,7 @@ class Sentinel: "name": process.name } logger.critical(msg) - elif int(process.timer.value) == -2: + elif int(process.timer.value) == TimerStatus.RECYCLED: logger.info(_("recycled worker %(name)s") % {"name": process.name}) else: logger.critical( @@ -348,11 +353,11 @@ class Sentinel: for p in self.pool: with p.timer.get_lock(): # Are you alive? - if not p.is_alive() or p.timer.value == 0: + if not p.is_alive() or p.timer.value == TimerStatus.TIMEOUT: self.reincarnate(p) continue # Decrement timer if work is being done - if p.timer.value > 0: + if p.timer.value > TimerStatus.TIMEOUT: p.timer.value -= cycle # Check Monitor if not self.monitor.is_alive(): diff --git a/django_q/enums.py b/django_q/enums.py new file mode 100644 index 0000000..8052eef --- /dev/null +++ b/django_q/enums.py @@ -0,0 +1,14 @@ +from enum import IntEnum + + +class TimerStatus(IntEnum): + """ + Sentinel values for the countdown timer a worker shares with the sentinel. + + Any positive value is the number of seconds left before the sentinel + considers the task timed out and reincarnates the worker. + """ + + TIMEOUT = 0 # task timed out, the worker has to be terminated + IDLE = -1 # waiting for work, or working without a timeout + RECYCLED = -2 # recycle limit reached, the worker stopped on its own diff --git a/django_q/tests/test_cached.py b/django_q/tests/test_cached.py index 08a8bb6..ab53d51 100644 --- a/django_q/tests/test_cached.py +++ b/django_q/tests/test_cached.py @@ -6,6 +6,7 @@ import pytest from django_q.brokers import get_broker from django_q.conf import Conf +from django_q.enums import TimerStatus from django_q.monitor import monitor from django_q.pusher import pusher from django_q.queues import Queue @@ -67,7 +68,7 @@ def test_cached(broker): assert task_queue.qsize() == task_count task_queue.put("STOP") result_queue = Queue() - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) assert result_queue.qsize() == task_count result_queue.put("STOP") monitor(result_queue) diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index f24d51e..d08c217 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -15,6 +15,7 @@ from django.utils import timezone from django_q.brokers import Broker, get_broker from django_q.cluster import Cluster, Sentinel from django_q.conf import Conf +from django_q.enums import TimerStatus from django_q.humanhash import DEFAULT_WORDLIST, uuid from django_q.models import Success, Task from django_q.monitor import monitor, save_task @@ -245,7 +246,7 @@ def test_cluster(broker): assert queue_size(broker=broker) == 0 # Test work task_queue.put("STOP") - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) assert task_queue.qsize() == 0 assert result_queue.qsize() == 1 # Test monitor @@ -272,7 +273,7 @@ def test_results(broker): pusher(task_queue, stop_event, broker=broker) task_queue.put("STOP") result_queue = Queue() - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) result_queue.put("STOP") monitor(result_queue) @@ -372,7 +373,7 @@ def test_enqueue(broker, admin_user): assert fetch_group("test_j", count=2, wait=10) is None # let a worker handle them result_queue = Queue() - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) assert result_queue.qsize() == task_count result_queue.put("STOP") # store the results @@ -541,7 +542,7 @@ def test_recycle(broker, monkeypatch, django_assert_num_queries): pusher(task_queue, stop_event, broker=broker) pusher(task_queue, stop_event, broker=broker) # worker should exit on recycle - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) # check if the work has been done assert result_queue.qsize() == 2 # save_limit test @@ -573,7 +574,7 @@ def test_save_limit_per_func(broker, monkeypatch): threading.Timer(3, stop_event.set).start() for i in range(3): pusher(task_queue, stop_event, broker=broker) - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker) assert start_event.is_set() assert s.status() == Conf.STOPPED @@ -618,7 +619,7 @@ def test_max_rss(broker, monkeypatch): # push the task pusher(task_queue, stop_event, broker=broker) # worker should exit on recycle - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) # check if the work has been done assert result_queue.qsize() == 1 # save_limit test @@ -654,7 +655,7 @@ def test_bad_secret(broker, monkeypatch): worker( task_queue, result_queue, - Value("f", -1), + Value("f", TimerStatus.IDLE), ) assert result_queue.qsize() == 0 broker.delete_queue() @@ -838,7 +839,7 @@ class TestSignals: event.set() pusher(task_queue, event, broker=broker) task_queue.put("STOP") - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) result_queue.put("STOP") monitor(result_queue, broker) broker.delete_queue() @@ -867,7 +868,7 @@ class TestSignals: event.set() pusher(task_queue, event, broker=broker) task_queue.put("STOP") - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) result_queue.put("STOP") monitor(result_queue, broker) broker.delete_queue() @@ -896,7 +897,7 @@ class TestSignals: event.set() pusher(task_queue, event, broker=broker) task_queue.put("STOP") - worker(task_queue, result_queue, Value("f", -1)) + worker(task_queue, result_queue, Value("f", TimerStatus.IDLE)) result_queue.put("STOP") monitor(result_queue, broker) broker.delete_queue() diff --git a/django_q/tests/test_scheduler.py b/django_q/tests/test_scheduler.py index fcad6f8..257b872 100644 --- a/django_q/tests/test_scheduler.py +++ b/django_q/tests/test_scheduler.py @@ -13,6 +13,7 @@ from django.utils.timezone import is_naive from django_q.brokers import Broker, get_broker from django_q.conf import Conf +from django_q.enums import TimerStatus from django_q.monitor import monitor from django_q.pusher import pusher from django_q.queues import Queue @@ -222,7 +223,7 @@ def test_scheduler(broker, monkeypatch): task_queue.put("STOP") # let a worker handle them result_queue = Queue() - worker(task_queue, result_queue, Value("b", -1)) + worker(task_queue, result_queue, Value("b", TimerStatus.IDLE)) assert result_queue.qsize() == 1 result_queue.put("STOP") # store the results diff --git a/django_q/worker.py b/django_q/worker.py index 293982b..c93f94d 100644 --- a/django_q/worker.py +++ b/django_q/worker.py @@ -17,6 +17,7 @@ except core.exceptions.AppRegistryNotReady: django.setup() from django_q.conf import Conf, error_reporter, logger, resource, setproctitle +from django_q.enums import TimerStatus from django_q.exceptions import TimeoutException from django_q.signals import post_execute_in_worker, post_spawn, pre_execute from django_q.timeout import TimeoutHandler @@ -54,11 +55,11 @@ def worker( setproctitle.setproctitle(f"qcluster {proc_name} idle") task_count = 0 if timeout is None: - timeout = -1 + timeout = TimerStatus.IDLE # Start reading the task queue for task in iter(task_queue.get, "STOP"): result = None - timer.value = -1 # Idle + timer.value = TimerStatus.IDLE task_count += 1 f = task["func"] @@ -92,7 +93,7 @@ def worker( pre_execute.send(sender="django_q", func=f, task=task) # execute the payload timer.value = timer_value # Busy - if timer.value != -1: + if timer.value != TimerStatus.IDLE: timer.value += 3 # Add buffer so that guard doesn't kill the process on timeout before it gets processed timeout_error = False @@ -122,15 +123,15 @@ def worker( result_queue.put(task) if timeout_error: # force destroy process due to timeout - timer.value = 0 + timer.value = TimerStatus.TIMEOUT break - timer.value = -1 # Idle + timer.value = TimerStatus.IDLE if setproctitle: setproctitle.setproctitle(f"qcluster {proc_name} idle") # Recycle if task_count == Conf.RECYCLE or rss_check(): - timer.value = -2 # Recycled + timer.value = TimerStatus.RECYCLED break logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name})