mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-06 23:38:11 +08:00
Adds group commands to task instances
* adds task.group_result() * adds task.group_count() * adds task.group_delete()
This commit is contained in:
+29
-18
@@ -1,6 +1,6 @@
|
|||||||
import logging
|
import logging
|
||||||
from django import get_version
|
|
||||||
|
|
||||||
|
from django import get_version
|
||||||
import importlib
|
import importlib
|
||||||
from django.core.urlresolvers import reverse
|
from django.core.urlresolvers import reverse
|
||||||
from django.utils.translation import ugettext_lazy as _
|
from django.utils.translation import ugettext_lazy as _
|
||||||
@@ -10,6 +10,7 @@ from django.dispatch import receiver
|
|||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
from picklefield import PickledObjectField
|
from picklefield import PickledObjectField
|
||||||
from picklefield.fields import dbsafe_decode
|
from picklefield.fields import dbsafe_decode
|
||||||
|
|
||||||
from django_q.signing import SignedPackage
|
from django_q.signing import SignedPackage
|
||||||
|
|
||||||
|
|
||||||
@@ -41,12 +42,20 @@ class Task(models.Model):
|
|||||||
values = Task.objects.filter(group=group_id).exclude(success=False).values_list('result', flat=True)
|
values = Task.objects.filter(group=group_id).exclude(success=False).values_list('result', flat=True)
|
||||||
return decode_results(values)
|
return decode_results(values)
|
||||||
|
|
||||||
|
def group_result(self, failures=False):
|
||||||
|
if self.group:
|
||||||
|
return self.get_result_group(self.group, failures)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_group_count(group_id, failures=False):
|
def get_group_count(group_id, failures=False):
|
||||||
if failures:
|
if failures:
|
||||||
return Failure.objects.filter(group=group_id).count()
|
return Failure.objects.filter(group=group_id).count()
|
||||||
return Task.objects.filter(group=group_id).count()
|
return Task.objects.filter(group=group_id).count()
|
||||||
|
|
||||||
|
def group_count(self, failures=False):
|
||||||
|
if self.group:
|
||||||
|
return self.get_group_count(self.group, failures)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def delete_group(group_id, objects=False):
|
def delete_group(group_id, objects=False):
|
||||||
group = Task.objects.filter(group=group_id)
|
group = Task.objects.filter(group=group_id)
|
||||||
@@ -54,6 +63,10 @@ class Task(models.Model):
|
|||||||
return group.delete()
|
return group.delete()
|
||||||
return group.update(group=None)
|
return group.update(group=None)
|
||||||
|
|
||||||
|
def group_delete(self, tasks=False):
|
||||||
|
if self.group:
|
||||||
|
return self.delete_group(self.group, tasks)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_task(task_id):
|
def get_task(task_id):
|
||||||
if len(task_id) == 32 and Task.objects.filter(id=task_id).exists():
|
if len(task_id) == 32 and Task.objects.filter(id=task_id).exists():
|
||||||
@@ -190,26 +203,26 @@ class Schedule(models.Model):
|
|||||||
|
|
||||||
|
|
||||||
class OrmQ(models.Model):
|
class OrmQ(models.Model):
|
||||||
key = models.CharField(max_length=100)
|
key = models.CharField(max_length=100)
|
||||||
payload = models.TextField()
|
payload = models.TextField()
|
||||||
lock = models.DateTimeField(null=True)
|
lock = models.DateTimeField(null=True)
|
||||||
|
|
||||||
def task(self):
|
def task(self):
|
||||||
return SignedPackage.loads(self.payload)
|
return SignedPackage.loads(self.payload)
|
||||||
|
|
||||||
def func(self):
|
def func(self):
|
||||||
return self.task()['func']
|
return self.task()['func']
|
||||||
|
|
||||||
def task_id(self):
|
def task_id(self):
|
||||||
return self.task()['id']
|
return self.task()['id']
|
||||||
|
|
||||||
def name(self):
|
def name(self):
|
||||||
return self.task()['name']
|
return self.task()['name']
|
||||||
|
|
||||||
class Meta:
|
class Meta:
|
||||||
app_label = 'django_q'
|
app_label = 'django_q'
|
||||||
verbose_name = _('Queued task')
|
verbose_name = _('Queued task')
|
||||||
verbose_name_plural = _('Queued tasks')
|
verbose_name_plural = _('Queued tasks')
|
||||||
|
|
||||||
|
|
||||||
# Backwards compatibility for Django 1.7
|
# Backwards compatibility for Django 1.7
|
||||||
@@ -218,5 +231,3 @@ def decode_results(values):
|
|||||||
# decode values in 1.7
|
# decode values in 1.7
|
||||||
return [dbsafe_decode(v) for v in values]
|
return [dbsafe_decode(v) for v in values]
|
||||||
return values
|
return values
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -215,13 +215,19 @@ def test_async(broker, admin_user):
|
|||||||
assert result(result_j.name) == result_j.result
|
assert result(result_j.name) == result_j.result
|
||||||
# groups
|
# groups
|
||||||
assert result_group('test_j')[0] == result_j.result
|
assert result_group('test_j')[0] == result_j.result
|
||||||
|
assert result_j.group_result()[0] == result_j.result
|
||||||
assert result_group('test_j', failures=True)[0] == result_j.result
|
assert result_group('test_j', failures=True)[0] == result_j.result
|
||||||
|
assert result_j.group_result(failures=True)[0] == result_j.result
|
||||||
assert fetch_group('test_j')[0].id == [result_j][0].id
|
assert fetch_group('test_j')[0].id == [result_j][0].id
|
||||||
assert fetch_group('test_j', failures=False)[0].id == [result_j][0].id
|
assert fetch_group('test_j', failures=False)[0].id == [result_j][0].id
|
||||||
assert count_group('test_j') == 1
|
assert count_group('test_j') == 1
|
||||||
|
assert result_j.group_count() == 1
|
||||||
assert count_group('test_j', failures=True) == 0
|
assert count_group('test_j', failures=True) == 0
|
||||||
|
assert result_j.group_count(failures=True) == 0
|
||||||
assert delete_group('test_j') == 1
|
assert delete_group('test_j') == 1
|
||||||
|
assert result_j.group_delete() == 0
|
||||||
assert delete_group('test_j', tasks=True) is None
|
assert delete_group('test_j', tasks=True) is None
|
||||||
|
assert result_j.group_delete(tasks=True) is None
|
||||||
# task k should not have been saved
|
# task k should not have been saved
|
||||||
assert fetch(k) is None
|
assert fetch(k) is None
|
||||||
broker.delete_queue()
|
broker.delete_queue()
|
||||||
|
|||||||
Reference in New Issue
Block a user