mirror of
https://github.com/django-q2/django-q2.git
synced 2026-09-15 13:37:56 +08:00
Added admin pages for failures and success.
SAVE_LIMIT can now be set to limit the amount of saved tasks results.
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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',),
|
||||
),
|
||||
]
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user