Fix async_task timeout parameter handling when cluster timeout is set to None (the default)

The cluster timeout configuration default value None is documented to mean
tasks never timeout out. Also the documentation states that the timeout
can be overridden for individual tasks.

With the old implementation the timeout parameter given to async_task was
not honored if the cluster timeout was set to None.

Fixes: #335
This commit is contained in:
Janne Rönkkö
2019-01-27 19:09:22 +02:00
parent 2650659da9
commit 0e2df88d92
2 changed files with 8 additions and 4 deletions
+6 -4
View File
@@ -172,7 +172,7 @@ class Sentinel(object):
else: else:
self.pool.remove(process) self.pool.remove(process)
self.spawn_worker() self.spawn_worker()
if self.timeout and int(process.timer.value) == 0: if process.timer.value == 0:
# only need to terminate on timeout, otherwise we risk destabilizing the queues # only need to terminate on timeout, otherwise we risk destabilizing the queues
process.terminate() process.terminate()
logger.warn(_("reincarnated worker {} after timeout").format(process.name)) logger.warn(_("reincarnated worker {} after timeout").format(process.name))
@@ -210,11 +210,11 @@ class Sentinel(object):
# Check Workers # Check Workers
for p in self.pool: for p in self.pool:
# Are you alive? # Are you alive?
if not p.is_alive() or (self.timeout and p.timer.value == 0): if not p.is_alive() or p.timer.value == 0:
self.reincarnate(p) self.reincarnate(p)
continue continue
# Decrement timer if work is being done # Decrement timer if work is being done
if self.timeout and p.timer.value > 0: if p.timer.value > 0:
p.timer.value -= cycle p.timer.value -= cycle
# Check Monitor # Check Monitor
if not self.monitor.is_alive(): if not self.monitor.is_alive():
@@ -347,6 +347,8 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
name = current_process().name name = current_process().name
logger.info(_('{} ready for work at {}').format(name, current_process().pid)) logger.info(_('{} ready for work at {}').format(name, current_process().pid))
task_count = 0 task_count = 0
if timeout is None:
timeout = -1
# 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
@@ -368,7 +370,7 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
# We're still going # We're still going
if not result: if not result:
db.close_old_connections() db.close_old_connections()
timer_value = task['kwargs'].pop('timeout', timeout or 0) timer_value = task['kwargs'].pop('timeout', timeout)
# signal execution # signal execution
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
+2
View File
@@ -247,6 +247,7 @@ def test_enqueue(broker, admin_user):
@pytest.mark.parametrize('cluster_config_timeout, async_task_kwargs', ( @pytest.mark.parametrize('cluster_config_timeout, async_task_kwargs', (
(1, {}), (1, {}),
(10, {'timeout': 1}), (10, {'timeout': 1}),
(None, {'timeout': 1}),
)) ))
def test_timeout(broker, cluster_config_timeout, async_task_kwargs): def test_timeout(broker, cluster_config_timeout, async_task_kwargs):
# set up the Sentinel # set up the Sentinel
@@ -269,6 +270,7 @@ def test_timeout(broker, cluster_config_timeout, async_task_kwargs):
(5, {}), (5, {}),
(10, {'timeout': 5}), (10, {'timeout': 5}),
(1, {'timeout': 5}), (1, {'timeout': 5}),
(None, {'timeout': 5}),
)) ))
def test_timeout_task_finishes(broker, cluster_config_timeout, async_task_kwargs): def test_timeout_task_finishes(broker, cluster_config_timeout, async_task_kwargs):
# set up the Sentinel # set up the Sentinel