mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 03:48:12 +08:00
Add option for acknowledging failed tasks (globally and per-task)
If a task fails with an exception, it is retried until it succeeds. This is contrary to what is said in the documentation: under the "Architecture" section, heading "Broker" it says that even when a task errors, it's still considered a successful delivery. Failed tasks never get acknowledged however, thereby being retried after the timeout period. See also issues #238 and #194. This patch adds an option to acknowledge failures, thereby closing issue #238. Issue #194 would require some more work. The default of this option is set to `False`, thereby maintaining backwards compatibility.
This commit is contained in:
@@ -17,7 +17,7 @@ from django_q.tasks import fetch, fetch_group, async, result, result_group, coun
|
||||
from django_q.models import Task, Success
|
||||
from django_q.conf import Conf
|
||||
from django_q.status import Stat
|
||||
from django_q.brokers import get_broker
|
||||
from django_q.brokers import get_broker, Broker
|
||||
from django_q.tests.tasks import multiply
|
||||
from django_q.queues import Queue
|
||||
|
||||
@@ -379,6 +379,57 @@ def test_update_failed(broker):
|
||||
assert saved_task.success is True
|
||||
assert saved_task.result == 'result'
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_acknowledge_failure_override():
|
||||
class VerifyAckMockBroker(Broker):
|
||||
def __init__(self, *args, **kwargs):
|
||||
super(VerifyAckMockBroker, self).__init__(*args, **kwargs)
|
||||
self.acknowledgements = {}
|
||||
|
||||
def acknowledge(self, task_id):
|
||||
count = self.acknowledgements.get(task_id, 0)
|
||||
self.acknowledgements[task_id] = count + 1
|
||||
|
||||
tag = uuid()
|
||||
task_fail_ack = {'id': tag[1],
|
||||
'name': tag[0],
|
||||
'ack_id': 'test_fail_ack_id',
|
||||
'ack_failure': True,
|
||||
'func': 'math.copysign',
|
||||
'args': (1, -1),
|
||||
'kwargs': {},
|
||||
'started': timezone.now(),
|
||||
'stopped': timezone.now(),
|
||||
'success': False,
|
||||
'result': None}
|
||||
|
||||
tag = uuid()
|
||||
task_fail_no_ack = task_fail_ack.copy()
|
||||
task_fail_no_ack.update({'id': tag[1],
|
||||
'name': tag[0],
|
||||
'ack_id': 'test_fail_no_ack_id'})
|
||||
del task_fail_no_ack['ack_failure']
|
||||
|
||||
tag = uuid()
|
||||
task_success_ack = task_fail_ack.copy()
|
||||
task_success_ack.update({'id': tag[1],
|
||||
'name': tag[0],
|
||||
'ack_id': 'test_success_ack_id',
|
||||
'success': True,})
|
||||
del task_success_ack['ack_failure']
|
||||
|
||||
result_queue = Queue()
|
||||
result_queue.put(task_fail_ack)
|
||||
result_queue.put(task_fail_no_ack)
|
||||
result_queue.put(task_success_ack)
|
||||
result_queue.put('STOP')
|
||||
broker = VerifyAckMockBroker(list_key='key')
|
||||
|
||||
monitor(result_queue, broker)
|
||||
|
||||
assert broker.acknowledgements.get('test_fail_ack_id') == 1
|
||||
assert broker.acknowledgements.get('test_fail_no_ack_id') is None
|
||||
assert broker.acknowledgements.get('test_success_ack_id') == 1
|
||||
|
||||
@pytest.mark.django_db
|
||||
def assert_result(task):
|
||||
|
||||
Reference in New Issue
Block a user