From 4c413fc65957c928f3a1ed9ccc58f81d8651c825 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Thu, 10 Sep 2015 13:31:37 +0200 Subject: [PATCH] Prevents crash on duplicate acknowledge --- django_q/brokers/ironmq.py | 11 +++++++++-- django_q/tests/test_brokers.py | 11 +++++++++++ 2 files changed, 20 insertions(+), 2 deletions(-) diff --git a/django_q/brokers/ironmq.py b/django_q/brokers/ironmq.py index fd47667..08bb860 100644 --- a/django_q/brokers/ironmq.py +++ b/django_q/brokers/ironmq.py @@ -1,3 +1,4 @@ +from requests.exceptions import HTTPError from django_q.conf import Conf from django_q.brokers import Broker from iron_mq import IronMQ @@ -31,13 +32,19 @@ class IronMQBroker(Broker): return self.connection.size() def delete_queue(self): - return self.connection.delete_queue()['msg'] + try: + return self.connection.delete_queue()['msg'] + except HTTPError: + return False def purge_queue(self): return self.connection.clear() def delete(self, task_id): - return self.connection.delete(task_id)['msg'] + try: + return self.connection.delete(task_id)['msg'] + except HTTPError: + return False def fail(self, task_id): self.delete(task_id) diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 68df44e..9990414 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -85,6 +85,13 @@ def test_disque(): task = broker.dequeue() assert task is not None broker.acknowledge(task[0]) + # test duplicate acknowledge + # broker.acknowledge(task[0]) + # + # this crashes Disque when followed by a JSCAN + # https://github.com/antirez/disque/issues/113 + # confirmed fix and merge is on the way + # # delete queue broker.enqueue('test') broker.enqueue('test') @@ -139,6 +146,8 @@ def test_ironmq(): task = broker.dequeue() assert task is not None broker.acknowledge(task[0]) + # duplicate acknowledge + broker.acknowledge(task[0]) # delete queue broker.enqueue('test') broker.enqueue('test') @@ -195,6 +204,8 @@ def test_sqs(): task = broker.dequeue() assert task is not None broker.acknowledge(task[0]) + # duplicate acknowledge + broker.acknowledge(task[0]) # delete queue broker.enqueue('test') broker.purge_queue()