adds fail method to broker

This commit is contained in:
Ilan Steemers
2015-09-03 16:34:00 +02:00
parent 7be533d9a6
commit 1ccfc66254
5 changed files with 17 additions and 0 deletions
+7
View File
@@ -55,6 +55,13 @@ class Broker(object):
"""
pass
def fail(self, task_id):
"""
Fails a task message
:param task_id:
:return:
"""
def ping(self):
"""
Checks whether the broker connection is available
+3
View File
@@ -33,6 +33,9 @@ class Sqs(Broker):
m.receipt_handle = task_id
return self.queue.delete_message(m)
def fail(self, task_id):
self.delete(task_id)
def delete_queue(self):
self.connection.delete_queue(self.queue)
+3
View File
@@ -27,6 +27,9 @@ class Disque(Broker):
def delete(self, task_id):
return self.connection.execute_command('DELJOB {}'.format(task_id))
def fail(self, task_id):
return self.delete(task_id)
def delete_queue(self):
jobs = self.connection.execute_command('JSCAN QUEUE {}'.format(self.list_key))[1]
if jobs:
+3
View File
@@ -32,6 +32,9 @@ class IronMQBroker(Broker):
def delete(self, task_id):
return self.connection.delete(task_id)['msg']
def fail(self, task_id):
self.delete(task_id)
def acknowledge(self, task_id):
return self.delete(task_id)
+1
View File
@@ -311,6 +311,7 @@ def pusher(task_queue, event, broker=None):
task = signing.SignedPackage.loads(task[1])
except (TypeError, signing.BadSignature) as e:
logger.error(e)
broker.fail(ack_id)
continue
task['ack_id'] = ack_id
task_queue.put(task)