diff --git a/django_q/admin.py b/django_q/admin.py index 5f90e4b..a39aac5 100644 --- a/django_q/admin.py +++ b/django_q/admin.py @@ -1,7 +1,6 @@ from django.contrib import admin -# Register your models here. -from .models import Task +from .models import Success, Failure class TaskAdmin(admin.ModelAdmin): @@ -9,7 +8,48 @@ class TaskAdmin(admin.ModelAdmin): u'name', 'func', 'started', - 'time_taken', - 'success' + 'time_taken' ) -admin.site.register(Task, TaskAdmin) + + def has_add_permission(self, request, obj=None): + """Don't allow adds""" + return False + + def get_queryset(self, request): + """Only show successes""" + qs = super(TaskAdmin, self).get_queryset(request) + return qs.filter(success=True) + + search_fields = ['name'] + readonly_fields = [] + + def get_readonly_fields(self, request, obj=None): + return list(self.readonly_fields) + \ + [field.name for field in obj._meta.fields] + + +def retry_failed(FailAdmin, request, queryset): + pass + + +retry_failed.short_description = "Resubmit selected tasks to Q" + + +class FailAdmin(admin.ModelAdmin): + list_display = ( + u'name', + 'func', + 'started', + 'result' + ) + actions = [retry_failed] + search_fields = ['name'] + readonly_fields = [] + + def get_readonly_fields(self, request, obj=None): + return list(self.readonly_fields) + \ + [field.name for field in obj._meta.fields] + + +admin.site.register(Success, TaskAdmin) +admin.site.register(Failure, FailAdmin) diff --git a/django_q/apps.py b/django_q/apps.py index b2c6424..8ee09ea 100644 --- a/django_q/apps.py +++ b/django_q/apps.py @@ -6,14 +6,28 @@ class SessionAdminConfig(AppConfig): name = 'django_q' verbose_name = "Django Q" - +""" +Sets the logging level for the app +""" try: LOG_LEVEL = settings.Q_LOG_LEVEL except AttributeError: LOG_LEVEL = "INFO" +""" +Using Django's secret key to sign task packages +""" try: SECRET_KEY = settings.SECRET_KEY except AttributeError: SECRET_KEY = 'omgicantbelieveyoudonthaveasecretkey' +""" +SAVE_LIMIT limits the amount of successful task executions saved to the database. +Set this to 0 for no limits. +Failures are not limited. +""" +try: + SAVE_LIMIT = settings.Q_SAVE_LIMIT +except AttributeError: + SAVE_LIMIT = 100 diff --git a/django_q/migrations/0001_initial.py b/django_q/migrations/0001_initial.py index eb57b2a..b91d43e 100644 --- a/django_q/migrations/0001_initial.py +++ b/django_q/migrations/0001_initial.py @@ -14,14 +14,33 @@ class Migration(migrations.Migration): migrations.CreateModel( name='Task', fields=[ - ('id', models.AutoField(verbose_name='ID', auto_created=True, serialize=False, primary_key=True)), + ('id', models.AutoField(serialize=False, auto_created=True, verbose_name='ID', primary_key=True)), ('name', models.CharField(max_length=100)), ('func', models.CharField(max_length=256)), - ('task', picklefield.fields.PickledObjectField(editable=False)), + ('args', picklefield.fields.PickledObjectField(editable=False)), + ('kwargs', picklefield.fields.PickledObjectField(editable=False)), ('result', picklefield.fields.PickledObjectField(editable=False)), ('started', models.DateTimeField()), ('stopped', models.DateTimeField()), ('success', models.BooleanField(default=True)), ], ), + migrations.CreateModel( + name='Failure', + fields=[ + ], + options={ + 'proxy': True, + }, + bases=('django_q.task',), + ), + migrations.CreateModel( + name='Success', + fields=[ + ], + options={ + 'proxy': True, + }, + bases=('django_q.task',), + ), ] diff --git a/django_q/models.py b/django_q/models.py index 0ff59f4..4224255 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -1,20 +1,51 @@ 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 = PickledObjectField() + args = PickledObjectField() + kwargs = PickledObjectField() result = PickledObjectField() started = models.DateTimeField() stopped = models.DateTimeField() success = models.BooleanField(default=True) + @staticmethod + def get_result(name): + return Task.objects.filter(name=name).values_list('result', flat=True) + def time_taken(self): return (self.stopped - self.started).total_seconds() class Meta: app_label = 'django_q' + + +class SuccessManager(models.Manager): + def get_queryset(self): + return super(SuccessManager, self).get_queryset().filter( + success=True) + + +class Success(Task): + objects = SuccessManager() + + class Meta: + app_label = 'django_q' + proxy = True + + +class FailureManager(models.Manager): + def get_queryset(self): + return super(FailureManager, self).get_queryset().filter( + success=False) + + +class Failure(Task): + objects = FailureManager() + + class Meta: + app_label = 'django_q' + proxy = True diff --git a/django_q/q.py b/django_q/q.py index ece9ed1..cf4ea5a 100644 --- a/django_q/q.py +++ b/django_q/q.py @@ -12,10 +12,10 @@ 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.apps import LOG_LEVEL, SECRET_KEY, SAVE_LIMIT from django_q.humanhash import uuid -from django_q.models import Task +from django_q.models import Task, Success prefix = 'django_q' q_list = '{}:q'.format(prefix) @@ -23,7 +23,7 @@ signer = Signer(SECRET_KEY) logger = logging.getLogger('django-q') coloredlogs.install(level=getattr(logging, LOG_LEVEL)) -r = redis.StrictRedis(decode_responses=True) +r = redis.StrictRedis() class Cluster(object): @@ -110,13 +110,7 @@ class Cluster(object): 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) + Cluster.save_task(task) logger.info("{} stopped".format(name)) @staticmethod @@ -134,12 +128,12 @@ class Cluster(object): 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 + task[0] = task[0].rsplit(":", 1)[0] + 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] @@ -169,6 +163,19 @@ class Cluster(object): # Replace it with a fresh one self.reincarnate(p.pid) + @staticmethod + def save_task(task): + if task[7] and 0 < SAVE_LIMIT < Success.objects.count(): + Success.objects.first().delete() + Task.objects.create(name=task[0], + func=task[1], + args=task[2], + kwargs=task[3], + started=task[4], + stopped=task[5], + result=task[6], + success=task[7]) + def stop(self): # Send the STOP signal to the pool self.running = False