From d06de5769b51232ab079895822cc5fa352d54532 Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Thu, 25 Jun 2015 19:48:38 +0200 Subject: [PATCH] added scheduled tasks with admin interface --- django_q/admin.py | 15 ++++++- django_q/core.py | 51 +++++++++++++++++++++- django_q/migrations/0001_initial.py | 33 +++++++++++--- django_q/models.py | 68 ++++++++++++++++++++++++++--- django_q/tests/tasks.py | 3 ++ requirements.txt | 1 + setup.py | 2 +- 7 files changed, 156 insertions(+), 17 deletions(-) diff --git a/django_q/admin.py b/django_q/admin.py index 1435910..5d9b9e0 100644 --- a/django_q/admin.py +++ b/django_q/admin.py @@ -2,7 +2,7 @@ from django.contrib import admin from django_q.core import async -from .models import Success, Failure +from .models import Success, Failure, Schedule class TaskAdmin(admin.ModelAdmin): @@ -59,6 +59,19 @@ class FailAdmin(admin.ModelAdmin): return list(self.readonly_fields) + \ [field.name for field in obj._meta.fields] +class ScheduleAdmin(admin.ModelAdmin): + list_display = ( + u'id', + 'func', + 'schedule_type', + 'repeats', + 'next_run', + 'result', + 'success' + ) + list_filter = ('next_run', 'schedule_type') + +admin.site.register(Schedule, ScheduleAdmin) admin.site.register(Success, TaskAdmin) admin.site.register(Failure, FailAdmin) diff --git a/django_q/core.py b/django_q/core.py index 108b459..26d815b 100644 --- a/django_q/core.py +++ b/django_q/core.py @@ -3,8 +3,10 @@ from __future__ import unicode_literals from __future__ import print_function from __future__ import division from __future__ import absolute_import +import ast from builtins import dict from builtins import range + from future import standard_library standard_library.install_aliases() @@ -28,6 +30,7 @@ except ImportError: # External import coloredlogs import redis +import arrow # Django from django.core import signing @@ -36,7 +39,7 @@ from django.utils import timezone # Local from .conf import LOG_LEVEL, SECRET_KEY, SAVE_LIMIT, WORKERS, COMPRESSED, PREFIX from .humanhash import uuid -from .models import Task, Success +from .models import Task, Success, Schedule SIGNAL_NAMES = dict((getattr(signal, n), n) for n in dir(signal) if n.startswith('SIG') and '_' not in n) @@ -197,6 +200,7 @@ class Sentinel(object): self.start_event.set() self.set_status(RUNNING) logger.info('Q Cluster-{} running.'.format(self.parent_pid)) + counter = 0 while True: for p in list(self.pool): if not p.is_alive(): @@ -206,6 +210,10 @@ class Sentinel(object): Stat(self).save() if self.stop_event.is_set(): break + counter += 1 + if counter > 15: + counter = 0 + scheduler() sleep(2) self.stop() @@ -438,3 +446,44 @@ class Stat(Status): except signing.BadSignature: continue return stats + + +def scheduler(): + for schedule in Schedule.objects.exclude(repeats=0).filter(next_run__lt=timezone.now()): + args = () + kwargs= {} + if schedule.kwargs: + try: + kwargs = eval('dict({})'.format(schedule.kwargs)) + except SyntaxError: + kwargs={} + if schedule.args: + args = ast.literal_eval(schedule.args) + if type(args) != tuple: + args = (args,) + if schedule.hook: + kwargs['hook'] = schedule.hook + schedule.task = async(schedule.func, *args, **kwargs) + if not schedule.schedule_type == schedule.ONCE: + next_run = arrow.get(schedule.next_run) + if schedule.schedule_type == schedule.HOURLY: + next_run = next_run.replace(hours=+1) + elif schedule.schedule_type == schedule.DAILY: + next_run = next_run.replace(days=+1) + elif schedule.schedule_type == schedule.WEEKLY: + next_run = next_run.replace(weeks=+1) + elif schedule.schedule_type == schedule.MONTHLY: + next_run = next_run.replace(months=+1) + elif schedule.schedule_type == schedule.QUARTERLY: + next_run = next_run.replace(months=+3) + elif schedule.schedule_type == schedule.YEARLY: + next_run = next_run.replace(years=+1) + schedule.next_run = next_run.datetime + schedule.repeats += -1 + else: + schedule.repeats = 0 + if not schedule.task: + logger.error('{} failed to create task from schedule {}').format(current_process().name, schedule.id) + else: + logger.info('{} created [{}] from schedule {}'.format(current_process().name, schedule.task, schedule.id)) + schedule.save() diff --git a/django_q/migrations/0001_initial.py b/django_q/migrations/0001_initial.py index 601d46e..4e4287d 100644 --- a/django_q/migrations/0001_initial.py +++ b/django_q/migrations/0001_initial.py @@ -2,6 +2,7 @@ from __future__ import unicode_literals from django.db import models, migrations +import django.utils.timezone import picklefield.fields @@ -11,19 +12,37 @@ class Migration(migrations.Migration): ] operations = [ + migrations.CreateModel( + name='Schedule', + fields=[ + ('id', models.AutoField(serialize=False, primary_key=True, auto_created=True, verbose_name='ID')), + ('func', models.CharField(max_length=256)), + ('hook', models.CharField(blank=True, max_length=256, null=True)), + ('args', models.CharField(blank=True, max_length=256, null=True)), + ('kwargs', models.CharField(blank=True, max_length=256, null=True)), + ('schedule_type', models.CharField(choices=[('O', 'Once'), ('H', 'Hourly'), ('D', 'Daily'), ('W', 'Weekly'), ('M', 'Monthly'), ('Q', 'Quarterly'), ('Y', 'Yearly')], max_length=1, verbose_name='Schedule Type', default='O')), + ('repeats', models.SmallIntegerField(verbose_name='Repeats', default=-1)), + ('next_run', models.DateTimeField(null=True, verbose_name='Next Run', default=django.utils.timezone.now)), + ('task', models.CharField(editable=False, max_length=100, null=True)), + ], + options={ + 'ordering': ['next_run'], + 'verbose_name': 'Scheduled task', + }, + ), migrations.CreateModel( name='Task', fields=[ ('id', models.AutoField(serialize=False, primary_key=True, auto_created=True, verbose_name='ID')), - ('name', models.CharField(max_length=100)), + ('name', models.CharField(editable=False, max_length=100)), ('func', models.CharField(max_length=256)), ('hook', models.CharField(max_length=256, null=True)), ('args', picklefield.fields.PickledObjectField(editable=False)), ('kwargs', picklefield.fields.PickledObjectField(editable=False)), - ('result', picklefield.fields.PickledObjectField(null=True, editable=False)), - ('started', models.DateTimeField()), - ('stopped', models.DateTimeField()), - ('success', models.BooleanField(default=True)), + ('result', picklefield.fields.PickledObjectField(editable=False, null=True)), + ('started', models.DateTimeField(editable=False)), + ('stopped', models.DateTimeField(editable=False)), + ('success', models.BooleanField(editable=False, default=True)), ], ), migrations.CreateModel( @@ -31,8 +50,8 @@ class Migration(migrations.Migration): fields=[ ], options={ - 'proxy': True, 'verbose_name': 'Failed task', + 'proxy': True, }, bases=('django_q.task',), ), @@ -41,8 +60,8 @@ class Migration(migrations.Migration): fields=[ ], options={ - 'proxy': True, 'verbose_name': 'Successful task', + 'proxy': True, }, bases=('django_q.task',), ), diff --git a/django_q/models.py b/django_q/models.py index a6d6ce8..7a7fe9c 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -1,22 +1,25 @@ import importlib import logging +from django.core.urlresolvers import reverse +from django.utils.translation import ugettext_lazy as _ from django.db import models from django.db.models.signals import pre_save from django.dispatch import receiver +from django.utils import timezone from picklefield import PickledObjectField class Task(models.Model): - name = models.CharField(max_length=100) + name = models.CharField(max_length=100, editable=False) func = models.CharField(max_length=256) hook = models.CharField(max_length=256, null=True) args = PickledObjectField() kwargs = PickledObjectField() result = PickledObjectField(null=True) - started = models.DateTimeField() - stopped = models.DateTimeField() - success = models.BooleanField(default=True) + started = models.DateTimeField(editable=False) + stopped = models.DateTimeField(editable=False) + success = models.BooleanField(default=True, editable=False) @staticmethod def get_result(name): @@ -30,6 +33,7 @@ class Task(models.Model): class Meta: app_label = 'django_q' + @receiver(pre_save, sender=Task) def call_hook(sender, instance, **kwargs): if instance.hook: @@ -40,7 +44,7 @@ def call_hook(sender, instance, **kwargs): f(instance) except Exception as e: logger = logging.getLogger('django-q') - logger.error('return hook failed on {}'.format(instance.name)) + logger.error(_('return hook failed on {}').format(instance.name)) logger.exception(e) @@ -55,7 +59,7 @@ class Success(Task): class Meta: app_label = 'django_q' - verbose_name = 'Successful task' + verbose_name = _('Successful task') proxy = True @@ -70,5 +74,55 @@ class Failure(Task): class Meta: app_label = 'django_q' - verbose_name = 'Failed task' + verbose_name = _('Failed task') proxy = True + + +class Schedule(models.Model): + func = models.CharField(max_length=256) + hook = models.CharField(max_length=256, null=True, blank=True) + args = models.CharField(max_length=256, null=True, blank=True) + kwargs = models.CharField(max_length=256, null=True, blank=True) + ONCE = 'O' + HOURLY = 'H' + DAILY = 'D' + WEEKLY = 'W' + MONTHLY = 'M' + QUARTERLY = 'Q' + YEARLY = 'Y' + TYPE = ( + (ONCE, _('Once')), + (HOURLY, _('Hourly')), + (DAILY, _('Daily')), + (WEEKLY, _('Weekly')), + (MONTHLY, _('Monthly')), + (QUARTERLY, _('Quarterly')), + (YEARLY, _('Yearly')), + ) + schedule_type = models.CharField(max_length=1, choices=TYPE, default=TYPE[0][0], verbose_name=_('Schedule Type')) + repeats = models.SmallIntegerField(default=-1, verbose_name=_('Repeats')) + next_run = models.DateTimeField(verbose_name=_('Next Run'), default=timezone.now, null=True) + task = models.CharField(max_length=100, editable=False, null=True) + + def result(self): + if Task.objects.filter(name=self.task).exists(): + task = Task.objects.get(name=self.task) + if task.success: + url = reverse('admin:django_q_success_change', args=(task.id,)) + else: + url = reverse('admin:django_q_failure_change', args=(task.id,)) + return '[{}]'.format(url, self.task) + + return None + + def success(self): + if Task.objects.filter(name=self.task).exists(): + return Task.objects.get(name=self.task).success + + success.boolean = True + result.allow_tags = True + + class Meta: + app_label = 'django_q' + verbose_name = _('Scheduled task') + ordering = ['next_run'] diff --git a/django_q/tests/tasks.py b/django_q/tests/tasks.py index 4258b14..7e328a6 100644 --- a/django_q/tests/tasks.py +++ b/django_q/tests/tasks.py @@ -3,6 +3,9 @@ def countdown(n): while n > 0: n -= 1 +def multiply(x, y): + return x * y + def count_letters(tup): total = 0 for word in tup: diff --git a/requirements.txt b/requirements.txt index 1ad3ac8..d891364 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,3 +1,4 @@ +arrow blessed==1.9.5 coloredlogs==1.0.1 django-picklefield==0.3.1 diff --git a/setup.py b/setup.py index 1028a88..838568e 100644 --- a/setup.py +++ b/setup.py @@ -35,7 +35,7 @@ setup( license='MIT', description='A multiprocessing task queue for Django', long_description=README, - install_requires=['django>=1.7', 'redis', 'coloredlogs', 'django-picklefield', 'blessed'], + install_requires=['django>=1.7', 'redis', 'coloredlogs', 'django-picklefield', 'blessed', 'arrow'], test_requires=['pytest', 'pytest-django', ], cmdclass={'test': PyTest}, classifiers=[