diff --git a/django_q/cluster.py b/django_q/cluster.py index d62c846..403b649 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -21,8 +21,7 @@ from django.utils.translation import ugettext_lazy as _ from django import db # Local -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 diff --git a/django_q/tasks.py b/django_q/tasks.py index e64600b..96a43fd 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,8 @@ from django.db import IntegrityError from django.utils import timezone # local -import time from django_q.signing import SignedPackage -import cluster +from django_q import cluster from django_q.conf import Conf, logger from django_q.models import Schedule, Task from django_q.humanhash import uuid @@ -112,14 +113,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 +129,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 SignedPackage.loads(r)['result'] - if (time.time() - start) * 1000 >= wait >= 0: + 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 +151,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,12 +172,12 @@ 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: @@ -186,9 +187,9 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non 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 +206,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,7 +222,7 @@ 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: @@ -237,9 +238,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 +254,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,12 +275,12 @@ 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: @@ -300,9 +301,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):