Merge pull request #48 from Koed00/dev

adds `catch_up` configuration option
This commit is contained in:
Ilan Steemers
2015-08-18 15:19:35 -07:00
3 changed files with 46 additions and 18 deletions
+18 -16
View File
@@ -436,20 +436,23 @@ def scheduler(list_key=Conf.Q_LIST):
# set up the next run time # set up the next run time
if not s.schedule_type == s.ONCE: if not s.schedule_type == s.ONCE:
next_run = arrow.get(s.next_run) next_run = arrow.get(s.next_run)
if s.schedule_type == s.MINUTES: while True:
next_run = next_run.replace(minutes=+(s.minutes or 1)) if s.schedule_type == s.MINUTES:
elif s.schedule_type == s.HOURLY: next_run = next_run.replace(minutes=+(s.minutes or 1))
next_run = next_run.replace(hours=+1) elif s.schedule_type == s.HOURLY:
elif s.schedule_type == s.DAILY: next_run = next_run.replace(hours=+1)
next_run = next_run.replace(days=+1) elif s.schedule_type == s.DAILY:
elif s.schedule_type == s.WEEKLY: next_run = next_run.replace(days=+1)
next_run = next_run.replace(weeks=+1) elif s.schedule_type == s.WEEKLY:
elif s.schedule_type == s.MONTHLY: next_run = next_run.replace(weeks=+1)
next_run = next_run.replace(months=+1) elif s.schedule_type == s.MONTHLY:
elif s.schedule_type == s.QUARTERLY: next_run = next_run.replace(months=+1)
next_run = next_run.replace(months=+3) elif s.schedule_type == s.QUARTERLY:
elif s.schedule_type == s.YEARLY: next_run = next_run.replace(months=+3)
next_run = next_run.replace(years=+1) elif s.schedule_type == s.YEARLY:
next_run = next_run.replace(years=+1)
if Conf.CATCH_UP or next_run > arrow.utcnow():
break
s.next_run = next_run.datetime s.next_run = next_run.datetime
s.repeats += -1 s.repeats += -1
# send it to the cluster # send it to the cluster
@@ -467,8 +470,7 @@ def scheduler(list_key=Conf.Q_LIST):
# default behavior is to delete a ONCE schedule # default behavior is to delete a ONCE schedule
if s.schedule_type == s.ONCE: if s.schedule_type == s.ONCE:
if s.repeats < 0: if s.repeats < 0:
s.delete() return s.delete()
return
# but not if it has a positive repeats # but not if it has a positive repeats
s.repeats = 0 s.repeats = 0
# save the schedule # save the schedule
+4
View File
@@ -77,6 +77,10 @@ class Conf(object):
# Global sync option to for debugging # Global sync option to for debugging
SYNC = conf.get('sync', False) SYNC = conf.get('sync', False)
# If set to False the scheduler won't execute tasks in the past.
# Instead it will reschedule the next run in the future. Defaults to True.
CATCH_UP = conf.get('catch_up', True)
# Use the secret key for package signing # Use the secret key for package signing
# Django itself should raise an error if it's not configured # Django itself should raise an error if it's not configured
SECRET_KEY = settings.SECRET_KEY SECRET_KEY = settings.SECRET_KEY
+24 -2
View File
@@ -1,10 +1,12 @@
from datetime import timedelta
from multiprocessing import Queue, Event, Value from multiprocessing import Queue, Event, Value
import pytest import pytest
import arrow import arrow
from django.utils import timezone from django.utils import timezone
from django_q.conf import redis_client from django_q.conf import redis_client, Conf
from django_q.cluster import pusher, worker, monitor, scheduler from django_q.cluster import pusher, worker, monitor, scheduler
from django_q.tasks import Schedule, fetch, schedule as create_schedule, queue_size from django_q.tasks import Schedule, fetch, schedule as create_schedule, queue_size
@@ -34,7 +36,7 @@ def test_scheduler(r):
# push it # push it
pusher(task_queue, stop_event, list_key=list_key, r=r) pusher(task_queue, stop_event, list_key=list_key, r=r)
assert task_queue.qsize() == 1 assert task_queue.qsize() == 1
assert queue_size(list_key=list_key,r=r) == 0 assert queue_size(list_key=list_key, r=r) == 0
task_queue.put('STOP') task_queue.put('STOP')
# let a worker handle them # let a worker handle them
result_queue = Queue() result_queue = Queue()
@@ -96,7 +98,27 @@ def test_scheduler(r):
kwargs='word="django"', kwargs='word="django"',
schedule_type=Schedule.DAILY schedule_type=Schedule.DAILY
) )
# scheduler
scheduler(list_key=list_key) scheduler(list_key=list_key)
# ONCE schedule should be deleted # ONCE schedule should be deleted
assert Schedule.objects.filter(pk=once_schedule.pk).exists() is False assert Schedule.objects.filter(pk=once_schedule.pk).exists() is False
# Catch up On
Conf.CATCH_UP = True
now = timezone.now()
schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='catch_up',
schedule_type=Schedule.HOURLY,
next_run=timezone.now() - timedelta(hours=12),
repeats=-1
)
scheduler(list_key=list_key)
schedule = Schedule.objects.get(pk=schedule.pk)
assert schedule.next_run < now
# Catch up off
Conf.CATCH_UP = False
scheduler(list_key=list_key)
schedule = Schedule.objects.get(pk=schedule.pk)
assert schedule.next_run > now
# Done
r.delete(list_key) r.delete(list_key)