diff --git a/django_q/cluster.py b/django_q/cluster.py index 80bbea2..3bc443b 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -176,7 +176,7 @@ class Sentinel(object): return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key, self.r) def spawn_worker(self): - self.spawn_process(worker, self.task_queue, self.result_queue, Value('b', -1)) + self.spawn_process(worker, self.task_queue, self.result_queue, Value('b', -1), self.timeout) def spawn_monitor(self): return self.spawn_process(monitor, self.result_queue) @@ -195,7 +195,7 @@ class Sentinel(object): else: self.pool.remove(process) self.spawn_worker() - if self.timeout and int(process.timer.value) >= self.timeout: + if self.timeout and int(process.timer.value) == 0: # only need to terminate on timeout, otherwise we risk destabilizing the queues process.terminate() logger.warn(_("reincarnated worker {} after timeout").format(process.name)) @@ -231,12 +231,12 @@ class Sentinel(object): # Check Workers for p in self.pool: # Are you alive? - if not p.is_alive() or (self.timeout and int(p.timer.value) >= self.timeout): + if not p.is_alive() or (self.timeout and int(p.timer.value) == 0): self.reincarnate(p) continue - # Increment timer if work is being done - if p.timer.value >= 0: - p.timer.value += 1 + # Decrement timer if work is being done + if p.timer.value > 0: + p.timer.value -= 1 # Check Monitor if not self.monitor.is_alive(): self.reincarnate(self.monitor) @@ -335,7 +335,7 @@ def monitor(result_queue): logger.info(_("{} stopped monitoring results").format(name)) -def worker(task_queue, result_queue, timer): +def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT): """ Takes a task from the task queue, tries to execute it and puts the result back in the result queue :type task_queue: multiprocessing.Queue @@ -370,7 +370,7 @@ def worker(task_queue, result_queue, timer): # We're still going if not result: # execute the payload - timer.value = 0 # Busy + timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy try: res = f(*task['args'], **task['kwargs']) result = (res, True) @@ -384,7 +384,7 @@ def worker(task_queue, result_queue, timer): timer.value = -1 # Idle # Recycle if task_count == Conf.RECYCLE: - timer.value = -2 + timer.value = -2 # Recycled break logger.info(_('{} stopped doing work').format(name)) diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index a6d9b2d..32fba17 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -213,6 +213,60 @@ def test_timeout(r): assert start_event.is_set() assert s.status() == Conf.STOPPED assert s.reincarnations == 1 + r.delete(list_key) + + +@pytest.mark.django_db +def test_timeout(r): + # set up the Sentinel + list_key = 'timeout_test:q' + async('django_q.tests.tasks.count_forever', list_key=list_key) + start_event = Event() + stop_event = Event() + # Set a timer to stop the Sentinel + threading.Timer(3, stop_event.set).start() + s = Sentinel(stop_event, start_event, list_key=list_key, timeout=1) + assert start_event.is_set() + assert s.status() == Conf.STOPPED + assert s.reincarnations == 1 + r.delete(list_key) + + +@pytest.mark.django_db +def test_timeout_override(r): + # set up the Sentinel + list_key = 'timeout_override_test:q' + async('django_q.tests.tasks.count_forever', list_key=list_key, timeout=1) + start_event = Event() + stop_event = Event() + # Set a timer to stop the Sentinel + threading.Timer(3, stop_event.set).start() + s = Sentinel(stop_event, start_event, list_key=list_key, timeout=10) + assert start_event.is_set() + assert s.status() == Conf.STOPPED + assert s.reincarnations == 1 + r.delete(list_key) + + +@pytest.mark.django_db +def test_recycle(r): + # set up the Sentinel + list_key = 'test_recycle_test:q' + async('django_q.tests.tasks.multiply', 2, 2, list_key=list_key) + async('django_q.tests.tasks.multiply', 2, 2, list_key=list_key) + async('django_q.tests.tasks.multiply', 2, 2, list_key=list_key) + start_event = Event() + stop_event = Event() + # override settings + Conf.RECYCLE = 2 + Conf.WORKERS = 1 + # Set a timer to stop the Sentinel + threading.Timer(3, stop_event.set).start() + s = Sentinel(stop_event, start_event, list_key=list_key) + assert start_event.is_set() + assert s.status() == Conf.STOPPED + assert s.reincarnations == 1 + r.delete(list_key) @pytest.mark.django_db diff --git a/docs/install.rst b/docs/install.rst index 3dd1c72..e7c3bac 100644 --- a/docs/install.rst +++ b/docs/install.rst @@ -65,11 +65,13 @@ recycle The number of tasks a worker will process before recycling . Useful to release memory resources on a regular basis. Defaults to ``500``. +.. _timeout: + timeout ~~~~~~~ The number of seconds a worker is allowed to spend on a task before it's terminated. Defaults to ``None``, meaning it will never time out. -Set this to something that makes sense for your project. +Set this to something that makes sense for your project. Can be overridden for individual tasks. compress ~~~~~~~~ diff --git a/docs/tasks.rst b/docs/tasks.rst index c16b842..c3fd630 100644 --- a/docs/tasks.rst +++ b/docs/tasks.rst @@ -78,7 +78,7 @@ When you are making individual calls to :func:`async` a lot though, it can help Reference --------- -.. py:function:: async(func, *args, hook=None, sync=False, redis=None, **kwargs) +.. py:function:: async(func, *args, hook=None, timeout=None, sync=False, redis=None, **kwargs) Puts a task in the cluster queue @@ -87,6 +87,7 @@ Reference :type func: object :param hook: Optional function to call after execution :type hook: object + :param int timeout: Overrides global cluster :ref:`timeout`. :param bool sync: If set to True, async will simulate a task execution :param redis: Optional redis connection :param kwargs: Keyword arguments for the task function