From 67180c1984f27e006021d13290d408f80e9a5b10 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Wed, 24 Jun 2015 15:29:16 +0200 Subject: [PATCH] Initial tests --- django_q/__init__.py | 23 +++++++++++- django_q/core.py | 23 ++++++++---- django_q/tests/test_q.py | 80 +++++++++++++++++++++++----------------- 3 files changed, 83 insertions(+), 43 deletions(-) diff --git a/django_q/__init__.py b/django_q/__init__.py index 53e495a..6f27633 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -1,7 +1,26 @@ from django_q.models import Task -from django_q.core import async, Cluster +from django_q.core import async default_app_config = 'django_q.apps.DjangoQConfig' + def result(name): - return Task.get_result(name) \ No newline at end of file + """ + Returns the result of the named task + :type name: str or unicode + :param name: the task name + :return: the result object of this task + :rtype: object or str + """ + return Task.get_result(name) + + +def get_task(name): + """ + Returns the processed task + :param name: the task name + :type name: str or unicode + :return: the full task object + :rtype: Task + """ + return Task.objects.get(name=name) diff --git a/django_q/core.py b/django_q/core.py index a09748c..679da63 100644 --- a/django_q/core.py +++ b/django_q/core.py @@ -56,7 +56,7 @@ def time_zone(value): return value class Cluster(object): - def __init__(self): + def __init__(self, list_key=Q_LIST): try: r.ping() except (): @@ -67,6 +67,7 @@ class Cluster(object): self.start_event = None self.stopped_event = None self.pid = current_process().pid + self.list_key = list_key signal.signal(signal.SIGTERM, self.sig_handler) signal.signal(signal.SIGINT, self.sig_handler) @@ -80,7 +81,7 @@ class Cluster(object): # Start Sentinel self.stop_event = Event() self.start_event = Event() - self.sentinel = Process(target=Sentinel, args=(self.stop_event, self.start_event)) + self.sentinel = Process(target=Sentinel, args=(self.stop_event, self.start_event, self.list_key)) self.sentinel.start() logger.info('Q Cluster-{} starting.'.format(self.pid)) return self.pid @@ -128,12 +129,13 @@ class Cluster(object): class Sentinel(object): - def __init__(self, stop_event, start_event): + def __init__(self, stop_event, start_event, list_key=Q_LIST): signal.signal(signal.SIGINT, signal.SIG_IGN) signal.signal(signal.SIGTERM, signal.SIG_DFL) self.pid = current_process().pid self.parent_pid = os.getppid() self.name = current_process().name + self.list_key = list_key self.status = None self.reincarnations = 0 self.tob = datetime.utcnow() @@ -163,7 +165,7 @@ class Sentinel(object): return p.pid def spawn_pusher(self): - return self.spawn_process(pusher, self.task_queue, self.event_out) + return self.spawn_process(pusher, self.task_queue, self.event_out, self.list_key) def spawn_worker(self): self.spawn_process(worker, self.task_queue, self.done_queue) @@ -231,10 +233,10 @@ class Sentinel(object): Stat(self, message).save() -def pusher(task_queue, e): +def pusher(task_queue, e, list_key=Q_LIST): logger.info('{} pushing tasks at {}'.format(current_process().name, current_process().pid)) while True: - task = r.blpop(Q_LIST, 1) + task = r.blpop(list_key, 1) if task: task = task[1] task_queue.put(task) @@ -309,14 +311,21 @@ def async(func, *args, **kwargs): """ Schedules a task with optional hook """ + # Check for hook if 'hook' in kwargs: hook = kwargs['hook'] del kwargs['hook'] else: hook = None + # Check for list_key override + if 'list_key' in kwargs: + list_key = kwargs['list_key'] + del kwargs['list_key'] + else: + list_key = Q_LIST task = {'name': uuid()[0], 'func': func, 'hook': hook, 'args': args, 'kwargs': kwargs, 'started': datetime.utcnow()} pack = SignedPackage.dumps(task) - r.rpush(Q_LIST, pack) + r.rpush(list_key, pack) logger.debug('Pushed {}'.format(pack)) return task['name'] diff --git a/django_q/tests/test_q.py b/django_q/tests/test_q.py index 497fbef..c76bef5 100644 --- a/django_q/tests/test_q.py +++ b/django_q/tests/test_q.py @@ -8,9 +8,9 @@ import pytest myPath = os.path.dirname(os.path.abspath(__file__)) sys.path.insert(0, myPath + '/../') -from django_q.core import Cluster, r, async, Q_LIST, pusher, worker, monitor +from django_q.core import Cluster, r, async, pusher, worker, monitor, Sentinel from django_q.humanhash import DEFAULT_WORDLIST -from django_q import result +from django_q import result, get_task class WordClass(object): @@ -50,11 +50,19 @@ def test_cluster_initial(): assert c.has_stopped +def test_sentinel(): + start_event = Event() + stop_event = Event() + stop_event.set() + Sentinel(stop_event, start_event, list_key='sentinel_test:q') + assert start_event.is_set() + @pytest.mark.django_db def test_cluster(): - task = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST) - task_count = r.llen(Q_LIST) - assert task_count >= 1 + list_key = 'cluster_test:q' + r.delete(list_key) + task = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, list_key=list_key) + assert r.llen(list_key) == 1 task_queue = Queue() assert task_queue.qsize() == 0 result_queue = Queue() @@ -62,9 +70,9 @@ def test_cluster(): event = Event() event.set() # Test push - pusher(task_queue, event) + pusher(task_queue, event, list_key=list_key) assert task_queue.qsize() == 1 - assert r.llen(Q_LIST) == task_count - 1 + assert r.llen(list_key) == 0 # Test work task_queue.put('STOP') worker(task_queue, result_queue) @@ -74,33 +82,51 @@ def test_cluster(): result_queue.put('STOP') monitor(result_queue) assert result_queue.qsize() == 0 - # check result assert result(task) == 1506 - + r.delete(list_key) @pytest.mark.django_db -def test_async(): - a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_q.assert_result') - b = async('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_q.assert_result') - c = Cluster() +def run_cluster(): + list_key = 'run_test:q' + r.delete(list_key) + c = Cluster(list_key=list_key) assert c.start() > 0 while not c.is_running: sleep(0.5) - assert isinstance(a, str) - assert isinstance(b, str) while c.stat.task_q_size > 0 and c.stat.done_q_size > 0: sleep(0.5) assert c.stop() is True - result_a = result(a) + r.delete(list_key) + +@pytest.mark.django_db +def blah_async(): + a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_q.assert_result') + b = async('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_q.assert_result') + # unknown argument + c = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany', + hook='django_q.tests.test_q.assert_bad_result') + # unknown function + d = async('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_q.assert_bad_result') + assert isinstance(a, str) + assert isinstance(b, str) + assert isinstance(c, str) + assert isinstance(d, str) + run_cluster() + result_a = get_task(a) assert result_a is not None assert result_a.success is True - assert result_a.result == 1506 - result_b = result(b) + assert result(a) == 1506 + result_b = get_task(b) assert result_b is not None assert result_b.success is True - assert result_b.result == 1506 - + assert result(b) == 1506 + result_c = get_task(c) + assert result_c is not None + assert result_c.success is False + result_d = get_task(d) + assert result_d is not None + assert result_d.success is False @pytest.mark.django_db @@ -110,20 +136,6 @@ def assert_result(task): assert task.result == 1506 -@pytest.mark.django_db -def broken_package(): - a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany', - hook='django_q.tests.test_q.assert_bad_result') - b = async('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_q.assert_bad_result') - assert isinstance(a, str) - assert isinstance(b, str) - sleep(5) - result_a = result(a) - assert result_a.success is False - result_b = result(b) - assert result_b.success is False - - @pytest.mark.django_db def assert_bad_result(task): assert task is not None