mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 16:58:13 +08:00
Remove arrow dependency, update packages, set up testing with compose (#17)
This commit is contained in:
@@ -173,11 +173,6 @@ def get_broker(list_key: str = Conf.PREFIX) -> Broker:
|
||||
m = importlib.import_module(module)
|
||||
broker = getattr(m, func)
|
||||
return broker(list_key=list_key)
|
||||
# disque
|
||||
elif Conf.DISQUE_NODES:
|
||||
from django_q.brokers import disque
|
||||
|
||||
return disque.Disque(list_key=list_key)
|
||||
# Iron MQ
|
||||
elif Conf.IRON_MQ:
|
||||
from django_q.brokers import ironmq
|
||||
|
||||
@@ -1,78 +0,0 @@
|
||||
import random
|
||||
|
||||
# External
|
||||
import redis
|
||||
|
||||
# Django
|
||||
from django.utils.translation import gettext_lazy as _
|
||||
from redis import Redis
|
||||
|
||||
from django_q.brokers import Broker
|
||||
from django_q.conf import Conf
|
||||
|
||||
|
||||
class Disque(Broker):
|
||||
def enqueue(self, task):
|
||||
retry = Conf.RETRY if Conf.RETRY > 0 else f"{Conf.RETRY} REPLICATE 1"
|
||||
return self.connection.execute_command(
|
||||
f"ADDJOB {self.list_key} {task} 500 RETRY {retry}"
|
||||
).decode()
|
||||
|
||||
def dequeue(self):
|
||||
tasks = self.connection.execute_command(
|
||||
f"GETJOB COUNT {Conf.BULK} TIMEOUT 1000 FROM {self.list_key}"
|
||||
)
|
||||
if tasks:
|
||||
return [(t[1].decode(), t[2].decode()) for t in tasks]
|
||||
|
||||
def queue_size(self):
|
||||
return self.connection.execute_command(f"QLEN {self.list_key}")
|
||||
|
||||
def acknowledge(self, task_id):
|
||||
command = "FASTACK" if Conf.DISQUE_FASTACK else "ACKJOB"
|
||||
return self.connection.execute_command(f"{command} {task_id}")
|
||||
|
||||
def ping(self) -> bool:
|
||||
return self.connection.execute_command("HELLO")[0] > 0
|
||||
|
||||
def delete(self, task_id):
|
||||
return self.connection.execute_command(f"DELJOB {task_id}")
|
||||
|
||||
def fail(self, task_id):
|
||||
return self.delete(task_id)
|
||||
|
||||
def delete_queue(self) -> int:
|
||||
jobs = self.connection.execute_command(f"JSCAN QUEUE {self.list_key}")[1]
|
||||
if jobs:
|
||||
job_ids = " ".join(jid.decode() for jid in jobs)
|
||||
self.connection.execute_command(f"DELJOB {job_ids}")
|
||||
return len(jobs)
|
||||
|
||||
def info(self) -> str:
|
||||
if not self._info:
|
||||
info = self.connection.info("server")
|
||||
self._info = f'Disque {info["disque_version"]}'
|
||||
return self._info
|
||||
|
||||
@staticmethod
|
||||
def get_connection(list_key: str = Conf.PREFIX) -> Redis:
|
||||
if not Conf.DISQUE_NODES:
|
||||
raise redis.exceptions.ConnectionError(_("No Disque nodes configured"))
|
||||
# randomize nodes
|
||||
random.shuffle(Conf.DISQUE_NODES)
|
||||
# find one that works
|
||||
for node in Conf.DISQUE_NODES:
|
||||
host, port = node.split(":")
|
||||
kwargs = {"host": host, "port": port}
|
||||
if Conf.DISQUE_AUTH:
|
||||
kwargs["password"] = Conf.DISQUE_AUTH
|
||||
redis_client = redis.Redis(**kwargs)
|
||||
redis_client.decode_responses = True
|
||||
try:
|
||||
redis_client.execute_command("HELLO")
|
||||
return redis_client
|
||||
except redis.exceptions.ConnectionError:
|
||||
continue
|
||||
raise redis.exceptions.ConnectionError(
|
||||
_("Could not connect to any Disque nodes")
|
||||
)
|
||||
+15
-23
@@ -6,13 +6,10 @@ import signal
|
||||
import socket
|
||||
import traceback
|
||||
import uuid
|
||||
from datetime import datetime
|
||||
from datetime import datetime, timedelta
|
||||
from multiprocessing import Event, Process, Value, current_process
|
||||
from time import sleep
|
||||
|
||||
# External
|
||||
import arrow
|
||||
|
||||
# Django
|
||||
from django import core, db
|
||||
from django.apps.registry import apps
|
||||
@@ -47,6 +44,8 @@ from django_q.signals import post_execute, pre_execute
|
||||
from django_q.signing import BadSignature, SignedPackage
|
||||
from django_q.status import Stat, Status
|
||||
|
||||
from .utils import add_months, add_years
|
||||
|
||||
|
||||
class Cluster:
|
||||
def __init__(self, broker: Broker = None):
|
||||
@@ -635,22 +634,22 @@ def scheduler(broker: Broker = None):
|
||||
q_options["hook"] = s.hook
|
||||
# set up the next run time
|
||||
if s.schedule_type != s.ONCE:
|
||||
next_run = arrow.get(s.next_run)
|
||||
next_run = s.next_run
|
||||
while True:
|
||||
if s.schedule_type == s.MINUTES:
|
||||
next_run = next_run.shift(minutes=+(s.minutes or 1))
|
||||
next_run = next_run + timedelta(minutes=(s.minutes or 1))
|
||||
elif s.schedule_type == s.HOURLY:
|
||||
next_run = next_run.shift(hours=+1)
|
||||
next_run = next_run + timedelta(hours=1)
|
||||
elif s.schedule_type == s.DAILY:
|
||||
next_run = next_run.shift(days=+1)
|
||||
next_run = next_run + timedelta(days=1)
|
||||
elif s.schedule_type == s.WEEKLY:
|
||||
next_run = next_run.shift(weeks=+1)
|
||||
next_run = next_run + timedelta(weeks=1)
|
||||
elif s.schedule_type == s.MONTHLY:
|
||||
next_run = next_run.shift(months=+1)
|
||||
next_run = add_months(next_run, 1)
|
||||
elif s.schedule_type == s.QUARTERLY:
|
||||
next_run = next_run.shift(months=+3)
|
||||
next_run = add_months(next_run, 3)
|
||||
elif s.schedule_type == s.YEARLY:
|
||||
next_run = next_run.shift(years=+1)
|
||||
next_run = add_years(next_run, 1)
|
||||
elif s.schedule_type == s.CRON:
|
||||
if not croniter:
|
||||
raise ImportError(
|
||||
@@ -658,18 +657,11 @@ def scheduler(broker: Broker = None):
|
||||
"Please install croniter to enable cron expressions"
|
||||
)
|
||||
)
|
||||
next_run = arrow.get(
|
||||
croniter(s.cron, localtime()).get_next()
|
||||
)
|
||||
if Conf.CATCH_UP or next_run > arrow.utcnow():
|
||||
next_run = croniter(s.cron, localtime()).get_next(datetime)
|
||||
if Conf.CATCH_UP or next_run > localtime():
|
||||
break
|
||||
# arrow always returns a tz aware datetime, and we don't want
|
||||
# this when we explicitly configured django with USE_TZ=False
|
||||
s.next_run = (
|
||||
next_run.datetime
|
||||
if settings.USE_TZ
|
||||
else next_run.datetime.replace(tzinfo=None)
|
||||
)
|
||||
|
||||
s.next_run = next_run
|
||||
s.repeats += -1
|
||||
# send it to the cluster
|
||||
scheduled_broker = broker
|
||||
|
||||
@@ -45,15 +45,6 @@ class Conf:
|
||||
|
||||
DJANGO_REDIS = conf.get("django_redis", None)
|
||||
|
||||
# Disque broker
|
||||
DISQUE_NODES = conf.get("disque_nodes", None)
|
||||
|
||||
# Optional Authentication
|
||||
DISQUE_AUTH = conf.get("disque_auth", None)
|
||||
|
||||
# Optional Fast acknowledge
|
||||
DISQUE_FASTACK = conf.get("disque_fastack", False)
|
||||
|
||||
# IronMQ broker
|
||||
IRON_MQ = conf.get("iron_mq", None)
|
||||
|
||||
|
||||
@@ -106,11 +106,16 @@ LOGGING = {
|
||||
|
||||
STATIC_URL = "/static/"
|
||||
|
||||
REDIS_HOST = os.environ.get("REDIS_HOST", "redis")
|
||||
|
||||
MONGO_HOST = os.environ.get("MONGO_HOST", "mongo")
|
||||
|
||||
|
||||
# Django Redis
|
||||
CACHES = {
|
||||
"default": {
|
||||
"BACKEND": "django_redis.cache.RedisCache",
|
||||
"LOCATION": "redis://127.0.0.1:6379/0",
|
||||
"LOCATION": f"redis://{REDIS_HOST}:6379/0",
|
||||
"OPTIONS": {
|
||||
"CLIENT_CLASS": "django_redis.client.DefaultClient",
|
||||
"PARSER_CLASS": "redis.connection.HiredisParser",
|
||||
@@ -125,4 +130,5 @@ Q_CLUSTER = {
|
||||
"testing": True,
|
||||
"log_level": "DEBUG",
|
||||
"django_redis": "default",
|
||||
"redis": f"redis://{REDIS_HOST}:6379/0"
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ import redis
|
||||
from django_q.brokers import Broker, get_broker
|
||||
from django_q.conf import Conf
|
||||
from django_q.humanhash import uuid
|
||||
from django_q.tests.settings import REDIS_HOST, MONGO_HOST
|
||||
|
||||
|
||||
def test_broker(monkeypatch):
|
||||
@@ -40,11 +41,11 @@ def test_redis(monkeypatch):
|
||||
broker = get_broker()
|
||||
assert broker.ping() is True
|
||||
assert broker.info() is not None
|
||||
monkeypatch.setattr(Conf, "REDIS", {"host": "127.0.0.1", "port": 7799})
|
||||
monkeypatch.setattr(Conf, "REDIS", {"host": REDIS_HOST, "port": 7799})
|
||||
broker = get_broker()
|
||||
with pytest.raises(Exception):
|
||||
broker.ping()
|
||||
monkeypatch.setattr(Conf, "REDIS", "redis://127.0.0.1:7799")
|
||||
monkeypatch.setattr(Conf, "REDIS", f"redis://{REDIS_HOST}:7799")
|
||||
broker = get_broker()
|
||||
with pytest.raises(Exception):
|
||||
broker.ping()
|
||||
@@ -58,68 +59,6 @@ def test_custom(monkeypatch):
|
||||
assert broker.__class__.__name__ == "Redis"
|
||||
|
||||
|
||||
def test_disque(monkeypatch):
|
||||
monkeypatch.setattr(Conf, "DISQUE_NODES", ["127.0.0.1:7711"])
|
||||
# check broker
|
||||
broker = get_broker(list_key="disque_test")
|
||||
assert broker.ping() is True
|
||||
assert broker.info() is not None
|
||||
# clear before we start
|
||||
broker.delete_queue()
|
||||
# async_task
|
||||
broker.enqueue("test")
|
||||
assert broker.queue_size() == 1
|
||||
# dequeue
|
||||
task = broker.dequeue()[0]
|
||||
assert task[1] == "test"
|
||||
broker.acknowledge(task[0])
|
||||
assert broker.queue_size() == 0
|
||||
# Retry test
|
||||
monkeypatch.setattr(Conf, "RETRY", 1)
|
||||
broker.enqueue("test")
|
||||
assert broker.queue_size() == 1
|
||||
broker.dequeue()
|
||||
assert broker.queue_size() == 0
|
||||
sleep(1.5)
|
||||
assert broker.queue_size() == 1
|
||||
task = broker.dequeue()[0]
|
||||
assert broker.queue_size() == 0
|
||||
broker.acknowledge(task[0])
|
||||
sleep(1.5)
|
||||
assert broker.queue_size() == 0
|
||||
# delete job
|
||||
task_id = broker.enqueue("test")
|
||||
broker.delete(task_id)
|
||||
assert broker.dequeue() is None
|
||||
# fail
|
||||
task_id = broker.enqueue("test")
|
||||
broker.fail(task_id)
|
||||
# bulk test
|
||||
for _ in range(5):
|
||||
broker.enqueue("test")
|
||||
monkeypatch.setattr(Conf, "BULK", 5)
|
||||
monkeypatch.setattr(Conf, "DISQUE_FASTACK", True)
|
||||
tasks = broker.dequeue()
|
||||
for task in tasks:
|
||||
assert task is not None
|
||||
broker.acknowledge(task[0])
|
||||
# test duplicate acknowledge
|
||||
broker.acknowledge(task[0])
|
||||
# delete queue
|
||||
broker.enqueue("test")
|
||||
broker.enqueue("test")
|
||||
broker.delete_queue()
|
||||
assert broker.queue_size() == 0
|
||||
# connection test
|
||||
monkeypatch.setattr(Conf, "DISQUE_NODES", ["127.0.0.1:7798", "127.0.0.1:7799"])
|
||||
with pytest.raises(redis.exceptions.ConnectionError):
|
||||
broker.get_connection()
|
||||
# connection test with no nodes
|
||||
monkeypatch.setattr(Conf, "DISQUE_NODES", None)
|
||||
with pytest.raises(redis.exceptions.ConnectionError):
|
||||
broker.get_connection()
|
||||
|
||||
|
||||
@pytest.mark.skipif(
|
||||
not os.getenv("IRON_MQ_TOKEN"), reason="requires IronMQ credentials"
|
||||
)
|
||||
@@ -312,7 +251,7 @@ def test_orm(monkeypatch):
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_mongo(monkeypatch):
|
||||
monkeypatch.setattr(Conf, "MONGO", {"host": "127.0.0.1", "port": 27017})
|
||||
monkeypatch.setattr(Conf, "MONGO", {"host": MONGO_HOST, "port": 27017})
|
||||
# check broker
|
||||
broker = get_broker(list_key="mongo_test")
|
||||
assert broker.ping() is True
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
from datetime import datetime
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
@@ -32,6 +33,7 @@ from django_q.tasks import (
|
||||
result_group,
|
||||
)
|
||||
from django_q.tests.tasks import TaskError, multiply
|
||||
from django_q.utils import add_months, add_years
|
||||
|
||||
|
||||
class WordClass:
|
||||
@@ -743,3 +745,43 @@ def assert_result(task):
|
||||
def assert_bad_result(task):
|
||||
assert task is not None
|
||||
assert task.success is False
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_add_months():
|
||||
# add some months
|
||||
initial_date = datetime(2020, 2, 2)
|
||||
new_date = add_months(initial_date, 3)
|
||||
assert new_date.year == 2020
|
||||
assert new_date.month == 5
|
||||
assert new_date.day == 2
|
||||
|
||||
# push to next year
|
||||
initial_date = datetime(2020, 11, 2)
|
||||
new_date = add_months(initial_date, 3)
|
||||
assert new_date.year == 2021
|
||||
assert new_date.month == 2
|
||||
assert new_date.day == 2
|
||||
|
||||
# last day of the month
|
||||
initial_date = datetime(2020, 1, 31)
|
||||
new_date = add_months(initial_date, 1)
|
||||
assert new_date.year == 2020
|
||||
assert new_date.month == 2
|
||||
assert new_date.day == 29
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_add_years():
|
||||
# add some months
|
||||
initial_date = datetime(2020, 2, 2)
|
||||
new_date = add_years(initial_date, 1)
|
||||
assert new_date.year == 2021
|
||||
assert new_date.month == 2
|
||||
assert new_date.day == 2
|
||||
|
||||
# test leap year
|
||||
initial_date = datetime(2020, 2, 29)
|
||||
new_date = add_years(initial_date, 1)
|
||||
assert new_date.year == 2021
|
||||
assert new_date.month == 2
|
||||
assert new_date.day == 28
|
||||
|
||||
@@ -3,7 +3,6 @@ from datetime import timedelta
|
||||
from multiprocessing import Event, Value
|
||||
from unittest import mock
|
||||
|
||||
import arrow
|
||||
import pytest
|
||||
from django.core.exceptions import ValidationError
|
||||
from django.db import IntegrityError
|
||||
@@ -128,7 +127,7 @@ def test_scheduler(broker, monkeypatch):
|
||||
assert schedule.repeats == 0
|
||||
assert schedule.last_run() is not None
|
||||
assert schedule.success() is True
|
||||
assert schedule.next_run < arrow.get(timezone.now()).shift(hours=+1)
|
||||
assert schedule.next_run < timezone.now() + timedelta(hours=1)
|
||||
task = fetch(schedule.task)
|
||||
assert task is not None
|
||||
assert task.success is True
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
import datetime
|
||||
from datetime import date
|
||||
from django.utils.timezone import make_aware
|
||||
import calendar
|
||||
|
||||
# credits: https://stackoverflow.com/a/4131114
|
||||
# Made them aware of timezone
|
||||
def add_months(d, months):
|
||||
month = d.month - 1 + months
|
||||
year = d.year + month // 12
|
||||
month = month % 12 + 1
|
||||
day = min(d.day, calendar.monthrange(year,month)[1])
|
||||
return d.replace(year=year, month=month, day=day)
|
||||
|
||||
# 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
|
||||
# Also made them aware of timezone
|
||||
def add_years(d, years):
|
||||
"""Return a date that's `years` years after the date (or datetime)
|
||||
object `d`. Return the same calendar date (month and day) in the
|
||||
destination year, if it exists, otherwise use the previous day
|
||||
(thus changing February 29 to February 28).
|
||||
|
||||
"""
|
||||
try:
|
||||
return d.replace(year = d.year + years)
|
||||
except ValueError:
|
||||
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)
|
||||
|
||||
Reference in New Issue
Block a user