diff --git a/django_q/__init__.py b/django_q/__init__.py index 4a9d8ad..ecb80d9 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -9,6 +9,6 @@ from .models import Task, Schedule, Success, Failure from .cluster import Cluster from .status import Stat -VERSION = (0, 6, 2) +VERSION = (0, 6, 3) default_app_config = 'django_q.apps.DjangoQConfig' diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index c247bb1..90c8d4b 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -153,6 +153,9 @@ def get_broker(list_key=Conf.PREFIX): if Conf.DISQUE_NODES: from brokers import disque return disque.Disque(list_key=list_key) + elif Conf.IRON_MQ: + from brokers import ironmq + return ironmq.IronMQBroker(list_key=list_key) # default to redis else: from brokers import redis_broker diff --git a/django_q/brokers/ironmq.py b/django_q/brokers/ironmq.py new file mode 100644 index 0000000..6388bf6 --- /dev/null +++ b/django_q/brokers/ironmq.py @@ -0,0 +1,44 @@ +from django_q.conf import Conf +from django_q.brokers import Broker +from iron_mq import IronMQ + + +class IronMQBroker(Broker): + + def enqueue(self, task): + return self.connection.post(task)['ids'][0] + + def dequeue(self): + timeout = Conf.RETRY or None + task = self.connection.get(timeout=timeout, wait=1)['messages'] + if task: + return task[0]['id'], task[0]['body'] + + def ping(self): + return self.connection.name == self.list_key + + def info(self): + return 'IronMQ' + + def queue_size(self): + return self.connection.size() + + def delete_queue(self): + return self.connection.delete_queue()['msg'] + + def purge_queue(self): + return self.connection.clear() + + 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) + + @staticmethod + def get_connection(list_key=Conf.PREFIX): + ironmq = IronMQ(name=None, **Conf.IRON_MQ) + return ironmq.queue(queue_name=list_key) \ No newline at end of file diff --git a/django_q/conf.py b/django_q/conf.py index ebd2842..6893702 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -37,6 +37,9 @@ class Conf(object): # Optional Authentication DISQUE_AUTH = conf.get('disque_auth', None) + # IronMQ broker + IRON_MQ = conf.get('iron_mq', None) + # Name of the cluster or site. For when you run multiple sites on one redis server PREFIX = conf.get('name', 'default') @@ -62,7 +65,7 @@ class Conf(object): WORKERS = 4 # Maximum number of tasks that each cluster can work on - QUEUE_LIMIT = conf.get('queue_limit', int(WORKERS)**2) + QUEUE_LIMIT = conf.get('queue_limit', int(WORKERS) ** 2) # Sets compression of redis packages COMPRESSED = conf.get('compress', False) diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 836e59f..09ce743 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -74,7 +74,7 @@ def test_disque(): broker.delete(task_id) assert broker.queue_size() == 0 # fail - task_id=broker.enqueue('test') + task_id = broker.enqueue('test') broker.fail(task_id) # delete queue broker.enqueue('test') @@ -83,3 +83,49 @@ def test_disque(): assert broker.queue_size() == 0 # back to django-redis Conf.DISQUE_NODES = None + + +@pytest.mark.skipif(not os.getenv('IRON_MQ_TOKEN'), + reason="requires IronMQ credentials") +def test_ironmq(): + Conf.IRON_MQ = {'host': os.getenv('IRON_MQ_HOST'), + 'token': os.getenv('IRON_MQ_TOKEN'), + 'project_id': os.getenv('IRON_MQ_PROJECT_ID')} + # check broker + broker = get_broker(list_key='djangoQ') + assert broker.ping() is True + assert broker.info() is not None + # clear before we start + broker.purge_queue() + # enqueue + broker.enqueue('test') + assert broker.queue_size() == 1 + # dequeue + task = broker.dequeue() + assert task[1] == 'test' + broker.acknowledge(task[0]) + assert broker.queue_size() == 0 + # Retry test + Conf.RETRY = 1 + broker.enqueue('test') + assert broker.queue_size() == 1 + assert broker.dequeue() is not None + sleep(1.5) + task = broker.dequeue() + assert len(task) > 0 + broker.acknowledge(task[0]) + sleep(1.5) + # delete job + task_id = broker.enqueue('test') + broker.delete(task_id) + assert broker.queue_size() == 0 + # fail + task_id = broker.enqueue('test') + broker.fail(task_id) + # delete queue + broker.enqueue('test') + broker.enqueue('test') + broker.purge_queue() + assert broker.queue_size() == 0 + # back to django-redis + Conf.IRON_MQ = None diff --git a/docs/conf.py b/docs/conf.py index bc9e9b6..191eb06 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -72,7 +72,7 @@ author = 'Ilan Steemers' # The short X.Y version. version = '0.6' # The full version, including alpha/beta/rc tags. -release = '0.6.2' +release = '0.6.3' # The language for content autogenerated by Sphinx. Refer to documentation # for a list of supported languages. diff --git a/requirements.in b/requirements.in index 1ab7ca1..c302143 100644 --- a/requirements.in +++ b/requirements.in @@ -6,3 +6,4 @@ hiredis redis psutil django-redis +iron-mq diff --git a/requirements.txt b/requirements.txt index 9fce959..4859074 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,9 +10,12 @@ django-picklefield==0.3.2 django-redis==4.2.0 future==0.15.0 hiredis==0.2.0 +iron-core==1.1.9 # via iron-mq +iron-mq==0.7 msgpack-python==0.4.6 # via django-redis psutil==3.2.1 -python-dateutil==2.4.2 # via arrow +python-dateutil==2.4.2 # via arrow, iron-core redis==2.10.3 +requests==2.7.0 # via iron-core six==1.9.0 # via python-dateutil wcwidth==0.1.4 # via blessed diff --git a/setup.py b/setup.py index d9ab8b0..ec25e3b 100644 --- a/setup.py +++ b/setup.py @@ -26,7 +26,7 @@ class PyTest(Command): setup( name='django-q', - version='0.6.2', + version='0.6.3', author='Ilan Steemers', author_email='koed00@gmail.com', keywords='django task queue worker redis disque multiprocessing',