From d6b9b10e3e2ac6a38c9613ffab68fa005ae253be Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Sun, 30 Aug 2015 20:26:30 +0200 Subject: [PATCH] Adds tests for brokers * test for redis and django-redis * test for disque TODO add disque to Travis --- django_q/brokers/__init__.py | 6 ++-- django_q/brokers/disque.py | 3 ++ django_q/tests/test_brokers.py | 61 ++++++++++++++++++++++++++++++++++ 3 files changed, 67 insertions(+), 3 deletions(-) create mode 100644 django_q/tests/test_brokers.py diff --git a/django_q/brokers/__init__.py b/django_q/brokers/__init__.py index 6ac89cf..9bdb7fa 100644 --- a/django_q/brokers/__init__.py +++ b/django_q/brokers/__init__.py @@ -6,7 +6,7 @@ class Broker(object): def __init__(self, list_key=Conf.Q_LIST): self.connection = self.get_connection() self.list_key = list_key - self.cache=self.get_cache() + self.cache = self.get_cache() def enqueue(self, task): pass @@ -27,7 +27,7 @@ class Broker(object): pass def set_stat(self, key, value, timeout): - key_list=self.cache.get(Conf.Q_STAT, []) + key_list = self.cache.get(Conf.Q_STAT, []) if key not in key_list: key_list.append(key) self.cache.set(Conf.Q_STAT, key_list) @@ -47,7 +47,7 @@ class Broker(object): stats.append(stat) else: key_list.remove(key) - self.cache.set(Conf.Q_STAT,key_list) + self.cache.set(Conf.Q_STAT, key_list) return stats @staticmethod diff --git a/django_q/brokers/disque.py b/django_q/brokers/disque.py index faff2c9..45f20d5 100644 --- a/django_q/brokers/disque.py +++ b/django_q/brokers/disque.py @@ -22,6 +22,9 @@ class Disque(Broker): def ping(self): return self.connection.ping() + def delete_queue(self, list_key=None): + raise NotImplementedError + @staticmethod def get_connection(): for node in Conf.DISQUE: diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py new file mode 100644 index 0000000..5c2020c --- /dev/null +++ b/django_q/tests/test_brokers.py @@ -0,0 +1,61 @@ +from time import sleep +import pytest +from django_q.conf import Conf +from django_q.brokers import get_broker, Broker + + +def test_broker(): + broker = Broker() + broker.enqueue('test') + broker.dequeue() + broker.queue_size() + broker.delete_queue() + broker.acknowledge('test') + broker.ping() + assert broker.get_stat('test_1') is None + broker.set_stat('test_1', 'test', 3) + assert broker.get_stat('test_1') == 'test' + assert broker.get_stats('test:*')[0] == 'test' + + +def test_redis(): + Conf.DJANGO_REDIS = None + broker = get_broker() + assert broker.ping() is True + Conf.REDIS = {'host': '127.0.0.1', 'port': 7712} + broker = get_broker() + with pytest.raises(Exception): + broker.ping() + + +def test_disque(): + Conf.DISQUE = ['127.0.0.1:7711'] + broker = get_broker() + assert broker.ping() is True + 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(1.5) + assert broker.queue_size() == 1 + task = broker.dequeue() + assert broker.queue_size() == 0 + broker.acknowledge(task[0]) + sleep(1.5) + assert broker.queue_size() == 0 + # errors + with pytest.raises(NotImplementedError): + broker.delete_queue() + Conf.DISQUE = ['127.0.0.1:7712', '127.0.0.1:7713'] + with pytest.raises(ConnectionError): + broker.get_connection() + # back to djangoredis + Conf.DJANGO_REDIS = 'default'