mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-24 21:18:11 +08:00
Merge pull request #297 from Eagllus/fix_199
Changing the path location where Django-Q is inserted.
This commit is contained in:
@@ -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)
|
||||
|
||||
|
||||
+7
-8
@@ -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)
|
||||
|
||||
+6
-6
@@ -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
|
||||
|
||||
|
||||
+49
-45
@@ -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']
|
||||
|
||||
Reference in New Issue
Block a user