Replacing format with f-strings

This commit is contained in:
Ilan
2020-02-17 16:59:12 +01:00
parent a1543c8d82
commit 0bca3cbd7a
7 changed files with 683 additions and 426 deletions
+178 -113
View File
@@ -4,27 +4,30 @@ from __future__ import division
from __future__ import print_function from __future__ import print_function
from __future__ import unicode_literals from __future__ import unicode_literals
from time import sleep
# external
import arrow
import ast import ast
# Standard # Standard
import importlib import importlib
import signal import signal
import socket import socket
import traceback import traceback
import uuid import uuid
from multiprocessing import Event, Process, Value, current_process
from time import sleep
# external
import arrow
# Django # Django
from django import db from django import db
from django.utils import timezone from django.utils import timezone
from django.utils.translation import gettext_lazy as _ from django.utils.translation import gettext_lazy as _
from multiprocessing import Event, Process, Value, current_process
# Local # Local
import django_q.tasks import django_q.tasks
from django_q.brokers import get_broker from django_q.brokers import get_broker
from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter
from django_q.humanhash import humanize
from django_q.models import Task, Success, Schedule from django_q.models import Task, Success, Schedule
from django_q.queues import Queue from django_q.queues import Queue
from django_q.signals import pre_execute from django_q.signals import pre_execute
@@ -49,10 +52,18 @@ class Cluster(object):
# Start Sentinel # Start Sentinel
self.stop_event = Event() self.stop_event = Event()
self.start_event = Event() self.start_event = Event()
self.sentinel = Process(target=Sentinel, self.sentinel = Process(
args=(self.stop_event, self.start_event, self.cluster_id, self.broker, self.timeout)) target=Sentinel,
args=(
self.stop_event,
self.start_event,
self.cluster_id,
self.broker,
self.timeout,
),
)
self.sentinel.start() self.sentinel.start()
logger.info(_('Q Cluster-{} starting.').format(self.cluster_id)) logger.info(_(f"Q Cluster {self.name} starting."))
while not self.start_event.is_set(): while not self.start_event.is_set():
sleep(0.1) sleep(0.1)
return self.pid return self.pid
@@ -60,17 +71,20 @@ class Cluster(object):
def stop(self): def stop(self):
if not self.sentinel.is_alive(): if not self.sentinel.is_alive():
return False return False
logger.info(_('Q Cluster-{} stopping.').format(self.cluster_id)) logger.info(_(f"Q Cluster {self.name} stopping."))
self.stop_event.set() self.stop_event.set()
self.sentinel.join() self.sentinel.join()
logger.info(_('Q Cluster-{} has stopped.').format(self.cluster_id)) logger.info(_(f"Q Cluster {self.name} has stopped."))
self.start_event = None self.start_event = None
self.stop_event = None self.stop_event = None
return True return True
def sig_handler(self, signum, frame): def sig_handler(self, signum, frame):
logger.debug(_('{} got signal {}').format(current_process().name, logger.debug(
Conf.SIGNAL_NAMES.get(signum, 'UNKNOWN'))) _(
f'{current_process().name} got signal {Conf.SIGNAL_NAMES.get(signum, "UNKNOWN")}'
)
)
self.stop() self.stop()
@property @property
@@ -79,6 +93,10 @@ class Cluster(object):
return Stat.get(pid=self.pid, cluster_id=self.cluster_id) return Stat.get(pid=self.pid, cluster_id=self.cluster_id)
return Status(pid=self.pid, cluster_id=self.cluster_id) return Status(pid=self.pid, cluster_id=self.cluster_id)
@property
def name(self):
return humanize(self.cluster_id.hex)
@property @property
def is_starting(self): def is_starting(self):
return self.stop_event and self.start_event and not self.start_event.is_set() return self.stop_event and self.start_event and not self.start_event.is_set()
@@ -89,7 +107,12 @@ class Cluster(object):
@property @property
def is_stopping(self): def is_stopping(self):
return self.stop_event and self.start_event and self.start_event.is_set() and self.stop_event.is_set() return (
self.stop_event
and self.start_event
and self.start_event.is_set()
and self.stop_event.is_set()
)
@property @property
def has_stopped(self): def has_stopped(self):
@@ -97,7 +120,15 @@ class Cluster(object):
class Sentinel(object): class Sentinel(object):
def __init__(self, stop_event, start_event, cluster_id, broker=None, timeout=Conf.TIMEOUT, start=True): def __init__(
self,
stop_event,
start_event,
cluster_id,
broker=None,
timeout=Conf.TIMEOUT,
start=True,
):
# Make sure we catch signals for the pool # Make sure we catch signals for the pool
signal.signal(signal.SIGINT, signal.SIG_IGN) signal.signal(signal.SIGINT, signal.SIG_IGN)
signal.signal(signal.SIGTERM, signal.SIG_DFL) signal.signal(signal.SIGTERM, signal.SIG_DFL)
@@ -113,7 +144,9 @@ class Sentinel(object):
self.pool_size = Conf.WORKERS self.pool_size = Conf.WORKERS
self.pool = [] self.pool = []
self.timeout = timeout self.timeout = timeout
self.task_queue = Queue(maxsize=Conf.QUEUE_LIMIT) if Conf.QUEUE_LIMIT else Queue() self.task_queue = (
Queue(maxsize=Conf.QUEUE_LIMIT) if Conf.QUEUE_LIMIT else Queue()
)
self.result_queue = Queue() self.result_queue = Queue()
self.event_out = Event() self.event_out = Event()
self.monitor = None self.monitor = None
@@ -155,7 +188,9 @@ class Sentinel(object):
return self.spawn_process(pusher, self.task_queue, self.event_out, self.broker) return self.spawn_process(pusher, self.task_queue, self.event_out, self.broker)
def spawn_worker(self): def spawn_worker(self):
self.spawn_process(worker, self.task_queue, self.result_queue, Value('f', -1), self.timeout) self.spawn_process(
worker, self.task_queue, self.result_queue, Value("f", -1), self.timeout
)
def spawn_monitor(self): def spawn_monitor(self):
return self.spawn_process(monitor, self.result_queue, self.broker) return self.spawn_process(monitor, self.result_queue, self.broker)
@@ -168,21 +203,21 @@ class Sentinel(object):
close_old_django_connections() close_old_django_connections()
if process == self.monitor: if process == self.monitor:
self.monitor = self.spawn_monitor() self.monitor = self.spawn_monitor()
logger.error(_("reincarnated monitor {} after sudden death").format(process.name)) logger.error(_(f"reincarnated monitor {process.name} after sudden death"))
elif process == self.pusher: elif process == self.pusher:
self.pusher = self.spawn_pusher() self.pusher = self.spawn_pusher()
logger.error(_("reincarnated pusher {} after sudden death").format(process.name)) logger.error(_(f"reincarnated pusher {process.name} after sudden death"))
else: else:
self.pool.remove(process) self.pool.remove(process)
self.spawn_worker() self.spawn_worker()
if process.timer.value == 0: if process.timer.value == 0:
# only need to terminate on timeout, otherwise we risk destabilizing the queues # only need to terminate on timeout, otherwise we risk destabilizing the queues
process.terminate() process.terminate()
logger.warn(_("reincarnated worker {} after timeout").format(process.name)) logger.warn(_(f"reincarnated worker {process.name} after timeout"))
elif int(process.timer.value) == -2: elif int(process.timer.value) == -2:
logger.info(_("recycled worker {}").format(process.name)) logger.info(_(f"recycled worker {process.name}"))
else: else:
logger.error(_("reincarnated worker {} after death").format(process.name)) logger.error(_(f"reincarnated worker {process.name} after death"))
self.reincarnations += 1 self.reincarnations += 1
@@ -201,10 +236,14 @@ class Sentinel(object):
set_cpu_affinity(Conf.CPU_AFFINITY, [w.pid for w in self.pool]) set_cpu_affinity(Conf.CPU_AFFINITY, [w.pid for w in self.pool])
def guard(self): def guard(self):
logger.info(_('{} guarding cluster at {}').format(current_process().name, self.pid)) logger.info(
_(
f"{current_process().name} guarding cluster {humanize(self.cluster_id.hex)}"
)
)
self.start_event.set() self.start_event.set()
Stat(self).save() Stat(self).save()
logger.info(_('Q Cluster-{} running.').format(self.parent_pid)) logger.info(_(f"Q Cluster {humanize(self.cluster_id.hex)} running."))
counter = 0 counter = 0
cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds
# Guard loop. Runs at least once # Guard loop. Runs at least once
@@ -238,7 +277,7 @@ class Sentinel(object):
def stop(self): def stop(self):
Stat(self).save() Stat(self).save()
name = current_process().name name = current_process().name
logger.info(_('{} stopping cluster processes').format(name)) logger.info(_(f"{name} stopping cluster processes"))
# Stopping pusher # Stopping pusher
self.event_out.set() self.event_out.set()
# Wait for it to stop # Wait for it to stop
@@ -247,7 +286,7 @@ class Sentinel(object):
Stat(self).save() Stat(self).save()
# Put poison pills in the queue # Put poison pills in the queue
for __ in range(len(self.pool)): for __ in range(len(self.pool)):
self.task_queue.put('STOP') self.task_queue.put("STOP")
self.task_queue.close() self.task_queue.close()
# wait for the task queue to empty # wait for the task queue to empty
self.task_queue.join_thread() self.task_queue.join_thread()
@@ -259,11 +298,11 @@ class Sentinel(object):
sleep(0.1) sleep(0.1)
Stat(self).save() Stat(self).save()
# Finally stop the monitor # Finally stop the monitor
self.result_queue.put('STOP') self.result_queue.put("STOP")
self.result_queue.close() self.result_queue.close()
# Wait for the result queue to empty # Wait for the result queue to empty
self.result_queue.join_thread() self.result_queue.join_thread()
logger.info(_('{} waiting for the monitor.').format(name)) logger.info(_(f"{name} waiting for the monitor."))
# Wait for everything to close or time out # Wait for everything to close or time out
count = 0 count = 0
if not self.timeout: if not self.timeout:
@@ -284,7 +323,7 @@ def pusher(task_queue, event, broker=None):
""" """
if not broker: if not broker:
broker = get_broker() broker = get_broker()
logger.info(_('{} pushing tasks at {}').format(current_process().name, current_process().pid)) logger.info(_(f"{current_process().name} pushing tasks at {current_process().pid}"))
while True: while True:
try: try:
task_set = broker.dequeue() task_set = broker.dequeue()
@@ -303,12 +342,12 @@ def pusher(task_queue, event, broker=None):
logger.error(e, traceback.format_exc()) logger.error(e, traceback.format_exc())
broker.fail(ack_id) broker.fail(ack_id)
continue continue
task['ack_id'] = ack_id task["ack_id"] = ack_id
task_queue.put(task) task_queue.put(task)
logger.debug(_('queueing from {}').format(broker.list_key)) logger.debug(_(f"queueing from {broker.list_key}"))
if event.is_set(): if event.is_set():
break break
logger.info(_("{} stopped pushing tasks").format(current_process().name)) logger.info(_(f"{current_process().name} stopped pushing tasks"))
def monitor(result_queue, broker=None): def monitor(result_queue, broker=None):
@@ -319,25 +358,25 @@ def monitor(result_queue, broker=None):
if not broker: if not broker:
broker = get_broker() broker = get_broker()
name = current_process().name name = current_process().name
logger.info(_("{} monitoring at {}").format(name, current_process().pid)) logger.info(_(f"{name} monitoring at {current_process().pid}"))
for task in iter(result_queue.get, 'STOP'): for task in iter(result_queue.get, "STOP"):
# save the result # save the result
if task.get('cached', False): if task.get("cached", False):
save_cached(task, broker) save_cached(task, broker)
else: else:
save_task(task, broker) save_task(task, broker)
# acknowledge result # acknowledge result
ack_id = task.pop('ack_id', False) ack_id = task.pop("ack_id", False)
if ack_id and (task['success'] or task.get('ack_failure', False)): if ack_id and (task["success"] or task.get("ack_failure", False)):
broker.acknowledge(ack_id) broker.acknowledge(ack_id)
# log the result # log the result
if task['success']: if task["success"]:
# log success # log success
logger.info(_("Processed [{}]").format(task['name'])) logger.info(_(f"Processed [{task['name']}]"))
else: else:
# log failure # log failure
logger.error(_("Failed [{}] - {}").format(task['name'], task['result'])) logger.error(_(f"Failed [{task['name']}] - {task['result']}"))
logger.info(_("{} stopped monitoring results").format(name)) logger.info(_(f"{name} stopped monitoring results"))
def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT): def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
@@ -348,22 +387,22 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
:type timer: multiprocessing.Value :type timer: multiprocessing.Value
""" """
name = current_process().name name = current_process().name
logger.info(_('{} ready for work at {}').format(name, current_process().pid)) logger.info(_(f"{name} ready for work at {current_process().pid}"))
task_count = 0 task_count = 0
if timeout is None: if timeout is None:
timeout = -1 timeout = -1
# Start reading the task queue # Start reading the task queue
for task in iter(task_queue.get, 'STOP'): for task in iter(task_queue.get, "STOP"):
result = None result = None
timer.value = -1 # Idle timer.value = -1 # Idle
task_count += 1 task_count += 1
# Get the function from the task # Get the function from the task
logger.info(_('{} processing [{}]').format(name, task['name'])) logger.info(_(f'{name} processing [{task["name"]}]'))
f = task['func'] f = task["func"]
# if it's not an instance try to get it from the string # if it's not an instance try to get it from the string
if not callable(task['func']): if not callable(task["func"]):
try: try:
module, func = f.rsplit('.', 1) module, func = f.rsplit(".", 1)
m = importlib.import_module(module) m = importlib.import_module(module)
f = getattr(m, func) f = getattr(m, func)
except (ValueError, ImportError, AttributeError) as e: except (ValueError, ImportError, AttributeError) as e:
@@ -373,30 +412,30 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
# We're still going # We're still going
if not result: if not result:
close_old_django_connections() close_old_django_connections()
timer_value = task.pop('timeout', timeout) timer_value = task.pop("timeout", timeout)
# signal execution # signal execution
pre_execute.send(sender="django_q", func=f, task=task) pre_execute.send(sender="django_q", func=f, task=task)
# execute the payload # execute the payload
timer.value = timer_value # Busy timer.value = timer_value # Busy
try: try:
res = f(*task['args'], **task['kwargs']) res = f(*task["args"], **task["kwargs"])
result = (res, True) result = (res, True)
except Exception as e: except Exception as e:
result = ('{} : {}'.format(e, traceback.format_exc()), False) result = (f"{e} : {traceback.format_exc()}", False)
if error_reporter: if error_reporter:
error_reporter.report() error_reporter.report()
with timer.get_lock(): with timer.get_lock():
# Process result # Process result
task['result'] = result[0] task["result"] = result[0]
task['success'] = result[1] task["success"] = result[1]
task['stopped'] = timezone.now() task["stopped"] = timezone.now()
result_queue.put(task) result_queue.put(task)
timer.value = -1 # Idle timer.value = -1 # Idle
# Recycle # Recycle
if task_count == Conf.RECYCLE: if task_count == Conf.RECYCLE:
timer.value = -2 # Recycled timer.value = -2 # Recycled
break break
logger.info(_('{} stopped doing work').format(name)) logger.info(_(f"{name} stopped doing work"))
def save_task(task, broker): def save_task(task, broker):
@@ -404,67 +443,77 @@ def save_task(task, broker):
Saves the task package to Django or the cache Saves the task package to Django or the cache
""" """
# SAVE LIMIT < 0 : Don't save success # SAVE LIMIT < 0 : Don't save success
if not task.get('save', Conf.SAVE_LIMIT >= 0) and task['success']: if not task.get("save", Conf.SAVE_LIMIT >= 0) and task["success"]:
return return
# enqueues next in a chain # enqueues next in a chain
if task.get('chain', None): if task.get("chain", None):
django_q.tasks.async_chain(task['chain'], group=task['group'], cached=task['cached'], sync=task['sync'], broker=broker) django_q.tasks.async_chain(
task["chain"],
group=task["group"],
cached=task["cached"],
sync=task["sync"],
broker=broker,
)
# SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning # SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning
close_old_django_connections() close_old_django_connections()
try: try:
if task['success'] and 0 < Conf.SAVE_LIMIT <= Success.objects.count(): if task["success"] and 0 < Conf.SAVE_LIMIT <= Success.objects.count():
Success.objects.last().delete() Success.objects.last().delete()
# check if this task has previous results # check if this task has previous results
if Task.objects.filter(id=task['id'], name=task['name']).exists(): if Task.objects.filter(id=task["id"], name=task["name"]).exists():
existing_task = Task.objects.get(id=task['id'], name=task['name']) existing_task = Task.objects.get(id=task["id"], name=task["name"])
# only update the result if it hasn't succeeded yet # only update the result if it hasn't succeeded yet
if not existing_task.success: if not existing_task.success:
existing_task.stopped = task['stopped'] existing_task.stopped = task["stopped"]
existing_task.result = task['result'] existing_task.result = task["result"]
existing_task.success = task['success'] existing_task.success = task["success"]
existing_task.save() existing_task.save()
else: else:
Task.objects.create(id=task['id'], Task.objects.create(
name=task['name'], id=task["id"],
func=task['func'], name=task["name"],
hook=task.get('hook'), func=task["func"],
args=task['args'], hook=task.get("hook"),
kwargs=task['kwargs'], args=task["args"],
started=task['started'], kwargs=task["kwargs"],
stopped=task['stopped'], started=task["started"],
result=task['result'], stopped=task["stopped"],
group=task.get('group'), result=task["result"],
success=task['success'] group=task.get("group"),
) success=task["success"],
)
except Exception as e: except Exception as e:
logger.error(e) logger.error(e)
def save_cached(task, broker): def save_cached(task, broker):
task_key = '{}:{}'.format(broker.list_key, task['id']) task_key = f'{broker.list_key}:{task["id"]}'
timeout = task['cached'] timeout = task["cached"]
if timeout is True: if timeout is True:
timeout = None timeout = None
try: try:
group = task.get('group', None) group = task.get("group", None)
iter_count = task.get('iter_count', 0) iter_count = task.get("iter_count", 0)
# if it's a group append to the group list # if it's a group append to the group list
if group: if group:
group_key = '{}:{}:keys'.format(broker.list_key, group) group_key = f"{broker.list_key}:{group}:keys"
group_list = broker.cache.get(group_key) or [] group_list = broker.cache.get(group_key) or []
# if it's an iter group, check if we are ready # if it's an iter group, check if we are ready
if iter_count and len(group_list) == iter_count - 1: if iter_count and len(group_list) == iter_count - 1:
group_args = '{}:{}:args'.format(broker.list_key, group) group_args = f"{broker.list_key}:{group}:args"
# collate the results into a Task result # collate the results into a Task result
results = [SignedPackage.loads(broker.cache.get(k))['result'] for k in group_list] results = [
results.append(task['result']) SignedPackage.loads(broker.cache.get(k))["result"]
task['result'] = results for k in group_list
task['id'] = group ]
task['args'] = SignedPackage.loads(broker.cache.get(group_args)) results.append(task["result"])
task.pop('iter_count', None) task["result"] = results
task.pop('group', None) task["id"] = group
if task.get('iter_cached', None): task["args"] = SignedPackage.loads(broker.cache.get(group_args))
task['cached'] = task.pop('iter_cached', None) task.pop("iter_count", None)
task.pop("group", None)
if task.get("iter_cached", None):
task["cached"] = task.pop("iter_cached", None)
save_cached(task, broker=broker) save_cached(task, broker=broker)
else: else:
save_task(task, broker) save_task(task, broker)
@@ -475,12 +524,16 @@ def save_cached(task, broker):
group_list.append(task_key) group_list.append(task_key)
broker.cache.set(group_key, group_list, timeout) broker.cache.set(group_key, group_list, timeout)
# async_task next in a chain # async_task next in a chain
if task.get('chain', None): if task.get("chain", None):
django_q.tasks.async_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker) django_q.tasks.async_chain(
task["chain"],
group=group,
cached=task["cached"],
sync=task["sync"],
broker=broker,
)
# save the task # save the task
broker.cache.set(task_key, broker.cache.set(task_key, SignedPackage.dumps(task), timeout)
SignedPackage.dumps(task),
timeout)
except Exception as e: except Exception as e:
logger.error(e) logger.error(e)
@@ -494,14 +547,18 @@ def scheduler(broker=None):
close_old_django_connections() close_old_django_connections()
try: try:
with db.transaction.atomic(): with db.transaction.atomic():
for s in Schedule.objects.select_for_update().exclude(repeats=0).filter(next_run__lt=timezone.now()): for s in (
Schedule.objects.select_for_update()
.exclude(repeats=0)
.filter(next_run__lt=timezone.now())
):
args = () args = ()
kwargs = {} kwargs = {}
# get args, kwargs and hook # get args, kwargs and hook
if s.kwargs: if s.kwargs:
try: try:
# eval should be safe here because dict() # eval should be safe here because dict()
kwargs = eval('dict({})'.format(s.kwargs)) kwargs = eval(f"dict({s.kwargs})")
except SyntaxError: except SyntaxError:
kwargs = {} kwargs = {}
if s.args: if s.args:
@@ -509,9 +566,9 @@ def scheduler(broker=None):
# single value won't eval to tuple, so: # single value won't eval to tuple, so:
if type(args) != tuple: if type(args) != tuple:
args = (args,) args = (args,)
q_options = kwargs.get('q_options', {}) q_options = kwargs.get("q_options", {})
if s.hook: if s.hook:
q_options['hook'] = s.hook q_options["hook"] = s.hook
# set up the next run time # set up the next run time
if not s.schedule_type == s.ONCE: if not s.schedule_type == s.ONCE:
next_run = arrow.get(s.next_run) next_run = arrow.get(s.next_run)
@@ -535,18 +592,23 @@ def scheduler(broker=None):
s.next_run = next_run.datetime s.next_run = next_run.datetime
s.repeats += -1 s.repeats += -1
# send it to the cluster # send it to the cluster
q_options['broker'] = broker q_options["broker"] = broker
q_options['group'] = q_options.get('group', s.name or s.id) q_options["group"] = q_options.get("group", s.name or s.id)
kwargs['q_options'] = q_options kwargs["q_options"] = q_options
s.task = django_q.tasks.async_task(s.func, *args, **kwargs) s.task = django_q.tasks.async_task(s.func, *args, **kwargs)
# log it # log it
if not s.task: if not s.task:
logger.error( logger.error(
_('{} failed to create a task from schedule [{}]').format(current_process().name, _(
s.name or s.id)) f"{current_process().name} failed to create a task from schedule [{s.name or s.id}]"
)
)
else: else:
logger.info( logger.info(
_('{} created a task from schedule [{}]').format(current_process().name, s.name or s.id)) _(
f"{current_process().name} created a task from schedule [{s.name or s.id}]"
)
)
# default behavior is to delete a ONCE schedule # default behavior is to delete a ONCE schedule
if s.schedule_type == s.ONCE: if s.schedule_type == s.ONCE:
if s.repeats < 0: if s.repeats < 0:
@@ -561,14 +623,15 @@ def scheduler(broker=None):
def close_old_django_connections(): def close_old_django_connections():
''' """
Close django connections unless running with sync=True. Close django connections unless running with sync=True.
''' """
if Conf.SYNC: if Conf.SYNC:
logger.warning( logger.warning(
'Preserving django database connections because sync=True. Beware ' "Preserving django database connections because sync=True. Beware "
'that tasks are now injected in the calling context/transactions ' "that tasks are now injected in the calling context/transactions "
'which may result in unexpected bahaviour.') "which may result in unexpected bahaviour."
)
else: else:
db.close_old_connections() db.close_old_connections()
@@ -583,11 +646,13 @@ def set_cpu_affinity(n, process_ids, actual=not Conf.TESTING):
""" """
# check if we have the psutil module # check if we have the psutil module
if not psutil: if not psutil:
logger.warning('Skipping cpu affinity because psutil was not found.') logger.warning("Skipping cpu affinity because psutil was not found.")
return return
# check if the platform supports cpu_affinity # check if the platform supports cpu_affinity
if actual and not hasattr(psutil.Process(process_ids[0]), 'cpu_affinity'): if actual and not hasattr(psutil.Process(process_ids[0]), "cpu_affinity"):
logger.warning('Faking cpu affinity because it is not supported on this platform') logger.warning(
"Faking cpu affinity because it is not supported on this platform"
)
actual = False actual = False
# get the available processors # get the available processors
cpu_list = list(range(psutil.cpu_count())) cpu_list = list(range(psutil.cpu_count()))
@@ -607,4 +672,4 @@ def set_cpu_affinity(n, process_ids, actual=not Conf.TESTING):
p = psutil.Process(pid) p = psutil.Process(pid)
if actual: if actual:
p.cpu_affinity(affinity) p.cpu_affinity(affinity)
logger.info(_('{} will use cpu {}').format(pid, affinity)) logger.info(_(f"{pid} will use cpu {affinity}"))
+63 -53
View File
@@ -1,15 +1,16 @@
import logging import logging
from copy import deepcopy
from signal import signal
from multiprocessing import cpu_count
# django
from django.utils.translation import gettext_lazy as _
from django.conf import settings
# external # external
import os import os
from copy import deepcopy
from multiprocessing import cpu_count
from signal import signal
import pkg_resources import pkg_resources
from django.conf import settings
# django
from django.utils.translation import gettext_lazy as _
# local # local
from django_q.queues import Queue from django_q.queues import Queue
@@ -25,64 +26,65 @@ class Conf(object):
""" """
Configuration class Configuration class
""" """
try: try:
conf = settings.Q_CLUSTER conf = settings.Q_CLUSTER
except AttributeError: except AttributeError:
conf = {} conf = {}
# Redis server configuration . Follows standard redis keywords # Redis server configuration . Follows standard redis keywords
REDIS = conf.get('redis', {}) REDIS = conf.get("redis", {})
# Support for Django-Redis connections # Support for Django-Redis connections
DJANGO_REDIS = conf.get('django_redis', None) DJANGO_REDIS = conf.get("django_redis", None)
# Disque broker # Disque broker
DISQUE_NODES = conf.get('disque_nodes', None) DISQUE_NODES = conf.get("disque_nodes", None)
# Optional Authentication # Optional Authentication
DISQUE_AUTH = conf.get('disque_auth', None) DISQUE_AUTH = conf.get("disque_auth", None)
# Optional Fast acknowledge # Optional Fast acknowledge
DISQUE_FASTACK = conf.get('disque_fastack', False) DISQUE_FASTACK = conf.get("disque_fastack", False)
# IronMQ broker # IronMQ broker
IRON_MQ = conf.get('iron_mq', None) IRON_MQ = conf.get("iron_mq", None)
# SQS broker # SQS broker
SQS = conf.get('sqs', None) SQS = conf.get("sqs", None)
# ORM broker # ORM broker
ORM = conf.get('orm', None) ORM = conf.get("orm", None)
# Custom broker class # Custom broker class
BROKER_CLASS = conf.get('broker_class', None) BROKER_CLASS = conf.get("broker_class", None)
# Database Poll # Database Poll
POLL = conf.get('poll', 0.2) POLL = conf.get("poll", 0.2)
# MongoDB broker # MongoDB broker
MONGO = conf.get('mongo', None) MONGO = conf.get("mongo", None)
MONGO_DB = conf.get('mongo_db', None) MONGO_DB = conf.get("mongo_db", None)
# Name of the cluster or site. For when you run multiple sites on one redis server # Name of the cluster or site. For when you run multiple sites on one redis server
PREFIX = conf.get('name', 'default') PREFIX = conf.get("name", "default")
# Log output level # Log output level
LOG_LEVEL = conf.get('log_level', 'INFO') LOG_LEVEL = conf.get("log_level", "INFO")
# Maximum number of successful tasks kept in the database. 0 saves everything. -1 saves none # Maximum number of successful tasks kept in the database. 0 saves everything. -1 saves none
# Failures are always saved # Failures are always saved
SAVE_LIMIT = conf.get('save_limit', 250) SAVE_LIMIT = conf.get("save_limit", 250)
# Guard loop sleep in seconds. Should be between 0 and 60 seconds. # Guard loop sleep in seconds. Should be between 0 and 60 seconds.
GUARD_CYCLE = conf.get('guard_cycle', 0.5) GUARD_CYCLE = conf.get("guard_cycle", 0.5)
# Disable the scheduler # Disable the scheduler
SCHEDULER = conf.get('scheduler', True) SCHEDULER = conf.get("scheduler", True)
# Number of workers in the pool. Default is cpu count if implemented, otherwise 4. # Number of workers in the pool. Default is cpu count if implemented, otherwise 4.
WORKERS = conf.get('workers', False) WORKERS = conf.get("workers", False)
if not WORKERS: if not WORKERS:
try: try:
WORKERS = cpu_count() WORKERS = cpu_count()
@@ -96,62 +98,62 @@ class Conf(object):
WORKERS = 4 WORKERS = 4
# Option to undaemonize the workers and allow them to spawn child processes # Option to undaemonize the workers and allow them to spawn child processes
DAEMONIZE_WORKERS = conf.get('daemonize_workers', True) DAEMONIZE_WORKERS = conf.get("daemonize_workers", True)
# Maximum number of tasks that each cluster can work on # Maximum number of tasks that each cluster can work on
QUEUE_LIMIT = conf.get('queue_limit', int(WORKERS) ** 2) QUEUE_LIMIT = conf.get("queue_limit", int(WORKERS) ** 2)
# Sets compression of redis packages # Sets compression of redis packages
COMPRESSED = conf.get('compress', False) COMPRESSED = conf.get("compress", False)
# Number of tasks each worker can handle before it gets recycled. Useful for releasing memory # Number of tasks each worker can handle before it gets recycled. Useful for releasing memory
RECYCLE = conf.get('recycle', 500) RECYCLE = conf.get("recycle", 500)
# Number of seconds to wait for a worker to finish. # Number of seconds to wait for a worker to finish.
TIMEOUT = conf.get('timeout', None) TIMEOUT = conf.get("timeout", None)
# Whether to acknowledge unsuccessful tasks. # Whether to acknowledge unsuccessful tasks.
# This causes failed tasks to be considered delivered, thereby removing them from # This causes failed tasks to be considered delivered, thereby removing them from
# the task queue. Defaults to False. # the task queue. Defaults to False.
ACK_FAILURES = conf.get('ack_failures', False) ACK_FAILURES = conf.get("ack_failures", False)
# Number of seconds to wait for acknowledgement before retrying a task # Number of seconds to wait for acknowledgement before retrying a task
# Only works with brokers that guarantee delivery. Defaults to 60 seconds. # Only works with brokers that guarantee delivery. Defaults to 60 seconds.
RETRY = conf.get('retry', 60) RETRY = conf.get("retry", 60)
# Sets the amount of tasks the cluster will try to pop off the broker. # Sets the amount of tasks the cluster will try to pop off the broker.
# If it supports bulk gets. # If it supports bulk gets.
BULK = conf.get('bulk', 1) BULK = conf.get("bulk", 1)
# The Django Admin label for this app # The Django Admin label for this app
LABEL = conf.get('label', 'Django Q') LABEL = conf.get("label", "Django Q")
# Sets the number of processors for each worker, defaults to all. # Sets the number of processors for each worker, defaults to all.
CPU_AFFINITY = conf.get('cpu_affinity', 0) CPU_AFFINITY = conf.get("cpu_affinity", 0)
# Global sync option to for debugging # Global sync option to for debugging
SYNC = conf.get('sync', False) SYNC = conf.get("sync", False)
# The Django cache to use # The Django cache to use
CACHE = conf.get('cache', 'default') CACHE = conf.get("cache", "default")
# Use the cache as result backend. Can be 'True' or an integer representing the global cache timeout. # Use the cache as result backend. Can be 'True' or an integer representing the global cache timeout.
# i.e 'cached: 60' , will make all results go the cache and expire in 60 seconds. # i.e 'cached: 60' , will make all results go the cache and expire in 60 seconds.
CACHED = conf.get('cached', False) CACHED = conf.get("cached", False)
# If set to False the scheduler won't execute tasks in the past. # If set to False the scheduler won't execute tasks in the past.
# Instead it will run once and reschedule the next run in the future. Defaults to True. # Instead it will run once and reschedule the next run in the future. Defaults to True.
CATCH_UP = conf.get('catch_up', True) CATCH_UP = conf.get("catch_up", True)
# Use the secret key for package signing # Use the secret key for package signing
# Django itself should raise an error if it's not configured # Django itself should raise an error if it's not configured
SECRET_KEY = settings.SECRET_KEY SECRET_KEY = settings.SECRET_KEY
# The redis stats key # The redis stats key
Q_STAT = 'django_q:{}:cluster'.format(PREFIX) Q_STAT = f"django_q:{PREFIX}:cluster"
# Optional error reporting setup # Optional error reporting setup
ERROR_REPORTER = conf.get('error_reporter', {}) ERROR_REPORTER = conf.get("error_reporter", {})
# OSX doesn't implement qsize because of missing sem_getvalue() # OSX doesn't implement qsize because of missing sem_getvalue()
try: try:
@@ -160,28 +162,33 @@ class Conf(object):
QSIZE = False QSIZE = False
# Getting the signal names # Getting the signal names
SIGNAL_NAMES = dict((getattr(signal, n), n) for n in dir(signal) if n.startswith('SIG') and '_' not in n) SIGNAL_NAMES = dict(
(getattr(signal, n), n)
for n in dir(signal)
if n.startswith("SIG") and "_" not in n
)
# Translators: Cluster status descriptions # Translators: Cluster status descriptions
STARTING = _('Starting') STARTING = _("Starting")
WORKING = _('Working') WORKING = _("Working")
IDLE = _("Idle") IDLE = _("Idle")
STOPPED = _('Stopped') STOPPED = _("Stopped")
STOPPING = _('Stopping') STOPPING = _("Stopping")
# to manage workarounds during testing # to manage workarounds during testing
TESTING = conf.get('testing', False) TESTING = conf.get("testing", False)
# logger # logger
logger = logging.getLogger('django-q') logger = logging.getLogger("django-q")
# Set up standard logging handler in case there is none # Set up standard logging handler in case there is none
if not logger.handlers: if not logger.handlers:
logger.setLevel(level=getattr(logging, Conf.LOG_LEVEL)) logger.setLevel(level=getattr(logging, Conf.LOG_LEVEL))
logger.propagate = False logger.propagate = False
formatter = logging.Formatter(fmt='%(asctime)s [Q] %(levelname)s %(message)s', formatter = logging.Formatter(
datefmt='%H:%M:%S') fmt="%(asctime)s [Q] %(levelname)s %(message)s", datefmt="%H:%M:%S"
)
handler = logging.StreamHandler() handler = logging.StreamHandler()
handler.setFormatter(formatter) handler.setFormatter(formatter)
logger.addHandler(handler) logger.addHandler(handler)
@@ -209,7 +216,8 @@ if Conf.ERROR_REPORTER:
# and instantiate an ErrorReporter using the provided config # and instantiate an ErrorReporter using the provided config
for name, conf in error_conf.items(): for name, conf in error_conf.items():
for entry in pkg_resources.iter_entry_points( for entry in pkg_resources.iter_entry_points(
'djangoq.errorreporters', name): "djangoq.errorreporters", name
):
Reporter = entry.load() Reporter = entry.load()
reporters.append(Reporter(**conf)) reporters.append(Reporter(**conf))
error_reporter = ErrorReporter(reporters) error_reporter = ErrorReporter(reporters)
@@ -221,9 +229,11 @@ else:
# get parent pid compatibility # get parent pid compatibility
def get_ppid(): def get_ppid():
if hasattr(os, 'getppid'): if hasattr(os, "getppid"):
return os.getppid() return os.getppid()
elif psutil: elif psutil:
return psutil.Process(os.getpid()).ppid() return psutil.Process(os.getpid()).ppid()
else: else:
raise OSError('Your OS does not support `os.getppid`. Please install `psutil` as an alternative provider.') raise OSError(
"Your OS does not support `os.getppid`. Please install `psutil` as an alternative provider."
)
+73 -55
View File
@@ -35,9 +35,15 @@ class Task(models.Model):
@staticmethod @staticmethod
def get_result_group(group_id, failures=False): def get_result_group(group_id, failures=False):
if failures: if failures:
values = Task.objects.filter(group=group_id).values_list('result', flat=True) values = Task.objects.filter(group=group_id).values_list(
"result", flat=True
)
else: else:
values = Task.objects.filter(group=group_id).exclude(success=False).values_list('result', flat=True) values = (
Task.objects.filter(group=group_id)
.exclude(success=False)
.values_list("result", flat=True)
)
return decode_results(values) return decode_results(values)
def group_result(self, failures=False): def group_result(self, failures=False):
@@ -86,76 +92,88 @@ class Task(models.Model):
return truncatechars(self.result, 100) return truncatechars(self.result, 100)
def __unicode__(self): def __unicode__(self):
return u'{}'.format(self.name or self.id) return f"{self.name or self.id}"
class Meta: class Meta:
app_label = 'django_q' app_label = "django_q"
ordering = ['-stopped'] ordering = ["-stopped"]
class SuccessManager(models.Manager): class SuccessManager(models.Manager):
def get_queryset(self): def get_queryset(self):
return super(SuccessManager, self).get_queryset().filter( return super(SuccessManager, self).get_queryset().filter(success=True)
success=True)
class Success(Task): class Success(Task):
objects = SuccessManager() objects = SuccessManager()
class Meta: class Meta:
app_label = 'django_q' app_label = "django_q"
verbose_name = _('Successful task') verbose_name = _("Successful task")
verbose_name_plural = _('Successful tasks') verbose_name_plural = _("Successful tasks")
ordering = ['-stopped'] ordering = ["-stopped"]
proxy = True proxy = True
class FailureManager(models.Manager): class FailureManager(models.Manager):
def get_queryset(self): def get_queryset(self):
return super(FailureManager, self).get_queryset().filter( return super(FailureManager, self).get_queryset().filter(success=False)
success=False)
class Failure(Task): class Failure(Task):
objects = FailureManager() objects = FailureManager()
class Meta: class Meta:
app_label = 'django_q' app_label = "django_q"
verbose_name = _('Failed task') verbose_name = _("Failed task")
verbose_name_plural = _('Failed tasks') verbose_name_plural = _("Failed tasks")
ordering = ['-stopped'] ordering = ["-stopped"]
proxy = True proxy = True
class Schedule(models.Model): class Schedule(models.Model):
name = models.CharField(max_length=100, null=True, blank=True) name = models.CharField(max_length=100, null=True, blank=True)
func = models.CharField(max_length=256, help_text='e.g. module.tasks.function') func = models.CharField(max_length=256, help_text="e.g. module.tasks.function")
hook = models.CharField(max_length=256, null=True, blank=True, help_text='e.g. module.tasks.result_function') hook = models.CharField(
max_length=256,
null=True,
blank=True,
help_text="e.g. module.tasks.result_function",
)
args = models.TextField(null=True, blank=True, help_text=_("e.g. 1, 2, 'John'")) args = models.TextField(null=True, blank=True, help_text=_("e.g. 1, 2, 'John'"))
kwargs = models.TextField(null=True, blank=True, help_text=_("e.g. x=1, y=2, name='John'")) kwargs = models.TextField(
ONCE = 'O' null=True, blank=True, help_text=_("e.g. x=1, y=2, name='John'")
MINUTES = 'I' )
HOURLY = 'H' ONCE = "O"
DAILY = 'D' MINUTES = "I"
WEEKLY = 'W' HOURLY = "H"
MONTHLY = 'M' DAILY = "D"
QUARTERLY = 'Q' WEEKLY = "W"
YEARLY = 'Y' MONTHLY = "M"
QUARTERLY = "Q"
YEARLY = "Y"
TYPE = ( TYPE = (
(ONCE, _('Once')), (ONCE, _("Once")),
(MINUTES, _('Minutes')), (MINUTES, _("Minutes")),
(HOURLY, _('Hourly')), (HOURLY, _("Hourly")),
(DAILY, _('Daily')), (DAILY, _("Daily")),
(WEEKLY, _('Weekly')), (WEEKLY, _("Weekly")),
(MONTHLY, _('Monthly')), (MONTHLY, _("Monthly")),
(QUARTERLY, _('Quarterly')), (QUARTERLY, _("Quarterly")),
(YEARLY, _('Yearly')), (YEARLY, _("Yearly")),
)
schedule_type = models.CharField(
max_length=1, choices=TYPE, default=TYPE[0][0], verbose_name=_("Schedule Type")
)
minutes = models.PositiveSmallIntegerField(
null=True, blank=True, help_text=_("Number of minutes for the Minutes type")
)
repeats = models.IntegerField(
default=-1, verbose_name=_("Repeats"), help_text=_("n = n times, -1 = forever")
)
next_run = models.DateTimeField(
verbose_name=_("Next Run"), default=timezone.now, null=True
) )
schedule_type = models.CharField(max_length=1, choices=TYPE, default=TYPE[0][0], verbose_name=_('Schedule Type'))
minutes = models.PositiveSmallIntegerField(null=True, blank=True,
help_text=_('Number of minutes for the Minutes type'))
repeats = models.IntegerField(default=-1, verbose_name=_('Repeats'), help_text=_('n = n times, -1 = forever'))
next_run = models.DateTimeField(verbose_name=_('Next Run'), default=timezone.now, null=True)
task = models.CharField(max_length=100, null=True, editable=False) task = models.CharField(max_length=100, null=True, editable=False)
def success(self): def success(self):
@@ -166,10 +184,10 @@ class Schedule(models.Model):
if self.task and Task.objects.filter(id=self.task): if self.task and Task.objects.filter(id=self.task):
task = Task.objects.get(id=self.task) task = Task.objects.get(id=self.task)
if task.success: if task.success:
url = reverse('admin:django_q_success_change', args=(task.id,)) url = reverse("admin:django_q_success_change", args=(task.id,))
else: else:
url = reverse('admin:django_q_failure_change', args=(task.id,)) url = reverse("admin:django_q_failure_change", args=(task.id,))
return format_html('<a href="{}">[{}]</a>'.format(url, task.name)) return format_html(f'<a href="{url}">[{task.name}]</a>')
return None return None
def __unicode__(self): def __unicode__(self):
@@ -179,10 +197,10 @@ class Schedule(models.Model):
last_run.allow_tags = True last_run.allow_tags = True
class Meta: class Meta:
app_label = 'django_q' app_label = "django_q"
verbose_name = _('Scheduled task') verbose_name = _("Scheduled task")
verbose_name_plural = _('Scheduled tasks') verbose_name_plural = _("Scheduled tasks")
ordering = ['next_run'] ordering = ["next_run"]
class OrmQ(models.Model): class OrmQ(models.Model):
@@ -194,23 +212,23 @@ class OrmQ(models.Model):
return SignedPackage.loads(self.payload) return SignedPackage.loads(self.payload)
def func(self): def func(self):
return self.task()['func'] return self.task()["func"]
def task_id(self): def task_id(self):
return self.task()['id'] return self.task()["id"]
def name(self): def name(self):
return self.task()['name'] return self.task()["name"]
class Meta: class Meta:
app_label = 'django_q' app_label = "django_q"
verbose_name = _('Queued task') verbose_name = _("Queued task")
verbose_name_plural = _('Queued tasks') verbose_name_plural = _("Queued tasks")
# Backwards compatibility for Django 1.7 # Backwards compatibility for Django 1.7
def decode_results(values): def decode_results(values):
if get_version().split('.')[1] == '7': if get_version().split(".")[1] == "7":
# decode values in 1.7 # decode values in 1.7
return [dbsafe_decode(v) for v in values] return [dbsafe_decode(v) for v in values]
return values return values
+161 -78
View File
@@ -24,20 +24,44 @@ def monitor(run_once=False, broker=None):
with term.fullscreen(), term.hidden_cursor(), term.cbreak(): with term.fullscreen(), term.hidden_cursor(), term.cbreak():
val = None val = None
start_width = int(term.width / 8) start_width = int(term.width / 8)
while val not in (u'q', u'Q',): while val not in ("q", "Q",):
col_width = int(term.width / 8) col_width = int(term.width / 8)
# In case of resize # In case of resize
if col_width != start_width: if col_width != start_width:
print(term.clear()) print(term.clear())
start_width = col_width start_width = col_width
print(term.move(0, 0) + term.black_on_green(term.center(_('Host'), width=col_width - 1))) print(
print(term.move(0, 1 * col_width) + term.black_on_green(term.center(_('Id'), width=col_width - 1))) term.move(0, 0)
print(term.move(0, 2 * col_width) + term.black_on_green(term.center(_('State'), width=col_width - 1))) + term.black_on_green(term.center(_("Host"), width=col_width - 1))
print(term.move(0, 3 * col_width) + term.black_on_green(term.center(_('Pool'), width=col_width - 1))) )
print(term.move(0, 4 * col_width) + term.black_on_green(term.center(_('TQ'), width=col_width - 1))) print(
print(term.move(0, 5 * col_width) + term.black_on_green(term.center(_('RQ'), width=col_width - 1))) term.move(0, 1 * col_width)
print(term.move(0, 6 * col_width) + term.black_on_green(term.center(_('RC'), width=col_width - 1))) + term.black_on_green(term.center(_("Id"), width=col_width - 1))
print(term.move(0, 7 * col_width) + term.black_on_green(term.center(_('Up'), width=col_width - 1))) )
print(
term.move(0, 2 * col_width)
+ term.black_on_green(term.center(_("State"), width=col_width - 1))
)
print(
term.move(0, 3 * col_width)
+ term.black_on_green(term.center(_("Pool"), width=col_width - 1))
)
print(
term.move(0, 4 * col_width)
+ term.black_on_green(term.center(_("TQ"), width=col_width - 1))
)
print(
term.move(0, 5 * col_width)
+ term.black_on_green(term.center(_("RQ"), width=col_width - 1))
)
print(
term.move(0, 6 * col_width)
+ term.black_on_green(term.center(_("RC"), width=col_width - 1))
)
print(
term.move(0, 7 * col_width)
+ term.black_on_green(term.center(_("Up"), width=col_width - 1))
)
i = 2 i = 2
stats = Stat.get_all(broker=broker) stats = Stat.get_all(broker=broker)
print(term.clear_eos()) print(term.clear_eos())
@@ -69,36 +93,83 @@ def monitor(run_once=False, broker=None):
uptime = (timezone.now() - stat.tob).total_seconds() uptime = (timezone.now() - stat.tob).total_seconds()
hours, remainder = divmod(uptime, 3600) hours, remainder = divmod(uptime, 3600)
minutes, seconds = divmod(remainder, 60) minutes, seconds = divmod(remainder, 60)
uptime = '%d:%02d:%02d' % (hours, minutes, seconds) uptime = "%d:%02d:%02d" % (hours, minutes, seconds)
# print to the terminal # print to the terminal
print(term.move(i, 0) + term.center(stat.host[:col_width - 1], width=col_width - 1)) print(
print(term.move(i, 1 * col_width) + term.center(str(stat.cluster_id)[-8:], width=col_width - 1)) term.move(i, 0)
print(term.move(i, 2 * col_width) + term.center(status, width=col_width - 1)) + term.center(stat.host[: col_width - 1], width=col_width - 1)
print(term.move(i, 3 * col_width) + term.center(workers, width=col_width - 1)) )
print(term.move(i, 4 * col_width) + term.center(tasks, width=col_width - 1)) print(
print(term.move(i, 5 * col_width) + term.center(results, width=col_width - 1)) term.move(i, 1 * col_width)
print(term.move(i, 6 * col_width) + term.center(stat.reincarnations, width=col_width - 1)) + term.center(str(stat.cluster_id)[-8:], width=col_width - 1)
print(term.move(i, 7 * col_width) + term.center(uptime, width=col_width - 1)) )
print(
term.move(i, 2 * col_width)
+ term.center(status, width=col_width - 1)
)
print(
term.move(i, 3 * col_width)
+ term.center(workers, width=col_width - 1)
)
print(
term.move(i, 4 * col_width)
+ term.center(tasks, width=col_width - 1)
)
print(
term.move(i, 5 * col_width)
+ term.center(results, width=col_width - 1)
)
print(
term.move(i, 6 * col_width)
+ term.center(stat.reincarnations, width=col_width - 1)
)
print(
term.move(i, 7 * col_width)
+ term.center(uptime, width=col_width - 1)
)
i += 1 i += 1
# bottom bar # bottom bar
i += 1 i += 1
queue_size = broker.queue_size() queue_size = broker.queue_size()
lock_size = broker.lock_size() lock_size = broker.lock_size()
if lock_size: if lock_size:
queue_size = '{}({})'.format(queue_size, lock_size) queue_size = f"{queue_size}({lock_size})"
print(term.move(i, 0) + term.white_on_cyan(term.center(broker.info(), width=col_width * 2))) print(
print(term.move(i, 2 * col_width) + term.black_on_cyan(term.center(_('Queued'), width=col_width))) term.move(i, 0)
print(term.move(i, 3 * col_width) + term.white_on_cyan(term.center(queue_size, width=col_width))) + term.white_on_cyan(term.center(broker.info(), width=col_width * 2))
print(term.move(i, 4 * col_width) + term.black_on_cyan(term.center(_('Success'), width=col_width))) )
print(term.move(i, 5 * col_width) + term.white_on_cyan( print(
term.center(models.Success.objects.count(), width=col_width))) term.move(i, 2 * col_width)
print(term.move(i, 6 * col_width) + term.black_on_cyan(term.center(_('Failures'), width=col_width))) + term.black_on_cyan(term.center(_("Queued"), width=col_width))
print(term.move(i, 7 * col_width) + term.white_on_cyan( )
term.center(models.Failure.objects.count(), width=col_width))) print(
term.move(i, 3 * col_width)
+ term.white_on_cyan(term.center(queue_size, width=col_width))
)
print(
term.move(i, 4 * col_width)
+ term.black_on_cyan(term.center(_("Success"), width=col_width))
)
print(
term.move(i, 5 * col_width)
+ term.white_on_cyan(
term.center(models.Success.objects.count(), width=col_width)
)
)
print(
term.move(i, 6 * col_width)
+ term.black_on_cyan(term.center(_("Failures"), width=col_width))
)
print(
term.move(i, 7 * col_width)
+ term.white_on_cyan(
term.center(models.Failure.objects.count(), width=col_width)
)
)
# for testing # for testing
if run_once: if run_once:
return Stat.get_all(broker=broker) return Stat.get_all(broker=broker)
print(term.move(i + 2, 0) + term.center(_('[Press q to quit]'))) print(term.move(i + 2, 0) + term.center(_("[Press q to quit]")))
val = term.inkey(timeout=1) val = term.inkey(timeout=1)
@@ -117,15 +188,19 @@ def info(broker=None):
reincarnations += cluster.reincarnations reincarnations += cluster.reincarnations
# calculate tasks pm and avg exec time # calculate tasks pm and avg exec time
tasks_per = 0 tasks_per = 0
per = _('day') per = _("day")
exec_time = 0 exec_time = 0
last_tasks = models.Success.objects.filter(stopped__gte=timezone.now() - timedelta(hours=24)) last_tasks = models.Success.objects.filter(
stopped__gte=timezone.now() - timedelta(hours=24)
)
tasks_per_day = last_tasks.count() tasks_per_day = last_tasks.count()
if tasks_per_day > 0: if tasks_per_day > 0:
# average execution time over the last 24 hours # average execution time over the last 24 hours
if not connection.vendor == 'sqlite': if not connection.vendor == "sqlite":
exec_time = last_tasks.aggregate(time_taken=Sum(F('stopped') - F('started'))) exec_time = last_tasks.aggregate(
exec_time = exec_time['time_taken'].total_seconds() / tasks_per_day time_taken=Sum(F("stopped") - F("started"))
)
exec_time = exec_time["time_taken"].total_seconds() / tasks_per_day
else: else:
# can't sum timedeltas on sqlite # can't sum timedeltas on sqlite
for t in last_tasks: for t in last_tasks:
@@ -134,58 +209,66 @@ def info(broker=None):
# tasks per second/minute/hour/day in the last 24 hours # tasks per second/minute/hour/day in the last 24 hours
if tasks_per_day > 24 * 60 * 60: if tasks_per_day > 24 * 60 * 60:
tasks_per = tasks_per_day / (24 * 60 * 60) tasks_per = tasks_per_day / (24 * 60 * 60)
per = _('second') per = _("second")
elif tasks_per_day > 24 * 60: elif tasks_per_day > 24 * 60:
tasks_per = tasks_per_day / (24 * 60) tasks_per = tasks_per_day / (24 * 60)
per = _('minute') per = _("minute")
elif tasks_per_day > 24: elif tasks_per_day > 24:
tasks_per = tasks_per_day / 24 tasks_per = tasks_per_day / 24
per = _('hour') per = _("hour")
else: else:
tasks_per = tasks_per_day tasks_per = tasks_per_day
# print to terminal # print to terminal
print(term.clear_eos()) print(term.clear_eos())
col_width = int(term.width / 6) col_width = int(term.width / 6)
print(term.black_on_green( print(
term.center( term.black_on_green(
_('-- {} {} on {} --').format(Conf.PREFIX.capitalize(), '.'.join(str(v) for v in VERSION), term.center(
broker.info())))) _(
print(term.cyan(_('Clusters')) + f'-- {Conf.PREFIX.capitalize()} { ".".join(str(v) for v in VERSION)} on {broker.info()} --'
term.move_x(1 * col_width) + )
term.white(str(clusters)) + )
term.move_x(2 * col_width) + )
term.cyan(_('Workers')) + )
term.move_x(3 * col_width) + print(
term.white(str(workers)) + term.cyan(_("Clusters"))
term.move_x(4 * col_width) + + term.move_x(1 * col_width)
term.cyan(_('Restarts')) + + term.white(str(clusters))
term.move_x(5 * col_width) + + term.move_x(2 * col_width)
term.white(str(reincarnations)) + term.cyan(_("Workers"))
) + term.move_x(3 * col_width)
print(term.cyan(_('Queued')) + + term.white(str(workers))
term.move_x(1 * col_width) + + term.move_x(4 * col_width)
term.white(str(broker.queue_size())) + + term.cyan(_("Restarts"))
term.move_x(2 * col_width) + + term.move_x(5 * col_width)
term.cyan(_('Successes')) + + term.white(str(reincarnations))
term.move_x(3 * col_width) + )
term.white(str(models.Success.objects.count())) + print(
term.move_x(4 * col_width) + term.cyan(_("Queued"))
term.cyan(_('Failures')) + + term.move_x(1 * col_width)
term.move_x(5 * col_width) + + term.white(str(broker.queue_size()))
term.white(str(models.Failure.objects.count())) + term.move_x(2 * col_width)
) + term.cyan(_("Successes"))
print(term.cyan(_('Schedules')) + + term.move_x(3 * col_width)
term.move_x(1 * col_width) + + term.white(str(models.Success.objects.count()))
term.white(str(models.Schedule.objects.count())) + + term.move_x(4 * col_width)
term.move_x(2 * col_width) + + term.cyan(_("Failures"))
term.cyan(_('Tasks/{}'.format(per))) + + term.move_x(5 * col_width)
term.move_x(3 * col_width) + + term.white(str(models.Failure.objects.count()))
term.white('{0:.2f}'.format(tasks_per)) + )
term.move_x(4 * col_width) + print(
term.cyan(_('Avg time')) + term.cyan(_("Schedules"))
term.move_x(5 * col_width) + + term.move_x(1 * col_width)
term.white('{0:.4f}'.format(exec_time)) + term.white(str(models.Schedule.objects.count()))
) + term.move_x(2 * col_width)
+ term.cyan(_(f"Tasks/{per}"))
+ term.move_x(3 * col_width)
+ term.white(f"{tasks_per:.2f}")
+ term.move_x(4 * col_width)
+ term.cyan(_("Avg time"))
+ term.move_x(5 * col_width)
+ term.white(f"{exec_time:.4f}")
)
return True return True
@@ -196,5 +279,5 @@ def get_ids():
for s in stat: for s in stat:
print(s.cluster_id) print(s.cluster_id)
else: else:
print('No clusters appear to be running.') print("No clusters appear to be running.")
return True return True
+9 -3
View File
@@ -14,16 +14,22 @@ def call_hook(sender, instance, **kwargs):
f = instance.hook f = instance.hook
if not callable(f): if not callable(f):
try: try:
module, func = f.rsplit('.', 1) module, func = f.rsplit(".", 1)
m = importlib.import_module(module) m = importlib.import_module(module)
f = getattr(m, func) f = getattr(m, func)
except (ValueError, ImportError, AttributeError): except (ValueError, ImportError, AttributeError):
logger.error(_('malformed return hook \'{}\' for [{}]').format(instance.hook, instance.name)) logger.error(
_(f"malformed return hook '{instance.hook}' for [{instance.name}]")
)
return return
try: try:
f(instance) f(instance)
except Exception as e: except Exception as e:
logger.error(_('return hook {} failed on [{}] because {}').format(instance.hook, instance.name, e)) logger.error(
_(
f"return hook {instance.hook} failed on [{instance.name}] because {str(e)}"
)
)
pre_enqueue = Signal(providing_args=["task"]) pre_enqueue = Signal(providing_args=["task"])
+6 -4
View File
@@ -28,7 +28,9 @@ class Stat(Status):
"""Status object for Cluster monitoring.""" """Status object for Cluster monitoring."""
def __init__(self, sentinel): def __init__(self, sentinel):
super(Stat, self).__init__(sentinel.parent_pid or sentinel.pid, cluster_id=sentinel.cluster_id) super(Stat, self).__init__(
sentinel.parent_pid or sentinel.pid, cluster_id=sentinel.cluster_id
)
self.broker = sentinel.broker or get_broker() self.broker = sentinel.broker or get_broker()
self.tob = sentinel.tob self.tob = sentinel.tob
self.reincarnations = sentinel.reincarnations self.reincarnations = sentinel.reincarnations
@@ -61,7 +63,7 @@ class Stat(Status):
:param cluster_id: cluster ID :param cluster_id: cluster ID
:return: redis key for the cluster statistic :return: redis key for the cluster statistic
""" """
return '{}:{}'.format(Conf.Q_STAT, cluster_id) return f"{Conf.Q_STAT}:{cluster_id}"
def save(self): def save(self):
try: try:
@@ -99,7 +101,7 @@ class Stat(Status):
if not broker: if not broker:
broker = get_broker() broker = get_broker()
stats = [] stats = []
packs = broker.get_stats('{}:*'.format(Conf.Q_STAT)) or [] packs = broker.get_stats(f"{Conf.Q_STAT}:*") or []
for pack in packs: for pack in packs:
try: try:
stats.append(SignedPackage.loads(pack)) stats.append(SignedPackage.loads(pack))
@@ -110,5 +112,5 @@ class Stat(Status):
def __getstate__(self): def __getstate__(self):
# Don't pickle the redis connection # Don't pickle the redis connection
state = dict(self.__dict__) state = dict(self.__dict__)
del state['broker'] del state["broker"]
return state return state
+193 -120
View File
@@ -8,6 +8,7 @@ from django.utils import timezone
from multiprocessing import Value from multiprocessing import Value
from django_q.brokers import get_broker from django_q.brokers import get_broker
# local # local
from django_q.conf import Conf, logger from django_q.conf import Conf, logger
from django_q.humanhash import uuid from django_q.humanhash import uuid
@@ -21,15 +22,30 @@ def async_task(func, *args, **kwargs):
"""Queue a task for the cluster.""" """Queue a task for the cluster."""
keywords = kwargs.copy() keywords = kwargs.copy()
opt_keys = ( opt_keys = (
'hook', 'group', 'save', 'sync', 'cached', 'ack_failure', 'iter_count', 'iter_cached', 'chain', 'broker', 'timeout') "hook",
q_options = keywords.pop('q_options', {}) "group",
"save",
"sync",
"cached",
"ack_failure",
"iter_count",
"iter_cached",
"chain",
"broker",
"timeout",
)
q_options = keywords.pop("q_options", {})
# get an id # get an id
tag = uuid() tag = uuid()
# build the task package # build the task package
task = {'id': tag[1], task = {
'name': keywords.pop('task_name', None) or q_options.pop('task_name', None) or tag[0], "id": tag[1],
'func': func, "name": keywords.pop("task_name", None)
'args': args} or q_options.pop("task_name", None)
or tag[0],
"func": func,
"args": args,
}
# push optionals # push optionals
for key in opt_keys: for key in opt_keys:
if q_options and key in q_options: if q_options and key in q_options:
@@ -37,28 +53,28 @@ def async_task(func, *args, **kwargs):
elif key in keywords: elif key in keywords:
task[key] = keywords.pop(key) task[key] = keywords.pop(key)
# don't serialize the broker # don't serialize the broker
broker = task.pop('broker', get_broker()) broker = task.pop("broker", get_broker())
# overrides # overrides
if 'cached' not in task and Conf.CACHED: if "cached" not in task and Conf.CACHED:
task['cached'] = Conf.CACHED task["cached"] = Conf.CACHED
if 'sync' not in task and Conf.SYNC: if "sync" not in task and Conf.SYNC:
task['sync'] = Conf.SYNC task["sync"] = Conf.SYNC
if 'ack_failure' not in task and Conf.ACK_FAILURES: if "ack_failure" not in task and Conf.ACK_FAILURES:
task['ack_failure'] = Conf.ACK_FAILURES task["ack_failure"] = Conf.ACK_FAILURES
# finalize # finalize
task['kwargs'] = keywords task["kwargs"] = keywords
task['started'] = timezone.now() task["started"] = timezone.now()
# signal it # signal it
pre_enqueue.send(sender="django_q", task=task) pre_enqueue.send(sender="django_q", task=task)
# sign it # sign it
pack = SignedPackage.dumps(task) pack = SignedPackage.dumps(task)
if task.get('sync', False): if task.get("sync", False):
return _sync(pack) return _sync(pack)
# push it # push it
enqueue_id = broker.enqueue(pack) enqueue_id = broker.enqueue(pack)
logger.info('Enqueued {}'.format(enqueue_id)) logger.info(f"Enqueued {enqueue_id}")
logger.debug('Pushed {}'.format(tag)) logger.debug(f"Pushed {tag}")
return task['id'] return task["id"]
def schedule(func, *args, **kwargs): def schedule(func, *args, **kwargs):
@@ -77,28 +93,29 @@ def schedule(func, *args, **kwargs):
:return: the schedule object. :return: the schedule object.
:rtype: Schedule :rtype: Schedule
""" """
name = kwargs.pop('name', None) name = kwargs.pop("name", None)
hook = kwargs.pop('hook', None) hook = kwargs.pop("hook", None)
schedule_type = kwargs.pop('schedule_type', Schedule.ONCE) schedule_type = kwargs.pop("schedule_type", Schedule.ONCE)
minutes = kwargs.pop('minutes', None) minutes = kwargs.pop("minutes", None)
repeats = kwargs.pop('repeats', -1) repeats = kwargs.pop("repeats", -1)
next_run = kwargs.pop('next_run', timezone.now()) next_run = kwargs.pop("next_run", timezone.now())
# check for name duplicates instead of am unique constraint # check for name duplicates instead of am unique constraint
if name and Schedule.objects.filter(name=name).exists(): if name and Schedule.objects.filter(name=name).exists():
raise IntegrityError("A schedule with the same name already exists.") raise IntegrityError("A schedule with the same name already exists.")
# create and return the schedule # create and return the schedule
return Schedule.objects.create(name=name, return Schedule.objects.create(
func=func, name=name,
hook=hook, func=func,
args=args, hook=hook,
kwargs=kwargs, args=args,
schedule_type=schedule_type, kwargs=kwargs,
minutes=minutes, schedule_type=schedule_type,
repeats=repeats, minutes=minutes,
next_run=next_run repeats=repeats,
) next_run=next_run,
)
def result(task_id, wait=0, cached=Conf.CACHED): def result(task_id, wait=0, cached=Conf.CACHED):
@@ -133,9 +150,9 @@ def result_cached(task_id, wait=0, broker=None):
broker = get_broker() broker = get_broker()
start = time() start = time()
while True: while True:
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) r = broker.cache.get(f"{broker.list_key}:{task_id}")
if r: if r:
return SignedPackage.loads(r)['result'] return SignedPackage.loads(r)["result"]
if (time() - start) * 1000 >= wait >= 0: if (time() - start) * 1000 >= wait >= 0:
break break
sleep(0.01) sleep(0.01)
@@ -156,7 +173,11 @@ def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHE
start = time() start = time()
if count: if count:
while True: while True:
if count_group(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0: if (
count_group(group_id) == count
or wait
and (time() - start) * 1000 >= wait >= 0
):
break break
sleep(0.01) sleep(0.01)
while True: while True:
@@ -177,17 +198,21 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non
start = time() start = time()
if count: if count:
while True: while True:
if count_group_cached(group_id) == count or wait and (time() - start) * 1000 >= wait > 0: if (
count_group_cached(group_id) == count
or wait
and (time() - start) * 1000 >= wait > 0
):
break break
sleep(0.01) sleep(0.01)
while True: while True:
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) group_list = broker.cache.get(f"{broker.list_key}:{group_id}:keys")
if group_list: if group_list:
result_list = [] result_list = []
for task_key in group_list: for task_key in group_list:
task = SignedPackage.loads(broker.cache.get(task_key)) task = SignedPackage.loads(broker.cache.get(task_key))
if task['success'] or failures: if task["success"] or failures:
result_list.append(task['result']) result_list.append(task["result"])
return result_list return result_list
if (time() - start) * 1000 >= wait >= 0: if (time() - start) * 1000 >= wait >= 0:
break break
@@ -226,19 +251,21 @@ def fetch_cached(task_id, wait=0, broker=None):
broker = get_broker() broker = get_broker()
start = time() start = time()
while True: while True:
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id)) r = broker.cache.get(f"{broker.list_key}:{task_id}")
if r: if r:
task = SignedPackage.loads(r) task = SignedPackage.loads(r)
t = Task(id=task['id'], t = Task(
name=task['name'], id=task["id"],
func=task['func'], name=task["name"],
hook=task.get('hook'), func=task["func"],
args=task['args'], hook=task.get("hook"),
kwargs=task['kwargs'], args=task["args"],
started=task['started'], kwargs=task["kwargs"],
stopped=task['stopped'], started=task["started"],
result=task['result'], stopped=task["stopped"],
success=task['success']) result=task["result"],
success=task["success"],
)
return t return t
if (time() - start) * 1000 >= wait >= 0: if (time() - start) * 1000 >= wait >= 0:
break break
@@ -259,7 +286,11 @@ def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED)
start = time() start = time()
if count: if count:
while True: while True:
if count_group(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0: if (
count_group(group_id) == count
or wait
and (time() - start) * 1000 >= wait >= 0
):
break break
sleep(0.01) sleep(0.01)
while True: while True:
@@ -280,27 +311,33 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None)
start = time() start = time()
if count: if count:
while True: while True:
if count_group_cached(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0: if (
count_group_cached(group_id) == count
or wait
and (time() - start) * 1000 >= wait >= 0
):
break break
sleep(0.01) sleep(0.01)
while True: while True:
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) group_list = broker.cache.get(f"{broker.list_key}:{group_id}:keys")
if group_list: if group_list:
task_list = [] task_list = []
for task_key in group_list: for task_key in group_list:
task = SignedPackage.loads(broker.cache.get(task_key)) task = SignedPackage.loads(broker.cache.get(task_key))
if task['success'] or failures: if task["success"] or failures:
t = Task(id=task['id'], t = Task(
name=task['name'], id=task["id"],
func=task['func'], name=task["name"],
hook=task.get('hook'), func=task["func"],
args=task['args'], hook=task.get("hook"),
kwargs=task['kwargs'], args=task["args"],
started=task['started'], kwargs=task["kwargs"],
stopped=task['stopped'], started=task["started"],
result=task['result'], stopped=task["stopped"],
group=task.get('group'), result=task["result"],
success=task['success']) group=task.get("group"),
success=task["success"],
)
task_list.append(t) task_list.append(t)
return task_list return task_list
if (time() - start) * 1000 >= wait >= 0: if (time() - start) * 1000 >= wait >= 0:
@@ -329,14 +366,14 @@ def count_group_cached(group_id, failures=False, broker=None):
""" """
if not broker: if not broker:
broker = get_broker() broker = get_broker()
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id)) group_list = broker.cache.get(f"{broker.list_key}:{group_id}:keys")
if group_list: if group_list:
if not failures: if not failures:
return len(group_list) return len(group_list)
failure_count = 0 failure_count = 0
for task_key in group_list: for task_key in group_list:
task = SignedPackage.loads(broker.cache.get(task_key)) task = SignedPackage.loads(broker.cache.get(task_key))
if not task['success']: if not task["success"]:
failure_count += 1 failure_count += 1
return failure_count return failure_count
@@ -362,7 +399,7 @@ def delete_group_cached(group_id, broker=None):
""" """
if not broker: if not broker:
broker = get_broker() broker = get_broker()
group_key = '{}:{}:keys'.format(broker.list_key, group_id) group_key = f"{broker.list_key}:{group_id}:keys"
group_list = broker.cache.get(group_key) group_list = broker.cache.get(group_key)
broker.cache.delete_many(group_list) broker.cache.delete_many(group_list)
broker.cache.delete(group_key) broker.cache.delete(group_key)
@@ -374,7 +411,7 @@ def delete_cached(task_id, broker=None):
""" """
if not broker: if not broker:
broker = get_broker() broker = get_broker()
return broker.cache.delete('{}:{}'.format(broker.list_key, task_id)) return broker.cache.delete(f"{broker.list_key}:{task_id}")
def queue_size(broker=None): def queue_size(broker=None):
@@ -398,17 +435,19 @@ def async_iter(func, args_iter, **kwargs):
iter_count = len(args_iter) iter_count = len(args_iter)
iter_group = uuid()[1] iter_group = uuid()[1]
# clean up the kwargs # clean up the kwargs
options = kwargs.get('q_options', kwargs) options = kwargs.get("q_options", kwargs)
options.pop('hook', None) options.pop("hook", None)
options['broker'] = options.get('broker', get_broker()) options["broker"] = options.get("broker", get_broker())
options['group'] = iter_group options["group"] = iter_group
options['iter_count'] = iter_count options["iter_count"] = iter_count
if options.get('cached', None): if options.get("cached", None):
options['iter_cached'] = options['cached'] options["iter_cached"] = options["cached"]
options['cached'] = True options["cached"] = True
# save the original arguments # save the original arguments
broker = options['broker'] broker = options["broker"]
broker.cache.set('{}:{}:args'.format(broker.list_key, iter_group), SignedPackage.dumps(args_iter)) broker.cache.set(
f"{broker.list_key}:{iter_group}:args", SignedPackage.dumps(args_iter)
)
for args in args_iter: for args in args_iter:
if not isinstance(args, tuple): if not isinstance(args, tuple):
args = (args,) args = (args,)
@@ -432,11 +471,11 @@ def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=No
args = task[1] args = task[1]
if len(task) > 2: if len(task) > 2:
kwargs = task[2] kwargs = task[2]
kwargs['chain'] = chain kwargs["chain"] = chain
kwargs['group'] = group kwargs["group"] = group
kwargs['cached'] = cached kwargs["cached"] = cached
kwargs['sync'] = sync kwargs["sync"] = sync
kwargs['broker'] = broker or get_broker() kwargs["broker"] = broker or get_broker()
async_task(task[0], *args, **kwargs) async_task(task[0], *args, **kwargs)
return group return group
@@ -446,11 +485,19 @@ class Iter(object):
An async task with iterable arguments An async task with iterable arguments
""" """
def __init__(self, func=None, args=None, kwargs=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None): def __init__(
self,
func=None,
args=None,
kwargs=None,
cached=Conf.CACHED,
sync=Conf.SYNC,
broker=None,
):
self.func = func self.func = func
self.args = args or [] self.args = args or []
self.kwargs = kwargs or {} self.kwargs = kwargs or {}
self.id = '' self.id = ""
self.broker = broker or get_broker() self.broker = broker or get_broker()
self.cached = cached self.cached = cached
self.sync = sync self.sync = sync
@@ -470,9 +517,9 @@ class Iter(object):
Start queueing the tasks to the worker cluster Start queueing the tasks to the worker cluster
:return: the task id :return: the task id
""" """
self.kwargs['cached'] = self.cached self.kwargs["cached"] = self.cached
self.kwargs['sync'] = self.sync self.kwargs["sync"] = self.sync
self.kwargs['broker'] = self.broker self.kwargs["broker"] = self.broker
self.id = async_iter(self.func, self.args, **self.kwargs) self.id = async_iter(self.func, self.args, **self.kwargs)
self.started = True self.started = True
return self.id return self.id
@@ -510,7 +557,7 @@ class Chain(object):
def __init__(self, chain=None, group=None, cached=Conf.CACHED, sync=Conf.SYNC): def __init__(self, chain=None, group=None, cached=Conf.CACHED, sync=Conf.SYNC):
self.chain = chain or [] self.chain = chain or []
self.group = group or '' self.group = group or ""
self.broker = get_broker() self.broker = get_broker()
self.cached = cached self.cached = cached
self.sync = sync self.sync = sync
@@ -533,8 +580,13 @@ class Chain(object):
Start queueing the chain to the worker cluster Start queueing the chain to the worker cluster
:return: the chain's group id :return: the chain's group id
""" """
self.group = async_chain(chain=self.chain[:], group=self.group, cached=self.cached, sync=self.sync, self.group = async_chain(
broker=self.broker) chain=self.chain[:],
group=self.group,
cached=self.cached,
sync=self.sync,
broker=self.broker,
)
self.started = True self.started = True
return self.group return self.group
@@ -545,7 +597,9 @@ class Chain(object):
:return: an unsorted list of results :return: an unsorted list of results
""" """
if self.started: if self.started:
return result_group(self.group, wait=wait, count=self.length(), cached=self.cached) return result_group(
self.group, wait=wait, count=self.length(), cached=self.cached
)
def fetch(self, failures=True, wait=0): def fetch(self, failures=True, wait=0):
""" """
@@ -555,7 +609,13 @@ class Chain(object):
:return: an unsorted list of task objects :return: an unsorted list of task objects
""" """
if self.started: if self.started:
return fetch_group(self.group, failures=failures, wait=wait, count=self.length(), cached=self.cached) return fetch_group(
self.group,
failures=failures,
wait=wait,
count=self.length(),
cached=self.cached,
)
def current(self): def current(self):
""" """
@@ -580,7 +640,7 @@ class AsyncTask(object):
""" """
def __init__(self, func, *args, **kwargs): def __init__(self, func, *args, **kwargs):
self.id = '' self.id = ""
self.started = False self.started = False
self.func = func self.func = func
self.args = args self.args = args
@@ -588,62 +648,62 @@ class AsyncTask(object):
@property @property
def broker(self): def broker(self):
return self._get_option('broker', None) return self._get_option("broker", None)
@broker.setter @broker.setter
def broker(self, value): def broker(self, value):
self._set_option('broker', value) self._set_option("broker", value)
@property @property
def sync(self): def sync(self):
return self._get_option('sync', None) return self._get_option("sync", None)
@sync.setter @sync.setter
def sync(self, value): def sync(self, value):
self._set_option('sync', value) self._set_option("sync", value)
@property @property
def save(self): def save(self):
return self._get_option('save', None) return self._get_option("save", None)
@save.setter @save.setter
def save(self, value): def save(self, value):
self._set_option('save', value) self._set_option("save", value)
@property @property
def hook(self): def hook(self):
return self._get_option('hook', None) return self._get_option("hook", None)
@hook.setter @hook.setter
def hook(self, value): def hook(self, value):
self._set_option('hook', value) self._set_option("hook", value)
@property @property
def group(self): def group(self):
return self._get_option('group', None) return self._get_option("group", None)
@group.setter @group.setter
def group(self, value): def group(self, value):
self._set_option('group', value) self._set_option("group", value)
@property @property
def cached(self): def cached(self):
return self._get_option('cached', Conf.CACHED) return self._get_option("cached", Conf.CACHED)
@cached.setter @cached.setter
def cached(self, value): def cached(self, value):
self._set_option('cached', value) self._set_option("cached", value)
def _set_option(self, key, value): def _set_option(self, key, value):
if 'q_options' in self.kwargs: if "q_options" in self.kwargs:
self.kwargs['q_options'][key] = value self.kwargs["q_options"][key] = value
else: else:
self.kwargs[key] = value self.kwargs[key] = value
self.started = False self.started = False
def _get_option(self, key, default=None): def _get_option(self, key, default=None):
if 'q_options' in self.kwargs: if "q_options" in self.kwargs:
return self.kwargs['q_options'].get(key, default) return self.kwargs["q_options"].get(key, default)
else: else:
return self.kwargs.get(key, default) return self.kwargs.get(key, default)
@@ -665,27 +725,40 @@ class AsyncTask(object):
def result_group(self, failures=False, wait=0, count=None): def result_group(self, failures=False, wait=0, count=None):
if self.started and self.group: if self.started and self.group:
return result_group(self.group, failures=failures, wait=wait, count=count, cached=self.cached) return result_group(
self.group,
failures=failures,
wait=wait,
count=count,
cached=self.cached,
)
def fetch_group(self, failures=True, wait=0, count=None): def fetch_group(self, failures=True, wait=0, count=None):
if self.started and self.group: if self.started and self.group:
return fetch_group(self.group, failures=failures, wait=wait, count=count, cached=self.cached) return fetch_group(
self.group,
failures=failures,
wait=wait,
count=count,
cached=self.cached,
)
def _sync(pack): def _sync(pack):
"""Simulate a package travelling through the cluster.""" """Simulate a package travelling through the cluster."""
from django_q.cluster import worker, monitor from django_q.cluster import worker, monitor
task_queue = Queue() task_queue = Queue()
result_queue = Queue() result_queue = Queue()
task = SignedPackage.loads(pack) task = SignedPackage.loads(pack)
task_queue.put(task) task_queue.put(task)
task_queue.put('STOP') task_queue.put("STOP")
worker(task_queue, result_queue, Value('f', -1)) worker(task_queue, result_queue, Value("f", -1))
result_queue.put('STOP') result_queue.put("STOP")
monitor(result_queue) monitor(result_queue)
task_queue.close() task_queue.close()
task_queue.join_thread() task_queue.join_thread()
result_queue.close() result_queue.close()
result_queue.join_thread() result_queue.join_thread()
return task['id'] return task["id"]