Merge pull request #138 from Koed00/dev

Adds task result update
This commit is contained in:
Ilan Steemers
2016-01-24 21:39:44 +01:00
4 changed files with 79 additions and 22 deletions
+22 -13
View File
@@ -408,19 +408,28 @@ def save_task(task, broker):
try: try:
if task['success'] and 0 < Conf.SAVE_LIMIT <= Success.objects.count(): if task['success'] and 0 < Conf.SAVE_LIMIT <= Success.objects.count():
Success.objects.last().delete() Success.objects.last().delete()
Task.objects.update_or_create(id=task['id'], # check if this task has previous results
name=task['name'], if Task.objects.filter(id=task['id'], name=task['name']).exists():
defaults={ existing_task = Task.objects.get(id=task['id'], name=task['name'])
'func': task['func'], # only update the result if it hasn't succeeded yet
'hook': task.get('hook'), if not existing_task.success:
'args': task['args'], existing_task.stopped = task['stopped']
'kwargs': task['kwargs'], existing_task.result = task['result']
'started': task['started'], existing_task.success = task['success']
'stopped': task['stopped'], existing_task.save()
'result': task['result'], else:
'group': task.get('group'), Task.objects.create(id=task['id'],
'success': task['success']} name=task['name'],
) func=task['func'],
hook=task.get('hook'),
args=task['args'],
kwargs=task['kwargs'],
started=task['started'],
stopped=task['stopped'],
result=task['result'],
group=task.get('group'),
success=task['success']
)
except Exception as e: except Exception as e:
logger.error(e) logger.error(e)
+50 -4
View File
@@ -1,16 +1,17 @@
import sys import sys
from multiprocessing import Queue, Event, Value
import threading import threading
from multiprocessing import Queue, Event, Value
from time import sleep from time import sleep
import os from django.utils import timezone
import os
import pytest import pytest
myPath = os.path.dirname(os.path.abspath(__file__)) myPath = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, myPath + '/../') sys.path.insert(0, myPath + '/../')
from django_q.cluster import Cluster, Sentinel, pusher, worker, monitor from django_q.cluster import Cluster, Sentinel, pusher, worker, monitor, save_task
from django_q.humanhash import DEFAULT_WORDLIST from django_q.humanhash import DEFAULT_WORDLIST, uuid
from django_q.tasks import fetch, fetch_group, async, result, result_group, count_group, delete_group, queue_size from django_q.tasks import fetch, fetch_group, async, result, result_group, count_group, delete_group, queue_size
from django_q.models import Task, Success from django_q.models import Task, Success
from django_q.conf import Conf from django_q.conf import Conf
@@ -340,6 +341,51 @@ def test_bad_secret(broker, monkeypatch):
broker.delete_queue() broker.delete_queue()
@pytest.mark.django_db
def test_update_failed(broker):
tag = uuid()
task = {'id': tag[1],
'name': tag[0],
'func': 'math.copysign',
'args': (1, -1),
'kwargs': {},
'started': timezone.now(),
'stopped': timezone.now(),
'success': False,
'result': None}
# initial save - no success
save_task(task, broker)
assert Task.objects.filter(id=task['id']).exists()
saved_task = Task.objects.get(id=task['id'])
assert saved_task.success is False
sleep(0.5)
# second save - no success
old_stopped = task['stopped']
task['stopped']=timezone.now()
save_task(task, broker)
saved_task = Task.objects.get(id=task['id'])
assert saved_task.stopped > old_stopped
# third save - success
task['stopped']=timezone.now()
task['result']='result'
task['success']=True
save_task(task, broker)
saved_task = Task.objects.get(id=task['id'])
assert saved_task.success is True
# fourth save - no success
task['result'] = None
task['success'] = False
task['stopped'] = old_stopped
save_task(task, broker)
# should not overwrite success
saved_task = Task.objects.get(id=task['id'])
assert saved_task.success is True
assert saved_task.result == 'result'
@pytest.mark.django_db @pytest.mark.django_db
def assert_result(task): def assert_result(task):
assert task is not None assert task is not None
+5 -2
View File
@@ -18,8 +18,11 @@ Some pointers:
* Don't set the :ref:`retry` timer to a lower or equal number than the task timeout. * Don't set the :ref:`retry` timer to a lower or equal number than the task timeout.
* Retry time includes time the task spends waiting in the clusters internal queue. * Retry time includes time the task spends waiting in the clusters internal queue.
* Don't set the :ref:`queue_limit` so high that tasks time out while waiting to be processed. * Don't set the :ref:`queue_limit` so high that tasks time out while waiting to be processed.
* In case a task is worked on twice, you will see a duplicate key error in the cluster logs. * In case a task is worked on twice, the task result will be updated with the latest results.
* Duplicate tasks do generate additional receipt messages, but the result is discarded in favor of the first result. * In some rare cases a non-atomic broker will re-queue a task after it has been acknowledged.
* If a task runs twice and a previous run has succeeded, the new result wil be discarded.
* Limiting the number of retries is handled globally in your actual broker's settings.
Support for more brokers is being worked on. Support for more brokers is being worked on.
+2 -3
View File
@@ -7,17 +7,16 @@
arrow==0.7.0 arrow==0.7.0
blessed==1.14.1 blessed==1.14.1
boto3==1.2.3 boto3==1.2.3
botocore==1.3.20 # via boto3 botocore==1.3.22 # via boto3
django-picklefield==0.3.2 django-picklefield==0.3.2
django-redis==4.3.0 django-redis==4.3.0
docutils==0.12 # via botocore docutils==0.12 # via botocore
future==0.15.2 future==0.15.2
futures==3.0.4 # via boto3
hiredis==0.2.0 hiredis==0.2.0
iron-core==1.2.0 # via iron-mq iron-core==1.2.0 # via iron-mq
iron-mq==0.8 iron-mq==0.8
jmespath==0.9.0 # via boto3, botocore jmespath==0.9.0 # via boto3, botocore
psutil==3.4.1 psutil==3.4.2
pymongo==3.2 pymongo==3.2
python-dateutil==2.4.2 # via arrow, botocore, iron-core python-dateutil==2.4.2 # via arrow, botocore, iron-core
redis==2.10.5 redis==2.10.5