diff --git a/django_q/__init__.py b/django_q/__init__.py index e69de29..14c0d51 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -0,0 +1,8 @@ +""" +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 +default_app_config = 'django_q.apps.SessionAdminConfig' diff --git a/django_q/q.py b/django_q/apps.py similarity index 65% rename from django_q/q.py rename to django_q/apps.py index 093cdbc..5a54994 100644 --- a/django_q/q.py +++ b/django_q/apps.py @@ -1,20 +1,20 @@ import importlib -from random import randint 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) @@ -26,16 +26,18 @@ logger = logging.getLogger('django-q') coloredlogs.install(level=logging.INFO) -def test(): - for i in range(20): - defer(u'testq.tasks.multiply', 5, i * randint(2, 100)) +class SessionAdminConfig(AppConfig): + name = 'django_q' + verbose_name = "Django Q" def defer(func, *args, **kwargs): - # [name, func, args, kwargs, started, finished, result] - pack = json.dumps([uuid()[0], func, args, kwargs, timezone.now()]) + # [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 class Worker(object): @@ -47,15 +49,12 @@ class Worker(object): self.stable = [] self.task_queue = Queue() self.done_queue = Queue() - self.fail_queue = Queue() # Spawn work horses for i in range(self.stable_size): self.spawn_horse() - # Spawn monitors - self.success_monitor_pid = None - self.spawn_success_monitor() - self.failure_monitor_pid = None - self.spawn_failure_monitor() + # Spawn monitor + self.monitor_pid = None + self.spawn_monitor() # Spawn pusher self.pusher_pid = None self.pusher_stop = Event() @@ -81,21 +80,18 @@ class Worker(object): 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, self.fail_queue) + self.spawn_process(self.horse, self.task_queue, self.done_queue) - def spawn_success_monitor(self): - self.success_monitor_pid = self.spawn_process(self.success_monitor, self.done_queue) - - def spawn_failure_monitor(self): - self.failure_monitor_pid = self.spawn_process(self.failure_monitor, self.fail_queue) + def spawn_monitor(self): + self.monitor_pid = self.spawn_process(self.monitor, self.done_queue) def reincarnate(self, pid): - if pid == self.success_monitor_pid: - self.spawn_success_monitor() - logger.warn("reincarnated success monitor after death of {}".format(pid)) - elif pid == self.failure_monitor_pid: - self.spawn_failure_monitor() - logger.warn("reincarnated failure monitor after death of {}".format(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)) @@ -110,29 +106,30 @@ class Worker(object): logger.debug('queueing {}'.format(task[1])) @staticmethod - def success_monitor(done_queue): + def monitor(done_queue): name = current_process().name - logger.info("{} monitoring successes at {}".format(name, current_process().pid)) + logger.info("{} monitoring at {}".format(name, current_process().pid)) for task in iter(done_queue.get, 'STOP'): - logger.info("Finished [{}:{}]".format(task[1], task[0])) - Task.objects.create(name=task[0], func=task[1], task=json.dumps(task), + 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]) + stopped=task[5], + result=result, + success=success) logger.info("{} stopped".format(name)) @staticmethod - def failure_monitor(fail_queue): - name = current_process().name - logger.info("{} monitoring failures at {}".format(name, current_process().pid)) - for task in iter(fail_queue.get, 'STOP'): - logger.error("Failure [{}:{} - {}]".format(task[1], task[0], task[6])) - Task.objects.create(name=task[0], func=task[1], task=json.dumps(task), - started=task[4], - stopped=task[5], success=False) - logger.info("{} stopped".format(name)) - - @staticmethod - def horse(task_queue, done_queue, fail_queue): + 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'): @@ -152,10 +149,12 @@ class Worker(object): f = getattr(m, func) result = f(*args, **kwargs) task.append(result) + task.append(True) done_queue.put(task) - except TypeError as e: + except Exception as e: task.append(e) - fail_queue.put(task) + task.append(False) + done_queue.put(task) logger.info('{} Stopped'.format(name)) def stable_boy(self): @@ -174,9 +173,7 @@ class Worker(object): logger.info('Stopping') # Wait for all the workers to finish the queue for p in self.stable: - if p.pid == self.failure_monitor_pid: - self.fail_queue.put('STOP') - elif p.pid == self.success_monitor_pid: + if p.pid == self.monitor_pid: self.done_queue.put('STOP') elif p.pid == self.pusher_pid: self.pusher_stop.set() diff --git a/django_q/management/commands/qworker.py b/django_q/management/commands/qworker.py index debf90f..ac4879c 100644 --- a/django_q/management/commands/qworker.py +++ b/django_q/management/commands/qworker.py @@ -1,5 +1,5 @@ from django.core.management.base import BaseCommand -from django_q.q import Worker +from django_q.apps import Worker class Command(BaseCommand): diff --git a/django_q/management/commands/testq.py b/django_q/management/commands/testq.py index 0da1e8b..8d07991 100644 --- a/django_q/management/commands/testq.py +++ b/django_q/management/commands/testq.py @@ -1,10 +1,11 @@ from django.core.management.base import BaseCommand -from django_q import q +from django_q.apps import defer class Command(BaseCommand): help = "My shiny new management command." def handle(self, *args, **options): - q.test() + for i in range(20): + defer('testq.tasks.multiply', 2, i) diff --git a/django_q/models.py b/django_q/models.py index da80f1b..7d68b5d 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -7,9 +7,13 @@ 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) started = models.DateTimeField() stopped = models.DateTimeField() success = models.BooleanField(default=True) def time_taken(self): return (self.stopped - self.started).total_seconds() + + class Meta: + app_label = 'django_q' diff --git a/django_q/tests.py b/django_q/tests.py deleted file mode 100644 index 7ce503c..0000000 --- a/django_q/tests.py +++ /dev/null @@ -1,3 +0,0 @@ -from django.test import TestCase - -# Create your tests here. diff --git a/django_q/tests/__init__.py b/django_q/tests/__init__.py new file mode 100644 index 0000000..dc155fb --- /dev/null +++ b/django_q/tests/__init__.py @@ -0,0 +1 @@ +__author__ = 'ilan' diff --git a/django_q/tests/settings.py b/django_q/tests/settings.py new file mode 100644 index 0000000..1fa5e46 --- /dev/null +++ b/django_q/tests/settings.py @@ -0,0 +1,105 @@ +""" +Django settings for testq project. + +Generated by 'django-admin startproject' using Django 1.8.2. + +For more information on this file, see +https://docs.djangoproject.com/en/1.8/topics/settings/ + +For the full list of settings and their values, see +https://docs.djangoproject.com/en/1.8/ref/settings/ +""" + +# Build paths inside the project like this: os.path.join(BASE_DIR, ...) +import os + +BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) + + +# Quick-start development settings - unsuitable for production +# See https://docs.djangoproject.com/en/1.8/howto/deployment/checklist/ + +# SECURITY WARNING: keep the secret key used in production secret! +SECRET_KEY = ')cqmpi+p@n&!u&fu@!m@9h&1bz9mwmstsahe)nf!ms+c$uc=x7' + +# SECURITY WARNING: don't run with debug turned on in production! +DEBUG = True + +ALLOWED_HOSTS = [] + + +# Application definition + +INSTALLED_APPS = ( + 'django.contrib.admin', + 'django.contrib.auth', + 'django.contrib.contenttypes', + 'django.contrib.sessions', + 'django.contrib.messages', + 'django.contrib.staticfiles', + 'django_extensions', + 'django_q' +) + +MIDDLEWARE_CLASSES = ( + 'django.contrib.sessions.middleware.SessionMiddleware', + 'django.middleware.common.CommonMiddleware', + 'django.middleware.csrf.CsrfViewMiddleware', + 'django.contrib.auth.middleware.AuthenticationMiddleware', + 'django.contrib.auth.middleware.SessionAuthenticationMiddleware', + 'django.contrib.messages.middleware.MessageMiddleware', + 'django.middleware.clickjacking.XFrameOptionsMiddleware', + 'django.middleware.security.SecurityMiddleware', +) + +ROOT_URLCONF = 'testq.urls' + +TEMPLATES = [ + { + 'BACKEND': 'django.template.backends.django.DjangoTemplates', + 'DIRS': [], + 'APP_DIRS': True, + 'OPTIONS': { + 'context_processors': [ + 'django.template.context_processors.debug', + 'django.template.context_processors.request', + 'django.contrib.auth.context_processors.auth', + 'django.contrib.messages.context_processors.messages', + ], + }, + }, +] + +WSGI_APPLICATION = 'testq.wsgi.application' + + +# Database +# https://docs.djangoproject.com/en/1.8/ref/settings/#databases + +DATABASES = { + 'default': { + 'ENGINE': 'django.db.backends.sqlite3', + 'NAME': os.path.join(BASE_DIR, 'db.sqlite3'), + } +} + + +# Internationalization +# https://docs.djangoproject.com/en/1.8/topics/i18n/ + +LANGUAGE_CODE = 'en-us' + +TIME_ZONE = 'UTC' + +USE_I18N = True + +USE_L10N = True + +USE_TZ = True + + +# Static files (CSS, JavaScript, Images) +# https://docs.djangoproject.com/en/1.8/howto/static-files/ + +STATIC_URL = '/static/' + diff --git a/django_q/tests/test_qworker.py b/django_q/tests/test_qworker.py new file mode 100644 index 0000000..8f25e7c --- /dev/null +++ b/django_q/tests/test_qworker.py @@ -0,0 +1,11 @@ +import pytest +from django_q.apps import Worker + + +@pytest.fixture +def qworker(): + return Worker() + + +def test_worker(qworker): + assert len(qworker.stable) == qworker.stable_size diff --git a/pytest.ini b/pytest.ini new file mode 100644 index 0000000..676f86b --- /dev/null +++ b/pytest.ini @@ -0,0 +1,2 @@ +[pytest] +DJANGO_SETTINGS_MODULE=.settings \ No newline at end of file