From 7fb7fca176eceb5112196b6620cd336a616124df Mon Sep 17 00:00:00 2001 From: Stan Triepels <1939656+GDay@users.noreply.github.com> Date: Sun, 2 Jul 2023 02:47:05 +0200 Subject: [PATCH] Move worker, scheduler, pusher and monitor to separate files (#100) --- django_q/admin.py | 32 +- django_q/brokers/orm.py | 8 +- django_q/cluster.py | 494 +------------ django_q/conf.py | 11 +- django_q/core_signing.py | 4 +- django_q/management/commands/qcluster.py | 5 +- django_q/management/commands/qinfo.py | 2 +- django_q/management/commands/qmemory.py | 2 +- django_q/management/commands/qmonitor.py | 2 +- .../0016_schedule_intended_date_kwarg.py | 1 + django_q/models.py | 21 +- django_q/monitor.py | 674 +++++------------- django_q/monitor_terminal.py | 510 +++++++++++++ django_q/pusher.py | 62 ++ django_q/scheduler.py | 124 ++++ django_q/signals.py | 1 + django_q/tasks.py | 3 +- django_q/tests/test_cached.py | 4 +- django_q/tests/test_cluster.py | 7 +- django_q/tests/test_monitor.py | 2 +- django_q/tests/test_scheduler.py | 22 +- django_q/utils.py | 22 +- django_q/worker.py | 118 +++ 23 files changed, 1127 insertions(+), 1004 deletions(-) create mode 100644 django_q/monitor_terminal.py create mode 100644 django_q/pusher.py create mode 100644 django_q/scheduler.py create mode 100644 django_q/worker.py diff --git a/django_q/admin.py b/django_q/admin.py index 9ba5719..00d194d 100644 --- a/django_q/admin.py +++ b/django_q/admin.py @@ -31,7 +31,15 @@ resubmit_task.short_description = _("Resubmit selected tasks to queue") class TaskAdmin(admin.ModelAdmin): """model admin for success tasks.""" - list_display = ("name", "group", "func", "cluster", "started", "stopped", "time_taken") + list_display = ( + "name", + "group", + "func", + "cluster", + "started", + "stopped", + "time_taken", + ) actions = [resubmit_task] def has_add_permission(self, request): @@ -55,7 +63,15 @@ class TaskAdmin(admin.ModelAdmin): class FailAdmin(admin.ModelAdmin): """model admin for failed tasks.""" - list_display = ("name", "group", "func", "cluster", "started", "stopped", "short_result") + list_display = ( + "name", + "group", + "func", + "cluster", + "started", + "stopped", + "short_result", + ) def has_add_permission(self, request): """Don't allow adds.""" @@ -132,7 +148,17 @@ class QueueAdmin(admin.ModelAdmin): """queue admin for ORM broker""" list_display = ("id", "key", "name", "group", "func", "lock", "task_id") - fields = ("key", "lock", "task_id", "name", "group", "func", "args", "kwargs", "q_options") + fields = ( + "key", + "lock", + "task_id", + "name", + "group", + "func", + "args", + "kwargs", + "q_options", + ) readonly_fields = fields[2:] def save_model(self, request, obj, form, change): diff --git a/django_q/brokers/orm.py b/django_q/brokers/orm.py index 9fcaf02..b8793f4 100644 --- a/django_q/brokers/orm.py +++ b/django_q/brokers/orm.py @@ -37,7 +37,9 @@ class ORM(Broker): def lock_size(self) -> int: return ( - self.get_connection().filter(key=self.list_key, lock__gt=timezone.now()).count() + self.get_connection() + .filter(key=self.list_key, lock__gt=timezone.now()) + .count() ) def purge_queue(self): @@ -62,7 +64,9 @@ class ORM(Broker): return package.pk def dequeue(self): - tasks = self.get_connection().filter(key=self.list_key, lock__lt=timezone.now())[ + tasks = self.get_connection().filter( + key=self.list_key, lock__lt=timezone.now() + )[ 0 : Conf.BULK # noqa: E203 ] if tasks: diff --git a/django_q/cluster.py b/django_q/cluster.py index fba613d..f3c1bce 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -1,11 +1,7 @@ # Standard -import ast -import pydoc import signal import socket -import traceback import uuid -from datetime import datetime, timedelta from multiprocessing import Event, Process, Value, current_process from time import sleep @@ -13,6 +9,11 @@ from time import sleep from django import core, db from django.apps.registry import apps +from django_q.monitor import monitor +from django_q.pusher import pusher +from django_q.scheduler import scheduler +from django_q.worker import worker + try: apps.check_apps_ready() except core.exceptions.AppRegistryNotReady: @@ -26,24 +27,14 @@ from django.utils.translation import gettext_lazy as _ # Local import django_q.tasks from django_q.brokers import Broker, get_broker -from django_q.conf import ( - Conf, - croniter, - error_reporter, - get_ppid, - logger, - psutil, - setproctitle, - resource, -) +from django_q.conf import Conf, get_ppid, logger, psutil, setproctitle from django_q.humanhash import humanize from django_q.models import Schedule, Success, Task from django_q.queues import Queue -from django_q.signals import post_execute, post_spawn, pre_execute from django_q.signing import BadSignature, SignedPackage from django_q.status import Stat, Status -from .utils import get_func_repr, localtime +from .utils import get_func_repr class Cluster: @@ -174,7 +165,7 @@ class Sentinel: def queue_name(self): # multi-queue: cluster name is (broker's) queue_name - return self.broker.list_key if self.broker else '--' + return self.broker.list_key if self.broker else "--" def start(self): self.broker.ping() @@ -248,20 +239,22 @@ class Sentinel: try: process_name = psutil.Process(process.pid).name() name_splits = process_name.split(" ") - task_name = name_splits[3] if len(name_splits) >= 4 and name_splits[2] == "processing" else "" + task_name = ( + name_splits[3] + if len(name_splits) >= 4 and name_splits[2] == "processing" + else "" + ) except psutil.NoSuchProcess: pass process.terminate() if task_name: - msg = ( - _("reincarnated worker %(name)s after timeout while processing task %(task_name)s") - % {"name": process.name, "task_name": task_name} - ) + msg = _( + "reincarnated worker %(name)s after timeout while processing task %(task_name)s" + ) % {"name": process.name, "task_name": task_name} else: - msg = ( - _("reincarnated worker %(name)s after timeout") - % {"name": process.name} - ) + msg = _("reincarnated worker %(name)s after timeout") % { + "name": process.name + } logger.critical(msg) elif int(process.timer.value) == -2: logger.info(_("recycled worker %(name)s") % {"name": process.name}) @@ -294,14 +287,18 @@ class Sentinel: _("%(name)s guarding cluster %(cluster_name)s") % { "name": current_process().name, - "cluster_name": humanize(self.cluster_id.hex) + f" [{self.queue_name()}]", + "cluster_name": humanize(self.cluster_id.hex) + + f" [{self.queue_name()}]", } ) self.start_event.set() Stat(self).save() logger.info( _("Q Cluster %(cluster_name)s running.") - % {"cluster_name": humanize(self.cluster_id.hex) + f" [{self.queue_name()}]"} + % { + "cluster_name": humanize(self.cluster_id.hex) + + f" [{self.queue_name()}]" + } ) counter = 0 cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds @@ -374,438 +371,6 @@ class Sentinel: Stat(self).save() -def pusher(task_queue: Queue, event: Event, broker: Broker = None): - """ - Pulls tasks of the broker and puts them in the task queue - :type broker: - :type task_queue: multiprocessing.Queue - :type event: multiprocessing.Event - """ - if not broker: - broker = get_broker() - proc_name = current_process().name - if setproctitle: - setproctitle.setproctitle(f"qcluster {proc_name} pusher") - logger.info( - _("%(name)s pushing tasks at %(id)s") - % {"name": proc_name, "id": current_process().pid} - ) - while True: - try: - task_set = broker.dequeue() - except Exception: - logger.exception("Failed to pull task from broker") - # broker probably crashed. Let the sentinel handle it. - sleep(10) - break - if task_set: - for task in task_set: - ack_id = task[0] - # unpack the task - try: - task = SignedPackage.loads(task[1]) - except (TypeError, BadSignature): - logger.exception("Failed to push task to queue") - broker.fail(ack_id) - continue - task["cluster"] = Conf.CLUSTER_NAME # save actual cluster name to orm task table - task["ack_id"] = ack_id - task_queue.put(task) - logger.debug( - _("queueing from %(list_key)s") % {"list_key": broker.list_key} - ) - if event.is_set(): - break - logger.info(_("%(name)s stopped pushing tasks") % {"name": current_process().name}) - - -def monitor(result_queue: Queue, broker: Broker = None): - """ - Gets finished tasks from the result queue and saves them to Django - :type broker: brokers.Broker - :type result_queue: multiprocessing.Queue - """ - if not broker: - broker = get_broker() - proc_name = current_process().name - if setproctitle: - setproctitle.setproctitle(f"qcluster {proc_name} monitor") - logger.info( - _("%(name)s monitoring at %(id)s") % {"name": proc_name, "id": current_process().pid} - ) - for task in iter(result_queue.get, "STOP"): - # save the result - if task.get("cached", False): - save_cached(task, broker) - else: - save_task(task, broker) - # acknowledge result - ack_id = task.pop("ack_id", False) - if ack_id and (task["success"] or task.get("ack_failure", False)): - broker.acknowledge(ack_id) - # signal execution done - post_execute.send(sender="django_q", task=task) - # log the result - info_name = get_func_repr(task["func"]) - if task["success"]: - # log success - logger.info( - _("Processed '%(info_name)s' (%(task_name)s)") - % {"info_name": info_name, "task_name": task["name"]} - ) - else: - # log failure - logger.error( - _("Failed '%(info_name)s' (%(task_name)s) - %(task_result)s") - % { - "info_name": info_name, - "task_name": task["name"], - "task_result": task["result"], - } - ) - logger.info(_("%(name)s stopped monitoring results") % {"name": proc_name}) - - -def worker( - task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT -): - """ - Takes a task from the task queue, tries to execute it and puts the result back in - the result queue - :param timeout: number of seconds wait for a worker to finish. - :type task_queue: multiprocessing.Queue - :type result_queue: multiprocessing.Queue - :type timer: multiprocessing.Value - """ - proc_name = current_process().name - logger.info( - _("%(proc_name)s ready for work at %(id)s") - % {"proc_name": proc_name, "id": current_process().pid} - ) - post_spawn.send(sender="django_q", proc_name=proc_name) - if setproctitle: - setproctitle.setproctitle(f"qcluster {proc_name} idle") - task_count = 0 - if timeout is None: - timeout = -1 - # Start reading the task queue - for task in iter(task_queue.get, "STOP"): - result = None - timer.value = -1 # Idle - task_count += 1 - f = task["func"] - - # Log task creation and set process name - # Get the function from the task - func_name = get_func_repr(f) - task_name = task["name"] - task_desc = ( - _("%(proc_name)s processing %(task_name)s '%(func_name)s'") - % { - "proc_name": proc_name, - "func_name": func_name, - "task_name": task_name, - } - ) - if "group" in task: - task_desc += f" [{task['group']}]" - logger.info(task_desc) - - if setproctitle: - proc_title = f"qcluster {proc_name} processing {task_name} '{func_name}'" - if "group" in task: - proc_title += f" [{task['group']}]" - setproctitle.setproctitle(proc_title) - - # if it's not an instance try to get it from the string - if not callable(f): - # locate() returns None if f cannot be loaded - f = pydoc.locate(f) - close_old_django_connections() - timer_value = task.pop("timeout", timeout) - # signal execution - pre_execute.send(sender="django_q", func=f, task=task) - # execute the payload - timer.value = timer_value # Busy - - try: - if f is None: - # raise a meaningfull error if task["func"] is not a valid function - raise ValueError(f"Function {task['func']} is not defined") - res = f(*task["args"], **task["kwargs"]) - result = (res, True) - except Exception as e: - result = (f"{e} : {traceback.format_exc()}", False) - if error_reporter: - error_reporter.report() - if task.get("sync", False): - raise - with timer.get_lock(): - # Process result - task["result"] = result[0] - task["success"] = result[1] - task["stopped"] = timezone.now() - result_queue.put(task) - timer.value = -1 # Idle - if setproctitle: - setproctitle.setproctitle(f"qcluster {proc_name} idle") - # Recycle - if task_count == Conf.RECYCLE or rss_check(): - timer.value = -2 # Recycled - break - logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name}) - - -def save_task(task, broker: Broker): - """ - Saves the task package to Django or the cache - :param task: the task package - :type broker: brokers.Broker - """ - # SAVE LIMIT < 0 : Don't save success - if not task.get("save", Conf.SAVE_LIMIT >= 0) and task["success"]: - return - # enqueues next in a chain - if task.get("chain", None): - 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 - close_old_django_connections() - - try: - filters = {} - if ( - Conf.SAVE_LIMIT_PER - and Conf.SAVE_LIMIT_PER in {"group", "name", "func"} - and Conf.SAVE_LIMIT_PER in task - ): - value = task[Conf.SAVE_LIMIT_PER] - if Conf.SAVE_LIMIT_PER == "func": - value = get_func_repr(value) - filters[Conf.SAVE_LIMIT_PER] = value - - with db.transaction.atomic(using=db.router.db_for_write(Success)): - last = Success.objects.filter(**filters).select_for_update().last() - if ( - task["success"] - and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count() - ): - last.delete() - - # check if this task has previous results - try: - existing_task = Task.objects.get(id=task["id"], name=task["name"]) - # only update the result if it hasn't succeeded yet - if not existing_task.success: - existing_task.stopped = task["stopped"] - existing_task.result = task["result"] - existing_task.success = task["success"] - existing_task.attempt_count = existing_task.attempt_count + 1 - existing_task.save() - - if ( - Conf.MAX_ATTEMPTS > 0 - and existing_task.attempt_count >= Conf.MAX_ATTEMPTS - ): - broker.acknowledge(task["ack_id"]) - - except Task.DoesNotExist: - # convert func to string - func = get_func_repr(task["func"]) - Task.objects.create( - id=task["id"], - name=task["name"], - func=func, - hook=task.get("hook"), - args=task["args"], - kwargs=task["kwargs"], - cluster=task.get("cluster"), - started=task["started"], - stopped=task["stopped"], - result=task["result"], - group=task.get("group"), - success=task["success"], - attempt_count=1, - ) - except Exception: - logger.exception("Could not save task result") - - -def save_cached(task, broker: Broker): - task_key = f'{broker.list_key}:{task["id"]}' - timeout = task["cached"] - if timeout is True: - timeout = None - try: - group = task.get("group", None) - iter_count = task.get("iter_count", 0) - # if it's a group append to the group list - if group: - group_key = f"{broker.list_key}:{group}:keys" - group_list = broker.cache.get(group_key) or [] - # if it's an iter group, check if we are ready - if iter_count and len(group_list) == iter_count - 1: - group_args = f"{broker.list_key}:{group}:args" - # collate the results into a Task result - 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"] = SignedPackage.loads(broker.cache.get(group_args)) - 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) - else: - save_task(task, broker) - broker.cache.delete_many(group_list) - broker.cache.delete_many([group_key, group_args]) - return - # save the group list - group_list.append(task_key) - broker.cache.set(group_key, group_list, timeout) - # async_task next in a chain - if task.get("chain", None): - django_q.tasks.async_chain( - task["chain"], - group=group, - cached=task["cached"], - sync=task["sync"], - broker=broker, - ) - # save the task - broker.cache.set(task_key, SignedPackage.dumps(task), timeout) - except Exception: - logger.exception("Could not save task result") - - -def scheduler(broker: Broker = None): - """ - Creates a task from a schedule at the scheduled time and schedules next run - """ - if not broker: - broker = get_broker() - close_old_django_connections() - try: - # Only default cluster will handler schedule with default(null) cluster - Q_default = db.models.Q(cluster__isnull=True) if Conf.CLUSTER_NAME == Conf.PREFIX else db.models.Q(pk__in=[]) - - with db.transaction.atomic(using=db.router.db_for_write(Schedule)): - for s in ( - Schedule.objects.select_for_update() - .exclude(repeats=0) - .filter(next_run__lt=timezone.now()) - .filter( - Q_default | db.models.Q(cluster=Conf.CLUSTER_NAME) - ) - ): - args = () - kwargs = {} - # get args, kwargs and hook - if s.kwargs: - try: - # first try the dict syntax - kwargs = ast.literal_eval(s.kwargs) - except (SyntaxError, ValueError): - # else use the kwargs syntax - try: - parsed_kwargs = ( - ast.parse(f"f({s.kwargs})").body[0].value.keywords - ) - kwargs = { - kwarg.arg: ast.literal_eval(kwarg.value) - for kwarg in parsed_kwargs - } - except (SyntaxError, ValueError): - kwargs = {} - if s.args: - args = ast.literal_eval(s.args) - # single value won't eval to tuple, so: - if type(args) != tuple: - args = (args,) - q_options = kwargs.get("q_options", {}) - if s.intended_date_kwarg: - kwargs[s.intended_date_kwarg] = s.next_run.isoformat() - if s.hook: - q_options["hook"] = s.hook - # set up the next run time - if s.schedule_type != s.ONCE: - next_run = s.next_run - while True: - next_run = s.calculate_next_run(next_run) - if Conf.CATCH_UP or next_run > localtime(): - break - - s.next_run = next_run - s.repeats += -1 - # send it to the cluster; any cluster name is allowed in multi-queue scenarios - # because `broker_name` is confusing, using `cluster` name is recommended and takes precedence - q_options["cluster"] = s.cluster or q_options.get("cluster", q_options.pop("broker_name", None)) - if q_options['cluster'] is None or q_options['cluster'] == Conf.CLUSTER_NAME: - q_options["broker"] = broker - q_options["group"] = q_options.get("group", s.name or s.id) - kwargs["q_options"] = q_options - s.task = django_q.tasks.async_task(s.func, *args, **kwargs) - # log it - if not s.task: - logger.error( - _( - "%(process_name)s failed to create a task from schedule " - "[%(schedule)s]" - ) - % { - "process_name": current_process().name, - "schedule": s.name or s.id, - } - ) - else: - logger.info( - _( - "%(process_name)s created task %(task_name)s from schedule " - "[%(schedule)s]" - ) - % { - "process_name": current_process().name, - "task_name": humanize(s.task), - "schedule": s.name or s.id, - } - ) - # default behavior is to delete a ONCE schedule - if s.schedule_type == s.ONCE: - if s.repeats < 0: - s.delete() - continue - # but not if it has a positive repeats - s.repeats = 0 - # save the schedule - s.save() - except Exception: - logger.exception("Could not create task from schedule") - - -def close_old_django_connections(): - """ - Close django connections unless running with sync=True. - """ - if Conf.SYNC: - logger.warning( - "Preserving django database connections because sync=True. Beware " - "that tasks are now injected in the calling context/transactions " - "which may result in unexpected behaviour." - ) - else: - db.close_old_connections() - - def set_cpu_affinity(n: int, process_ids: list, actual: bool = not Conf.TESTING): """ Sets the cpu affinity for the supplied processes. @@ -846,12 +411,3 @@ def set_cpu_affinity(n: int, process_ids: list, actual: bool = not Conf.TESTING) _("%(pid)s will use cpu %(affinity)s") % {"pid": pid, "affinity": affinity} ) - - -def rss_check(): - if Conf.MAX_RSS: - if resource: - return resource.getrusage(resource.RUSAGE_SELF).ru_maxrss >= Conf.MAX_RSS - elif psutil: - return psutil.Process().memory_info().rss >= Conf.MAX_RSS * 1024 - return False diff --git a/django_q/conf.py b/django_q/conf.py index a9877c7..ab7380e 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -44,15 +44,18 @@ class Conf: conf = {} _Q_CLUSTER_NAME = os.getenv("Q_CLUSTER_NAME") - if _Q_CLUSTER_NAME and _Q_CLUSTER_NAME != conf.get("name") and \ - _Q_CLUSTER_NAME != conf.get("cluster_name"): + if ( + _Q_CLUSTER_NAME + and _Q_CLUSTER_NAME != conf.get("name") + and _Q_CLUSTER_NAME != conf.get("cluster_name") + ): conf["cluster_name"] = _Q_CLUSTER_NAME alt_conf = conf.pop("ALT_CLUSTERS") if isinstance(alt_conf, dict): alt_conf = alt_conf.get(_Q_CLUSTER_NAME) if isinstance(alt_conf, dict): - alt_conf.pop('name', None) - alt_conf.pop('cluster_name', None) + alt_conf.pop("name", None) + alt_conf.pop("cluster_name", None) conf.update(alt_conf) # Redis server configuration . Follows standard redis keywords diff --git a/django_q/core_signing.py b/django_q/core_signing.py index c88f399..c55c4d1 100644 --- a/django_q/core_signing.py +++ b/django_q/core_signing.py @@ -37,7 +37,9 @@ def loads( """ # TimestampSigner.unsign() returns str but base64 and zlib compression # operate on bytes. - base64d = force_bytes(TimestampSigner(key=key, salt=salt).unsign(s, max_age=max_age)) + base64d = force_bytes( + TimestampSigner(key=key, salt=salt).unsign(s, max_age=max_age) + ) decompress = False if base64d[:1] == b".": # It's compressed; uncompress it first diff --git a/django_q/management/commands/qcluster.py b/django_q/management/commands/qcluster.py index 2330c01..bfedcca 100644 --- a/django_q/management/commands/qcluster.py +++ b/django_q/management/commands/qcluster.py @@ -1,8 +1,9 @@ +import os + from django.core.management.base import BaseCommand from django.utils.translation import gettext as _ from django_q.cluster import Cluster -import os class Command(BaseCommand): @@ -23,7 +24,7 @@ class Command(BaseCommand): dest="cluster_name", default=None, help="Set alternative cluster name instead of the name in Q_CLUSTER settings (for multi-queue setup). " - "On Linux you should set name through `Q_CLUSTER_NAME=cluster_name python manage.py qcluster` instead." + "On Linux you should set name through `Q_CLUSTER_NAME=cluster_name python manage.py qcluster` instead.", ) def handle(self, *args, **options): diff --git a/django_q/management/commands/qinfo.py b/django_q/management/commands/qinfo.py index 2efe131..5032ce7 100644 --- a/django_q/management/commands/qinfo.py +++ b/django_q/management/commands/qinfo.py @@ -3,7 +3,7 @@ from django.utils.translation import gettext as _ from django_q import VERSION from django_q.conf import Conf -from django_q.monitor import get_ids, info +from django_q.monitor_terminal import get_ids, info class Command(BaseCommand): diff --git a/django_q/management/commands/qmemory.py b/django_q/management/commands/qmemory.py index e7f84ce..7f7afbc 100644 --- a/django_q/management/commands/qmemory.py +++ b/django_q/management/commands/qmemory.py @@ -1,7 +1,7 @@ from django.core.management.base import BaseCommand from django.utils.translation import gettext as _ -from django_q.monitor import memory +from django_q.monitor_terminal import memory class Command(BaseCommand): diff --git a/django_q/management/commands/qmonitor.py b/django_q/management/commands/qmonitor.py index ca925a0..183b8a1 100644 --- a/django_q/management/commands/qmonitor.py +++ b/django_q/management/commands/qmonitor.py @@ -1,7 +1,7 @@ from django.core.management.base import BaseCommand from django.utils.translation import gettext as _ -from django_q.monitor import monitor +from django_q.monitor_terminal import monitor class Command(BaseCommand): diff --git a/django_q/migrations/0016_schedule_intended_date_kwarg.py b/django_q/migrations/0016_schedule_intended_date_kwarg.py index 234f1f9..6b3992e 100644 --- a/django_q/migrations/0016_schedule_intended_date_kwarg.py +++ b/django_q/migrations/0016_schedule_intended_date_kwarg.py @@ -1,6 +1,7 @@ # Generated by Django 4.1.2 on 2023-01-15 22:34 from django.db import migrations, models + import django_q.models diff --git a/django_q/models.py b/django_q/models.py index e754776..4305667 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -8,19 +8,19 @@ from django.db import models from django.template.defaultfilters import truncatechars from django.urls import reverse from django.utils import timezone -from django.utils.timezone import is_aware -from django.utils.html import format_html -from django.utils.translation import gettext_lazy as _ from django.utils.functional import cached_property +from django.utils.html import format_html +from django.utils.timezone import is_aware +from django.utils.translation import gettext_lazy as _ # External from picklefield import PickledObjectField from picklefield.fields import dbsafe_decode # Local -from django_q.conf import croniter, Conf +from django_q.conf import Conf, croniter from django_q.signing import SignedPackage -from django_q.utils import localtime, add_months, add_years +from django_q.utils import add_months, add_years, localtime from .utils import get_func_repr @@ -218,8 +218,11 @@ class Schedule(models.Model): ) task = models.CharField(max_length=100, null=True, editable=False) cluster = models.CharField( - max_length=100, default=None, null=True, blank=True, - help_text=_("Name of the target cluster") + max_length=100, + default=None, + null=True, + blank=True, + help_text=_("Name of the target cluster"), ) intended_date_kwarg = models.CharField( max_length=100, @@ -309,7 +312,9 @@ class Schedule(models.Model): class OrmQ(models.Model): key = models.CharField(max_length=100, help_text=_("Name of the target cluster")) payload = models.TextField() - lock = models.DateTimeField(null=True, help_text=_("Prevent any cluster from pulling until")) + lock = models.DateTimeField( + null=True, help_text=_("Prevent any cluster from pulling until") + ) @cached_property def task(self): diff --git a/django_q/monitor.py b/django_q/monitor.py index 1c89f6e..da7e3db 100644 --- a/django_q/monitor.py +++ b/django_q/monitor.py @@ -1,510 +1,198 @@ -from datetime import timedelta +from multiprocessing.process import current_process +from multiprocessing.queues import Queue -# django -from django.db import connection -from django.db.models import F, Sum -from django.utils import timezone -from django.utils.translation import gettext as _ +from django import db +from django.utils.translation import gettext_lazy as _ -from django_q import VERSION, models -from django_q.brokers import get_broker +import django_q.tasks +from django_q.brokers import Broker, get_broker +from django_q.conf import Conf, logger, setproctitle +from django_q.models import Success, Task +from django_q.signals import post_execute +from django_q.signing import SignedPackage +from django_q.utils import close_old_django_connections, get_func_repr -# local -from django_q.conf import Conf -from django_q.status import Stat - -# optional try: - import psutil -except ImportError: - psutil = None + import setproctitle +except ModuleNotFoundError: + setproctitle = None -def get_process_mb(pid): - try: - process = psutil.Process(pid) - mb_used = round(process.memory_info().rss / 1024**2, 2) - except psutil.NoSuchProcess: - mb_used = "NO_PROCESS_FOUND" - return mb_used - - -BLESSED_INSTALL_MESSAGE = ( - "Blessed is not installed. Please install blessed to use this: " - "https://pypi.org/project/blessed/" -) - - -def monitor(run_once=False, broker=None): +def monitor(result_queue: Queue, broker: Broker = None): + """ + Gets finished tasks from the result queue and saves them to Django + :type broker: brokers.Broker + :type result_queue: multiprocessing.Queue + """ if not broker: broker = get_broker() - try: - from blessed import Terminal + proc_name = current_process().name + if setproctitle: + setproctitle.setproctitle(f"qcluster {proc_name} monitor") + logger.info( + _("%(name)s monitoring at %(id)s") + % {"name": proc_name, "id": current_process().pid} + ) + for task in iter(result_queue.get, "STOP"): + # save the result + if task.get("cached", False): + save_cached(task, broker) + else: + save_task(task, broker) + # acknowledge result + ack_id = task.pop("ack_id", False) + if ack_id and (task["success"] or task.get("ack_failure", False)): + broker.acknowledge(ack_id) + # signal execution done + post_execute.send(sender="django_q", task=task) + # log the result + info_name = get_func_repr(task["func"]) + if task["success"]: + # log success + logger.info( + _("Processed '%(info_name)s' (%(task_name)s)") + % {"info_name": info_name, "task_name": task["name"]} + ) + else: + # log failure + logger.error( + _("Failed '%(info_name)s' (%(task_name)s) - %(task_result)s") + % { + "info_name": info_name, + "task_name": task["name"], + "task_result": task["result"], + } + ) + logger.info(_("%(name)s stopped monitoring results") % {"name": proc_name}) - term = Terminal() - except ImportError: - print(BLESSED_INSTALL_MESSAGE) + +def save_task(task, broker: Broker): + """ + Saves the task package to Django or the cache + :param task: the task package + :type broker: brokers.Broker + """ + # SAVE LIMIT < 0 : Don't save success + if not task.get("save", Conf.SAVE_LIMIT >= 0) and task["success"]: return + # enqueues next in a chain + if task.get("chain", None): + 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 + close_old_django_connections() - broker.ping() - with term.fullscreen(), term.hidden_cursor(), term.cbreak(): - val = None - start_width = int(term.width / 8) - while val not in ( - "q", - "Q", + try: + filters = {} + if ( + Conf.SAVE_LIMIT_PER + and Conf.SAVE_LIMIT_PER in {"group", "name", "func"} + and Conf.SAVE_LIMIT_PER in task ): - col_width = int(term.width / 8) - # In case of resize - if col_width != start_width: - print(term.clear()) - start_width = col_width - print( - term.move(0, 0) - + term.black_on_green(term.center(_("Host"), width=col_width - 1)) + value = task[Conf.SAVE_LIMIT_PER] + if Conf.SAVE_LIMIT_PER == "func": + value = get_func_repr(value) + filters[Conf.SAVE_LIMIT_PER] = value + + with db.transaction.atomic(using=db.router.db_for_write(Success)): + last = Success.objects.filter(**filters).select_for_update().last() + if ( + task["success"] + and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count() + ): + last.delete() + + # check if this task has previous results + try: + existing_task = Task.objects.get(id=task["id"], name=task["name"]) + # only update the result if it hasn't succeeded yet + if not existing_task.success: + existing_task.stopped = task["stopped"] + existing_task.result = task["result"] + existing_task.success = task["success"] + existing_task.attempt_count = existing_task.attempt_count + 1 + existing_task.save() + + if ( + Conf.MAX_ATTEMPTS > 0 + and existing_task.attempt_count >= Conf.MAX_ATTEMPTS + ): + broker.acknowledge(task["ack_id"]) + + except Task.DoesNotExist: + # convert func to string + func = get_func_repr(task["func"]) + Task.objects.create( + id=task["id"], + name=task["name"], + func=func, + hook=task.get("hook"), + args=task["args"], + kwargs=task["kwargs"], + cluster=task.get("cluster"), + started=task["started"], + stopped=task["stopped"], + result=task["result"], + group=task.get("group"), + success=task["success"], + attempt_count=1, ) - print( - term.move(0, 1 * col_width) - + term.black_on_green(term.center(_("Id"), 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 - stats = Stat.get_all(broker=broker) - print(term.clear_eos()) - for stat in stats: - status = stat.status - # color status - if stat.status == Conf.WORKING: - status = term.green(str(Conf.WORKING)) - elif stat.status == Conf.STOPPING: - status = term.yellow(str(Conf.STOPPING)) - elif stat.status == Conf.STOPPED: - status = term.red(str(Conf.STOPPED)) - elif stat.status == Conf.IDLE: - status = str(Conf.IDLE) - # color q's - tasks = str(stat.task_q_size) - if stat.task_q_size > 0: - tasks = term.cyan(str(stat.task_q_size)) - if Conf.QUEUE_LIMIT and stat.task_q_size == Conf.QUEUE_LIMIT: - tasks = term.green(str(stat.task_q_size)) - results = stat.done_q_size - if results > 0: - results = term.cyan(str(results)) - # color workers - workers = len(stat.workers) - if workers < Conf.WORKERS: - workers = term.yellow(str(workers)) - # format uptime - uptime = (timezone.now() - stat.tob).total_seconds() - hours, remainder = divmod(uptime, 3600) - minutes, seconds = divmod(remainder, 60) - uptime = "%d:%02d:%02d" % (hours, minutes, seconds) - # print to the terminal - print( - term.move(i, 0) - + term.center(stat.host[: col_width - 1], width=col_width - 1) - ) - print( - term.move(i, 1 * col_width) - + term.center(str(stat.cluster_id)[-8:], 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 - # bottom bar - i += 1 - queue_size = broker.queue_size() - lock_size = broker.lock_size() - if 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( - term.move(i, 2 * col_width) - + term.black_on_cyan(term.center(_("Queued"), 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 - if run_once: - return Stat.get_all(broker=broker) - print(term.move(i + 2, 0) + term.center(_("[Press q to quit]"))) - val = term.inkey(timeout=1) + except Exception: + logger.exception("Could not save task result") -def info(broker=None): - if not broker: - broker = get_broker() +def save_cached(task, broker: Broker): + task_key = f'{broker.list_key}:{task["id"]}' + timeout = task["cached"] + if timeout is True: + timeout = None try: - from blessed import Terminal - - term = Terminal() - except ImportError: - print(BLESSED_INSTALL_MESSAGE) - return - - broker.ping() - stat = Stat.get_all(broker=broker) - # general stats - clusters = len(stat) - workers = 0 - reincarnations = 0 - for cluster in stat: - workers += len(cluster.workers) - reincarnations += cluster.reincarnations - # calculate tasks pm and avg exec time - tasks_per = 0 - per = _("day") - exec_time = 0 - last_tasks = models.Success.objects.filter( - stopped__gte=timezone.now() - timedelta(hours=24) - ) - tasks_per_day = last_tasks.count() - if tasks_per_day > 0: - # average execution time over the last 24 hours - if connection.vendor != "sqlite": - exec_time = last_tasks.aggregate( - time_taken=Sum(F("stopped") - F("started")) - ) - exec_time = exec_time["time_taken"].total_seconds() / tasks_per_day - else: - # can't sum timedeltas on sqlite - for t in last_tasks: - exec_time += t.time_taken() - exec_time = exec_time / tasks_per_day - # tasks per second/minute/hour/day in the last 24 hours - if tasks_per_day > 24 * 60 * 60: - tasks_per = tasks_per_day / (24 * 60 * 60) - per = _("second") - elif tasks_per_day > 24 * 60: - tasks_per = tasks_per_day / (24 * 60) - per = _("minute") - elif tasks_per_day > 24: - tasks_per = tasks_per_day / 24 - per = _("hour") - else: - tasks_per = tasks_per_day - # print to terminal - print(term.clear_eos()) - col_width = int(term.width / 6) - print( - term.black_on_green( - term.center( - _("-- %(prefix)s %(version)s on %(info)s --") - % { - "prefix": Conf.PREFIX.capitalize(), - "version": ".".join(str(v) for v in VERSION), - "info": broker.info(), - } - ) - ) - ) - print( - term.cyan(_("Clusters")) - + 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) - + term.white(str(workers)) - + term.move_x(4 * col_width) - + term.cyan(_("Restarts")) - + term.move_x(5 * col_width) - + term.white(str(reincarnations)) - ) - print( - term.cyan(_("Queued")) - + term.move_x(1 * col_width) - + term.white(str(broker.queue_size())) - + term.move_x(2 * col_width) - + term.cyan(_("Successes")) - + term.move_x(3 * col_width) - + term.white(str(models.Success.objects.count())) - + term.move_x(4 * col_width) - + term.cyan(_("Failures")) - + term.move_x(5 * col_width) - + term.white(str(models.Failure.objects.count())) - ) - print( - term.cyan(_("Schedules")) - + term.move_x(1 * col_width) - + term.white(str(models.Schedule.objects.count())) - + term.move_x(2 * col_width) - + term.cyan(_("Tasks/%(per)s") % {"per": 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 - - -def memory(run_once=False, workers=False, broker=None): - if not broker: - broker = get_broker() - try: - from blessed import Terminal - - term = Terminal() - except ImportError: - print(BLESSED_INSTALL_MESSAGE) - return - broker.ping() - if not psutil: - print(term.clear_eos()) - print( - term.white_on_red( - 'Cannot start "qmemory" command. Missing "psutil" library.' - ) - ) - return - with term.fullscreen(), term.hidden_cursor(), term.cbreak(): - MEMORY_AVAILABLE_LOWEST_PERCENTAGE = 100.0 - MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT = timezone.now() - cols = 8 - val = None - start_width = int(term.width / cols) - while val not in ["q", "Q"]: - col_width = int(term.width / cols) - # In case of resize - if col_width != start_width: - print(term.clear()) - start_width = col_width - # sentinel, monitor and workers memory usage - print( - term.move(0, 0 * col_width) - + term.black_on_green(term.center(_("Host"), width=col_width - 1)) - ) - print( - term.move(0, 1 * col_width) - + term.black_on_green(term.center(_("Id"), width=col_width - 1)) - ) - print( - term.move(0, 2 * col_width) - + term.black_on_green( - term.center(_("Available (%)"), width=col_width - 1) + group = task.get("group", None) + iter_count = task.get("iter_count", 0) + # if it's a group append to the group list + if group: + group_key = f"{broker.list_key}:{group}:keys" + group_list = broker.cache.get(group_key) or [] + # if it's an iter group, check if we are ready + if iter_count and len(group_list) == iter_count - 1: + group_args = f"{broker.list_key}:{group}:args" + # collate the results into a Task result + 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"] = SignedPackage.loads(broker.cache.get(group_args)) + 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) + else: + save_task(task, broker) + broker.cache.delete_many(group_list) + broker.cache.delete_many([group_key, group_args]) + return + # save the group list + group_list.append(task_key) + broker.cache.set(group_key, group_list, timeout) + # async_task next in a chain + if task.get("chain", None): + django_q.tasks.async_chain( + task["chain"], + group=group, + cached=task["cached"], + sync=task["sync"], + broker=broker, ) - ) - print( - term.move(0, 3 * col_width) - + term.black_on_green( - term.center(_("Available (MB)"), width=col_width - 1) - ) - ) - print( - term.move(0, 4 * col_width) - + term.black_on_green(term.center(_("Total (MB)"), width=col_width - 1)) - ) - print( - term.move(0, 5 * col_width) - + term.black_on_green( - term.center(_("Sentinel (MB)"), width=col_width - 1) - ) - ) - print( - term.move(0, 6 * col_width) - + term.black_on_green( - term.center(_("Monitor (MB)"), width=col_width - 1) - ) - ) - print( - term.move(0, 7 * col_width) - + term.black_on_green( - term.center(_("Workers (MB)"), width=col_width - 1) - ) - ) - row = 2 - stats = Stat.get_all(broker=broker) - print(term.clear_eos()) - for stat in stats: - # memory available (%) - memory_available_percentage = round( - psutil.virtual_memory().available - * 100 - / psutil.virtual_memory().total, - 2, - ) - # memory available (MB) - memory_available = round( - psutil.virtual_memory().available / 1024**2, 2 - ) - if memory_available_percentage < MEMORY_AVAILABLE_LOWEST_PERCENTAGE: - MEMORY_AVAILABLE_LOWEST_PERCENTAGE = memory_available_percentage - MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT = timezone.now() - print( - term.move(row, 0 * col_width) - + term.center(stat.host[: col_width - 1], width=col_width - 1) - ) - print( - term.move(row, 1 * col_width) - + term.center(str(stat.cluster_id)[-8:], width=col_width - 1) - ) - print( - term.move(row, 2 * col_width) - + term.center(memory_available_percentage, width=col_width - 1) - ) - print( - term.move(row, 3 * col_width) - + term.center(memory_available, width=col_width - 1) - ) - print( - term.move(row, 4 * col_width) - + term.center( - round(psutil.virtual_memory().total / 1024**2, 2), - width=col_width - 1, - ) - ) - print( - term.move(row, 5 * col_width) - + term.center(get_process_mb(stat.sentinel), width=col_width - 1) - ) - print( - term.move(row, 6 * col_width) - + term.center( - get_process_mb(getattr(stat, "monitor", None)), - width=col_width - 1, - ) - ) - workers_mb = 0 - for worker_pid in stat.workers: - result = get_process_mb(worker_pid) - if isinstance(result, str): - result = 0 - workers_mb += result - print( - term.move(row, 7 * col_width) - + term.center( - workers_mb or "NO_PROCESSES_FOUND", width=col_width - 1 - ) - ) - row += 1 - # each worker's memory usage - if workers: - row += 2 - col_width = int(term.width / (1 + Conf.WORKERS)) - print( - term.move(row, 0 * col_width) - + term.black_on_cyan(term.center(_("Id"), width=col_width - 1)) - ) - for worker_num in range(Conf.WORKERS): - print( - term.move(row, (worker_num + 1) * col_width) - + term.black_on_cyan( - term.center( - "Worker #{} (MB)".format(worker_num + 1), - width=col_width - 1, - ) - ) - ) - row += 2 - for stat in stats: - print( - term.move(row, 0 * col_width) - + term.center(str(stat.cluster_id)[-8:], width=col_width - 1) - ) - for idx, worker_pid in enumerate(stat.workers): - mb_used = get_process_mb(worker_pid) - print( - term.move(row, (idx + 1) * col_width) - + term.center(mb_used, width=col_width - 1) - ) - row += 1 - row += 1 - print( - term.move(row, 0) - + _("Available lowest (): %(memory_percent)s ((at)s)") - % { - "memory_percent": str(MEMORY_AVAILABLE_LOWEST_PERCENTAGE), - "at": MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT.strftime( - "%Y-%m-%d %H:%M:%S+00:00" - ), - } - ) - # for testing - if run_once: - return Stat.get_all(broker=broker) - print(term.move(row + 2, 0) + term.center(_("[Press q to quit]"))) - val = term.inkey(timeout=1) - - -def get_ids(): - # prints id (PID) of running clusters - stat = Stat.get_all() - if stat: - for s in stat: - print(s.cluster_id) - else: - print(_("No clusters appear to be running.")) - return True + # save the task + broker.cache.set(task_key, SignedPackage.dumps(task), timeout) + except Exception: + logger.exception("Could not save task result") diff --git a/django_q/monitor_terminal.py b/django_q/monitor_terminal.py new file mode 100644 index 0000000..1c89f6e --- /dev/null +++ b/django_q/monitor_terminal.py @@ -0,0 +1,510 @@ +from datetime import timedelta + +# django +from django.db import connection +from django.db.models import F, Sum +from django.utils import timezone +from django.utils.translation import gettext as _ + +from django_q import VERSION, models +from django_q.brokers import get_broker + +# local +from django_q.conf import Conf +from django_q.status import Stat + +# optional +try: + import psutil +except ImportError: + psutil = None + + +def get_process_mb(pid): + try: + process = psutil.Process(pid) + mb_used = round(process.memory_info().rss / 1024**2, 2) + except psutil.NoSuchProcess: + mb_used = "NO_PROCESS_FOUND" + return mb_used + + +BLESSED_INSTALL_MESSAGE = ( + "Blessed is not installed. Please install blessed to use this: " + "https://pypi.org/project/blessed/" +) + + +def monitor(run_once=False, broker=None): + if not broker: + broker = get_broker() + try: + from blessed import Terminal + + term = Terminal() + except ImportError: + print(BLESSED_INSTALL_MESSAGE) + return + + broker.ping() + with term.fullscreen(), term.hidden_cursor(), term.cbreak(): + val = None + start_width = int(term.width / 8) + while val not in ( + "q", + "Q", + ): + col_width = int(term.width / 8) + # In case of resize + if col_width != start_width: + print(term.clear()) + start_width = col_width + print( + term.move(0, 0) + + term.black_on_green(term.center(_("Host"), width=col_width - 1)) + ) + print( + term.move(0, 1 * col_width) + + term.black_on_green(term.center(_("Id"), 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 + stats = Stat.get_all(broker=broker) + print(term.clear_eos()) + for stat in stats: + status = stat.status + # color status + if stat.status == Conf.WORKING: + status = term.green(str(Conf.WORKING)) + elif stat.status == Conf.STOPPING: + status = term.yellow(str(Conf.STOPPING)) + elif stat.status == Conf.STOPPED: + status = term.red(str(Conf.STOPPED)) + elif stat.status == Conf.IDLE: + status = str(Conf.IDLE) + # color q's + tasks = str(stat.task_q_size) + if stat.task_q_size > 0: + tasks = term.cyan(str(stat.task_q_size)) + if Conf.QUEUE_LIMIT and stat.task_q_size == Conf.QUEUE_LIMIT: + tasks = term.green(str(stat.task_q_size)) + results = stat.done_q_size + if results > 0: + results = term.cyan(str(results)) + # color workers + workers = len(stat.workers) + if workers < Conf.WORKERS: + workers = term.yellow(str(workers)) + # format uptime + uptime = (timezone.now() - stat.tob).total_seconds() + hours, remainder = divmod(uptime, 3600) + minutes, seconds = divmod(remainder, 60) + uptime = "%d:%02d:%02d" % (hours, minutes, seconds) + # print to the terminal + print( + term.move(i, 0) + + term.center(stat.host[: col_width - 1], width=col_width - 1) + ) + print( + term.move(i, 1 * col_width) + + term.center(str(stat.cluster_id)[-8:], 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 + # bottom bar + i += 1 + queue_size = broker.queue_size() + lock_size = broker.lock_size() + if 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( + term.move(i, 2 * col_width) + + term.black_on_cyan(term.center(_("Queued"), 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 + if run_once: + return Stat.get_all(broker=broker) + print(term.move(i + 2, 0) + term.center(_("[Press q to quit]"))) + val = term.inkey(timeout=1) + + +def info(broker=None): + if not broker: + broker = get_broker() + try: + from blessed import Terminal + + term = Terminal() + except ImportError: + print(BLESSED_INSTALL_MESSAGE) + return + + broker.ping() + stat = Stat.get_all(broker=broker) + # general stats + clusters = len(stat) + workers = 0 + reincarnations = 0 + for cluster in stat: + workers += len(cluster.workers) + reincarnations += cluster.reincarnations + # calculate tasks pm and avg exec time + tasks_per = 0 + per = _("day") + exec_time = 0 + last_tasks = models.Success.objects.filter( + stopped__gte=timezone.now() - timedelta(hours=24) + ) + tasks_per_day = last_tasks.count() + if tasks_per_day > 0: + # average execution time over the last 24 hours + if connection.vendor != "sqlite": + exec_time = last_tasks.aggregate( + time_taken=Sum(F("stopped") - F("started")) + ) + exec_time = exec_time["time_taken"].total_seconds() / tasks_per_day + else: + # can't sum timedeltas on sqlite + for t in last_tasks: + exec_time += t.time_taken() + exec_time = exec_time / tasks_per_day + # tasks per second/minute/hour/day in the last 24 hours + if tasks_per_day > 24 * 60 * 60: + tasks_per = tasks_per_day / (24 * 60 * 60) + per = _("second") + elif tasks_per_day > 24 * 60: + tasks_per = tasks_per_day / (24 * 60) + per = _("minute") + elif tasks_per_day > 24: + tasks_per = tasks_per_day / 24 + per = _("hour") + else: + tasks_per = tasks_per_day + # print to terminal + print(term.clear_eos()) + col_width = int(term.width / 6) + print( + term.black_on_green( + term.center( + _("-- %(prefix)s %(version)s on %(info)s --") + % { + "prefix": Conf.PREFIX.capitalize(), + "version": ".".join(str(v) for v in VERSION), + "info": broker.info(), + } + ) + ) + ) + print( + term.cyan(_("Clusters")) + + 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) + + term.white(str(workers)) + + term.move_x(4 * col_width) + + term.cyan(_("Restarts")) + + term.move_x(5 * col_width) + + term.white(str(reincarnations)) + ) + print( + term.cyan(_("Queued")) + + term.move_x(1 * col_width) + + term.white(str(broker.queue_size())) + + term.move_x(2 * col_width) + + term.cyan(_("Successes")) + + term.move_x(3 * col_width) + + term.white(str(models.Success.objects.count())) + + term.move_x(4 * col_width) + + term.cyan(_("Failures")) + + term.move_x(5 * col_width) + + term.white(str(models.Failure.objects.count())) + ) + print( + term.cyan(_("Schedules")) + + term.move_x(1 * col_width) + + term.white(str(models.Schedule.objects.count())) + + term.move_x(2 * col_width) + + term.cyan(_("Tasks/%(per)s") % {"per": 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 + + +def memory(run_once=False, workers=False, broker=None): + if not broker: + broker = get_broker() + try: + from blessed import Terminal + + term = Terminal() + except ImportError: + print(BLESSED_INSTALL_MESSAGE) + return + broker.ping() + if not psutil: + print(term.clear_eos()) + print( + term.white_on_red( + 'Cannot start "qmemory" command. Missing "psutil" library.' + ) + ) + return + with term.fullscreen(), term.hidden_cursor(), term.cbreak(): + MEMORY_AVAILABLE_LOWEST_PERCENTAGE = 100.0 + MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT = timezone.now() + cols = 8 + val = None + start_width = int(term.width / cols) + while val not in ["q", "Q"]: + col_width = int(term.width / cols) + # In case of resize + if col_width != start_width: + print(term.clear()) + start_width = col_width + # sentinel, monitor and workers memory usage + print( + term.move(0, 0 * col_width) + + term.black_on_green(term.center(_("Host"), width=col_width - 1)) + ) + print( + term.move(0, 1 * col_width) + + term.black_on_green(term.center(_("Id"), width=col_width - 1)) + ) + print( + term.move(0, 2 * col_width) + + term.black_on_green( + term.center(_("Available (%)"), width=col_width - 1) + ) + ) + print( + term.move(0, 3 * col_width) + + term.black_on_green( + term.center(_("Available (MB)"), width=col_width - 1) + ) + ) + print( + term.move(0, 4 * col_width) + + term.black_on_green(term.center(_("Total (MB)"), width=col_width - 1)) + ) + print( + term.move(0, 5 * col_width) + + term.black_on_green( + term.center(_("Sentinel (MB)"), width=col_width - 1) + ) + ) + print( + term.move(0, 6 * col_width) + + term.black_on_green( + term.center(_("Monitor (MB)"), width=col_width - 1) + ) + ) + print( + term.move(0, 7 * col_width) + + term.black_on_green( + term.center(_("Workers (MB)"), width=col_width - 1) + ) + ) + row = 2 + stats = Stat.get_all(broker=broker) + print(term.clear_eos()) + for stat in stats: + # memory available (%) + memory_available_percentage = round( + psutil.virtual_memory().available + * 100 + / psutil.virtual_memory().total, + 2, + ) + # memory available (MB) + memory_available = round( + psutil.virtual_memory().available / 1024**2, 2 + ) + if memory_available_percentage < MEMORY_AVAILABLE_LOWEST_PERCENTAGE: + MEMORY_AVAILABLE_LOWEST_PERCENTAGE = memory_available_percentage + MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT = timezone.now() + print( + term.move(row, 0 * col_width) + + term.center(stat.host[: col_width - 1], width=col_width - 1) + ) + print( + term.move(row, 1 * col_width) + + term.center(str(stat.cluster_id)[-8:], width=col_width - 1) + ) + print( + term.move(row, 2 * col_width) + + term.center(memory_available_percentage, width=col_width - 1) + ) + print( + term.move(row, 3 * col_width) + + term.center(memory_available, width=col_width - 1) + ) + print( + term.move(row, 4 * col_width) + + term.center( + round(psutil.virtual_memory().total / 1024**2, 2), + width=col_width - 1, + ) + ) + print( + term.move(row, 5 * col_width) + + term.center(get_process_mb(stat.sentinel), width=col_width - 1) + ) + print( + term.move(row, 6 * col_width) + + term.center( + get_process_mb(getattr(stat, "monitor", None)), + width=col_width - 1, + ) + ) + workers_mb = 0 + for worker_pid in stat.workers: + result = get_process_mb(worker_pid) + if isinstance(result, str): + result = 0 + workers_mb += result + print( + term.move(row, 7 * col_width) + + term.center( + workers_mb or "NO_PROCESSES_FOUND", width=col_width - 1 + ) + ) + row += 1 + # each worker's memory usage + if workers: + row += 2 + col_width = int(term.width / (1 + Conf.WORKERS)) + print( + term.move(row, 0 * col_width) + + term.black_on_cyan(term.center(_("Id"), width=col_width - 1)) + ) + for worker_num in range(Conf.WORKERS): + print( + term.move(row, (worker_num + 1) * col_width) + + term.black_on_cyan( + term.center( + "Worker #{} (MB)".format(worker_num + 1), + width=col_width - 1, + ) + ) + ) + row += 2 + for stat in stats: + print( + term.move(row, 0 * col_width) + + term.center(str(stat.cluster_id)[-8:], width=col_width - 1) + ) + for idx, worker_pid in enumerate(stat.workers): + mb_used = get_process_mb(worker_pid) + print( + term.move(row, (idx + 1) * col_width) + + term.center(mb_used, width=col_width - 1) + ) + row += 1 + row += 1 + print( + term.move(row, 0) + + _("Available lowest (): %(memory_percent)s ((at)s)") + % { + "memory_percent": str(MEMORY_AVAILABLE_LOWEST_PERCENTAGE), + "at": MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT.strftime( + "%Y-%m-%d %H:%M:%S+00:00" + ), + } + ) + # for testing + if run_once: + return Stat.get_all(broker=broker) + print(term.move(row + 2, 0) + term.center(_("[Press q to quit]"))) + val = term.inkey(timeout=1) + + +def get_ids(): + # prints id (PID) of running clusters + stat = Stat.get_all() + if stat: + for s in stat: + print(s.cluster_id) + else: + print(_("No clusters appear to be running.")) + return True diff --git a/django_q/pusher.py b/django_q/pusher.py new file mode 100644 index 0000000..503274e --- /dev/null +++ b/django_q/pusher.py @@ -0,0 +1,62 @@ +from multiprocessing import Event +from multiprocessing.process import current_process +from multiprocessing.queues import Queue +from time import sleep + +from django.utils.translation import gettext_lazy as _ + +from django_q.brokers import Broker, get_broker +from django_q.conf import Conf, logger +from django_q.signing import BadSignature, SignedPackage + +try: + import setproctitle +except ModuleNotFoundError: + setproctitle = None + + +def pusher(task_queue: Queue, event: Event, broker: Broker = None): + """ + Pulls tasks of the broker and puts them in the task queue + :type broker: + :type task_queue: multiprocessing.Queue + :type event: multiprocessing.Event + """ + if not broker: + broker = get_broker() + proc_name = current_process().name + if setproctitle: + setproctitle.setproctitle(f"qcluster {proc_name} pusher") + logger.info( + _("%(name)s pushing tasks at %(id)s") + % {"name": proc_name, "id": current_process().pid} + ) + while True: + try: + task_set = broker.dequeue() + except Exception: + logger.exception("Failed to pull task from broker") + # broker probably crashed. Let the sentinel handle it. + sleep(10) + break + if task_set: + for task in task_set: + ack_id = task[0] + # unpack the task + try: + task = SignedPackage.loads(task[1]) + except (TypeError, BadSignature): + logger.exception("Failed to push task to queue") + broker.fail(ack_id) + continue + task[ + "cluster" + ] = Conf.CLUSTER_NAME # save actual cluster name to orm task table + task["ack_id"] = ack_id + task_queue.put(task) + logger.debug( + _("queueing from %(list_key)s") % {"list_key": broker.list_key} + ) + if event.is_set(): + break + logger.info(_("%(name)s stopped pushing tasks") % {"name": current_process().name}) diff --git a/django_q/scheduler.py b/django_q/scheduler.py new file mode 100644 index 0000000..10aaa98 --- /dev/null +++ b/django_q/scheduler.py @@ -0,0 +1,124 @@ +import ast +from multiprocessing.process import current_process + +from django import db +from django.utils import timezone +from django.utils.translation import gettext_lazy as _ + +import django_q.tasks +from django_q.brokers import Broker, get_broker +from django_q.conf import Conf, logger +from django_q.humanhash import humanize +from django_q.models import Schedule +from django_q.utils import close_old_django_connections, localtime + + +def scheduler(broker: Broker = None): + """ + Creates a task from a schedule at the scheduled time and schedules next run + """ + if not broker: + broker = get_broker() + close_old_django_connections() + try: + # Only default cluster will handler schedule with default(null) cluster + Q_default = ( + db.models.Q(cluster__isnull=True) + if Conf.CLUSTER_NAME == Conf.PREFIX + else db.models.Q(pk__in=[]) + ) + + with db.transaction.atomic(using=db.router.db_for_write(Schedule)): + for s in ( + Schedule.objects.select_for_update() + .exclude(repeats=0) + .filter(next_run__lt=timezone.now()) + .filter(Q_default | db.models.Q(cluster=Conf.CLUSTER_NAME)) + ): + args = () + kwargs = {} + # get args, kwargs and hook + if s.kwargs: + try: + # first try the dict syntax + kwargs = ast.literal_eval(s.kwargs) + except (SyntaxError, ValueError): + # else use the kwargs syntax + try: + parsed_kwargs = ( + ast.parse(f"f({s.kwargs})").body[0].value.keywords + ) + kwargs = { + kwarg.arg: ast.literal_eval(kwarg.value) + for kwarg in parsed_kwargs + } + except (SyntaxError, ValueError): + kwargs = {} + if s.args: + args = ast.literal_eval(s.args) + # single value won't eval to tuple, so: + if type(args) != tuple: + args = (args,) + q_options = kwargs.get("q_options", {}) + if s.intended_date_kwarg: + kwargs[s.intended_date_kwarg] = s.next_run.isoformat() + if s.hook: + q_options["hook"] = s.hook + # set up the next run time + if s.schedule_type != s.ONCE: + next_run = s.next_run + while True: + next_run = s.calculate_next_run(next_run) + if Conf.CATCH_UP or next_run > localtime(): + break + + s.next_run = next_run + s.repeats += -1 + # send it to the cluster; any cluster name is allowed in multi-queue scenarios + # because `broker_name` is confusing, using `cluster` name is recommended and takes precedence + q_options["cluster"] = s.cluster or q_options.get( + "cluster", q_options.pop("broker_name", None) + ) + if ( + q_options["cluster"] is None + or q_options["cluster"] == Conf.CLUSTER_NAME + ): + q_options["broker"] = broker + q_options["group"] = q_options.get("group", s.name or s.id) + kwargs["q_options"] = q_options + s.task = django_q.tasks.async_task(s.func, *args, **kwargs) + # log it + if not s.task: + logger.error( + _( + "%(process_name)s failed to create a task from schedule " + "[%(schedule)s]" + ) + % { + "process_name": current_process().name, + "schedule": s.name or s.id, + } + ) + else: + logger.info( + _( + "%(process_name)s created task %(task_name)s from schedule " + "[%(schedule)s]" + ) + % { + "process_name": current_process().name, + "task_name": humanize(s.task), + "schedule": s.name or s.id, + } + ) + # default behavior is to delete a ONCE schedule + if s.schedule_type == s.ONCE: + if s.repeats < 0: + s.delete() + continue + # but not if it has a positive repeats + s.repeats = 0 + # save the schedule + s.save() + except Exception: + logger.exception("Could not create task from schedule") diff --git a/django_q/signals.py b/django_q/signals.py index bb171a3..1476089 100644 --- a/django_q/signals.py +++ b/django_q/signals.py @@ -31,6 +31,7 @@ def call_hook(sender, instance, **kwargs): % {"hook": instance.hook, "name": instance.name, "error": str(e)} ) + # args: proc_name post_spawn = Signal() diff --git a/django_q/tasks.py b/django_q/tasks.py index 5ffdd94..e256429 100644 --- a/django_q/tasks.py +++ b/django_q/tasks.py @@ -763,7 +763,8 @@ class AsyncTask: def _sync(pack): """Simulate a package travelling through the cluster.""" - from django_q.cluster import monitor, worker + from django_q.monitor import monitor + from django_q.worker import worker task_queue = Queue() result_queue = Queue() diff --git a/django_q/tests/test_cached.py b/django_q/tests/test_cached.py index ad480c6..092a7ee 100644 --- a/django_q/tests/test_cached.py +++ b/django_q/tests/test_cached.py @@ -3,8 +3,9 @@ from multiprocessing import Event, Value import pytest from django_q.brokers import get_broker -from django_q.cluster import monitor, pusher, worker from django_q.conf import Conf +from django_q.monitor import monitor +from django_q.pusher import pusher from django_q.queues import Queue from django_q.tasks import ( AsyncTask, @@ -21,6 +22,7 @@ from django_q.tasks import ( result, result_group, ) +from django_q.worker import worker @pytest.fixture diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index da240e8..7ea7e4c 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -12,10 +12,12 @@ import pytest from django.utils import timezone from django_q.brokers import Broker, get_broker -from django_q.cluster import Cluster, Sentinel, monitor, pusher, save_task, worker +from django_q.cluster import Cluster, Sentinel from django_q.conf import Conf from django_q.humanhash import DEFAULT_WORDLIST, uuid from django_q.models import Success, Task +from django_q.monitor import monitor, save_task +from django_q.pusher import pusher from django_q.queues import Queue from django_q.signals import post_execute, pre_enqueue, pre_execute from django_q.status import Stat @@ -29,8 +31,9 @@ from django_q.tasks import ( result, result_group, ) -from django_q.tests.tasks import multiply, TaskError +from django_q.tests.tasks import TaskError, multiply from django_q.utils import add_months, add_years +from django_q.worker import worker myPath = os.path.dirname(os.path.abspath(__file__)) sys.path.insert(0, myPath + "/../") diff --git a/django_q/tests/test_monitor.py b/django_q/tests/test_monitor.py index a7a7980..10a81f5 100644 --- a/django_q/tests/test_monitor.py +++ b/django_q/tests/test_monitor.py @@ -5,7 +5,7 @@ import pytest from django_q.brokers import get_broker from django_q.cluster import Cluster from django_q.conf import Conf -from django_q.monitor import get_ids, info, monitor +from django_q.monitor_terminal import get_ids, info, monitor from django_q.status import Stat from django_q.tasks import async_task diff --git a/django_q/tests/test_scheduler.py b/django_q/tests/test_scheduler.py index cb6546d..edecf53 100644 --- a/django_q/tests/test_scheduler.py +++ b/django_q/tests/test_scheduler.py @@ -3,8 +3,8 @@ from datetime import datetime, timedelta from multiprocessing import Event, Value from unittest import mock -import pytest import django +import pytest from django.core.exceptions import ValidationError from django.db import IntegrityError from django.test import override_settings @@ -12,9 +12,11 @@ from django.utils import timezone from django.utils.timezone import is_naive from django_q.brokers import Broker, get_broker -from django_q.cluster import localtime, monitor, pusher, scheduler, worker from django_q.conf import Conf +from django_q.monitor import monitor +from django_q.pusher import pusher from django_q.queues import Queue +from django_q.scheduler import scheduler from django_q.tasks import Schedule, fetch from django_q.tasks import schedule as create_schedule from django_q.tests.settings import BASE_DIR @@ -22,7 +24,8 @@ from django_q.tests.testing_utilities.multiple_database_routers import ( TestingMultipleAppsDatabaseRouter, TestingReplicaDatabaseRouter, ) -from django_q.utils import add_months +from django_q.utils import add_months, localtime +from django_q.worker import worker if django.VERSION < (4, 0): # pytz is the default in django 3.2. Remove when no support for 3.2 @@ -85,7 +88,7 @@ def test_scheduler_daylight_saving_time_daily(broker, monkeypatch): # 28th of March 2021 is the day when sunlight saving starts (at 2 am) monkeypatch.setattr(Conf, "TIME_ZONE", "Europe/Amsterdam") - tz = ZoneInfo('Europe/Amsterdam') + tz = ZoneInfo("Europe/Amsterdam") broker.list_key = "scheduler_test:q" # Let's start a schedule at 1 am on the 27th of March. This is in AMS timezone. # So, 2021-03-27 00:00:00 when saved (due to TZ being Amsterdam and saved in UTC) @@ -181,7 +184,6 @@ def test_scheduler_daylight_saving_time_daily(broker, monkeypatch): assert str(next_run) == "2021-11-01 01:00:00+01:00" - @pytest.mark.django_db def test_scheduler(broker, monkeypatch): broker.list_key = "scheduler_test:q" @@ -407,7 +409,7 @@ def test_scheduler(broker, monkeypatch): def test_intended_schedule_kwarg(broker, monkeypatch): broker.list_key = "scheduler_test:q" broker.delete_queue() - run_date = timezone.now()-timedelta(hours=1) + run_date = timezone.now() - timedelta(hours=1) schedule = create_schedule( "math.copysign", 1, @@ -417,10 +419,10 @@ def test_intended_schedule_kwarg(broker, monkeypatch): schedule_type=Schedule.HOURLY, repeats=1, next_run=run_date, - intended_date_kwarg='intended_date', + intended_date_kwarg="intended_date", ) assert schedule.last_run() is None - assert schedule.intended_date_kwarg == 'intended_date' + assert schedule.intended_date_kwarg == "intended_date" # run scheduler scheduler(broker=broker) # set up the workflow @@ -431,8 +433,8 @@ def test_intended_schedule_kwarg(broker, monkeypatch): pusher(task_queue, stop_event, broker=broker) assert task_queue.qsize() == 1 task = task_queue.get() - assert 'intended_date' in task['kwargs'] - assert task['kwargs']['intended_date'] == run_date.isoformat() + assert "intended_date" in task["kwargs"] + assert task["kwargs"]["intended_date"] == run_date.isoformat() @override_settings( diff --git a/django_q/utils.py b/django_q/utils.py index aed3c29..ee9e8aa 100644 --- a/django_q/utils.py +++ b/django_q/utils.py @@ -1,13 +1,13 @@ -from datetime import datetime import calendar import inspect -from datetime import date +from datetime import date, datetime import django -from django.utils import timezone +from django import db from django.conf import settings +from django.utils import timezone -from django_q.conf import Conf +from django_q.conf import Conf, logger if django.VERSION < (4, 0): # pytz is the default in django 3.2. Remove when no support for 3.2 @@ -72,3 +72,17 @@ def localtime(value=None) -> datetime: return datetime.now() else: return value + + +def close_old_django_connections(): + """ + Close django connections unless running with sync=True. + """ + if Conf.SYNC: + logger.warning( + "Preserving django database connections because sync=True. Beware " + "that tasks are now injected in the calling context/transactions " + "which may result in unexpected behaviour." + ) + else: + db.close_old_connections() diff --git a/django_q/worker.py b/django_q/worker.py new file mode 100644 index 0000000..fc627fc --- /dev/null +++ b/django_q/worker.py @@ -0,0 +1,118 @@ +import pydoc +import traceback +from multiprocessing import Value +from multiprocessing.process import current_process +from multiprocessing.queues import Queue + +from django.utils import timezone +from django.utils.translation import gettext_lazy as _ + +from django_q.conf import Conf, error_reporter, logger, resource, setproctitle +from django_q.signals import post_spawn, pre_execute +from django_q.utils import close_old_django_connections, get_func_repr + +try: + import psutil +except ImportError: + psutil = None + +try: + import setproctitle +except ModuleNotFoundError: + setproctitle = None + + +def worker( + task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT +): + """ + Takes a task from the task queue, tries to execute it and puts the result back in + the result queue + :param timeout: number of seconds wait for a worker to finish. + :type task_queue: multiprocessing.Queue + :type result_queue: multiprocessing.Queue + :type timer: multiprocessing.Value + """ + proc_name = current_process().name + logger.info( + _("%(proc_name)s ready for work at %(id)s") + % {"proc_name": proc_name, "id": current_process().pid} + ) + post_spawn.send(sender="django_q", proc_name=proc_name) + if setproctitle: + setproctitle.setproctitle(f"qcluster {proc_name} idle") + task_count = 0 + if timeout is None: + timeout = -1 + # Start reading the task queue + for task in iter(task_queue.get, "STOP"): + result = None + timer.value = -1 # Idle + task_count += 1 + f = task["func"] + + # Log task creation and set process name + # Get the function from the task + func_name = get_func_repr(f) + task_name = task["name"] + task_desc = _("%(proc_name)s processing %(task_name)s '%(func_name)s'") % { + "proc_name": proc_name, + "func_name": func_name, + "task_name": task_name, + } + if "group" in task: + task_desc += f" [{task['group']}]" + logger.info(task_desc) + + if setproctitle: + proc_title = f"qcluster {proc_name} processing {task_name} '{func_name}'" + if "group" in task: + proc_title += f" [{task['group']}]" + setproctitle.setproctitle(proc_title) + + # if it's not an instance try to get it from the string + if not callable(f): + # locate() returns None if f cannot be loaded + f = pydoc.locate(f) + close_old_django_connections() + timer_value = task.pop("timeout", timeout) + # signal execution + pre_execute.send(sender="django_q", func=f, task=task) + # execute the payload + timer.value = timer_value # Busy + + try: + if f is None: + # raise a meaningfull error if task["func"] is not a valid function + raise ValueError(f"Function {task['func']} is not defined") + res = f(*task["args"], **task["kwargs"]) + result = (res, True) + except Exception as e: + result = (f"{e} : {traceback.format_exc()}", False) + if error_reporter: + error_reporter.report() + if task.get("sync", False): + raise + with timer.get_lock(): + # Process result + task["result"] = result[0] + task["success"] = result[1] + task["stopped"] = timezone.now() + result_queue.put(task) + timer.value = -1 # Idle + if setproctitle: + setproctitle.setproctitle(f"qcluster {proc_name} idle") + # Recycle + if task_count == Conf.RECYCLE or rss_check(): + timer.value = -2 # Recycled + break + logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name}) + + +def rss_check(): + if Conf.MAX_RSS: + if resource: + return resource.getrusage(resource.RUSAGE_SELF).ru_maxrss >= Conf.MAX_RSS + elif psutil: + return psutil.Process().memory_info().rss >= Conf.MAX_RSS * 1024 + return False