diff --git a/django_q/cluster.py b/django_q/cluster.py index 027b6a5..c351ee5 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -478,7 +478,12 @@ def save_task(task, broker: Broker): existing_task.stopped = task["stopped"] existing_task.result = task["result"] existing_task.success = task["success"] + existing_task.attempt_count = existing_task.attempt_count + 1 existing_task.save() + + if 0 < Conf.ATTEMPT_COUNT == existing_task.attempt_count: + broker.acknowledge(task['ack_id']) + else: Task.objects.create( id=task["id"], diff --git a/django_q/conf.py b/django_q/conf.py index 2e9ddc2..964364f 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -164,6 +164,9 @@ class Conf: # Optional error reporting setup ERROR_REPORTER = conf.get("error_reporter", {}) + # Optional attempt count. set to 0 for infinite attempts + ATTEMPT_COUNT = conf.get('attempt_count', 0) + # OSX doesn't implement qsize because of missing sem_getvalue() try: QSIZE = Queue().qsize() == 0 diff --git a/django_q/models.py b/django_q/models.py index 7037f12..7a05693 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -29,6 +29,7 @@ class Task(models.Model): started = models.DateTimeField(editable=False) stopped = models.DateTimeField(editable=False) success = models.BooleanField(default=True, editable=False) + attempt_count = models.IntegerField(default=0) @staticmethod def get_result(task_id):