diff --git a/django_q/cluster.py b/django_q/cluster.py index 4dd61f6..d62c846 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -21,12 +21,12 @@ from django.utils.translation import ugettext_lazy as _ from django import db # Local -import signing 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 +297,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 +456,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 +479,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..e64600b 100644 --- a/django_q/tasks.py +++ b/django_q/tasks.py @@ -7,7 +7,7 @@ from django.utils import timezone # local import time -import signing +from django_q.signing import SignedPackage import cluster from django_q.conf import Conf, logger from django_q.models import Schedule, Task @@ -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 @@ -132,7 +132,7 @@ def result_cached(task_id, wait=0, broker=None): while True: r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) if r: - return signing.SignedPackage.loads(r)['result'] + return SignedPackage.loads(r)['result'] if (time.time() - start) * 1000 >= wait >= 0: break time.sleep(0.01) @@ -182,7 +182,7 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non 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 @@ -225,7 +225,7 @@ def fetch_cached(task_id, wait=0, broker=None): 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'], @@ -285,7 +285,7 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None) 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'], @@ -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,) @@ -674,7 +674,7 @@ def _sync(pack): """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))