diff --git a/django_q/__init__.py b/django_q/__init__.py index 092c838..5ccb120 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -1,9 +1,9 @@ -import os -import sys +# import os +# import sys import django -myPath = os.path.dirname(os.path.abspath(__file__)) -sys.path.insert(0, myPath) +# myPath = os.path.dirname(os.path.abspath(__file__)) +# sys.path.insert(0, myPath) VERSION = (0, 9, 2) diff --git a/django_q/cluster.py b/django_q/cluster.py index 4dd61f6..403b649 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -21,12 +21,11 @@ from django.utils.translation import ugettext_lazy as _ from django import db # Local -import signing -import tasks - +from django_q import tasks from django_q.compat import range from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter, rollbar from django_q.models import Task, Success, Schedule +from django_q.signing import SignedPackage, BadSignature from django_q.status import Stat, Status from django_q.brokers import get_broker from django_q.signals import pre_execute @@ -297,8 +296,8 @@ def pusher(task_queue, event, broker=None): ack_id = task[0] # unpack the task try: - task = signing.SignedPackage.loads(task[1]) - except (TypeError, signing.BadSignature) as e: + task = SignedPackage.loads(task[1]) + except (TypeError, BadSignature) as e: logger.error(e) broker.fail(ack_id) continue @@ -456,11 +455,11 @@ def save_cached(task, broker): if iter_count and len(group_list) == iter_count - 1: group_args = '{}:{}:args'.format(broker.list_key, group) # collate the results into a Task result - results = [signing.SignedPackage.loads(broker.cache.get(k))['result'] for k in group_list] + results = [SignedPackage.loads(broker.cache.get(k))['result'] for k in group_list] results.append(task['result']) task['result'] = results task['id'] = group - task['args'] = signing.SignedPackage.loads(broker.cache.get(group_args)) + task['args'] = SignedPackage.loads(broker.cache.get(group_args)) task.pop('iter_count', None) task.pop('group', None) if task.get('iter_cached', None): @@ -479,7 +478,7 @@ def save_cached(task, broker): tasks.async_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker) # save the task broker.cache.set(task_key, - signing.SignedPackage.dumps(task), + SignedPackage.dumps(task), timeout) except Exception as e: logger.error(e) diff --git a/django_q/status.py b/django_q/status.py index db03825..51f32ee 100644 --- a/django_q/status.py +++ b/django_q/status.py @@ -2,7 +2,7 @@ import socket from django.utils import timezone from django_q.brokers import get_broker from django_q.conf import Conf, logger -import signing +from django_q.signing import SignedPackage, BadSignature class Status(object): @@ -64,7 +64,7 @@ class Stat(Status): def save(self): try: - self.broker.set_stat(self.key, signing.SignedPackage.dumps(self, True), 3) + self.broker.set_stat(self.key, SignedPackage.dumps(self, True), 3) except Exception as e: logger.error(e) @@ -83,8 +83,8 @@ class Stat(Status): pack = broker.get_stat(Stat.get_key(cluster_id)) if pack: try: - return signing.SignedPackage.loads(pack) - except signing.BadSignature: + return SignedPackage.loads(pack) + except BadSignature: return None return Status(cluster_id) @@ -101,8 +101,8 @@ class Stat(Status): packs = broker.get_stats('{}:*'.format(Conf.Q_STAT)) or [] for pack in packs: try: - stats.append(signing.SignedPackage.loads(pack)) - except signing.BadSignature: + stats.append(SignedPackage.loads(pack)) + except BadSignature: continue return stats diff --git a/django_q/tasks.py b/django_q/tasks.py index e64bfc9..bf94a2a 100644 --- a/django_q/tasks.py +++ b/django_q/tasks.py @@ -1,4 +1,6 @@ """Provides task functionality.""" +# Standard +from time import sleep, time from multiprocessing import Value # django @@ -6,9 +8,7 @@ from django.db import IntegrityError from django.utils import timezone # local -import time -import signing -import cluster +from django_q.signing import SignedPackage from django_q.conf import Conf, logger from django_q.models import Schedule, Task from django_q.humanhash import uuid @@ -48,7 +48,7 @@ def async(func, *args, **kwargs): # signal it pre_enqueue.send(sender="django_q", task=task) # sign it - pack = signing.SignedPackage.dumps(task) + pack = SignedPackage.dumps(task) if task.get('sync', False): return _sync(pack) # push it @@ -112,14 +112,14 @@ def result(task_id, wait=0, cached=Conf.CACHED): """ if cached: return result_cached(task_id, wait) - start = time.time() + start = time() while True: r = Task.get_result(task_id) if r: return r - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def result_cached(task_id, wait=0, broker=None): @@ -128,14 +128,14 @@ def result_cached(task_id, wait=0, broker=None): """ if not broker: broker = get_broker() - start = time.time() + start = time() while True: r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) if r: - return signing.SignedPackage.loads(r)['result'] - if (time.time() - start) * 1000 >= wait >= 0: + return SignedPackage.loads(r)['result'] + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHED): @@ -150,19 +150,19 @@ def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHE """ if cached: return result_group_cached(group_id, failures, wait, count) - start = time.time() + start = time() if count: while True: - if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0: + if count_group(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) while True: r = Task.get_result_group(group_id, failures) if r: return r - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def result_group_cached(group_id, failures=False, wait=0, count=None, broker=None): @@ -171,24 +171,24 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non """ if not broker: broker = get_broker() - start = time.time() + start = time() if count: while True: - if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait > 0: + if count_group_cached(group_id) == count or wait and (time() - start) * 1000 >= wait > 0: break - time.sleep(0.01) + sleep(0.01) while True: group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) if group_list: result_list = [] for task_key in group_list: - task = signing.SignedPackage.loads(broker.cache.get(task_key)) + task = SignedPackage.loads(broker.cache.get(task_key)) if task['success'] or failures: result_list.append(task['result']) return result_list - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def fetch(task_id, wait=0, cached=Conf.CACHED): @@ -205,14 +205,14 @@ def fetch(task_id, wait=0, cached=Conf.CACHED): """ if cached: return fetch_cached(task_id, wait) - start = time.time() + start = time() while True: t = Task.get_task(task_id) if t: return t - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def fetch_cached(task_id, wait=0, broker=None): @@ -221,11 +221,11 @@ def fetch_cached(task_id, wait=0, broker=None): """ if not broker: broker = get_broker() - start = time.time() + start = time() while True: r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) if r: - task = signing.SignedPackage.loads(r) + task = SignedPackage.loads(r) t = Task(id=task['id'], name=task['name'], func=task['func'], @@ -237,9 +237,9 @@ def fetch_cached(task_id, wait=0, broker=None): result=task['result'], success=task['success']) return t - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED): @@ -253,19 +253,19 @@ def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED) """ if cached: return fetch_group_cached(group_id, failures, wait, count) - start = time.time() + start = time() if count: while True: - if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0: + if count_group(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) while True: r = Task.get_task_group(group_id, failures) if r: return r - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None): @@ -274,18 +274,18 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None) """ if not broker: broker = get_broker() - start = time.time() + start = time() if count: while True: - if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0: + if count_group_cached(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) while True: group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) if group_list: task_list = [] for task_key in group_list: - task = signing.SignedPackage.loads(broker.cache.get(task_key)) + task = SignedPackage.loads(broker.cache.get(task_key)) if task['success'] or failures: t = Task(id=task['id'], name=task['name'], @@ -300,9 +300,9 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None) success=task['success']) task_list.append(t) return task_list - if (time.time() - start) * 1000 >= wait >= 0: + if (time() - start) * 1000 >= wait >= 0: break - time.sleep(0.01) + sleep(0.01) def count_group(group_id, failures=False, cached=Conf.CACHED): @@ -332,7 +332,7 @@ def count_group_cached(group_id, failures=False, broker=None): return len(group_list) failure_count = 0 for task_key in group_list: - task = signing.SignedPackage.loads(broker.cache.get(task_key)) + task = SignedPackage.loads(broker.cache.get(task_key)) if not task['success']: failure_count += 1 return failure_count @@ -405,7 +405,7 @@ def async_iter(func, args_iter, **kwargs): options['cached'] = True # save the original arguments broker = options['broker'] - broker.cache.set('{}:{}:args'.format(broker.list_key, iter_group), signing.SignedPackage.dumps(args_iter)) + broker.cache.set('{}:{}:args'.format(broker.list_key, iter_group), SignedPackage.dumps(args_iter)) for args in args_iter: if type(args) is not tuple: args = (args,) @@ -671,13 +671,17 @@ class Async(object): def _sync(pack): + # Python 2.6 is unable to handle this import on top of the file + # because it creates a circular dependency between tasks and cluster + from django_q.cluster import worker, monitor + """Simulate a package travelling through the cluster.""" task_queue = Queue() result_queue = Queue() - task = signing.SignedPackage.loads(pack) + task = SignedPackage.loads(pack) task_queue.put(task) task_queue.put('STOP') - cluster.worker(task_queue, result_queue, Value('f', -1)) + worker(task_queue, result_queue, Value('f', -1)) result_queue.put('STOP') - cluster.monitor(result_queue) + monitor(result_queue) return task['id']