From 061775b3bb535eb8e84dedf036acc9fd606d325e Mon Sep 17 00:00:00 2001 From: Ilan Steemers Date: Wed, 16 Sep 2015 15:47:39 +0200 Subject: [PATCH] Adds group commands to task instances * adds task.group_result() * adds task.group_count() * adds task.group_delete() --- django_q/models.py | 47 +++++++++++++++++++++------------- django_q/tests/test_cluster.py | 6 +++++ 2 files changed, 35 insertions(+), 18 deletions(-) diff --git a/django_q/models.py b/django_q/models.py index e3037d3..203bcd5 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -1,6 +1,6 @@ import logging -from django import get_version +from django import get_version import importlib from django.core.urlresolvers import reverse from django.utils.translation import ugettext_lazy as _ @@ -10,6 +10,7 @@ from django.dispatch import receiver from django.utils import timezone from picklefield import PickledObjectField from picklefield.fields import dbsafe_decode + 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) return decode_results(values) + def group_result(self, failures=False): + if self.group: + return self.get_result_group(self.group, failures) + @staticmethod def get_group_count(group_id, failures=False): if failures: return Failure.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 def delete_group(group_id, objects=False): group = Task.objects.filter(group=group_id) @@ -54,6 +63,10 @@ class Task(models.Model): return group.delete() return group.update(group=None) + def group_delete(self, tasks=False): + if self.group: + return self.delete_group(self.group, tasks) + @staticmethod def get_task(task_id): 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): - key = models.CharField(max_length=100) - payload = models.TextField() - lock = models.DateTimeField(null=True) + key = models.CharField(max_length=100) + payload = models.TextField() + lock = models.DateTimeField(null=True) - def task(self): - return SignedPackage.loads(self.payload) + def task(self): + return SignedPackage.loads(self.payload) - def func(self): - return self.task()['func'] + def func(self): + return self.task()['func'] - def task_id(self): - return self.task()['id'] + def task_id(self): + return self.task()['id'] - def name(self): - return self.task()['name'] + def name(self): + return self.task()['name'] - class Meta: - app_label = 'django_q' - verbose_name = _('Queued task') - verbose_name_plural = _('Queued tasks') + class Meta: + app_label = 'django_q' + verbose_name = _('Queued task') + verbose_name_plural = _('Queued tasks') # Backwards compatibility for Django 1.7 @@ -218,5 +231,3 @@ def decode_results(values): # decode values in 1.7 return [dbsafe_decode(v) for v in values] return values - - diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index ceabb72..bc21f97 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -215,13 +215,19 @@ def test_async(broker, admin_user): assert result(result_j.name) == result_j.result # groups 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_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', failures=False)[0].id == [result_j][0].id assert count_group('test_j') == 1 + assert result_j.group_count() == 1 assert count_group('test_j', failures=True) == 0 + assert result_j.group_count(failures=True) == 0 assert delete_group('test_j') == 1 + assert result_j.group_delete() == 0 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 assert fetch(k) is None broker.delete_queue()