Compare commits

...
10 Commits
45 changed files with 802 additions and 397 deletions
+2
View File
@@ -0,0 +1,2 @@
# flake8, black, isort
b1d000d007f3f77069719523268a0c6256dc0860
-37
View File
@@ -1,37 +0,0 @@
categories:
-
label: breaking
title: Breaking
-
label: feature
title: New
-
label: bug
title: "Bug Fixes"
-
label: dependencies
title: "Dependency Updates"
-
label: security
title: Security
name-template: v$NEXT_PATCH_VERSION
tag-template: v$NEXT_PATCH_VERSION
template: |
# Changes
$CHANGES
version-resolver:
major:
labels:
- breaking
- major
minor:
labels:
- feature
- minor
patch:
labels:
- bug
- dependencies
- security
- patch
default: patch
+2 -2
View File
@@ -22,8 +22,8 @@ jobs:
- name: Install dependencies - name: Install dependencies
run: | run: |
apt-get update sudo apt-get update
apt-get -y install gettext sudo apt-get -y install gettext
python -m pip install pip setuptools django poetry python -m pip install pip setuptools django poetry
# compile messages to get .mo files # compile messages to get .mo files
django-admin compilemessages django-admin compilemessages
-14
View File
@@ -1,14 +0,0 @@
name: Update release draft
on:
push:
branches:
- master
jobs:
update_release_draft:
runs-on: ubuntu-latest
steps:
- uses: release-drafter/release-drafter@v5
with:
config-name: release-drafter.yml
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
+5 -1
View File
@@ -73,7 +73,11 @@ jobs:
- name: Upload to coveralls - name: Upload to coveralls
run: | run: |
python -m pip install --upgrade pip python -m pip install --upgrade pip
python -m pip install coveralls python -m pip install coveralls flake8 black
coveralls --service=github --finish coveralls --service=github --finish
env: env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
- name: Check flake8/black
run: |
flake8 .
black --check .
+20
View File
@@ -2,6 +2,26 @@
## [Unreleased](https://github.com/GDay/django-q2/tree/HEAD) ## [Unreleased](https://github.com/GDay/django-q2/tree/HEAD)
## [v1.4.7](https://github.com/GDay/django-q2/tree/v1.4.7) (2022-12-21)
**Merged pull requests:**
- Fix: handling exceptions inside job function https://github.com/GDay/django-q2/pull/51
- Fix: Daylight saving time issue with scheduler https://github.com/GDay/django-q2/pull/47
- Chore: Fix badge and add download badge https://github.com/GDay/django-q2/pull/52
- Chore: Remove release drafter https://github.com/GDay/django-q2/pull/53
## [v1.4.6](https://github.com/GDay/django-q2/tree/v1.4.6) (2022-11-30)
**Merged pull requests:**
- Fix: Log exceptions with logger.exception https://github.com/GDay/django-q2/pull/42
- Chore: flake8, isort, black https://github.com/GDay/django-q2/pull/40
## [v1.4.5](https://github.com/GDay/django-q2/tree/v1.4.5) (2022-11-13)
- Fix release workflow
## [v1.4.4](https://github.com/GDay/django-q2/tree/v1.4.4) (2022-11-13) ## [v1.4.4](https://github.com/GDay/django-q2/tree/v1.4.4) (2022-11-13)
**Merged pull requests:** **Merged pull requests:**
+4 -2
View File
@@ -1,7 +1,7 @@
A multiprocessing distributed task queue for Django A multiprocessing distributed task queue for Django
--------------------------------------------------- ---------------------------------------------------
|image0| |image1| |docs| |image0| |image1| |docs| |downloads|
:: ::
@@ -245,4 +245,6 @@ Acknowledgements
.. |docs| image:: https://readthedocs.org/projects/docs/badge/?version=latest .. |docs| image:: https://readthedocs.org/projects/docs/badge/?version=latest
:alt: Documentation Status :alt: Documentation Status
:scale: 100 :scale: 100
:target: https://django-q.readthedocs.org/ :target: https://django-q2.readthedocs.org/
.. |downloads| image:: https://img.shields.io/pypi/dm/django-q2
:target: https://img.shields.io/pypi/dm/django-q2
+2 -2
View File
@@ -1,7 +1,7 @@
VERSION = (1, 3, 9)
import django import django
VERSION = (1, 4, 7)
if django.VERSION < (3, 2): if django.VERSION < (3, 2):
default_app_config = "django_q.apps.DjangoQConfig" default_app_config = "django_q.apps.DjangoQConfig"
+17 -7
View File
@@ -1,9 +1,9 @@
"""Admin module for Django.""" """Admin module for Django."""
from django.contrib import admin
from django.db.models.expressions import OuterRef, Subquery
from django.urls import reverse from django.urls import reverse
from django.utils.html import format_html from django.utils.html import format_html
from django.contrib import admin
from django.utils.translation import gettext_lazy as _ from django.utils.translation import gettext_lazy as _
from django.db.models.expressions import OuterRef, Subquery
from django_q.conf import Conf, croniter from django_q.conf import Conf, croniter
from django_q.models import Failure, OrmQ, Schedule, Success, Task from django_q.models import Failure, OrmQ, Schedule, Success, Task
@@ -82,18 +82,27 @@ class ScheduleAdmin(admin.ModelAdmin):
readonly_fields = ("cron",) readonly_fields = ("cron",)
list_filter = ("next_run", "schedule_type", "cluster") list_filter = ("next_run", "schedule_type", "cluster")
search_fields = ("name", "func",) search_fields = (
"name",
"func",
)
list_display_links = ("id", "name") list_display_links = ("id", "name")
def get_queryset(self, request): def get_queryset(self, request):
qs = super().get_queryset(request) qs = super().get_queryset(request)
task_query = Task.objects.filter(id=OuterRef('task')).values('id', 'name', 'success') task_query = Task.objects.filter(id=OuterRef("task")).values(
qs = qs.annotate(task_id=Subquery(task_query.values('id')), task_name=Subquery(task_query.values('name')), "id", "name", "success"
task_success=Subquery(task_query.values('success'))) )
qs = qs.annotate(
task_id=Subquery(task_query.values("id")),
task_name=Subquery(task_query.values("name")),
task_success=Subquery(task_query.values("success")),
)
return qs return qs
def get_success(self, obj): def get_success(self, obj):
return obj.task_success return obj.task_success
get_success.boolean = True get_success.boolean = True
get_success.short_description = _("success") get_success.short_description = _("success")
@@ -105,6 +114,7 @@ class ScheduleAdmin(admin.ModelAdmin):
url = reverse("admin:django_q_failure_change", args=(obj.task_id,)) url = reverse("admin:django_q_failure_change", args=(obj.task_id,))
return format_html(f'<a href="{url}">[{obj.task_name}]</a>') return format_html(f'<a href="{url}">[{obj.task_name}]</a>')
return None return None
get_last_run.allow_tags = True get_last_run.allow_tags = True
get_last_run.short_description = _("last_run") get_last_run.short_description = _("last_run")
@@ -112,7 +122,7 @@ class ScheduleAdmin(admin.ModelAdmin):
class QueueAdmin(admin.ModelAdmin): class QueueAdmin(admin.ModelAdmin):
"""queue admin for ORM broker""" """queue admin for ORM broker"""
list_display = ("id", "key", "name", "group", "func", "lock", "task_id") list_display = ("id", "key", "name", "group", "func", "lock", "task_id")
def save_model(self, request, obj, form, change): def save_model(self, request, obj, form, change):
obj.save(using=Conf.ORM) obj.save(using=Conf.ORM)
+1 -1
View File
@@ -9,4 +9,4 @@ class DjangoQConfig(AppConfig):
default_auto_field = "django.db.models.AutoField" default_auto_field = "django.db.models.AutoField"
def ready(self): def ready(self):
from django_q.signals import call_hook from django_q.signals import call_hook # noqa: F401
+2 -1
View File
@@ -39,7 +39,8 @@ class Sqs(Broker):
raise ValueError("receive_message_wait_time_seconds should be int") raise ValueError("receive_message_wait_time_seconds should be int")
if wait_time_second > 20: if wait_time_second > 20:
raise ValueError( raise ValueError(
"receive_message_wait_time_seconds is invalid. Reason: Must be >= 0 and <= 20" "receive_message_wait_time_seconds is invalid. Reason: Must be >= 0"
" and <= 20"
) )
params.update({"WaitTimeSeconds": wait_time_second}) params.update({"WaitTimeSeconds": wait_time_second})
+3 -2
View File
@@ -62,7 +62,7 @@ class ORM(Broker):
def dequeue(self): def dequeue(self):
tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[ tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[
0 : Conf.BULK 0 : Conf.BULK # noqa: E203
] ]
if tasks: if tasks:
task_list = [] task_list = []
@@ -73,7 +73,8 @@ class ORM(Broker):
.update(lock=timezone.now()) .update(lock=timezone.now())
): ):
task_list.append((task.pk, task.payload)) task_list.append((task.pk, task.payload))
# else don't process, as another cluster has been faster than us on that task # else don't process, as another cluster has been faster than us on
# that task
return task_list return task_list
# empty queue, spare the cpu # empty queue, spare the cpu
sleep(Conf.POLL) sleep(Conf.POLL)
+133 -81
View File
@@ -6,6 +6,7 @@ import socket
import traceback import traceback
import uuid import uuid
from datetime import datetime, timedelta from datetime import datetime, timedelta
from pytz import timezone as pytz_timezone
from multiprocessing import Event, Process, Value, current_process from multiprocessing import Event, Process, Value, current_process
from time import sleep from time import sleep
@@ -43,7 +44,7 @@ from django_q.signals import post_execute, pre_execute
from django_q.signing import BadSignature, SignedPackage from django_q.signing import BadSignature, SignedPackage
from django_q.status import Stat, Status from django_q.status import Stat, Status
from .utils import add_months, add_years, get_func_repr from .utils import get_func_repr, localtime
class Cluster: class Cluster:
@@ -74,7 +75,7 @@ class Cluster:
), ),
) )
self.sentinel.start() self.sentinel.start()
logger.info(_("Q Cluster %(name)s starting.") % {'name': self.name}) logger.info(_("Q Cluster %(name)s starting.") % {"name": self.name})
while not self.start_event.is_set(): while not self.start_event.is_set():
sleep(0.1) sleep(0.1)
return self.pid return self.pid
@@ -82,19 +83,21 @@ class Cluster:
def stop(self) -> bool: def stop(self) -> bool:
if not self.sentinel.is_alive(): if not self.sentinel.is_alive():
return False return False
logger.info(_("Q Cluster %(name)s stopping.") % {'name': self.name}) logger.info(_("Q Cluster %(name)s stopping.") % {"name": self.name})
self.stop_event.set() self.stop_event.set()
self.sentinel.join() self.sentinel.join()
logger.info(_("Q Cluster %(name)s has stopped.") % {'name': self.name}) logger.info(_("Q Cluster %(name)s has stopped.") % {"name": self.name})
self.start_event = None self.start_event = None
self.stop_event = None self.stop_event = None
return True return True
def sig_handler(self, signum, frame): def sig_handler(self, signum, frame):
logger.debug( logger.debug(
_( _("%(name)s got signal %(signal)s")
'%(name)s got signal %(signal)s' % {
) % {'name': current_process().name, 'signal': Conf.SIGNAL_NAMES.get(signum, "UNKNOWN")} "name": current_process().name,
"signal": Conf.SIGNAL_NAMES.get(signum, "UNKNOWN"),
}
) )
self.stop() self.stop()
@@ -216,21 +219,34 @@ class Sentinel:
db.connections.close_all() db.connections.close_all()
if process == self.monitor: if process == self.monitor:
self.monitor = self.spawn_monitor() self.monitor = self.spawn_monitor()
logger.error(_("reincarnated monitor %(name)s after sudden death") % {'name': process.name}) logger.error(
_("reincarnated monitor %(name)s after sudden death")
% {"name": process.name}
)
elif process == self.pusher: elif process == self.pusher:
self.pusher = self.spawn_pusher() self.pusher = self.spawn_pusher()
logger.error(_("reincarnated pusher %(name)s after sudden death") % {'name': process.name}) logger.error(
_("reincarnated pusher %(name)s after sudden death")
% {"name": process.name}
)
else: else:
self.pool.remove(process) self.pool.remove(process)
self.spawn_worker() self.spawn_worker()
if process.timer.value == 0: if process.timer.value == 0:
# only need to terminate on timeout, otherwise we risk destabilizing the queues # only need to terminate on timeout, otherwise we risk destabilizing
# the queues
process.terminate() process.terminate()
logger.warning(_("reincarnated worker %(name)s after timeout") % {'name': process.name}) logger.warning(
_("reincarnated worker %(name)s after timeout")
% {"name": process.name}
)
elif int(process.timer.value) == -2: elif int(process.timer.value) == -2:
logger.info(_("recycled worker %(name)s") % {'name': process.name}) logger.info(_("recycled worker %(name)s") % {"name": process.name})
else: else:
logger.error(_("reincarnated worker %(name)s after death") % {'name': process.name}) logger.error(
_("reincarnated worker %(name)s after death")
% {"name": process.name}
)
self.reincarnations += 1 self.reincarnations += 1
@@ -252,13 +268,18 @@ class Sentinel:
def guard(self): def guard(self):
logger.info( logger.info(
_( _("%(name)s guarding cluster %(cluster_name)s")
"%(name)s guarding cluster %(cluster_name)s" % {
) % {'name': current_process().name, 'cluster_name': humanize(self.cluster_id.hex)} "name": current_process().name,
"cluster_name": humanize(self.cluster_id.hex),
}
) )
self.start_event.set() self.start_event.set()
Stat(self).save() Stat(self).save()
logger.info(_("Q Cluster %(cluster_name)s running.") % {'cluster_name': humanize(self.cluster_id.hex)}) logger.info(
_("Q Cluster %(cluster_name)s running.")
% {"cluster_name": humanize(self.cluster_id.hex)}
)
counter = 0 counter = 0
cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds
# Guard loop. Runs at least once # Guard loop. Runs at least once
@@ -292,7 +313,7 @@ class Sentinel:
def stop(self): def stop(self):
Stat(self).save() Stat(self).save()
name = current_process().name name = current_process().name
logger.info(_("%(name)s stopping cluster processes") % {'name': name}) logger.info(_("%(name)s stopping cluster processes") % {"name": name})
# Stopping pusher # Stopping pusher
self.event_out.set() self.event_out.set()
# Wait for it to stop # Wait for it to stop
@@ -317,7 +338,7 @@ class Sentinel:
self.result_queue.close() self.result_queue.close()
# Wait for the result queue to empty # Wait for the result queue to empty
self.result_queue.join_thread() self.result_queue.join_thread()
logger.info(_("%(name)s waiting for the monitor.") % {'name': name}) logger.info(_("%(name)s waiting for the monitor.") % {"name": name})
# Wait for everything to close or time out # Wait for everything to close or time out
count = 0 count = 0
if not self.timeout: if not self.timeout:
@@ -339,12 +360,15 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
""" """
if not broker: if not broker:
broker = get_broker() broker = get_broker()
logger.info(_("%(process_name)s pushing tasks at %(id)s") % {'process_name': current_process().name, 'id': current_process().pid}) logger.info(
_("%(process_name)s pushing tasks at %(id)s")
% {"process_name": current_process().name, "id": current_process().pid}
)
while True: while True:
try: try:
task_set = broker.dequeue() task_set = broker.dequeue()
except Exception as e: except Exception as e:
logger.error(e, traceback.format_exc()) logger.exception("Failed to pull task from broker")
# broker probably crashed. Let the sentinel handle it. # broker probably crashed. Let the sentinel handle it.
sleep(10) sleep(10)
break break
@@ -355,15 +379,17 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
try: try:
task = SignedPackage.loads(task[1]) task = SignedPackage.loads(task[1])
except (TypeError, BadSignature) as e: except (TypeError, BadSignature) as e:
logger.error(e, traceback.format_exc()) logger.exception("Failed to push task to queue")
broker.fail(ack_id) broker.fail(ack_id)
continue continue
task["ack_id"] = ack_id task["ack_id"] = ack_id
task_queue.put(task) task_queue.put(task)
logger.debug(_("queueing from %(list_key)s") % {'list_key': broker.list_key}) logger.debug(
_("queueing from %(list_key)s") % {"list_key": broker.list_key}
)
if event.is_set(): if event.is_set():
break break
logger.info(_("%(name)s stopped pushing tasks") % {'name': current_process().name}) logger.info(_("%(name)s stopped pushing tasks") % {"name": current_process().name})
def monitor(result_queue: Queue, broker: Broker = None): def monitor(result_queue: Queue, broker: Broker = None):
@@ -375,7 +401,9 @@ def monitor(result_queue: Queue, broker: Broker = None):
if not broker: if not broker:
broker = get_broker() broker = get_broker()
name = current_process().name name = current_process().name
logger.info(_("%(name)s monitoring at %(id)s") % {'name': name, 'id': current_process().pid}) logger.info(
_("%(name)s monitoring at %(id)s") % {"name": name, "id": current_process().pid}
)
for task in iter(result_queue.get, "STOP"): for task in iter(result_queue.get, "STOP"):
# save the result # save the result
if task.get("cached", False): if task.get("cached", False):
@@ -389,28 +417,42 @@ def monitor(result_queue: Queue, broker: Broker = None):
# signal execution done # signal execution done
post_execute.send(sender="django_q", task=task) post_execute.send(sender="django_q", task=task)
# log the result # log the result
info_name = get_func_repr(task['func']) info_name = get_func_repr(task["func"])
if task["success"]: if task["success"]:
# log success # log success
logger.info(_("Processed '%(info_name)s' (%(task_name)s)") % {'info_name': info_name, 'task_name': task['name']}) logger.info(
_("Processed '%(info_name)s' (%(task_name)s)")
% {"info_name": info_name, "task_name": task["name"]}
)
else: else:
# log failure # log failure
logger.error(_("Failed '%(info_name)s' (%(task_name)s) - %(task_result)s") % {'info_name': info_name, 'task_name': task['name'], 'task_result': task['result']}) logger.error(
logger.info(_("%(name)s stopped monitoring results") % {'name': name}) _("Failed '%(info_name)s' (%(task_name)s) - %(task_result)s")
% {
"info_name": info_name,
"task_name": task["name"],
"task_result": task["result"],
}
)
logger.info(_("%(name)s stopped monitoring results") % {"name": name})
def worker( def worker(
task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT
): ):
""" """
Takes a task from the task queue, tries to execute it and puts the result back in the result queue Takes a task from the task queue, tries to execute it and puts the result back in
the result queue
:param timeout: number of seconds wait for a worker to finish. :param timeout: number of seconds wait for a worker to finish.
:type task_queue: multiprocessing.Queue :type task_queue: multiprocessing.Queue
:type result_queue: multiprocessing.Queue :type result_queue: multiprocessing.Queue
:type timer: multiprocessing.Value :type timer: multiprocessing.Value
""" """
proc_name = current_process().name proc_name = current_process().name
logger.info(_("%(proc_name)s ready for work at %(id)s") % {'proc_name': proc_name, 'id': current_process().pid}) logger.info(
_("%(proc_name)s ready for work at %(id)s")
% {"proc_name": proc_name, "id": current_process().pid}
)
task_count = 0 task_count = 0
if timeout is None: if timeout is None:
timeout = -1 timeout = -1
@@ -422,7 +464,14 @@ def worker(
# Get the function from the task # Get the function from the task
func = task["func"] func = task["func"]
func_name = get_func_repr(func) func_name = get_func_repr(func)
logger.info(_("%(proc_name)s processing '%(func_name)s' (%(task_name)s)") % {'proc_name': proc_name, 'func_name': func_name, 'task_name': task['name']}) logger.info(
_("%(proc_name)s processing '%(func_name)s' (%(task_name)s)")
% {
"proc_name": proc_name,
"func_name": func_name,
"task_name": task["name"],
}
)
f = task["func"] f = task["func"]
# if it's not an instance try to get it from the string # if it's not an instance try to get it from the string
if not callable(task["func"]): if not callable(task["func"]):
@@ -436,12 +485,12 @@ def worker(
try: try:
res = f(*task["args"], **task["kwargs"]) res = f(*task["args"], **task["kwargs"])
result = (res, True) result = (res, True)
except Exception: except Exception as e:
result = (_("Could not process '%(func_name)s'. Check the location of the function and the args/kwargs.") % {'func_name': func_name}, False) result = (f"{e} : {traceback.format_exc()}", False)
if error_reporter: if error_reporter:
error_reporter.report() error_reporter.report()
if task.get("sync", False): if task.get("sync", False):
raise Exception(result) raise
with timer.get_lock(): with timer.get_lock():
# Process result # Process result
task["result"] = result[0] task["result"] = result[0]
@@ -453,7 +502,8 @@ def worker(
if task_count == Conf.RECYCLE or rss_check(): if task_count == Conf.RECYCLE or rss_check():
timer.value = -2 # Recycled timer.value = -2 # Recycled
break break
logger.info(_("%(proc_name)s stopped doing work") % {'proc_name': proc_name}) logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name})
def save_task(task, broker: Broker): def save_task(task, broker: Broker):
""" """
@@ -478,16 +528,27 @@ def save_task(task, broker: Broker):
try: try:
filters = {} filters = {}
if Conf.SAVE_LIMIT_PER and Conf.SAVE_LIMIT_PER in {"group", "name", "func"} and Conf.SAVE_LIMIT_PER in task: if (
Conf.SAVE_LIMIT_PER
and Conf.SAVE_LIMIT_PER in {"group", "name", "func"}
and Conf.SAVE_LIMIT_PER in task
):
value = task[Conf.SAVE_LIMIT_PER] value = task[Conf.SAVE_LIMIT_PER]
if Conf.SAVE_LIMIT_PER == "func": if Conf.SAVE_LIMIT_PER == "func":
value = get_func_repr(value) value = get_func_repr(value)
filters[Conf.SAVE_LIMIT_PER] = value filters[Conf.SAVE_LIMIT_PER] = value
database_to_use = {"using": Conf.ORM if Conf.ORM else Schedule.objects.db} if not Conf.HAS_REPLICA else {} database_to_use = (
{"using": Conf.ORM if Conf.ORM else Schedule.objects.db}
if not Conf.HAS_REPLICA
else {}
)
with db.transaction.atomic(**database_to_use): with db.transaction.atomic(**database_to_use):
last = Success.objects.filter(**filters).select_for_update().last() last = Success.objects.filter(**filters).select_for_update().last()
if task["success"] and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count(): if (
task["success"]
and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count()
):
last.delete() last.delete()
# check if this task has previous results # check if this task has previous results
@@ -588,7 +649,11 @@ def scheduler(broker: Broker = None):
broker = get_broker() broker = get_broker()
close_old_django_connections() close_old_django_connections()
try: try:
database_to_use = {"using": Conf.ORM if Conf.ORM else Schedule.objects.db} if not Conf.HAS_REPLICA else {} database_to_use = (
{"using": Conf.ORM if Conf.ORM else Schedule.objects.db}
if not Conf.HAS_REPLICA
else {}
)
with db.transaction.atomic(**database_to_use): with db.transaction.atomic(**database_to_use):
for s in ( for s in (
Schedule.objects.select_for_update() Schedule.objects.select_for_update()
@@ -608,8 +673,13 @@ def scheduler(broker: Broker = None):
except (SyntaxError, ValueError): except (SyntaxError, ValueError):
# else use the kwargs syntax # else use the kwargs syntax
try: try:
parsed_kwargs = ast.parse(f"f({s.kwargs})").body[0].value.keywords parsed_kwargs = (
kwargs = {kwarg.arg: ast.literal_eval(kwarg.value) for kwarg in parsed_kwargs} ast.parse(f"f({s.kwargs})").body[0].value.keywords
)
kwargs = {
kwarg.arg: ast.literal_eval(kwarg.value)
for kwarg in parsed_kwargs
}
except (SyntaxError, ValueError): except (SyntaxError, ValueError):
kwargs = {} kwargs = {}
if s.args: if s.args:
@@ -624,32 +694,7 @@ def scheduler(broker: Broker = None):
if s.schedule_type != s.ONCE: if s.schedule_type != s.ONCE:
next_run = s.next_run next_run = s.next_run
while True: while True:
if s.schedule_type == s.MINUTES: next_run = s.calculate_next_run(next_run)
next_run = next_run + timedelta(minutes=(s.minutes or 1))
elif s.schedule_type == s.HOURLY:
next_run = next_run + timedelta(hours=1)
elif s.schedule_type == s.DAILY:
next_run = next_run + timedelta(days=1)
elif s.schedule_type == s.WEEKLY:
next_run = next_run + timedelta(weeks=1)
elif s.schedule_type == s.BIWEEKLY:
next_run = next_run + timedelta(weeks=2)
elif s.schedule_type == s.MONTHLY:
next_run = add_months(next_run, 1)
elif s.schedule_type == s.BIMONTHLY:
next_run = add_months(next_run, 2)
elif s.schedule_type == s.QUARTERLY:
next_run = add_months(next_run, 3)
elif s.schedule_type == s.YEARLY:
next_run = add_years(next_run, 1)
elif s.schedule_type == s.CRON:
if not croniter:
raise ImportError(
_(
"Please install croniter to enable cron expressions"
)
)
next_run = croniter(s.cron, localtime()).get_next(datetime)
if Conf.CATCH_UP or next_run > localtime(): if Conf.CATCH_UP or next_run > localtime():
break break
@@ -659,7 +704,8 @@ def scheduler(broker: Broker = None):
scheduled_broker = broker scheduled_broker = broker
try: try:
scheduled_broker = get_broker(q_options["broker_name"]) scheduled_broker = get_broker(q_options["broker_name"])
except: # invalid broker_name or non existing broker with broker_name except: # noqa: E722
# invalid broker_name or non existing broker with broker_name
pass pass
q_options["broker"] = scheduled_broker q_options["broker"] = scheduled_broker
q_options["group"] = q_options.get("group", s.name or s.id) q_options["group"] = q_options.get("group", s.name or s.id)
@@ -669,14 +715,24 @@ def scheduler(broker: Broker = None):
if not s.task: if not s.task:
logger.error( logger.error(
_( _(
"%(process_name)s failed to create a task from schedule [%(schedule)s]" "%(process_name)s failed to create a task from schedule "
) % {'process_name': current_process().name, 'schedule': s.name or s.id} "[%(schedule)s]"
)
% {
"process_name": current_process().name,
"schedule": s.name or s.id,
}
) )
else: else:
logger.info( logger.info(
_( _(
"%(process_name)s created a task from schedule [%(schedule)s]" "%(process_name)s created a task from schedule "
) % {'process_name': current_process().name, 'schedule': s.name or s.id} "[%(schedule)s]"
)
% {
"process_name": current_process().name,
"schedule": s.name or s.id,
}
) )
# 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:
@@ -741,7 +797,10 @@ def set_cpu_affinity(n: int, process_ids: list, actual: bool = not Conf.TESTING)
p = psutil.Process(pid) p = psutil.Process(pid)
if actual: if actual:
p.cpu_affinity(affinity) p.cpu_affinity(affinity)
logger.info(_("%(pid)s will use cpu %(affinity)s") % {'pid': pid, 'affinity': affinity}) logger.info(
_("%(pid)s will use cpu %(affinity)s")
% {"pid": pid, "affinity": affinity}
)
def rss_check(): def rss_check():
@@ -751,10 +810,3 @@ def rss_check():
elif psutil: elif psutil:
return psutil.Process().memory_info().rss >= Conf.MAX_RSS * 1024 return psutil.Process().memory_info().rss >= Conf.MAX_RSS * 1024
return False return False
def localtime() -> datetime:
"""Override for timezone.localtime to deal with naive times and local times"""
if settings.USE_TZ:
return timezone.localtime()
return datetime.now()
+28 -11
View File
@@ -73,7 +73,8 @@ class Conf:
# Log output level # Log output level
LOG_LEVEL = conf.get("log_level", "INFO") LOG_LEVEL = conf.get("log_level", "INFO")
# Maximum number of successful tasks kept in the database. 0 saves everything. -1 saves none # Maximum number of successful tasks kept in the database. 0 saves everything.
# -1 saves none
# Failures are always saved # Failures are always saved
SAVE_LIMIT = conf.get("save_limit", 250) SAVE_LIMIT = conf.get("save_limit", 250)
@@ -82,7 +83,13 @@ class Conf:
# Verify SAVE_LIMIT_PER is valid # Verify SAVE_LIMIT_PER is valid
if SAVE_LIMIT_PER not in ["group", "name", "func", None]: if SAVE_LIMIT_PER not in ["group", "name", "func", None]:
warn(_("SAVE_LIMIT_PER (%(option)s) is not a valid option. Options are: 'group', 'name', 'func' and None. Default is None.") % {'option': SAVE_LIMIT_PER}) warn(
_(
"SAVE_LIMIT_PER (%(option)s) is not a valid option. Options are: "
"'group', 'name', 'func' and None. Default is None."
)
% {"option": SAVE_LIMIT_PER}
)
# Guard loop sleep in seconds. Should be between 0 and 60 seconds. # Guard loop sleep in seconds. Should be between 0 and 60 seconds.
GUARD_CYCLE = conf.get("guard_cycle", 0.5) GUARD_CYCLE = conf.get("guard_cycle", 0.5)
@@ -113,11 +120,12 @@ class Conf:
# Sets compression of redis packages # Sets compression of redis packages
COMPRESSED = conf.get("compress", False) COMPRESSED = conf.get("compress", False)
# Number of tasks each worker can handle before it gets recycled. Useful for releasing memory # Number of tasks each worker can handle before it gets recycled.
# Useful for releasing memory
RECYCLE = conf.get("recycle", 500) RECYCLE = conf.get("recycle", 500)
# The maximum resident set size in kilobytes before a worker will recycle. Useful for limiting memory usage # The maximum resident set size in kilobytes before a worker will recycle.
# Not available on all platforms # Useful for limiting memory usage. Not available on all platforms
MAX_RSS = conf.get("max_rss", None) MAX_RSS = conf.get("max_rss", None)
# Number of seconds to wait for a worker to finish. # Number of seconds to wait for a worker to finish.
@@ -135,9 +143,10 @@ class Conf:
# Verify if retry and timeout settings are correct # Verify if retry and timeout settings are correct
if not TIMEOUT or (TIMEOUT > RETRY): if not TIMEOUT or (TIMEOUT > RETRY):
warn( warn(
"""Retry and timeout are misconfigured. Set retry larger than timeout, "Retry and timeout are misconfigured. Set retry larger than timeout,"
failure to do so will cause the tasks to be retriggered before completion. "failure to do so will cause the tasks to be retriggered before completion."
See https://django-q2.readthedocs.io/en/master/configure.html#retry for details.""" "See https://django-q2.readthedocs.io/en/master/configure.html#retry "
"for details."
) )
# Sets the amount of tasks the cluster will try to pop off the broker. # Sets the amount of tasks the cluster will try to pop off the broker.
@@ -156,12 +165,14 @@ class Conf:
# The Django cache to use # The Django cache to use
CACHE = conf.get("cache", "default") CACHE = conf.get("cache", "default")
# Use the cache as result backend. Can be 'True' or an integer representing the global cache timeout. # Use the cache as result backend. Can be 'True' or an integer representing the
# global cache timeout.
# i.e 'cached: 60' , will make all results go the cache and expire in 60 seconds. # i.e 'cached: 60' , will make all results go the cache and expire in 60 seconds.
CACHED = conf.get("cached", False) CACHED = conf.get("cached", False)
# If set to False the scheduler won't execute tasks in the past. # If set to False the scheduler won't execute tasks in the past.
# Instead it will run once and reschedule the next run in the future. Defaults to True. # Instead it will run once and reschedule the next run in the future. Defaults to
# True.
CATCH_UP = conf.get("catch_up", True) CATCH_UP = conf.get("catch_up", True)
# Use the secret key for package signing # Use the secret key for package signing
@@ -200,6 +211,11 @@ class Conf:
# to manage workarounds during testing # to manage workarounds during testing
TESTING = conf.get("testing", False) TESTING = conf.get("testing", False)
# Timezone for next_run, overrules Django timezone
TIME_ZONE = None
if settings.USE_TZ:
TIME_ZONE = conf.get("time_zone", settings.TIME_ZONE)
# logger # logger
logger = logging.getLogger("django-q") logger = logging.getLogger("django-q")
@@ -257,5 +273,6 @@ def get_ppid():
return psutil.Process(os.getpid()).ppid() return psutil.Process(os.getpid()).ppid()
else: else:
raise OSError( raise OSError(
"Your OS does not support `os.getppid`. Please install `psutil` as an alternative provider." "Your OS does not support `os.getppid`. Please install `psutil` as an "
"alternative provider."
) )
+1
View File
@@ -6,6 +6,7 @@ from django.core.signing import BadSignature, JSONSerializer, SignatureExpired
from django.core.signing import Signer as Sgnr from django.core.signing import Signer as Sgnr
from django.core.signing import TimestampSigner as TsS from django.core.signing import TimestampSigner as TsS
from django.core.signing import b64_decode, dumps from django.core.signing import b64_decode, dumps
try: try:
from django.core.signing import base62 from django.core.signing import base62
except ImportError: except ImportError:
+9 -3
View File
@@ -337,12 +337,18 @@ class HumanHasher:
# Split `bytes` into `target` segments. # Split `bytes` into `target` segments.
seg_size = length // target seg_size = length // target
segments = [bytes[i * seg_size : (i + 1) * seg_size] for i in range(target)] # fmt: off
segments = [
bytes[i * seg_size : (i + 1) * seg_size] for i in range(target) # noqa: E203 E501
]
# fmt: on
# Catch any left-over bytes in the last segment. # Catch any left-over bytes in the last segment.
segments[-1].extend(bytes[target * seg_size :]) segments[-1].extend(bytes[target * seg_size :]) # noqa: E203 E501
# Use a simple XOR checksum-like function for compression. # Use a simple XOR checksum-like function for compression.
checksum = lambda bytes: reduce(operator.xor, bytes, 0) def checksum(bytes):
return reduce(operator.xor, bytes, 0)
checksums = list(map(checksum, segments)) checksums = list(map(checksum, segments))
return checksums return checksums
+118 -37
View File
@@ -5,61 +5,142 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = []
]
operations = [ operations = [
migrations.CreateModel( migrations.CreateModel(
name='Schedule', name="Schedule",
fields=[ fields=[
('id', models.AutoField(verbose_name='ID', auto_created=True, serialize=False, primary_key=True)), (
('func', models.CharField(max_length=256, help_text='e.g. module.tasks.function')), "id",
('hook', models.CharField(null=True, blank=True, max_length=256, help_text='e.g. module.tasks.result_function')), models.AutoField(
('args', models.CharField(null=True, blank=True, max_length=256, help_text="e.g. 1, 2, 'John'")), verbose_name="ID",
('kwargs', models.CharField(null=True, blank=True, max_length=256, help_text="e.g. x=1, y=2, name='John'")), auto_created=True,
('schedule_type', models.CharField(verbose_name='Schedule Type', choices=[('O', 'Once'), ('H', 'Hourly'), ('D', 'Daily'), ('W', 'Weekly'), ('M', 'Monthly'), ('Q', 'Quarterly'), ('Y', 'Yearly')], default='O', max_length=1)), serialize=False,
('repeats', models.SmallIntegerField(verbose_name='Repeats', default=-1, help_text='n = n times, -1 = forever')), primary_key=True,
('next_run', models.DateTimeField(verbose_name='Next Run', default=django.utils.timezone.now, null=True)), ),
('task', models.CharField(editable=False, null=True, max_length=100)), ),
(
"func",
models.CharField(
max_length=256, help_text="e.g. module.tasks.function"
),
),
(
"hook",
models.CharField(
null=True,
blank=True,
max_length=256,
help_text="e.g. module.tasks.result_function",
),
),
(
"args",
models.CharField(
null=True,
blank=True,
max_length=256,
help_text="e.g. 1, 2, 'John'",
),
),
(
"kwargs",
models.CharField(
null=True,
blank=True,
max_length=256,
help_text="e.g. x=1, y=2, name='John'",
),
),
(
"schedule_type",
models.CharField(
verbose_name="Schedule Type",
choices=[
("O", "Once"),
("H", "Hourly"),
("D", "Daily"),
("W", "Weekly"),
("M", "Monthly"),
("Q", "Quarterly"),
("Y", "Yearly"),
],
default="O",
max_length=1,
),
),
(
"repeats",
models.SmallIntegerField(
verbose_name="Repeats",
default=-1,
help_text="n = n times, -1 = forever",
),
),
(
"next_run",
models.DateTimeField(
verbose_name="Next Run",
default=django.utils.timezone.now,
null=True,
),
),
("task", models.CharField(editable=False, null=True, max_length=100)),
], ],
options={ options={
'verbose_name': 'Scheduled task', "verbose_name": "Scheduled task",
'ordering': ['next_run'], "ordering": ["next_run"],
}, },
), ),
migrations.CreateModel( migrations.CreateModel(
name='Task', name="Task",
fields=[ fields=[
('id', models.AutoField(verbose_name='ID', auto_created=True, serialize=False, primary_key=True)), (
('name', models.CharField(editable=False, max_length=100)), "id",
('func', models.CharField(max_length=256)), models.AutoField(
('hook', models.CharField(null=True, max_length=256)), verbose_name="ID",
('args', picklefield.fields.PickledObjectField(editable=False, null=True)), auto_created=True,
('kwargs', picklefield.fields.PickledObjectField(editable=False, null=True)), serialize=False,
('result', picklefield.fields.PickledObjectField(editable=False, null=True)), primary_key=True,
('started', models.DateTimeField(editable=False)), ),
('stopped', models.DateTimeField(editable=False)), ),
('success', models.BooleanField(editable=False, default=True)), ("name", models.CharField(editable=False, max_length=100)),
("func", models.CharField(max_length=256)),
("hook", models.CharField(null=True, max_length=256)),
(
"args",
picklefield.fields.PickledObjectField(editable=False, null=True),
),
(
"kwargs",
picklefield.fields.PickledObjectField(editable=False, null=True),
),
(
"result",
picklefield.fields.PickledObjectField(editable=False, null=True),
),
("started", models.DateTimeField(editable=False)),
("stopped", models.DateTimeField(editable=False)),
("success", models.BooleanField(editable=False, default=True)),
], ],
), ),
migrations.CreateModel( migrations.CreateModel(
name='Failure', name="Failure",
fields=[ fields=[],
],
options={ options={
'verbose_name': 'Failed task', "verbose_name": "Failed task",
'proxy': True, "proxy": True,
}, },
bases=('django_q.task',), bases=("django_q.task",),
), ),
migrations.CreateModel( migrations.CreateModel(
name='Success', name="Success",
fields=[ fields=[],
],
options={ options={
'verbose_name': 'Successful task', "verbose_name": "Successful task",
'proxy': True, "proxy": True,
}, },
bases=('django_q.task',), bases=("django_q.task",),
), ),
] ]
+11 -7
View File
@@ -4,18 +4,22 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0001_initial'), ("django_q", "0001_initial"),
] ]
operations = [ operations = [
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='args', name="args",
field=models.TextField(help_text="e.g. 1, 2, 'John'", blank=True, null=True), field=models.TextField(
help_text="e.g. 1, 2, 'John'", blank=True, null=True
),
), ),
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='kwargs', name="kwargs",
field=models.TextField(help_text="e.g. x=1, y=2, name='John'", blank=True, null=True), field=models.TextField(
help_text="e.g. x=1, y=2, name='John'", blank=True, null=True
),
), ),
] ]
+24 -12
View File
@@ -4,29 +4,41 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0002_auto_20150630_1624'), ("django_q", "0002_auto_20150630_1624"),
] ]
operations = [ operations = [
migrations.AlterModelOptions( migrations.AlterModelOptions(
name='failure', name="failure",
options={'verbose_name_plural': 'Failed tasks', 'verbose_name': 'Failed task'}, options={
"verbose_name_plural": "Failed tasks",
"verbose_name": "Failed task",
},
), ),
migrations.AlterModelOptions( migrations.AlterModelOptions(
name='schedule', name="schedule",
options={'verbose_name_plural': 'Scheduled tasks', 'ordering': ['next_run'], 'verbose_name': 'Scheduled task'}, options={
"verbose_name_plural": "Scheduled tasks",
"ordering": ["next_run"],
"verbose_name": "Scheduled task",
},
), ),
migrations.AlterModelOptions( migrations.AlterModelOptions(
name='success', name="success",
options={'verbose_name_plural': 'Successful tasks', 'verbose_name': 'Successful task'}, options={
"verbose_name_plural": "Successful tasks",
"verbose_name": "Successful task",
},
), ),
migrations.RemoveField( migrations.RemoveField(
model_name='task', model_name="task",
name='id', name="id",
), ),
migrations.AddField( migrations.AddField(
model_name='task', model_name="task",
name='id', name="id",
field=models.CharField(max_length=32, primary_key=True, editable=False, serialize=False), field=models.CharField(
max_length=32, primary_key=True, editable=False, serialize=False
),
), ),
] ]
+16 -8
View File
@@ -1,23 +1,31 @@
from django.db import migrations, models from django.db import migrations
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0003_auto_20150708_1326'), ("django_q", "0003_auto_20150708_1326"),
] ]
operations = [ operations = [
migrations.AlterModelOptions( migrations.AlterModelOptions(
name='failure', name="failure",
options={'verbose_name_plural': 'Failed tasks', 'verbose_name': 'Failed task', 'ordering': ['-stopped']}, options={
"verbose_name_plural": "Failed tasks",
"verbose_name": "Failed task",
"ordering": ["-stopped"],
},
), ),
migrations.AlterModelOptions( migrations.AlterModelOptions(
name='success', name="success",
options={'verbose_name_plural': 'Successful tasks', 'verbose_name': 'Successful task', 'ordering': ['-stopped']}, options={
"verbose_name_plural": "Successful tasks",
"verbose_name": "Successful task",
"ordering": ["-stopped"],
},
), ),
migrations.AlterModelOptions( migrations.AlterModelOptions(
name='task', name="task",
options={'ordering': ['-stopped']}, options={"ordering": ["-stopped"]},
), ),
] ]
@@ -4,18 +4,18 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0004_auto_20150710_1043'), ("django_q", "0004_auto_20150710_1043"),
] ]
operations = [ operations = [
migrations.AddField( migrations.AddField(
model_name='schedule', model_name="schedule",
name='name', name="name",
field=models.CharField(max_length=100, null=True), field=models.CharField(max_length=100, null=True),
), ),
migrations.AddField( migrations.AddField(
model_name='task', model_name="task",
name='group', name="group",
field=models.CharField(max_length=100, null=True, editable=False), field=models.CharField(max_length=100, null=True, editable=False),
), ),
] ]
+25 -7
View File
@@ -4,18 +4,36 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0005_auto_20150718_1506'), ("django_q", "0005_auto_20150718_1506"),
] ]
operations = [ operations = [
migrations.AddField( migrations.AddField(
model_name='schedule', model_name="schedule",
name='minutes', name="minutes",
field=models.PositiveSmallIntegerField(help_text='Number of minutes for the Minutes type', blank=True, null=True), field=models.PositiveSmallIntegerField(
help_text="Number of minutes for the Minutes type",
blank=True,
null=True,
),
), ),
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='schedule_type', name="schedule_type",
field=models.CharField(max_length=1, choices=[('O', 'Once'), ('I', 'Minutes'), ('H', 'Hourly'), ('D', 'Daily'), ('W', 'Weekly'), ('M', 'Monthly'), ('Q', 'Quarterly'), ('Y', 'Yearly')], default='O', verbose_name='Schedule Type'), field=models.CharField(
max_length=1,
choices=[
("O", "Once"),
("I", "Minutes"),
("H", "Hourly"),
("D", "Daily"),
("W", "Weekly"),
("M", "Monthly"),
("Q", "Quarterly"),
("Y", "Yearly"),
],
default="O",
verbose_name="Schedule Type",
),
), ),
] ]
+16 -8
View File
@@ -4,21 +4,29 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0006_auto_20150805_1817'), ("django_q", "0006_auto_20150805_1817"),
] ]
operations = [ operations = [
migrations.CreateModel( migrations.CreateModel(
name='OrmQ', name="OrmQ",
fields=[ fields=[
('id', models.AutoField(primary_key=True, auto_created=True, verbose_name='ID', serialize=False)), (
('key', models.CharField(max_length=100)), "id",
('payload', models.TextField()), models.AutoField(
('lock', models.DateTimeField(null=True)), primary_key=True,
auto_created=True,
verbose_name="ID",
serialize=False,
),
),
("key", models.CharField(max_length=100)),
("payload", models.TextField()),
("lock", models.DateTimeField(null=True)),
], ],
options={ options={
'verbose_name_plural': 'Queued tasks', "verbose_name_plural": "Queued tasks",
'verbose_name': 'Queued task', "verbose_name": "Queued task",
}, },
), ),
] ]
@@ -4,13 +4,13 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0007_ormq'), ("django_q", "0007_ormq"),
] ]
operations = [ operations = [
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='name', name="name",
field=models.CharField(blank=True, max_length=100, null=True), field=models.CharField(blank=True, max_length=100, null=True),
), ),
] ]
@@ -4,13 +4,17 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0008_auto_20160224_1026'), ("django_q", "0008_auto_20160224_1026"),
] ]
operations = [ operations = [
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='repeats', name="repeats",
field=models.IntegerField(default=-1, help_text='n = n times, -1 = forever', verbose_name='Repeats'), field=models.IntegerField(
default=-1,
help_text="n = n times, -1 = forever",
verbose_name="Repeats",
),
), ),
] ]
+16 -10
View File
@@ -5,23 +5,29 @@ from django.db import migrations
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0009_auto_20171009_0915'), ("django_q", "0009_auto_20171009_0915"),
] ]
operations = [ operations = [
migrations.AlterField( migrations.AlterField(
model_name='task', model_name="task",
name='args', name="args",
field=picklefield.fields.PickledObjectField(editable=False, null=True, protocol=-1), field=picklefield.fields.PickledObjectField(
editable=False, null=True, protocol=-1
),
), ),
migrations.AlterField( migrations.AlterField(
model_name='task', model_name="task",
name='kwargs', name="kwargs",
field=picklefield.fields.PickledObjectField(editable=False, null=True, protocol=-1), field=picklefield.fields.PickledObjectField(
editable=False, null=True, protocol=-1
),
), ),
migrations.AlterField( migrations.AlterField(
model_name='task', model_name="task",
name='result', name="result",
field=picklefield.fields.PickledObjectField(editable=False, null=True, protocol=-1), field=picklefield.fields.PickledObjectField(
editable=False, null=True, protocol=-1
),
), ),
] ]
+24 -7
View File
@@ -6,18 +6,35 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0010_auto_20200610_0856'), ("django_q", "0010_auto_20200610_0856"),
] ]
operations = [ operations = [
migrations.AddField( migrations.AddField(
model_name='schedule', model_name="schedule",
name='cron', name="cron",
field=models.CharField(blank=True, help_text='Cron expression', max_length=100, null=True), field=models.CharField(
blank=True, help_text="Cron expression", max_length=100, null=True
),
), ),
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='schedule_type', name="schedule_type",
field=models.CharField(choices=[('O', 'Once'), ('I', 'Minutes'), ('H', 'Hourly'), ('D', 'Daily'), ('W', 'Weekly'), ('M', 'Monthly'), ('Q', 'Quarterly'), ('Y', 'Yearly'), ('C', 'Cron')], default='O', max_length=1, verbose_name='Schedule Type'), field=models.CharField(
choices=[
("O", "Once"),
("I", "Minutes"),
("H", "Hourly"),
("D", "Daily"),
("W", "Weekly"),
("M", "Monthly"),
("Q", "Quarterly"),
("Y", "Yearly"),
("C", "Cron"),
],
default="O",
max_length=1,
verbose_name="Schedule Type",
),
), ),
] ]
+10 -4
View File
@@ -8,13 +8,19 @@ import django_q.models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0011_auto_20200628_1055'), ("django_q", "0011_auto_20200628_1055"),
] ]
operations = [ operations = [
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='cron', name="cron",
field=models.CharField(blank=True, help_text='Cron expression', max_length=100, null=True, validators=[django_q.models.validate_cron]), field=models.CharField(
blank=True,
help_text="Cron expression",
max_length=100,
null=True,
validators=[django_q.models.validate_cron],
),
), ),
] ]
@@ -6,13 +6,13 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0012_auto_20200702_1608'), ("django_q", "0012_auto_20200702_1608"),
] ]
operations = [ operations = [
migrations.AddField( migrations.AddField(
model_name='task', model_name="task",
name='attempt_count', name="attempt_count",
field=models.IntegerField(default=0), field=models.IntegerField(default=0),
), ),
] ]
+3 -3
View File
@@ -6,13 +6,13 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0013_task_attempt_count'), ("django_q", "0013_task_attempt_count"),
] ]
operations = [ operations = [
migrations.AddField( migrations.AddField(
model_name='schedule', model_name="schedule",
name='cluster', name="cluster",
field=models.CharField(blank=True, default=None, max_length=100, null=True), field=models.CharField(blank=True, default=None, max_length=100, null=True),
), ),
] ]
@@ -6,13 +6,30 @@ from django.db import migrations, models
class Migration(migrations.Migration): class Migration(migrations.Migration):
dependencies = [ dependencies = [
('django_q', '0014_schedule_cluster'), ("django_q", "0014_schedule_cluster"),
] ]
operations = [ operations = [
migrations.AlterField( migrations.AlterField(
model_name='schedule', model_name="schedule",
name='schedule_type', name="schedule_type",
field=models.CharField(choices=[('O', 'Once'), ('I', 'Minutes'), ('H', 'Hourly'), ('D', 'Daily'), ('W', 'Weekly'), ('BW', 'Biweekly'), ('M', 'Monthly'), ('BM', 'Bimonthly'), ('Q', 'Quarterly'), ('Y', 'Yearly'), ('C', 'Cron')], default='O', max_length=2, verbose_name='Schedule Type'), field=models.CharField(
choices=[
("O", "Once"),
("I", "Minutes"),
("H", "Hourly"),
("D", "Daily"),
("W", "Weekly"),
("BW", "Biweekly"),
("M", "Monthly"),
("BM", "Bimonthly"),
("Q", "Quarterly"),
("Y", "Yearly"),
("C", "Cron"),
],
default="O",
max_length=2,
verbose_name="Schedule Type",
),
), ),
] ]
+59 -2
View File
@@ -1,3 +1,5 @@
from datetime import datetime, timedelta
# Django # Django
from django import get_version from django import get_version
from django.core.exceptions import ValidationError from django.core.exceptions import ValidationError
@@ -5,6 +7,7 @@ from django.db import models
from django.template.defaultfilters import truncatechars from django.template.defaultfilters import truncatechars
from django.urls import reverse from django.urls import reverse
from django.utils import timezone from django.utils import timezone
from django.utils.timezone import is_aware
from django.utils.html import format_html from django.utils.html import format_html
from django.utils.translation import gettext_lazy as _ from django.utils.translation import gettext_lazy as _
@@ -13,8 +16,10 @@ from picklefield import PickledObjectField
from picklefield.fields import dbsafe_decode from picklefield.fields import dbsafe_decode
# Local # Local
from django_q.conf import croniter from django_q.conf import croniter, Conf
from django_q.signing import SignedPackage from django_q.signing import SignedPackage
from django_q.utils import localtime, add_months, add_years
from .utils import get_func_repr from .utils import get_func_repr
@@ -207,6 +212,59 @@ class Schedule(models.Model):
task = models.CharField(max_length=100, null=True, editable=False) task = models.CharField(max_length=100, null=True, editable=False)
cluster = models.CharField(max_length=100, default=None, null=True, blank=True) cluster = models.CharField(max_length=100, default=None, null=True, blank=True)
def calculate_next_run(self, next_run=None):
# next run is always in UTC
next_run = next_run or self.next_run
if self.schedule_type == self.CRON:
if not croniter:
raise ImportError(
_("Please install croniter to enable cron expressions")
)
return croniter(self.cron, localtime()).get_next(datetime)
if self.schedule_type == self.MINUTES:
add = timedelta(minutes=(self.minutes or 1))
elif self.schedule_type == self.HOURLY:
add = timedelta(hours=1)
elif self.schedule_type == self.DAILY:
add = timedelta(days=1)
elif self.schedule_type == self.WEEKLY:
add = timedelta(weeks=1)
elif self.schedule_type == self.BIWEEKLY:
add = timedelta(weeks=2)
elif self.schedule_type == self.MONTHLY:
add = timedelta(days=(add_months(next_run, 1) - next_run).days)
elif self.schedule_type == self.BIMONTHLY:
add = timedelta(days=(add_months(next_run, 2) - next_run).days)
elif self.schedule_type == self.QUARTERLY:
add = timedelta(days=(add_months(next_run, 3) - next_run).days)
elif self.schedule_type == self.YEARLY:
add = timedelta(days=(add_years(next_run, 1) - next_run).days)
# add normal timedelta, we will correct this later based on timezone
next_run += add
# DST differencers don't matter with minutes, hourly or yearly, so skip those
if self.schedule_type not in [self.MINUTES, self.HOURLY, self.YEARLY]:
# Get localtimes and then remove the tzinfo, so we can get the actual difference
current_next_run = localtime(next_run - add).replace(tzinfo=None)
new_next_run = localtime(next_run).replace(tzinfo=None)
# get the difference between them, this should be (-)1 or (-)0.5 hour
# based on DST active or not
extra_diff = (new_next_run - current_next_run) - add
# if we have one positive hour difference, then subtract it, so we are even
# and vice versa. In most cases, this will be 0, as there won't be a
# timezone diff
if extra_diff > timedelta(hours=0):
next_run -= extra_diff
else:
next_run += extra_diff
return next_run
def success(self): def success(self):
if self.task and Task.objects.filter(id=self.task): if self.task and Task.objects.filter(id=self.task):
return Task.objects.get(id=self.task).success return Task.objects.get(id=self.task).success
@@ -229,7 +287,6 @@ class Schedule(models.Model):
last_run.allow_tags = True last_run.allow_tags = True
last_run.short_description = _("last_run") last_run.short_description = _("last_run")
class Meta: class Meta:
app_label = "django_q" app_label = "django_q"
verbose_name = _("Scheduled task") verbose_name = _("Scheduled task")
+32 -19
View File
@@ -19,28 +19,30 @@ try:
except ImportError: except ImportError:
psutil = None psutil = None
# optional
try:
from blessed import Terminal
except ImportError:
pass
def get_process_mb(pid): def get_process_mb(pid):
try: try:
process = psutil.Process(pid) process = psutil.Process(pid)
mb_used = round(process.memory_info().rss / 1024 ** 2, 2) mb_used = round(process.memory_info().rss / 1024**2, 2)
except psutil.NoSuchProcess: except psutil.NoSuchProcess:
mb_used = "NO_PROCESS_FOUND" mb_used = "NO_PROCESS_FOUND"
return mb_used return mb_used
BLESSED_INSTALL_MESSAGE = "Blessed is not installed. Please install blessed to use this: https://pypi.org/project/blessed/"
BLESSED_INSTALL_MESSAGE = (
"Blessed is not installed. Please install blessed to use this: "
"https://pypi.org/project/blessed/"
)
def monitor(run_once=False, broker=None): def monitor(run_once=False, broker=None):
if not broker: if not broker:
broker = get_broker() broker = get_broker()
try: try:
from blessed import Terminal
term = Terminal() term = Terminal()
except: except ImportError:
print(BLESSED_INSTALL_MESSAGE) print(BLESSED_INSTALL_MESSAGE)
return return
@@ -204,8 +206,10 @@ def info(broker=None):
if not broker: if not broker:
broker = get_broker() broker = get_broker()
try: try:
from blessed import Terminal
term = Terminal() term = Terminal()
except: except ImportError:
print(BLESSED_INSTALL_MESSAGE) print(BLESSED_INSTALL_MESSAGE)
return return
@@ -256,9 +260,12 @@ def info(broker=None):
print( print(
term.black_on_green( term.black_on_green(
term.center( term.center(
_( _("-- %(prefix)s %(version)s on %(info)s --")
'-- %(prefix)s %(version)s on %(info)s --' % {
) % {'prefix': Conf.PREFIX.capitalize(), 'version': ".".join(str(v) for v in VERSION), 'info': broker.info()} "prefix": Conf.PREFIX.capitalize(),
"version": ".".join(str(v) for v in VERSION),
"info": broker.info(),
}
) )
) )
) )
@@ -293,7 +300,7 @@ def info(broker=None):
+ term.move_x(1 * col_width) + term.move_x(1 * col_width)
+ term.white(str(models.Schedule.objects.count())) + term.white(str(models.Schedule.objects.count()))
+ term.move_x(2 * col_width) + term.move_x(2 * col_width)
+ term.cyan(_("Tasks/%(per)s") % {'per': per}) + term.cyan(_("Tasks/%(per)s") % {"per": per})
+ term.move_x(3 * col_width) + term.move_x(3 * col_width)
+ term.white(f"{tasks_per:.2f}") + term.white(f"{tasks_per:.2f}")
+ term.move_x(4 * col_width) + term.move_x(4 * col_width)
@@ -308,8 +315,10 @@ def memory(run_once=False, workers=False, broker=None):
if not broker: if not broker:
broker = get_broker() broker = get_broker()
try: try:
from blessed import Terminal
term = Terminal() term = Terminal()
except: except ImportError:
print(BLESSED_INSTALL_MESSAGE) print(BLESSED_INSTALL_MESSAGE)
return return
broker.ping() broker.ping()
@@ -389,7 +398,7 @@ def memory(run_once=False, workers=False, broker=None):
) )
# memory available (MB) # memory available (MB)
memory_available = round( memory_available = round(
psutil.virtual_memory().available / 1024 ** 2, 2 psutil.virtual_memory().available / 1024**2, 2
) )
if memory_available_percentage < MEMORY_AVAILABLE_LOWEST_PERCENTAGE: if memory_available_percentage < MEMORY_AVAILABLE_LOWEST_PERCENTAGE:
MEMORY_AVAILABLE_LOWEST_PERCENTAGE = memory_available_percentage MEMORY_AVAILABLE_LOWEST_PERCENTAGE = memory_available_percentage
@@ -413,7 +422,7 @@ def memory(run_once=False, workers=False, broker=None):
print( print(
term.move(row, 4 * col_width) term.move(row, 4 * col_width)
+ term.center( + term.center(
round(psutil.virtual_memory().total / 1024 ** 2, 2), round(psutil.virtual_memory().total / 1024**2, 2),
width=col_width - 1, width=col_width - 1,
) )
) )
@@ -475,9 +484,13 @@ def memory(run_once=False, workers=False, broker=None):
row += 1 row += 1
print( print(
term.move(row, 0) term.move(row, 0)
+ _("Available lowest (): %(memory_percent)s ((at)s)") % { 'memory_percent': str(MEMORY_AVAILABLE_LOWEST_PERCENTAGE), 'at': MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT.strftime( + _("Available lowest (): %(memory_percent)s ((at)s)")
"%Y-%m-%d %H:%M:%S+00:00" % {
)} "memory_percent": str(MEMORY_AVAILABLE_LOWEST_PERCENTAGE),
"at": MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT.strftime(
"%Y-%m-%d %H:%M:%S+00:00"
),
}
) )
# for testing # for testing
if run_once: if run_once:
+2 -1
View File
@@ -1,5 +1,6 @@
""" """
The code is derived from https://github.com/althonos/pronto/commit/3384010dfb4fc7c66a219f59276adef3288a886b The code is derived from
https://github.com/althonos/pronto/commit/3384010dfb4fc7c66a219f59276adef3288a886b
""" """
import multiprocessing import multiprocessing
import multiprocessing.queues import multiprocessing.queues
+4 -4
View File
@@ -19,16 +19,16 @@ def call_hook(sender, instance, **kwargs):
f = getattr(m, func) f = getattr(m, func)
except (ValueError, ImportError, AttributeError): except (ValueError, ImportError, AttributeError):
logger.error( logger.error(
_("malformed return hook '%(hook)s' for [%(name)s]") % {'hook': instance.hook, 'name': instance.name} _("malformed return hook '%(hook)s' for [%(name)s]")
% {"hook": instance.hook, "name": instance.name}
) )
return return
try: try:
f(instance) f(instance)
except Exception as e: except Exception as e:
logger.error( logger.error(
_( _("return hook %(hook)s failed on [%(name)s] because %(error)s")
"return hook %(hook)s failed on [%(name)s] because %(error)s" % {"hook": instance.hook, "name": instance.name, "error": str(e)}
) % {'hook': instance.hook, 'name': instance.name, 'error': str(e)}
) )
+4 -2
View File
@@ -600,7 +600,8 @@ class Chain:
def result(self, wait=0): def result(self, wait=0):
""" """
return the full list of results from the chain when it finishes. blocks until timeout. return the full list of results from the chain when it finishes. blocks until
timeout.
:param int wait: how many milliseconds to wait for a result :param int wait: how many milliseconds to wait for a result
:return: an unsorted list of results :return: an unsorted list of results
""" """
@@ -611,7 +612,8 @@ class Chain:
def fetch(self, failures=True, wait=0): def fetch(self, failures=True, wait=0):
""" """
get the task result objects from the chain when it finishes. blocks until timeout. get the task result objects from the chain when it finishes. blocks until
timeout.
:param failures: include failed tasks :param failures: include failed tasks
:param int wait: how many milliseconds to wait for a result :param int wait: how many milliseconds to wait for a result
:return: an unsorted list of task objects :return: an unsorted list of task objects
+2 -4
View File
@@ -1,7 +1,5 @@
import os import os
import django
BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
@@ -77,7 +75,7 @@ DATABASES = {
LANGUAGE_CODE = "en-us" LANGUAGE_CODE = "en-us"
TIME_ZONE = "UTC" TIME_ZONE = "Europe/Amsterdam"
USE_I18N = True USE_I18N = True
@@ -130,5 +128,5 @@ Q_CLUSTER = {
"testing": True, "testing": True,
"log_level": "DEBUG", "log_level": "DEBUG",
"django_redis": "default", "django_redis": "default",
"redis": f"redis://{REDIS_HOST}:6379/0" "redis": f"redis://{REDIS_HOST}:6379/0",
} }
+1 -2
View File
@@ -2,12 +2,11 @@ import os
from time import sleep from time import sleep
import pytest import pytest
import redis
from django_q.brokers import Broker, get_broker from django_q.brokers import Broker, get_broker
from django_q.conf import Conf from django_q.conf import Conf
from django_q.humanhash import uuid from django_q.humanhash import uuid
from django_q.tests.settings import REDIS_HOST, MONGO_HOST from django_q.tests.settings import MONGO_HOST, REDIS_HOST
def test_broker(monkeypatch): def test_broker(monkeypatch):
+12 -12
View File
@@ -1,8 +1,8 @@
from datetime import datetime
import os import os
import sys import sys
import threading import threading
import uuid as uuidlib import uuid as uuidlib
from datetime import datetime
from math import copysign from math import copysign
from multiprocessing import Event, Value from multiprocessing import Event, Value
from time import sleep from time import sleep
@@ -11,9 +11,6 @@ from typing import Optional
import pytest import pytest
from django.utils import timezone from django.utils import timezone
myPath = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, myPath + "/../")
from django_q.brokers import Broker, get_broker from django_q.brokers import Broker, get_broker
from django_q.cluster import Cluster, Sentinel, monitor, pusher, save_task, worker from django_q.cluster import Cluster, Sentinel, monitor, pusher, save_task, worker
from django_q.conf import Conf from django_q.conf import Conf
@@ -32,9 +29,12 @@ from django_q.tasks import (
result, result,
result_group, result_group,
) )
from django_q.tests.tasks import TaskError, multiply from django_q.tests.tasks import multiply, TaskError
from django_q.utils import add_months, add_years from django_q.utils import add_months, add_years
myPath = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, myPath + "/../")
class WordClass: class WordClass:
def __init__(self): def __init__(self):
@@ -64,7 +64,7 @@ def test_sync(broker):
@pytest.mark.django_db @pytest.mark.django_db
def test_sync_raise_exception(broker): def test_sync_raise_exception(broker):
with pytest.raises(Exception): with pytest.raises(TaskError):
async_task("django_q.tests.tasks.raise_exception", broker=broker, sync=True) async_task("django_q.tests.tasks.raise_exception", broker=broker, sync=True)
@@ -409,6 +409,7 @@ def test_recycle(broker, monkeypatch):
assert Success.objects.count() == Conf.SAVE_LIMIT assert Success.objects.count() == Conf.SAVE_LIMIT
broker.delete_queue() broker.delete_queue()
@pytest.mark.django_db @pytest.mark.django_db
def test_save_limit_per_func(broker, monkeypatch): def test_save_limit_per_func(broker, monkeypatch):
# set up the Sentinel # set up the Sentinel
@@ -442,15 +443,14 @@ def test_save_limit_per_func(broker, monkeypatch):
# run monitor # run monitor
monitor(result_queue) monitor(result_queue)
assert Success.objects.count() == 3 assert Success.objects.count() == 3
assert set(Success.objects.filter().values_list('func', flat=True)) == { assert set(Success.objects.filter().values_list("func", flat=True)) == {
'django_q.tests.tasks.countdown', "django_q.tests.tasks.countdown",
'django_q.tests.tasks.hello', "django_q.tests.tasks.hello",
'django_q.tests.tasks.multiply', "django_q.tests.tasks.multiply",
} }
broker.delete_queue() broker.delete_queue()
@pytest.mark.django_db @pytest.mark.django_db
def test_max_rss(broker, monkeypatch): def test_max_rss(broker, monkeypatch):
# set up the Sentinel # set up the Sentinel
@@ -538,7 +538,6 @@ def test_attempt_count(broker, monkeypatch):
assert saved_task.attempt_count == 1 assert saved_task.attempt_count == 1
sleep(0.5) sleep(0.5)
# second save # second save
old_stopped = task["stopped"]
task["stopped"] = timezone.now() task["stopped"] = timezone.now()
save_task(task, broker) save_task(task, broker)
saved_task = Task.objects.get(id=task["id"]) saved_task = Task.objects.get(id=task["id"])
@@ -770,6 +769,7 @@ def test_add_months():
assert new_date.month == 2 assert new_date.month == 2
assert new_date.day == 29 assert new_date.day == 29
@pytest.mark.django_db @pytest.mark.django_db
def test_add_years(): def test_add_years():
# add some months # add some months
+76 -12
View File
@@ -1,5 +1,6 @@
import os import os
from datetime import timedelta import pytz
from datetime import datetime, timedelta
from multiprocessing import Event, Value from multiprocessing import Event, Value
from unittest import mock from unittest import mock
@@ -11,7 +12,7 @@ from django.utils import timezone
from django.utils.timezone import is_naive from django.utils.timezone import is_naive
from django_q.brokers import Broker, get_broker from django_q.brokers import Broker, get_broker
from django_q.cluster import monitor, pusher, scheduler, worker, localtime from django_q.cluster import localtime, monitor, pusher, scheduler, worker
from django_q.conf import Conf from django_q.conf import Conf
from django_q.queues import Queue from django_q.queues import Queue
from django_q.tasks import Schedule, fetch from django_q.tasks import Schedule, fetch
@@ -21,12 +22,15 @@ from django_q.tests.testing_utilities.multiple_database_routers import (
TestingMultipleAppsDatabaseRouter, TestingMultipleAppsDatabaseRouter,
TestingReplicaDatabaseRouter, TestingReplicaDatabaseRouter,
) )
from django_q.utils import add_months, add_years from django_q.utils import add_months
@pytest.fixture @pytest.fixture
def broker(monkeypatch) -> Broker: def broker(monkeypatch) -> Broker:
"""Patches the Conf object setting the DJANGO_REDIS attribute allowing a default redis configuration.""" """
Patches the Conf object setting the DJANGO_REDIS attribute allowing a default
redis configuration.
"""
monkeypatch.setattr(Conf, "DJANGO_REDIS", "default") monkeypatch.setattr(Conf, "DJANGO_REDIS", "default")
return get_broker() return get_broker()
@@ -66,7 +70,7 @@ REPLICA_DATABASES = {
} }
MULTIPLE_APPS_DATABASE_ROUTERS = [ MULTIPLE_APPS_DATABASE_ROUTERS = [
f"{TestingMultipleAppsDatabaseRouter.__module__}.{TestingMultipleAppsDatabaseRouter.__name__}" f"{TestingMultipleAppsDatabaseRouter.__module__}.{TestingMultipleAppsDatabaseRouter.__name__}" # noqa: E501
] ]
MULTIPLE_APPS_DATABASES = { MULTIPLE_APPS_DATABASES = {
"default": { "default": {
@@ -80,6 +84,63 @@ MULTIPLE_APPS_DATABASES = {
} }
@pytest.mark.django_db
def test_scheduler_daylight_saving_time_daily(broker, monkeypatch):
# Set up a startdate in the Amsterdam timezone (without dst 1 hour ahead). The
# 28th of March 2021 is the day when sunlight saving starts (at 2 am)
monkeypatch.setattr(Conf, "TIME_ZONE", "Europe/Amsterdam")
tz = pytz.timezone('Europe/Amsterdam')
broker.list_key = "scheduler_test:q"
# Let's start a schedule at 1 am on the 27th of March. This is in AMS timezone.
# So, 2021-03-27 00:00:00 when saved (due to TZ being Amsterdam and saved in UTC)
start_date = datetime(2021, 3, 27, 1, 0, 0)
# Create schedule with the next run date on the start date. It will move one day
# forward when we run the scheduler
schedule = create_schedule(
"math.copysign",
1,
-1,
name="test math",
schedule_type=Schedule.DAILY,
next_run=start_date,
)
# Run scheduler so we get the next run date
scheduler(broker=broker)
schedule.refresh_from_db()
# It's now the day after exactly at midnight UTC
next_run = schedule.next_run
assert str(next_run) == "2021-03-28 00:00:00+00:00"
# In the Amsterdam timezone, it's 1 hour over midnight (+01)
next_run = next_run.astimezone(tz)
assert str(next_run) == "2021-03-28 01:00:00+01:00"
# Run scheduler so we get the next run date
scheduler(broker=broker)
schedule.refresh_from_db()
next_run = schedule.next_run
assert str(next_run) == "2021-03-28 23:00:00+00:00"
next_run = next_run.astimezone(tz)
# In the Amsterdam timezone, it's 1 hour over midnight (+02)
assert str(next_run) == "2021-03-29 01:00:00+02:00"
# Run scheduler so we get the next run date
scheduler(broker=broker)
schedule.refresh_from_db()
next_run = schedule.next_run
assert str(next_run) == "2021-03-29 23:00:00+00:00"
next_run = next_run.astimezone(tz)
assert str(next_run) == "2021-03-30 01:00:00+02:00"
@pytest.mark.django_db @pytest.mark.django_db
def test_scheduler(broker, monkeypatch): def test_scheduler(broker, monkeypatch):
broker.list_key = "scheduler_test:q" broker.list_key = "scheduler_test:q"
@@ -234,7 +295,7 @@ def test_scheduler(broker, monkeypatch):
"django_q.tests.tasks.word_multiply", "django_q.tests.tasks.word_multiply",
2, 2,
word="catch_up", word="catch_up",
schedule_type=Schedule.BIMONTHLY schedule_type=Schedule.BIMONTHLY,
) )
scheduler(broker=broker) scheduler(broker=broker)
schedule = Schedule.objects.get(pk=schedule.pk) schedule = Schedule.objects.get(pk=schedule.pk)
@@ -245,7 +306,7 @@ def test_scheduler(broker, monkeypatch):
"django_q.tests.tasks.word_multiply", "django_q.tests.tasks.word_multiply",
2, 2,
word="catch_up", word="catch_up",
schedule_type=Schedule.BIWEEKLY schedule_type=Schedule.BIWEEKLY,
) )
scheduler(broker=broker) scheduler(broker=broker)
schedule = Schedule.objects.get(pk=schedule.pk) schedule = Schedule.objects.get(pk=schedule.pk)
@@ -311,7 +372,8 @@ def test_scheduler_atomic_transaction_must_specify_a_database_when_no_replicas_a
""" """
GIVEN a environment without a read replica database GIVEN a environment without a read replica database
WHEN the scheduler is called WHEN the scheduler is called
THEN the transaction atomic must be called using the configured database in the Conf.ORM settings. THEN the transaction atomic must be called using the configured database in the
Conf.ORM settings.
""" """
broker = orm_no_replica_broker broker = orm_no_replica_broker
with mock.patch("django_q.cluster.db") as mocked_db: with mock.patch("django_q.cluster.db") as mocked_db:
@@ -324,13 +386,14 @@ def test_scheduler_atomic_transaction_must_specify_a_database_when_no_replicas_a
DATABASE_ROUTERS=REPLICA_DATABASE_ROUTERS, DATABASES=REPLICA_DATABASES DATABASE_ROUTERS=REPLICA_DATABASE_ROUTERS, DATABASES=REPLICA_DATABASES
) )
@pytest.mark.django_db @pytest.mark.django_db
def test_scheduler_atomic_transaction_must_specify_no_database_when_read_write_replicas_are_used( def test_scheduler_atomic_must_specify_no_db_when_read_write_replicas_are_used(
orm_replica_broker: Broker, orm_replica_broker: Broker,
): ):
""" """
GIVEN a environment with a read/write configured replica database GIVEN a environment with a read/write configured replica database
WHEN the scheduler is called WHEN the scheduler is called
THEN the transaction must be called without a specific database, thus letting the database router pick. THEN the transaction must be called without a specific database, thus letting the
database router pick.
""" """
with mock.patch("django_q.cluster.db") as mocked_db: with mock.patch("django_q.cluster.db") as mocked_db:
scheduler(broker=orm_replica_broker) scheduler(broker=orm_replica_broker)
@@ -342,13 +405,14 @@ def test_scheduler_atomic_transaction_must_specify_no_database_when_read_write_r
DATABASE_ROUTERS=MULTIPLE_APPS_DATABASE_ROUTERS, DATABASES=MULTIPLE_APPS_DATABASES DATABASE_ROUTERS=MULTIPLE_APPS_DATABASE_ROUTERS, DATABASES=MULTIPLE_APPS_DATABASES
) )
@pytest.mark.django_db @pytest.mark.django_db
def test_scheduler_atomic_transaction_must_specify_the_database_based_on_router_redirection( def test_scheduler_atomic_must_specify_the_database_based_on_router_redirection(
orm_no_replica_broker: Broker, orm_no_replica_broker: Broker,
): ):
""" """
GIVEN a environment without a read replica database GIVEN a environment without a read replica database
WHEN the scheduler is called WHEN the scheduler is called
THEN the transaction atomic must be called using the configured database in the Conf.ORM settings. THEN the transaction atomic must be called using the configured database in the
Conf.ORM settings.
""" """
broker = orm_no_replica_broker broker = orm_no_replica_broker
with mock.patch("django_q.cluster.db") as mocked_db: with mock.patch("django_q.cluster.db") as mocked_db:
+24 -9
View File
@@ -1,6 +1,13 @@
from datetime import datetime
import pytz
import calendar
import inspect import inspect
from datetime import date from datetime import date
import calendar
from django.utils import timezone
from django.conf import settings
from django_q.conf import Conf
# credits: https://stackoverflow.com/a/4131114 # credits: https://stackoverflow.com/a/4131114
# Made them aware of timezone # Made them aware of timezone
@@ -8,21 +15,21 @@ def add_months(d, months):
month = d.month - 1 + months month = d.month - 1 + months
year = d.year + month // 12 year = d.year + month // 12
month = month % 12 + 1 month = month % 12 + 1
day = min(d.day, calendar.monthrange(year,month)[1]) day = min(d.day, calendar.monthrange(year, month)[1])
return d.replace(year=year, month=month, day=day) return d.replace(year=year, month=month, day=day)
# credits: https://stackoverflow.com/a/15743908 # credits: https://stackoverflow.com/a/15743908
# Changed the last line to make it a little easier to read and changed it to move February 29 to 28 next year # Changed the last line to make it a little easier to read and changed it to move
# Also made them aware of timezone # February 29 to 28 next year.
def add_years(d, years): def add_years(d, years):
"""Return a date that's `years` years after the date (or datetime) """Return a date that's `years` years after the date (or datetime)
object `d`. Return the same calendar date (month and day) in the object `d`. Return the same calendar date (month and day) in the
destination year, if it exists, otherwise use the previous day destination year, if it exists, otherwise use the previous day
(thus changing February 29 to February 28). (thus changing February 29 to February 28).
""" """
try: try:
return d.replace(year = d.year + years) return d.replace(year=d.year + years)
except ValueError: except ValueError:
new_date = d + (date(d.year + years, 3, 1) - date(d.year, 3, 1)) new_date = d + (date(d.year + years, 3, 1) - date(d.year, 3, 1))
return d.replace(year=new_date.year, month=new_date.month, day=new_date.day) return d.replace(year=new_date.year, month=new_date.month, day=new_date.day)
@@ -32,11 +39,19 @@ def get_func_repr(func):
# convert func to string # convert func to string
if inspect.isfunction(func): if inspect.isfunction(func):
return f"{func.__module__}.{func.__name__}" return f"{func.__module__}.{func.__name__}"
elif inspect.ismethod(func) and hasattr(func.__self__, '__name__'): elif inspect.ismethod(func) and hasattr(func.__self__, "__name__"):
return ( return (
f"{func.__self__.__module__}." f"{func.__self__.__module__}." f"{func.__self__.__name__}.{func.__name__}"
f"{func.__self__.__name__}.{func.__name__}"
) )
else: else:
return str(func) return str(func)
def localtime(value=None) -> datetime:
"""Override for timezone.localtime to deal with naive times and local times"""
if settings.USE_TZ:
return timezone.localtime(value=value, timezone=pytz.timezone(Conf.TIME_ZONE))
if value is None:
return datetime.now()
else:
return value
+44 -43
View File
@@ -16,12 +16,10 @@
import os import os
import sys import sys
import sphinx_rtd_theme
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 + "/../")
os.environ['DJANGO_SETTINGS_MODULE'] = 'django_q.tests.settings' os.environ["DJANGO_SETTINGS_MODULE"] = "django_q.tests.settings"
nitpick_ignore = [('py:class', 'datetime')] nitpick_ignore = [("py:class", "datetime")]
# If extensions (or modules to document with autodoc) are in another directory, # If extensions (or modules to document with autodoc) are in another directory,
# add these directories to sys.path here. If the directory is relative to the # add these directories to sys.path here. If the directory is relative to the
@@ -37,50 +35,54 @@ nitpick_ignore = [('py:class', 'datetime')]
# extensions coming with Sphinx (named 'sphinx.ext.*') or your custom # extensions coming with Sphinx (named 'sphinx.ext.*') or your custom
# ones. # ones.
extensions = [ extensions = [
'sphinx_rtd_theme', "sphinx_rtd_theme",
'sphinx.ext.todo', "sphinx.ext.todo",
'sphinx.ext.intersphinx', "sphinx.ext.intersphinx",
# 'sphinx.ext.autodoc' # 'sphinx.ext.autodoc'
] ]
intersphinx_mapping = {'python': ('https://docs.python.org/3.8', None), intersphinx_mapping = {
'django': ('https://docs.djangoproject.com/en/2.2/', "python": ("https://docs.python.org/3.8", None),
'https://docs.djangoproject.com/en/2.2/_objects/')} "django": (
"https://docs.djangoproject.com/en/2.2/",
"https://docs.djangoproject.com/en/2.2/_objects/",
),
}
# Add any paths that contain templates here, relative to this directory. # Add any paths that contain templates here, relative to this directory.
templates_path = ['_templates'] templates_path = ["_templates"]
# The suffix(es) of source filenames. # The suffix(es) of source filenames.
# You can specify multiple suffix as a list of string: # You can specify multiple suffix as a list of string:
# source_suffix = ['.rst', '.md'] # source_suffix = ['.rst', '.md']
source_suffix = '.rst' source_suffix = ".rst"
# The encoding of source files. # The encoding of source files.
# source_encoding = 'utf-8-sig' # source_encoding = 'utf-8-sig'
# The master toctree document. # The master toctree document.
master_doc = 'index' master_doc = "index"
# General information about the project. # General information about the project.
project = 'Django Q2' project = "Django Q2"
copyright = '2015-2021, Ilan Steemers - 2022, Stan Triepels' copyright = "2015-2021, Ilan Steemers - 2022, Stan Triepels"
author = 'Ilan Steemers, Stan Triepels' author = "Ilan Steemers, Stan Triepels"
# The version info for the project you're documenting, acts as replacement for # The version info for the project you're documenting, acts as replacement for
# |version| and |release|, also used in various other places throughout the # |version| and |release|, also used in various other places throughout the
# built documents. # built documents.
# #
# The short X.Y version. # The short X.Y version.
version = '1.4' version = "1.4"
# The full version, including alpha/beta/rc tags. # The full version, including alpha/beta/rc tags.
release = '1.4.4' release = "1.4.7"
# The language for content autogenerated by Sphinx. Refer to documentation # The language for content autogenerated by Sphinx. Refer to documentation
# for a list of supported languages. # for a list of supported languages.
# #
# This is also used if you do content translation via gettext catalogs. # This is also used if you do content translation via gettext catalogs.
# Usually you set "language" from the command line for these cases. # Usually you set "language" from the command line for these cases.
language = 'en' language = "en"
# There are two options for replacing |today|: either, you set today to some # There are two options for replacing |today|: either, you set today to some
# non-false value, then it is used: # non-false value, then it is used:
@@ -90,7 +92,7 @@ language = 'en'
# List of patterns, relative to source directory, that match files and # List of patterns, relative to source directory, that match files and
# directories to ignore when looking for source files. # directories to ignore when looking for source files.
exclude_patterns = ['_build'] exclude_patterns = ["_build"]
# The reST default role (used for this markup: `text`) to use for all # The reST default role (used for this markup: `text`) to use for all
# documents. # documents.
@@ -108,7 +110,7 @@ add_module_names = False
# show_authors = False # show_authors = False
# The name of the Pygments (syntax highlighting) style to use. # The name of the Pygments (syntax highlighting) style to use.
pygments_style = 'sphinx' pygments_style = "sphinx"
# A list of ignored prefixes for module index sorting. # A list of ignored prefixes for module index sorting.
# modindex_common_prefix = [] # modindex_common_prefix = []
@@ -124,7 +126,7 @@ todo_include_todos = True
# The theme to use for HTML and HTML Help pages. See the documentation for # The theme to use for HTML and HTML Help pages. See the documentation for
# a list of builtin themes. # a list of builtin themes.
html_theme = 'sphinx_rtd_theme' html_theme = "sphinx_rtd_theme"
# Theme options are theme-specific and customize the look and feel of a theme # Theme options are theme-specific and customize the look and feel of a theme
# further. For a list of options available for each theme, see the # further. For a list of options available for each theme, see the
@@ -136,11 +138,11 @@ html_theme_options = {
# 'github_banner': True, # 'github_banner': True,
} }
html_sidebars = { html_sidebars = {
'**': [ "**": [
'about.html', "about.html",
'navigation.html', "navigation.html",
'relations.html', "relations.html",
'searchbox.html', "searchbox.html",
] ]
} }
# Add any paths that contain custom themes here, relative to this directory. # Add any paths that contain custom themes here, relative to this directory.
@@ -161,12 +163,12 @@ html_sidebars = {
# The name of an image file (within the static path) to use as favicon of the # The name of an image file (within the static path) to use as favicon of the
# docs. This file should be a Windows icon file (.ico) being 16x16 or 32x32 # docs. This file should be a Windows icon file (.ico) being 16x16 or 32x32
# pixels large. # pixels large.
html_favicon = '_static/favicon.ico' html_favicon = "_static/favicon.ico"
# Add any paths that contain custom static files (such as style sheets) here, # Add any paths that contain custom static files (such as style sheets) here,
# relative to this directory. They are copied after the builtin static files, # relative to this directory. They are copied after the builtin static files,
# so a file named "default.css" will overwrite the builtin "default.css". # so a file named "default.css" will overwrite the builtin "default.css".
html_static_path = ['_static'] html_static_path = ["_static"]
# Add any extra paths that contain custom files (such as robots.txt or # Add any extra paths that contain custom files (such as robots.txt or
# .htaccess) here, relative to this directory. These files are copied # .htaccess) here, relative to this directory. These files are copied
@@ -229,20 +231,17 @@ html_static_path = ['_static']
# html_search_scorer = 'scorer.js' # html_search_scorer = 'scorer.js'
# Output file base name for HTML help builder. # Output file base name for HTML help builder.
htmlhelp_basename = 'DjangoQ2doc' htmlhelp_basename = "DjangoQ2doc"
# -- Options for LaTeX output --------------------------------------------- # -- Options for LaTeX output ---------------------------------------------
latex_elements = { latex_elements = {
# The paper size ('letterpaper' or 'a4paper'). # The paper size ('letterpaper' or 'a4paper').
# 'papersize': 'letterpaper', # 'papersize': 'letterpaper',
# The font size ('10pt', '11pt' or '12pt'). # The font size ('10pt', '11pt' or '12pt').
# 'pointsize': '10pt', # 'pointsize': '10pt',
# Additional stuff for the LaTeX preamble. # Additional stuff for the LaTeX preamble.
# 'preamble': '', # 'preamble': '',
# Latex figure (float) alignment # Latex figure (float) alignment
# 'figure_align': 'htbp', # 'figure_align': 'htbp',
} }
@@ -251,8 +250,7 @@ latex_elements = {
# (source start file, target name, title, # (source start file, target name, title,
# author, documentclass [howto, manual, or own class]). # author, documentclass [howto, manual, or own class]).
latex_documents = [ latex_documents = [
(master_doc, 'DjangoQ2.tex', 'Django Q2 Documentation', (master_doc, "DjangoQ2.tex", "Django Q2 Documentation", "Ilan Steemers", "manual"),
'Ilan Steemers', 'manual'),
] ]
# The name of an image file (relative to this directory) to place at the top of # The name of an image file (relative to this directory) to place at the top of
@@ -280,10 +278,7 @@ latex_documents = [
# One entry per manual page. List of tuples # One entry per manual page. List of tuples
# (source start file, name, description, authors, manual section). # (source start file, name, description, authors, manual section).
man_pages = [ man_pages = [(master_doc, "djangoq2", "Django Q2 Documentation", [author], 1)]
(master_doc, 'djangoq2', 'Django Q2 Documentation',
[author], 1)
]
# If true, show URL addresses after external links. # If true, show URL addresses after external links.
@@ -296,9 +291,15 @@ man_pages = [
# (source start file, target name, title, author, # (source start file, target name, title, author,
# dir menu entry, description, category) # dir menu entry, description, category)
texinfo_documents = [ texinfo_documents = [
(master_doc, 'DjangoQ2', 'Django Q2 Documentation', (
author, 'DjangoQ2', 'A multiprocessing distributed task queue for Django.', master_doc,
'Miscellaneous'), "DjangoQ2",
"Django Q2 Documentation",
author,
"DjangoQ2",
"A multiprocessing distributed task queue for Django.",
"Miscellaneous",
),
] ]
# Documents to append as an appendix to all manuals. # Documents to append as an appendix to all manuals.
+7
View File
@@ -70,6 +70,13 @@ Set this to something that makes sense for your project. Can be overridden for i
See :ref:`retry` for details how to set values for timeout and retry. See :ref:`retry` for details how to set values for timeout and retry.
.. _time_zone:
time_zone
~~~~~~~
The timezone that is used for task scheduling. Use this if you are having issue with DST. The scheduler uses UTC to calculate the next date and will therefore ignore any DST changes. This will cause 1 hour or 0.5 hour changes in the schedule when time is moved one hour ahead or back. Defaults to `settings.TIME_ZONE` if `USE_TZ` is enabled.
.. _ack_failures: .. _ack_failures:
ack_failures ack_failures
+1 -1
View File
@@ -1,6 +1,6 @@
[tool.poetry] [tool.poetry]
name = "django-q2" name = "django-q2"
version = "1.4.4" version = "1.4.7"
packages = [ packages = [
{ include = "django_q" }, { include = "django_q" },
] ]
+2
View File
@@ -0,0 +1,2 @@
[flake8]
max-line-length=88