From c8d2927dd297f39237b0f3cc63e1808f43289e15 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Wed, 17 Jun 2015 13:33:46 +0200 Subject: [PATCH] restructured app and added task signing --- django_q/__init__.py | 5 +- django_q/apps.py | 184 +--------------- .../commands/{qworker.py => qcluster.py} | 4 +- django_q/management/commands/testq.py | 4 +- django_q/models.py | 5 +- django_q/q.py | 198 ++++++++++++++++++ django_q/tests/tasks.py | 0 django_q/tests/test_qworker.py | 4 +- 8 files changed, 218 insertions(+), 186 deletions(-) rename django_q/management/commands/{qworker.py => qcluster.py} (75%) create mode 100644 django_q/q.py create mode 100644 django_q/tests/tasks.py diff --git a/django_q/__init__.py b/django_q/__init__.py index 14c0d51..d695d7b 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -3,6 +3,7 @@ A multiprocessing task queue application for Django Author: Ilan Steemers (koed00@gmail.com Github: https://github.com/Koed00/django-q """ -from .apps import defer -from .apps import Worker + +from django_q.q import Cluster, async + default_app_config = 'django_q.apps.SessionAdminConfig' diff --git a/django_q/apps.py b/django_q/apps.py index 5a54994..b2c6424 100644 --- a/django_q/apps.py +++ b/django_q/apps.py @@ -1,187 +1,19 @@ -import importlib -from multiprocessing import Queue, Process, Event, current_process, cpu_count -import sys -import signal -import logging -from time import sleep - from django.apps import AppConfig -import jsonpickle as json -import coloredlogs -import redis - from django.conf import settings -from django.utils import timezone - -from .models import Task -from .humanhash import uuid - -r = redis.StrictRedis(decode_responses=True) -secret = settings.SECRET_KEY -prefix = 'django_q' -q_list = '{}:q'.format(prefix) - -logger = logging.getLogger('django-q') -coloredlogs.install(level=logging.INFO) - class SessionAdminConfig(AppConfig): name = 'django_q' verbose_name = "Django Q" -def defer(func, *args, **kwargs): - # [name, func, args, kwargs, started, finished, result, success] - name = uuid()[0] - pack = json.dumps([name, func, args, kwargs, timezone.now()]) - r.rpush(q_list, pack) - logger.debug('Pushed {}'.format(pack)) - return name +try: + LOG_LEVEL = settings.Q_LOG_LEVEL +except AttributeError: + LOG_LEVEL = "INFO" +try: + SECRET_KEY = settings.SECRET_KEY +except AttributeError: + SECRET_KEY = 'omgicantbelieveyoudonthaveasecretkey' -class Worker(object): - def __init__(self): - signal.signal(signal.SIGTERM, self.sig_handler) - signal.signal(signal.SIGINT, self.sig_handler) - self.running = True - self.stable_size = cpu_count() - self.stable = [] - self.task_queue = Queue() - self.done_queue = Queue() - # Spawn work horses - for i in range(self.stable_size): - self.spawn_horse() - # Spawn monitor - self.monitor_pid = None - self.spawn_monitor() - # Spawn pusher - self.pusher_pid = None - self.pusher_stop = Event() - self.spawn_pusher() - # Monitor process health - while self.running: - self.stable_boy() - sleep(1) - - def spawn_process(self, target, *args): - # This is just for PyCharm to not crash. Ignore it. - if not hasattr(sys.stdin, 'close'): - def dummy_close(): - pass - - sys.stdin.close = dummy_close - p = Process(target=target, args=args) - self.stable.append(p) - p.start() - return p.pid - - def spawn_pusher(self): - self.pusher_pid = self.spawn_process(self.pusher, self.task_queue, self.pusher_stop) - - def spawn_horse(self): - self.spawn_process(self.horse, self.task_queue, self.done_queue) - - def spawn_monitor(self): - self.monitor_pid = self.spawn_process(self.monitor, self.done_queue) - - def reincarnate(self, pid): - if pid == self.monitor_pid: - self.spawn_monitor() - logger.warn("reincarnated monitor after death of {}".format(pid)) - elif pid == self.pusher_pid: - self.spawn_pusher() - logger.warn("reincarnated pusher after death of {}".format(pid)) - else: - self.spawn_horse() - logger.warn("reincarnated work horse after death of {}".format(pid)) - - @staticmethod - def pusher(task_queue, e): - logger.info('{} pushing tasks at {}'.format(current_process().name, current_process().pid)) - while not e.is_set(): - task = r.blpop(q_list, 1) - if task: - task_queue.put(task[1]) - logger.debug('queueing {}'.format(task[1])) - - @staticmethod - def monitor(done_queue): - name = current_process().name - logger.info("{} monitoring at {}".format(name, current_process().pid)) - for task in iter(done_queue.get, 'STOP'): - name = task[0] - func = task[1] - result = task[6] - success = task[7] - if success: - logger.info("Finished [{}:{}]".format(func, name)) - result = json.dumps(result) - else: - logger.error("Failed [{}:{}] - {}".format(func, name, result)) - Task.objects.create(name=name, - func=func, - task=json.dumps(task), - started=task[4], - stopped=task[5], - result=result, - success=success) - logger.info("{} stopped".format(name)) - - @staticmethod - def horse(task_queue, done_queue): - name = current_process().name - logger.info('{} ready for work at {}'.format(name, current_process().pid)) - for pack in iter(task_queue.get, 'STOP'): - try: - task = json.loads(pack) - except TypeError as e: - logger.error(e) - continue - func = task[1] - module, func = func.rsplit('.', 1) - args = task[2] - kwargs = task[3] - logger.info('{} processing [{}:{}]'.format(name, func, task[0])) - task.append(timezone.now()) - try: - m = importlib.import_module(module) - f = getattr(m, func) - result = f(*args, **kwargs) - task.append(result) - task.append(True) - done_queue.put(task) - except Exception as e: - task.append(e) - task.append(False) - done_queue.put(task) - logger.info('{} Stopped'.format(name)) - - def stable_boy(self): - # Check if all the horses are alive - for p in list(self.stable): - if not p.is_alive(): - # Be humane - p.terminate() - self.stable.remove(p) - # Replace it with a fresh one - self.reincarnate(p.pid) - - def stop(self): - # Send the STOP signal to the stable - self.running = False - logger.info('Stopping') - # Wait for all the workers to finish the queue - for p in self.stable: - if p.pid == self.monitor_pid: - self.done_queue.put('STOP') - elif p.pid == self.pusher_pid: - self.pusher_stop.set() - else: - self.task_queue.put('STOP') - p.join() - - logger.info('Goodbye. Have a wonderful time.') - - def sig_handler(self, signum, frame): - self.stop() diff --git a/django_q/management/commands/qworker.py b/django_q/management/commands/qcluster.py similarity index 75% rename from django_q/management/commands/qworker.py rename to django_q/management/commands/qcluster.py index ac4879c..1016a52 100644 --- a/django_q/management/commands/qworker.py +++ b/django_q/management/commands/qcluster.py @@ -1,9 +1,9 @@ from django.core.management.base import BaseCommand -from django_q.apps import Worker +from django_q import Cluster class Command(BaseCommand): help = "My shiny new management command." def handle(self, *args, **options): - q = Worker() + q = Cluster() diff --git a/django_q/management/commands/testq.py b/django_q/management/commands/testq.py index 8d07991..793a1cd 100644 --- a/django_q/management/commands/testq.py +++ b/django_q/management/commands/testq.py @@ -1,6 +1,6 @@ from django.core.management.base import BaseCommand -from django_q.apps import defer +from django_q import async class Command(BaseCommand): @@ -8,4 +8,4 @@ class Command(BaseCommand): def handle(self, *args, **options): for i in range(20): - defer('testq.tasks.multiply', 2, i) + async('testq.tasks.multiply', 2, i) diff --git a/django_q/models.py b/django_q/models.py index 7d68b5d..0ff59f4 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -1,13 +1,14 @@ from django.db import models import socket from django.utils.translation import ugettext_lazy as _ +from picklefield import PickledObjectField class Task(models.Model): name = models.CharField(max_length=100) func = models.CharField(max_length=256) - task = models.TextField(null=True) - result = models.TextField(null=True) + task = PickledObjectField() + result = PickledObjectField() started = models.DateTimeField() stopped = models.DateTimeField() success = models.BooleanField(default=True) diff --git a/django_q/q.py b/django_q/q.py new file mode 100644 index 0000000..ece9ed1 --- /dev/null +++ b/django_q/q.py @@ -0,0 +1,198 @@ +import importlib +import logging +import pickle +import signal +from multiprocessing import cpu_count, Queue, Event, Process, current_process +import sys +from time import sleep + +import coloredlogs + +from django.utils import timezone +import redis + +from django.core.signing import Signer, BadSignature +from django_q.apps import LOG_LEVEL, SECRET_KEY + +from django_q.humanhash import uuid +from django_q.models import Task + +prefix = 'django_q' +q_list = '{}:q'.format(prefix) +signer = Signer(SECRET_KEY) +logger = logging.getLogger('django-q') +coloredlogs.install(level=getattr(logging, LOG_LEVEL)) + +r = redis.StrictRedis(decode_responses=True) + + +class Cluster(object): + def __init__(self): + signal.signal(signal.SIGTERM, self.sig_handler) + signal.signal(signal.SIGINT, self.sig_handler) + try: + r.ping() + except (): + logger.error('Can not connect to Redis server') + return + self.running = True + self.pool_size = cpu_count() + self.pool = [] + self.task_queue = Queue() + self.done_queue = Queue() + # Spawn workers + for i in range(self.pool_size): + self.spawn_worker() + # Spawn monitor + self.monitor_pid = None + self.spawn_monitor() + # Spawn pusher + self.pusher_pid = None + self.pusher_stop = Event() + self.spawn_pusher() + # Monitor process health + while self.running: + self.medic() + sleep(1) + + def spawn_process(self, target, *args): + # This is just for PyCharm to not crash. Ignore it. + if not hasattr(sys.stdin, 'close'): + def dummy_close(): + pass + + sys.stdin.close = dummy_close + p = Process(target=target, args=args) + self.pool.append(p) + p.start() + return p.pid + + def spawn_pusher(self): + self.pusher_pid = self.spawn_process(self.pusher, self.task_queue, self.pusher_stop) + + def spawn_worker(self): + self.spawn_process(self.worker, self.task_queue, self.done_queue) + + def spawn_monitor(self): + self.monitor_pid = self.spawn_process(self.monitor, self.done_queue) + + def reincarnate(self, pid): + if pid == self.monitor_pid: + self.spawn_monitor() + logger.warn("reincarnated monitor after death of {}".format(pid)) + elif pid == self.pusher_pid: + self.spawn_pusher() + logger.warn("reincarnated pusher after death of {}".format(pid)) + else: + self.spawn_worker() + logger.warn("reincarnated work worker after death of {}".format(pid)) + + @staticmethod + def pusher(task_queue, e): + logger.info('{} pushing tasks at {}'.format(current_process().name, current_process().pid)) + while not e.is_set(): + task = r.blpop(q_list, 1) + if task: + task = task[1] + task_queue.put(task) + logger.debug('queueing {}'.format(task)) + + @staticmethod + def monitor(done_queue): + name = current_process().name + logger.info("{} monitoring at {}".format(name, current_process().pid)) + for task in iter(done_queue.get, 'STOP'): + name = task[0] + func = task[1] + result = task[6] + success = task[7] + if success: + logger.info("Finished [{}:{}]".format(func, name)) + else: + logger.error("Failed [{}:{}] - {}".format(func, name, result)) + Task.objects.create(name=name, + func=func, + task=task, + started=task[4], + stopped=task[5], + result=task[6], + success=success) + logger.info("{} stopped".format(name)) + + @staticmethod + def worker(task_queue, done_queue): + name = current_process().name + logger.info('{} ready for work at {}'.format(name, current_process().pid)) + for pack in iter(task_queue.get, 'STOP'): + # unpickle the task + try: + task = pickle.loads(pack) + except TypeError as e: + logger.error(e) + continue + # check signature + try: + task[0] = signer.unsign(task[0]) + except BadSignature as e: + logger.error("Bad signature on task.") + task.append(timezone.now()) + task.append(e) + task.append(False) + done_queue.put(task) + continue + func = task[1] + module, func = func.rsplit('.', 1) + args = task[2] + kwargs = task[3] + logger.info('{} processing [{}:{}]'.format(name, func, task[0])) + task.append(timezone.now()) + try: + m = importlib.import_module(module) + f = getattr(m, func) + result = f(*args, **kwargs) + task.append(result) + task.append(True) + done_queue.put(task) + except Exception as e: + task.append(e) + task.append(False) + done_queue.put(task) + logger.info('{} Stopped'.format(name)) + + def medic(self): + # Check if all the workers are alive + for p in list(self.pool): + if not p.is_alive(): + # Be humane + p.terminate() + self.pool.remove(p) + # Replace it with a fresh one + self.reincarnate(p.pid) + + def stop(self): + # Send the STOP signal to the pool + self.running = False + logger.info('Stopping') + # Wait for all the workers to finish the queue + for p in self.pool: + if p.pid == self.monitor_pid: + self.done_queue.put('STOP') + elif p.pid == self.pusher_pid: + self.pusher_stop.set() + else: + self.task_queue.put('STOP') + p.join() + + logger.info('Goodbye. Have a wonderful time.') + + def sig_handler(self, signum, frame): + self.stop() + + +def async(func, *args, **kwargs): + # [name, func, args, kwargs, started, finished, result, success] + name = signer.sign(uuid()[0]) + pack = pickle.dumps([name, func, args, kwargs, timezone.now()]) + r.rpush(q_list, pack) + logger.debug('Pushed {}'.format(pack)) + return name diff --git a/django_q/tests/tasks.py b/django_q/tests/tasks.py new file mode 100644 index 0000000..e69de29 diff --git a/django_q/tests/test_qworker.py b/django_q/tests/test_qworker.py index 8f25e7c..bab5a9a 100644 --- a/django_q/tests/test_qworker.py +++ b/django_q/tests/test_qworker.py @@ -1,10 +1,10 @@ import pytest -from django_q.apps import Worker +from django_q import Cluster @pytest.fixture def qworker(): - return Worker() + return Cluster() def test_worker(qworker):