From 3f75e7fa05c55741b2369112068f28f130bc188a Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Fri, 4 Sep 2015 14:49:17 +0200 Subject: [PATCH] removed SQS and IronMQ brokers for now Both brokers are not atomic and cause all kinds of problems during testing. Needs more time. Meanwhile I want to release the pluggable broker backend. --- README.rst | 2 +- django_q/brokers/__init__.py | 9 +--- django_q/brokers/aws_sqs.py | 63 ----------------------- django_q/brokers/disque.py | 4 +- django_q/brokers/iron_mq.py | 44 ---------------- django_q/brokers/redis_broker.py | 3 ++ django_q/conf.py | 31 ----------- django_q/tests/tasks.py | 9 ---- django_q/tests/test_brokers.py | 88 ++++++++------------------------ django_q/tests/test_cluster.py | 10 ++-- django_q/tests/test_config.py | 10 ---- setup.py | 2 +- 12 files changed, 36 insertions(+), 239 deletions(-) delete mode 100644 django_q/brokers/aws_sqs.py delete mode 100644 django_q/brokers/iron_mq.py delete mode 100644 django_q/tests/test_config.py diff --git a/README.rst b/README.rst index 9f8550b..dde7c86 100644 --- a/README.rst +++ b/README.rst @@ -20,7 +20,7 @@ Features - Django Admin integration - PaaS compatible with multiple instances - Multi cluster monitor -- Redis broker +- Redis broker and Disque broker - Python 2 and 3 Requirements diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index 4dffaf0..c247bb1 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -149,15 +149,10 @@ def get_broker(list_key=Conf.PREFIX): :type list_key: str :return: """ - if Conf.IRONMQ: - from brokers import iron_mq - return iron_mq.IronMQBroker(list_key=list_key) - elif Conf.DISQUE_NODES: + # disque + if Conf.DISQUE_NODES: from brokers import disque return disque.Disque(list_key=list_key) - elif Conf.SQS: - from brokers import aws_sqs - return aws_sqs.Sqs(list_key=list_key) # default to redis else: from brokers import redis_broker diff --git a/django_q/brokers/aws_sqs.py b/django_q/brokers/aws_sqs.py deleted file mode 100644 index 2b9e254..0000000 --- a/django_q/brokers/aws_sqs.py +++ /dev/null @@ -1,63 +0,0 @@ -import os -from django_q.conf import Conf -from django_q.brokers import Broker -import boto.sqs -from boto.sqs.message import Message - - -class Sqs(Broker): - def __init__(self, list_key=Conf.PREFIX): - super(Sqs, self).__init__(list_key) - self.queue = self.get_queue() - - def enqueue(self, task): - m = Message() - m.set_body(task) - self.queue.write(m) - return m.id - - def dequeue(self): - rs = self.queue.get_messages(visibility_timeout=Conf.RETRY or 30) - if rs: - m = rs[0] - return m.receipt_handle, m.get_body() - - def acknowledge(self, task_id): - return self.delete(task_id) - - def queue_size(self): - return self.queue.count() - - def delete(self, task_id): - m = Message() - 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) - - def purge_queue(self): - self.queue.purge() - - def ping(self): - try: - self.connection.get_all_queues() - return True - except Exception as e: - raise e - - def info(self): - return 'AWS SQS' - - @staticmethod - def get_connection(list_key=Conf.PREFIX): - conn = boto.sqs.connect_to_region(Conf.SQS['region'], - aws_access_key_id=Conf.SQS['aws_access_key_id'], - aws_secret_access_key=Conf.SQS['aws_secret_access_key']) - return conn - - def get_queue(self): - return self.connection.create_queue(self.list_key) diff --git a/django_q/brokers/disque.py b/django_q/brokers/disque.py index 8ae163e..7e43ccd 100644 --- a/django_q/brokers/disque.py +++ b/django_q/brokers/disque.py @@ -33,7 +33,9 @@ class Disque(Broker): def delete_queue(self): jobs = self.connection.execute_command('JSCAN QUEUE {}'.format(self.list_key))[1] if jobs: - self.connection.execute_command('DELJOB {}'.format(' '.join(map(str, jobs)))) + job_ids = ' '.join(jid.decode() for jid in jobs) + self.connection.execute_command('DELJOB {}'.format(job_ids)) + return len(jobs) def info(self): info = self.connection.info('server') diff --git a/django_q/brokers/iron_mq.py b/django_q/brokers/iron_mq.py deleted file mode 100644 index d59d2f4..0000000 --- a/django_q/brokers/iron_mq.py +++ /dev/null @@ -1,44 +0,0 @@ -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()['msg'] - - 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.IRONMQ) - return ironmq.queue(queue_name=list_key) diff --git a/django_q/brokers/redis_broker.py b/django_q/brokers/redis_broker.py index 8ef157c..68289aa 100644 --- a/django_q/brokers/redis_broker.py +++ b/django_q/brokers/redis_broker.py @@ -27,6 +27,9 @@ class Redis(Broker): def delete_queue(self): return self.connection.delete(self.list_key) + def purge_queue(self): + return self.connection.ltrim(self.list_key, 1, 0) + def ping(self): try: return self.connection.ping() diff --git a/django_q/conf.py b/django_q/conf.py index 620c9d8..024d9fc 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -8,7 +8,6 @@ from django.conf import settings # external import os -import redis # optional try: @@ -38,14 +37,6 @@ class Conf(object): # Optional Authentication DISQUE_AUTH = conf.get('disque_auth', None) - # Amazon SQS broker - SQS = conf.get('sqs', None) - - # IronMQ broker - IRONMQ = conf.get('ironmq', None) - if IRONMQ and os.environ.get('IRONMQ_TOKEN'): - IRONMQ['token'] = os.environ['IRONMQ_TOKEN'] - # Name of the cluster or site. For when you run multiple sites on one redis server PREFIX = conf.get('name', 'default') @@ -142,28 +133,6 @@ if not logger.handlers: logger.addHandler(handler) -# Django-redis support -if Conf.DJANGO_REDIS: - try: - import django_redis - except ImportError: - django_redis = None - - -def get_redis_client(): - """ - Returns a connection from redis-py or django-redis - :return: a redis client - """ - if Conf.DJANGO_REDIS and django_redis: - return django_redis.get_redis_connection(Conf.DJANGO_REDIS) - return redis.StrictRedis(**Conf.REDIS) - - -# redis client -redis_client = get_redis_client() - - # get parent pid compatibility def get_ppid(): if hasattr(os, 'getppid'): diff --git a/django_q/tests/tasks.py b/django_q/tests/tasks.py index 9c8f8ad..70eb782 100644 --- a/django_q/tests/tasks.py +++ b/django_q/tests/tasks.py @@ -1,7 +1,3 @@ -# simple countdown, returns nothing -from time import sleep - - def countdown(n): while n > 0: n -= 1 @@ -26,11 +22,6 @@ def word_multiply(x, word=''): return len(word) * x -def count_forever(): - while True: - sleep(0.5) - - def get_task_name(task): return task.name diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index b0bbfca..feef253 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -11,9 +11,12 @@ def test_broker(): broker.enqueue('test') broker.dequeue() broker.queue_size() + broker.purge_queue() + broker.delete('id') broker.delete_queue() broker.acknowledge('test') broker.ping() + broker.info() assert broker.get_stat('test_1') is None broker.set_stat('test_1', 'test', 3) assert broker.get_stat('test_1') == 'test' @@ -24,6 +27,7 @@ def test_redis(): Conf.DJANGO_REDIS = None broker = get_broker() assert broker.ping() is True + assert broker.info() is not None Conf.REDIS = {'host': '127.0.0.1', 'port': 7712} broker = get_broker() with pytest.raises(Exception): @@ -32,15 +36,18 @@ def test_redis(): Conf.DJANGO_REDIS = 'default' -@pytest.mark.skipif(not os.getenv('DISQUE', None), - reason="No disque server configured") def test_disque(): Conf.DISQUE_NODES = ['127.0.0.1:7711'] + Conf.DISQUE_AUTH = 'foobared' broker = get_broker(list_key='disque_test') assert broker.ping() is True + assert broker.info() is not None + # clear before we start broker.delete_queue() + # enqueue broker.enqueue('test') assert broker.queue_size() == 1 + # dequeue task = broker.dequeue() assert task[1] == 'test' broker.acknowledge(task[0]) @@ -58,76 +65,21 @@ def test_disque(): broker.acknowledge(task[0]) sleep(1.5) assert broker.queue_size() == 0 + # connection test Conf.DISQUE_NODES = ['127.0.0.1:7712', '127.0.0.1:7713'] with pytest.raises(redis.exceptions.ConnectionError): broker.get_connection() + # 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.delete_queue() assert broker.queue_size() == 0 # back to django-redis Conf.DISQUE_NODES = None - - -@pytest.mark.skipif(not os.getenv('AWS_ACCESS_KEY_ID'), - reason="requires AWS SQS credentials") -def test_sqs(): - Conf.SQS = {'aws_access_key_id': os.getenv('AWS_ACCESS_KEY_ID'), - 'aws_secret_access_key': os.getenv('AWS_SECRET_ACCESS_KEY'), - 'region': os.getenv('SQS_REGION', 'eu-west-1')} - broker = get_broker(list_key='sqs_test') - assert broker.ping() is True - if broker.queue_size() > 0: - broker.purge_queue() - broker.enqueue('test') - task = broker.dequeue() - assert task[1] == 'test' - broker.acknowledge(task[0]) - assert broker.queue_size() == 0 - # Retry test - Conf.RETRY = 1 - broker.enqueue('test') - broker.dequeue() - assert broker.queue_size() == 0 - sleep(2) - # task should re-queue - assert broker.queue_size() == 1 - task = broker.dequeue() - assert task[1] == 'test' - assert broker.acknowledge(task[0]) is True - assert broker.queue_size() == 0 - broker.delete_queue() - # back to defaults - Conf.SQS = None - - -@pytest.mark.skipif(not os.getenv('IRONMQ_TOKEN'), - reason="requires IronMQ credentials") -def test_ironmq(): - Conf.IRONMQ = {'host': 'mq-aws-eu-west-1.iron.io', - 'token': os.getenv('IRONMQ_TOKEN'), - 'project_id': os.getenv('IRONMQ_PROJECT')} - broker = get_broker(list_key='djangoQ_test') - assert broker.ping() is True - broker.delete_queue() - broker.enqueue('test') - assert broker.queue_size() == 1 - 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 - broker.dequeue() - assert broker.queue_size() == 0 - sleep(2) - assert broker.queue_size() == 1 - task = broker.dequeue() - assert broker.queue_size() == 0 - broker.acknowledge(task[0]) - sleep(2) - assert broker.queue_size() == 0 - broker.delete_queue() - assert broker.queue_size() == 0 - # back to django-redis - Conf.IRONMQ = None diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index f35c3c3..16d131e 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -226,7 +226,8 @@ def test_async(broker, admin_user): def test_timeout(broker): # set up the Sentinel broker.list_key = 'timeout_test:q' - async('django_q.tests.tasks.count_forever',broker=broker) + broker.purge_queue() + async('django_q.tests.tasks.count_forever', broker=broker) start_event = Event() stop_event = Event() # Set a timer to stop the Sentinel @@ -242,12 +243,13 @@ def test_timeout(broker): def test_timeout(broker): # set up the Sentinel broker.list_key = 'timeout_test:q' + broker.purge_queue() async('django_q.tests.tasks.count_forever', broker=broker) start_event = Event() stop_event = Event() # Set a timer to stop the Sentinel threading.Timer(3, stop_event.set).start() - s = Sentinel(stop_event, start_event,broker=broker, timeout=1) + s = Sentinel(stop_event, start_event, broker=broker, timeout=1) assert start_event.is_set() assert s.status() == Conf.STOPPED assert s.reincarnations == 1 @@ -284,7 +286,7 @@ def test_recycle(broker): Conf.WORKERS = 1 # set a timer to stop the Sentinel threading.Timer(3, stop_event.set).start() - s = Sentinel(stop_event, start_event,broker=broker) + s = Sentinel(stop_event, start_event, broker=broker) assert start_event.is_set() assert s.status() == Conf.STOPPED assert s.reincarnations == 1 @@ -310,7 +312,7 @@ def test_recycle(broker): @pytest.mark.django_db def test_bad_secret(broker, monkeypatch): - broker.list_key='test_bad_secret:q' + broker.list_key = 'test_bad_secret:q' async('math.copysign', 1, -1, broker=broker) stop_event = Event() stop_event.set() diff --git a/django_q/tests/test_config.py b/django_q/tests/test_config.py deleted file mode 100644 index 10ade3d..0000000 --- a/django_q/tests/test_config.py +++ /dev/null @@ -1,10 +0,0 @@ -import pytest - -from django_q import conf - - -def test_django_redis(): - conf.Conf.DJANGO_REDIS = None - assert conf.redis_client.ping() is True - conf.Conf.DJANGO_REDIS = 'default' - assert conf.redis_client.ping() is True diff --git a/setup.py b/setup.py index 35b0afc..1631e60 100644 --- a/setup.py +++ b/setup.py @@ -29,7 +29,7 @@ setup( version='0.6.0', author='Ilan Steemers', author_email='koed00@gmail.com', - keywords='django task queue worker redis multiprocessing', + keywords='django task queue worker redis disque multiprocessing', packages=['django_q'], include_package_data=True, url='https://django-q.readthedocs.org',