Compare commits

...
Author SHA1 Message Date
Stan a8f3bfb832 Migrate timer status to enum instead of int 2026-08-13 16:53:42 +02:00
6 changed files with 46 additions and 23 deletions
+10 -5
View File
@@ -31,6 +31,7 @@ from django_q.conf import (
psutil, psutil,
setproctitle, setproctitle,
) )
from django_q.enums import TimerStatus
from django_q.humanhash import humanize from django_q.humanhash import humanize
from django_q.monitor import monitor from django_q.monitor import monitor
from django_q.pusher import pusher from django_q.pusher import pusher
@@ -234,7 +235,11 @@ class Sentinel:
def spawn_worker(self): def spawn_worker(self):
self.spawn_process( 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: def spawn_monitor(self) -> Process:
@@ -271,7 +276,7 @@ class Sentinel:
self.pool.remove(process) self.pool.remove(process)
self.spawn_worker() self.spawn_worker()
if process.timer.value == 0: if process.timer.value == TimerStatus.TIMEOUT:
# only need to terminate on timeout, otherwise we risk destabilizing # only need to terminate on timeout, otherwise we risk destabilizing
# the queues # the queues
task_name = "" task_name = ""
@@ -296,7 +301,7 @@ class Sentinel:
"name": process.name "name": process.name
} }
logger.critical(msg) 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}) logger.info(_("recycled worker %(name)s") % {"name": process.name})
else: else:
logger.critical( logger.critical(
@@ -348,11 +353,11 @@ class Sentinel:
for p in self.pool: for p in self.pool:
with p.timer.get_lock(): with p.timer.get_lock():
# Are you alive? # 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) self.reincarnate(p)
continue continue
# Decrement timer if work is being done # Decrement timer if work is being done
if p.timer.value > 0: if p.timer.value > TimerStatus.TIMEOUT:
p.timer.value -= cycle p.timer.value -= cycle
# Check Monitor # Check Monitor
if not self.monitor.is_alive(): if not self.monitor.is_alive():
+14
View File
@@ -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
+2 -1
View File
@@ -6,6 +6,7 @@ import pytest
from django_q.brokers import get_broker from django_q.brokers import get_broker
from django_q.conf import Conf from django_q.conf import Conf
from django_q.enums import TimerStatus
from django_q.monitor import monitor from django_q.monitor import monitor
from django_q.pusher import pusher from django_q.pusher import pusher
from django_q.queues import Queue from django_q.queues import Queue
@@ -67,7 +68,7 @@ def test_cached(broker):
assert task_queue.qsize() == task_count assert task_queue.qsize() == task_count
task_queue.put("STOP") task_queue.put("STOP")
result_queue = Queue() 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 assert result_queue.qsize() == task_count
result_queue.put("STOP") result_queue.put("STOP")
monitor(result_queue) monitor(result_queue)
+11 -10
View File
@@ -15,6 +15,7 @@ from django.utils import timezone
from django_q.brokers import Broker, get_broker from django_q.brokers import Broker, get_broker
from django_q.cluster import Cluster, Sentinel from django_q.cluster import Cluster, Sentinel
from django_q.conf import Conf from django_q.conf import Conf
from django_q.enums import TimerStatus
from django_q.humanhash import DEFAULT_WORDLIST, uuid from django_q.humanhash import DEFAULT_WORDLIST, uuid
from django_q.models import Success, Task from django_q.models import Success, Task
from django_q.monitor import monitor, save_task from django_q.monitor import monitor, save_task
@@ -245,7 +246,7 @@ def test_cluster(broker):
assert queue_size(broker=broker) == 0 assert queue_size(broker=broker) == 0
# Test work # Test work
task_queue.put("STOP") 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 task_queue.qsize() == 0
assert result_queue.qsize() == 1 assert result_queue.qsize() == 1
# Test monitor # Test monitor
@@ -272,7 +273,7 @@ def test_results(broker):
pusher(task_queue, stop_event, broker=broker) pusher(task_queue, stop_event, broker=broker)
task_queue.put("STOP") task_queue.put("STOP")
result_queue = Queue() result_queue = Queue()
worker(task_queue, result_queue, Value("f", -1)) worker(task_queue, result_queue, Value("f", TimerStatus.IDLE))
result_queue.put("STOP") result_queue.put("STOP")
monitor(result_queue) monitor(result_queue)
@@ -372,7 +373,7 @@ def test_enqueue(broker, admin_user):
assert fetch_group("test_j", count=2, wait=10) is None assert fetch_group("test_j", count=2, wait=10) is None
# let a worker handle them # let a worker handle them
result_queue = Queue() 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 assert result_queue.qsize() == task_count
result_queue.put("STOP") result_queue.put("STOP")
# store the results # 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)
pusher(task_queue, stop_event, broker=broker) pusher(task_queue, stop_event, broker=broker)
# worker should exit on recycle # 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 # check if the work has been done
assert result_queue.qsize() == 2 assert result_queue.qsize() == 2
# save_limit test # save_limit test
@@ -573,7 +574,7 @@ def test_save_limit_per_func(broker, monkeypatch):
threading.Timer(3, stop_event.set).start() threading.Timer(3, stop_event.set).start()
for i in range(3): for i in range(3):
pusher(task_queue, stop_event, broker=broker) 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) s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker)
assert start_event.is_set() assert start_event.is_set()
assert s.status() == Conf.STOPPED assert s.status() == Conf.STOPPED
@@ -618,7 +619,7 @@ def test_max_rss(broker, monkeypatch):
# push the task # push the task
pusher(task_queue, stop_event, broker=broker) pusher(task_queue, stop_event, broker=broker)
# worker should exit on recycle # 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 # check if the work has been done
assert result_queue.qsize() == 1 assert result_queue.qsize() == 1
# save_limit test # save_limit test
@@ -654,7 +655,7 @@ def test_bad_secret(broker, monkeypatch):
worker( worker(
task_queue, task_queue,
result_queue, result_queue,
Value("f", -1), Value("f", TimerStatus.IDLE),
) )
assert result_queue.qsize() == 0 assert result_queue.qsize() == 0
broker.delete_queue() broker.delete_queue()
@@ -838,7 +839,7 @@ class TestSignals:
event.set() event.set()
pusher(task_queue, event, broker=broker) pusher(task_queue, event, broker=broker)
task_queue.put("STOP") 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") result_queue.put("STOP")
monitor(result_queue, broker) monitor(result_queue, broker)
broker.delete_queue() broker.delete_queue()
@@ -867,7 +868,7 @@ class TestSignals:
event.set() event.set()
pusher(task_queue, event, broker=broker) pusher(task_queue, event, broker=broker)
task_queue.put("STOP") 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") result_queue.put("STOP")
monitor(result_queue, broker) monitor(result_queue, broker)
broker.delete_queue() broker.delete_queue()
@@ -896,7 +897,7 @@ class TestSignals:
event.set() event.set()
pusher(task_queue, event, broker=broker) pusher(task_queue, event, broker=broker)
task_queue.put("STOP") 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") result_queue.put("STOP")
monitor(result_queue, broker) monitor(result_queue, broker)
broker.delete_queue() broker.delete_queue()
+2 -1
View File
@@ -13,6 +13,7 @@ from django.utils.timezone import is_naive
from django_q.brokers import Broker, get_broker from django_q.brokers import Broker, get_broker
from django_q.conf import Conf from django_q.conf import Conf
from django_q.enums import TimerStatus
from django_q.monitor import monitor from django_q.monitor import monitor
from django_q.pusher import pusher from django_q.pusher import pusher
from django_q.queues import Queue from django_q.queues import Queue
@@ -222,7 +223,7 @@ def test_scheduler(broker, monkeypatch):
task_queue.put("STOP") task_queue.put("STOP")
# let a worker handle them # let a worker handle them
result_queue = Queue() 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 assert result_queue.qsize() == 1
result_queue.put("STOP") result_queue.put("STOP")
# store the results # store the results
+7 -6
View File
@@ -17,6 +17,7 @@ except core.exceptions.AppRegistryNotReady:
django.setup() django.setup()
from django_q.conf import Conf, error_reporter, logger, resource, setproctitle 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.exceptions import TimeoutException
from django_q.signals import post_execute_in_worker, post_spawn, pre_execute from django_q.signals import post_execute_in_worker, post_spawn, pre_execute
from django_q.timeout import TimeoutHandler from django_q.timeout import TimeoutHandler
@@ -54,11 +55,11 @@ def worker(
setproctitle.setproctitle(f"qcluster {proc_name} idle") setproctitle.setproctitle(f"qcluster {proc_name} idle")
task_count = 0 task_count = 0
if timeout is None: if timeout is None:
timeout = -1 timeout = TimerStatus.IDLE
# Start reading the task queue # Start reading the task queue
for task in iter(task_queue.get, "STOP"): for task in iter(task_queue.get, "STOP"):
result = None result = None
timer.value = -1 # Idle timer.value = TimerStatus.IDLE
task_count += 1 task_count += 1
f = task["func"] f = task["func"]
@@ -92,7 +93,7 @@ def worker(
pre_execute.send(sender="django_q", func=f, task=task) pre_execute.send(sender="django_q", func=f, task=task)
# execute the payload # execute the payload
timer.value = timer_value # Busy 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 timer.value += 3 # Add buffer so that guard doesn't kill the process on timeout before it gets processed
timeout_error = False timeout_error = False
@@ -122,15 +123,15 @@ def worker(
result_queue.put(task) result_queue.put(task)
if timeout_error: if timeout_error:
# force destroy process due to timeout # force destroy process due to timeout
timer.value = 0 timer.value = TimerStatus.TIMEOUT
break break
timer.value = -1 # Idle timer.value = TimerStatus.IDLE
if setproctitle: if setproctitle:
setproctitle.setproctitle(f"qcluster {proc_name} idle") setproctitle.setproctitle(f"qcluster {proc_name} idle")
# Recycle # Recycle
if task_count == Conf.RECYCLE or rss_check(): if task_count == Conf.RECYCLE or rss_check():
timer.value = -2 # Recycled timer.value = TimerStatus.RECYCLED
break break
logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name}) logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name})