Build improvements (#569)

* Black linting

* Adding Black to dev dependencies and upping minimal python to 3.6.2 for compatibility

* Updating packages

* Removing pip-tools input and exporting requirements with poetry

* Deleting old setup files and test runner

* Trying 1.3.7

* Looser extras requirements to prevent conflicts

* Sorted imports with isort

* Added iSort to dev dependencies

* Fixes localtime for naive setups
This commit is contained in:
Ilan Steemers
2021-05-30 17:10:07 +02:00
committed by GitHub
parent dc95da3a6b
commit 15155c7a99
45 changed files with 1497 additions and 4194 deletions
+1 -1
View File
@@ -1,4 +1,4 @@
VERSION = (1, 3, 6)
VERSION = (1, 3, 7)
default_app_config = "django_q.apps.DjangoQConfig"
+4 -4
View File
@@ -3,7 +3,7 @@ from django.contrib import admin
from django.utils.translation import gettext_lazy as _
from django_q.conf import Conf, croniter
from django_q.models import Success, Failure, Schedule, OrmQ
from django_q.models import Failure, OrmQ, Schedule, Success
from django_q.tasks import async_task
@@ -60,7 +60,7 @@ class FailAdmin(admin.ModelAdmin):
class ScheduleAdmin(admin.ModelAdmin):
""" model admin for schedules """
"""model admin for schedules"""
list_display = (
"id",
@@ -84,7 +84,7 @@ class ScheduleAdmin(admin.ModelAdmin):
class QueueAdmin(admin.ModelAdmin):
""" queue admin for ORM broker """
"""queue admin for ORM broker"""
list_display = ("id", "key", "task_id", "name", "func", "lock")
@@ -100,7 +100,7 @@ class QueueAdmin(admin.ModelAdmin):
def has_add_permission(self, request):
"""Don't allow adds."""
return False
list_filter = ("key",)
+1 -1
View File
@@ -1,7 +1,7 @@
import importlib
from typing import Optional
from django.core.cache import caches, InvalidCacheBackendError
from django.core.cache import InvalidCacheBackendError, caches
from django_q.conf import Conf
+4 -2
View File
@@ -38,7 +38,9 @@ class Sqs(Broker):
if not isinstance(wait_time_second, int):
raise ValueError("receive_message_wait_time_seconds should be int")
if wait_time_second > 20:
raise ValueError("receive_message_wait_time_seconds is invalid. Reason: Must be >= 0 and <= 20")
raise ValueError(
"receive_message_wait_time_seconds is invalid. Reason: Must be >= 0 and <= 20"
)
params.update({"WaitTimeSeconds": wait_time_second})
tasks = self.queue.receive_messages(**params)
@@ -80,7 +82,7 @@ class Sqs(Broker):
config["region_name"] = config["aws_region"]
del config["aws_region"]
if 'receive_message_wait_time_seconds' in config:
if "receive_message_wait_time_seconds" in config:
del config["receive_message_wait_time_seconds"]
return Session(**config)
+2 -1
View File
@@ -2,10 +2,11 @@ import random
# External
import redis
from redis 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
+30 -16
View File
@@ -6,6 +6,7 @@ import signal
import socket
import traceback
import uuid
from datetime import datetime
from multiprocessing import Event, Process, Value, current_process
from time import sleep
@@ -13,13 +14,14 @@ from time import sleep
import arrow
# Django
from django import db, core
from django import core, db
from django.apps.registry import apps
try:
apps.check_apps_ready()
except core.exceptions.AppRegistryNotReady:
import django
django.setup()
from django.conf import settings
@@ -28,21 +30,21 @@ from django.utils.translation import gettext_lazy as _
# Local
import django_q.tasks
from django_q.brokers import get_broker, Broker
from django_q.brokers import Broker, get_broker
from django_q.conf import (
Conf,
croniter,
error_reporter,
get_ppid,
logger,
psutil,
get_ppid,
error_reporter,
croniter,
resource,
)
from django_q.humanhash import humanize
from django_q.models import Task, Success, Schedule
from django_q.models import Schedule, Success, Task
from django_q.queues import Queue
from django_q.signals import pre_execute
from django_q.signing import SignedPackage, BadSignature
from django_q.signing import BadSignature, SignedPackage
from django_q.status import Stat, Status
@@ -485,8 +487,11 @@ def save_task(task, broker: Broker):
existing_task.attempt_count = existing_task.attempt_count + 1
existing_task.save()
if Conf.MAX_ATTEMPTS > 0 and existing_task.attempt_count >= Conf.MAX_ATTEMPTS:
broker.acknowledge(task['ack_id'])
if (
Conf.MAX_ATTEMPTS > 0
and existing_task.attempt_count >= Conf.MAX_ATTEMPTS
):
broker.acknowledge(task["ack_id"])
else:
func = task["func"]
@@ -495,8 +500,8 @@ def save_task(task, broker: Broker):
func = f"{func.__module__}.{func.__name__}"
elif inspect.ismethod(func):
func = (
f'{func.__self__.__module__}.'
f'{func.__self__.__name__}.{func.__name__}'
f"{func.__self__.__module__}."
f"{func.__self__.__name__}.{func.__name__}"
)
Task.objects.create(
id=task["id"],
@@ -510,7 +515,7 @@ def save_task(task, broker: Broker):
result=task["result"],
group=task.get("group"),
success=task["success"],
attempt_count=1
attempt_count=1,
)
except Exception as e:
logger.error(e)
@@ -582,7 +587,9 @@ def scheduler(broker: Broker = None):
Schedule.objects.select_for_update()
.exclude(repeats=0)
.filter(next_run__lt=timezone.now())
.filter(db.models.Q(cluster__isnull=True) | db.models.Q(cluster=Conf.PREFIX))
.filter(
db.models.Q(cluster__isnull=True) | db.models.Q(cluster=Conf.PREFIX)
)
):
args = ()
kwargs = {}
@@ -627,7 +634,7 @@ def scheduler(broker: Broker = None):
)
)
next_run = arrow.get(
croniter(s.cron, timezone.localtime()).get_next()
croniter(s.cron, localtime()).get_next()
)
if Conf.CATCH_UP or next_run > arrow.utcnow():
break
@@ -643,7 +650,7 @@ def scheduler(broker: Broker = None):
scheduled_broker = broker
try:
scheduled_broker = get_broker(q_options["broker_name"])
except: # invalid broker_name or non existing broker with broker_name
except: # invalid broker_name or non existing broker with broker_name
pass
q_options["broker"] = scheduled_broker
q_options["group"] = q_options.get("group", s.name or s.id)
@@ -734,4 +741,11 @@ def rss_check():
return resource.getrusage(resource.RUSAGE_SELF).ru_maxrss >= Conf.MAX_RSS
elif psutil:
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()
+9 -7
View File
@@ -135,12 +135,13 @@ class Conf:
RETRY = conf.get("retry", 60)
# Verify if retry and timeout settings are correct
if not TIMEOUT or (TIMEOUT > RETRY):
warn("""Retry and timeout are misconfigured. Set retry larger than timeout,
if not TIMEOUT or (TIMEOUT > RETRY):
warn(
"""Retry and timeout are misconfigured. Set retry larger than timeout,
failure to do so will cause the tasks to be retriggered before completion.
See https://django-q.readthedocs.io/en/latest/configure.html#retry for details.""")
See https://django-q.readthedocs.io/en/latest/configure.html#retry for details."""
)
# Sets the amount of tasks the cluster will try to pop off the broker.
# If it supports bulk gets.
BULK = conf.get("bulk", 1)
@@ -176,7 +177,7 @@ class Conf:
ERROR_REPORTER = conf.get("error_reporter", {})
# Optional attempt count. set to 0 for infinite attempts
MAX_ATTEMPTS = conf.get('max_attempts', 0)
MAX_ATTEMPTS = conf.get("max_attempts", 0)
# OSX doesn't implement qsize because of missing sem_getvalue()
try:
@@ -200,7 +201,8 @@ class Conf:
# to manage workarounds during testing
TESTING = conf.get("testing", False)
# logger
logger = logging.getLogger("django-q")
+4 -9
View File
@@ -2,15 +2,10 @@ import datetime
import time
import zlib
from django.core.signing import (
BadSignature,
SignatureExpired,
b64_decode,
JSONSerializer,
Signer as Sgnr,
TimestampSigner as TsS,
dumps,
)
from django.core.signing import BadSignature, JSONSerializer, SignatureExpired
from django.core.signing import Signer as Sgnr
from django.core.signing import TimestampSigner as TsS
from django.core.signing import b64_decode, dumps
from django.utils import baseconv
from django.utils.crypto import constant_time_compare
from django.utils.encoding import force_bytes, force_str
+266 -45
View File
@@ -4,50 +4,269 @@ humanhash: Human-readable representations of digests.
The simplest ways to use this module are the :func:`humanize` and :func:`uuid`
functions. For tighter control over the output, see :class:`HumanHasher`.
"""
from argparse import ArgumentError
import operator
import uuid as uuidlib
from argparse import ArgumentError
from functools import reduce
DEFAULT_WORDLIST = (
'ack', 'alabama', 'alanine', 'alaska', 'alpha', 'angel', 'apart', 'april',
'arizona', 'arkansas', 'artist', 'asparagus', 'aspen', 'august', 'autumn',
'avocado', 'bacon', 'bakerloo', 'batman', 'beer', 'berlin', 'beryllium',
'black', 'blossom', 'blue', 'bluebird', 'bravo', 'bulldog', 'burger',
'butter', 'california', 'carbon', 'cardinal', 'carolina', 'carpet', 'cat',
'ceiling', 'charlie', 'chicken', 'coffee', 'cola', 'cold', 'colorado',
'comet', 'connecticut', 'crazy', 'cup', 'dakota', 'december', 'delaware',
'delta', 'diet', 'don', 'double', 'early', 'earth', 'east', 'echo',
'edward', 'eight', 'eighteen', 'eleven', 'emma', 'enemy', 'equal',
'failed', 'fanta', 'fifteen', 'fillet', 'finch', 'fish', 'five', 'fix',
'floor', 'florida', 'football', 'four', 'fourteen', 'foxtrot', 'freddie',
'friend', 'fruit', 'gee', 'georgia', 'glucose', 'golf', 'green', 'grey',
'hamper', 'happy', 'harry', 'hawaii', 'helium', 'high', 'hot', 'hotel',
'hydrogen', 'idaho', 'illinois', 'india', 'indigo', 'ink', 'iowa',
'island', 'item', 'jersey', 'jig', 'johnny', 'juliet', 'july', 'jupiter',
'kansas', 'kentucky', 'kilo', 'king', 'kitten', 'lactose', 'lake', 'lamp',
'lemon', 'leopard', 'lima', 'lion', 'lithium', 'london', 'louisiana',
'low', 'magazine', 'magnesium', 'maine', 'mango', 'march', 'mars',
'maryland', 'massachusetts', 'may', 'mexico', 'michigan', 'mike',
'minnesota', 'mirror', 'mississippi', 'missouri', 'mobile', 'mockingbird',
'monkey', 'montana', 'moon', 'mountain', 'muppet', 'music', 'nebraska',
'neptune', 'network', 'nevada', 'nine', 'nineteen', 'nitrogen', 'north',
'november', 'nuts', 'october', 'ohio', 'oklahoma', 'one', 'orange',
'oranges', 'oregon', 'oscar', 'oven', 'oxygen', 'papa', 'paris', 'pasta',
'pennsylvania', 'pip', 'pizza', 'pluto', 'potato', 'princess', 'purple',
'quebec', 'queen', 'quiet', 'red', 'river', 'robert', 'robin', 'romeo',
'rugby', 'sad', 'salami', 'saturn', 'september', 'seven', 'seventeen',
'shade', 'sierra', 'single', 'sink', 'six', 'sixteen', 'skylark', 'snake',
'social', 'sodium', 'solar', 'south', 'spaghetti', 'speaker', 'spring',
'stairway', 'steak', 'stream', 'summer', 'sweet', 'table', 'tango', 'ten',
'tennessee', 'tennis', 'texas', 'thirteen', 'three', 'timing', 'triple',
'twelve', 'twenty', 'two', 'uncle', 'undress', 'uniform', 'uranus', 'utah',
'vegan', 'venus', 'vermont', 'victor', 'video', 'violet', 'virginia',
'washington', 'west', 'whiskey', 'white', 'william', 'winner', 'winter',
'wisconsin', 'wolfram', 'wyoming', 'xray', 'yankee', 'yellow', 'zebra',
'zulu')
"ack",
"alabama",
"alanine",
"alaska",
"alpha",
"angel",
"apart",
"april",
"arizona",
"arkansas",
"artist",
"asparagus",
"aspen",
"august",
"autumn",
"avocado",
"bacon",
"bakerloo",
"batman",
"beer",
"berlin",
"beryllium",
"black",
"blossom",
"blue",
"bluebird",
"bravo",
"bulldog",
"burger",
"butter",
"california",
"carbon",
"cardinal",
"carolina",
"carpet",
"cat",
"ceiling",
"charlie",
"chicken",
"coffee",
"cola",
"cold",
"colorado",
"comet",
"connecticut",
"crazy",
"cup",
"dakota",
"december",
"delaware",
"delta",
"diet",
"don",
"double",
"early",
"earth",
"east",
"echo",
"edward",
"eight",
"eighteen",
"eleven",
"emma",
"enemy",
"equal",
"failed",
"fanta",
"fifteen",
"fillet",
"finch",
"fish",
"five",
"fix",
"floor",
"florida",
"football",
"four",
"fourteen",
"foxtrot",
"freddie",
"friend",
"fruit",
"gee",
"georgia",
"glucose",
"golf",
"green",
"grey",
"hamper",
"happy",
"harry",
"hawaii",
"helium",
"high",
"hot",
"hotel",
"hydrogen",
"idaho",
"illinois",
"india",
"indigo",
"ink",
"iowa",
"island",
"item",
"jersey",
"jig",
"johnny",
"juliet",
"july",
"jupiter",
"kansas",
"kentucky",
"kilo",
"king",
"kitten",
"lactose",
"lake",
"lamp",
"lemon",
"leopard",
"lima",
"lion",
"lithium",
"london",
"louisiana",
"low",
"magazine",
"magnesium",
"maine",
"mango",
"march",
"mars",
"maryland",
"massachusetts",
"may",
"mexico",
"michigan",
"mike",
"minnesota",
"mirror",
"mississippi",
"missouri",
"mobile",
"mockingbird",
"monkey",
"montana",
"moon",
"mountain",
"muppet",
"music",
"nebraska",
"neptune",
"network",
"nevada",
"nine",
"nineteen",
"nitrogen",
"north",
"november",
"nuts",
"october",
"ohio",
"oklahoma",
"one",
"orange",
"oranges",
"oregon",
"oscar",
"oven",
"oxygen",
"papa",
"paris",
"pasta",
"pennsylvania",
"pip",
"pizza",
"pluto",
"potato",
"princess",
"purple",
"quebec",
"queen",
"quiet",
"red",
"river",
"robert",
"robin",
"romeo",
"rugby",
"sad",
"salami",
"saturn",
"september",
"seven",
"seventeen",
"shade",
"sierra",
"single",
"sink",
"six",
"sixteen",
"skylark",
"snake",
"social",
"sodium",
"solar",
"south",
"spaghetti",
"speaker",
"spring",
"stairway",
"steak",
"stream",
"summer",
"sweet",
"table",
"tango",
"ten",
"tennessee",
"tennis",
"texas",
"thirteen",
"three",
"timing",
"triple",
"twelve",
"twenty",
"two",
"uncle",
"undress",
"uniform",
"uranus",
"utah",
"vegan",
"venus",
"vermont",
"victor",
"video",
"violet",
"virginia",
"washington",
"west",
"whiskey",
"white",
"william",
"winner",
"winter",
"wisconsin",
"wolfram",
"wyoming",
"xray",
"yankee",
"yellow",
"zebra",
"zulu",
)
class HumanHasher:
@@ -70,7 +289,7 @@ class HumanHasher:
raise ArgumentError("Wordlist must have exactly 256 items")
self.wordlist = wordlist
def humanize(self, hexdigest, words=4, separator='-'):
def humanize(self, hexdigest, words=4, separator="-"):
"""
Humanize a given hexadecimal digest.
@@ -84,7 +303,10 @@ class HumanHasher:
"""
# Gets a list of byte values between 0-255.
bytes = [int(x, 16) for x in list(map(''.join, list(zip(hexdigest[::2], hexdigest[1::2]))))]
bytes = [
int(x, 16)
for x in list(map("".join, list(zip(hexdigest[::2], hexdigest[1::2]))))
]
# Compress an arbitrary number of bytes to `words`.
compressed = self.compress(bytes, words)
# Map the compressed byte values through the word list.
@@ -115,10 +337,9 @@ class HumanHasher:
# Split `bytes` into `target` segments.
seg_size = length // target
segments = [bytes[i * seg_size:(i + 1) * seg_size]
for i in range(target)]
segments = [bytes[i * seg_size : (i + 1) * seg_size] for i in range(target)]
# Catch any left-over bytes in the last segment.
segments[-1].extend(bytes[target * seg_size:])
segments[-1].extend(bytes[target * seg_size :])
# Use a simple XOR checksum-like function for compression.
checksum = lambda bytes: reduce(operator.xor, bytes, 0)
@@ -134,7 +355,7 @@ class HumanHasher:
as :meth:`humanize` (they'll be passed straight through).
"""
digest = str(uuidlib.uuid4()).replace('-', '')
digest = str(uuidlib.uuid4()).replace("-", "")
return self.humanize(digest, **params), digest
+5 -5
View File
@@ -10,15 +10,15 @@ class Command(BaseCommand):
def add_arguments(self, parser):
parser.add_argument(
'--run-once',
action='store_true',
dest='run_once',
"--run-once",
action="store_true",
dest="run_once",
default=False,
help='Run once and then stop.',
help="Run once and then stop.",
)
def handle(self, *args, **options):
q = Cluster()
q.start()
if options.get('run_once', False):
if options.get("run_once", False):
q.stop()
+1 -1
View File
@@ -3,7 +3,7 @@ from django.utils.translation import gettext as _
from django_q import VERSION
from django_q.conf import Conf
from django_q.monitor import info, get_ids
from django_q.monitor import get_ids, info
class Command(BaseCommand):
+1 -1
View File
@@ -27,5 +27,5 @@ class Command(BaseCommand):
def handle(self, *args, **options):
memory(
run_once=options.get("run_once", False),
workers=options.get("workers", False)
workers=options.get("workers", False),
)
+2 -2
View File
@@ -1,6 +1,6 @@
from django.db import models, migrations
import picklefield.fields
import django.utils.timezone
import picklefield.fields
from django.db import migrations, models
class Migration(migrations.Migration):
@@ -1,4 +1,4 @@
from django.db import models, migrations
from django.db import migrations, models
class Migration(migrations.Migration):
@@ -1,4 +1,4 @@
from django.db import models, migrations
from django.db import migrations, models
class Migration(migrations.Migration):
@@ -1,4 +1,4 @@
from django.db import models, migrations
from django.db import migrations, models
class Migration(migrations.Migration):
@@ -1,4 +1,4 @@
from django.db import models, migrations
from django.db import migrations, models
class Migration(migrations.Migration):
@@ -1,4 +1,4 @@
from django.db import models, migrations
from django.db import migrations, models
class Migration(migrations.Migration):
+1 -1
View File
@@ -1,4 +1,4 @@
from django.db import models, migrations
from django.db import migrations, models
class Migration(migrations.Migration):
@@ -1,5 +1,5 @@
from django.db import migrations
import picklefield.fields
from django.db import migrations
class Migration(migrations.Migration):
@@ -1,6 +1,7 @@
# Generated by Django 3.0.8 on 2020-07-02 16:08
from django.db import migrations, models
import django_q.models
+58 -18
View File
@@ -5,15 +5,16 @@ from blessed import Terminal
# django
from django.db import connection
from django.db.models import Sum, F
from django.db.models import F, Sum
from django.utils import timezone
from django.utils.translation import gettext as _
from django_q import VERSION, models
from django_q.brokers import get_broker
# local
from django_q.conf import Conf
from django_q.status import Stat
from django_q.brokers import get_broker
from django_q import models, VERSION
# optional
try:
@@ -27,7 +28,7 @@ def get_process_mb(pid):
process = psutil.Process(pid)
mb_used = round(process.memory_info().rss / 1024 ** 2, 2)
except psutil.NoSuchProcess:
mb_used = 'NO_PROCESS_FOUND'
mb_used = "NO_PROCESS_FOUND"
return mb_used
@@ -39,7 +40,10 @@ def monitor(run_once=False, broker=None):
with term.fullscreen(), term.hidden_cursor(), term.cbreak():
val = None
start_width = int(term.width / 8)
while val not in ("q", "Q",):
while val not in (
"q",
"Q",
):
col_width = int(term.width / 8)
# In case of resize
if col_width != start_width:
@@ -294,7 +298,11 @@ def memory(run_once=False, workers=False, broker=None):
broker.ping()
if not psutil:
print(term.clear_eos())
print(term.white_on_red("Cannot start \"qmemory\" command. Missing \"psutil\" library."))
print(
term.white_on_red(
'Cannot start "qmemory" command. Missing "psutil" library.'
)
)
return
with term.fullscreen(), term.hidden_cursor(), term.cbreak():
MEMORY_AVAILABLE_LOWEST_PERCENTAGE = 100.0
@@ -319,11 +327,15 @@ def memory(run_once=False, workers=False, broker=None):
)
print(
term.move(0, 2 * col_width)
+ term.black_on_green(term.center(_("Available (%)"), width=col_width - 1))
+ term.black_on_green(
term.center(_("Available (%)"), width=col_width - 1)
)
)
print(
term.move(0, 3 * col_width)
+ term.black_on_green(term.center(_("Available (MB)"), width=col_width - 1))
+ term.black_on_green(
term.center(_("Available (MB)"), width=col_width - 1)
)
)
print(
term.move(0, 4 * col_width)
@@ -331,24 +343,37 @@ def memory(run_once=False, workers=False, broker=None):
)
print(
term.move(0, 5 * col_width)
+ term.black_on_green(term.center(_("Sentinel (MB)"), width=col_width - 1))
+ term.black_on_green(
term.center(_("Sentinel (MB)"), width=col_width - 1)
)
)
print(
term.move(0, 6 * col_width)
+ term.black_on_green(term.center(_("Monitor (MB)"), width=col_width - 1))
+ term.black_on_green(
term.center(_("Monitor (MB)"), width=col_width - 1)
)
)
print(
term.move(0, 7 * col_width)
+ term.black_on_green(term.center(_("Workers (MB)"), width=col_width - 1))
+ term.black_on_green(
term.center(_("Workers (MB)"), width=col_width - 1)
)
)
row = 2
stats = Stat.get_all(broker=broker)
print(term.clear_eos())
for stat in stats:
# memory available (%)
memory_available_percentage = round(psutil.virtual_memory().available * 100 / psutil.virtual_memory().total, 2)
memory_available_percentage = round(
psutil.virtual_memory().available
* 100
/ psutil.virtual_memory().total,
2,
)
# memory available (MB)
memory_available = round(psutil.virtual_memory().available / 1024 ** 2, 2)
memory_available = round(
psutil.virtual_memory().available / 1024 ** 2, 2
)
if memory_available_percentage < MEMORY_AVAILABLE_LOWEST_PERCENTAGE:
MEMORY_AVAILABLE_LOWEST_PERCENTAGE = memory_available_percentage
MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT = timezone.now()
@@ -370,7 +395,10 @@ def memory(run_once=False, workers=False, broker=None):
)
print(
term.move(row, 4 * col_width)
+ term.center(round(psutil.virtual_memory().total / 1024 ** 2, 2), width=col_width - 1)
+ term.center(
round(psutil.virtual_memory().total / 1024 ** 2, 2),
width=col_width - 1,
)
)
print(
term.move(row, 5 * col_width)
@@ -378,7 +406,10 @@ def memory(run_once=False, workers=False, broker=None):
)
print(
term.move(row, 6 * col_width)
+ term.center(get_process_mb(getattr(stat, 'monitor', None)), width=col_width - 1)
+ term.center(
get_process_mb(getattr(stat, "monitor", None)),
width=col_width - 1,
)
)
workers_mb = 0
for worker_pid in stat.workers:
@@ -388,7 +419,9 @@ def memory(run_once=False, workers=False, broker=None):
workers_mb += result
print(
term.move(row, 7 * col_width)
+ term.center(workers_mb or 'NO_PROCESSES_FOUND', width=col_width - 1)
+ term.center(
workers_mb or "NO_PROCESSES_FOUND", width=col_width - 1
)
)
row += 1
# each worker's memory usage
@@ -402,7 +435,12 @@ def memory(run_once=False, workers=False, broker=None):
for worker_num in range(Conf.WORKERS):
print(
term.move(row, (worker_num + 1) * col_width)
+ term.black_on_cyan(term.center("Worker #{} (MB)".format(worker_num + 1), width=col_width - 1))
+ term.black_on_cyan(
term.center(
"Worker #{} (MB)".format(worker_num + 1),
width=col_width - 1,
)
)
)
row += 2
for stat in stats:
@@ -422,7 +460,9 @@ def memory(run_once=False, workers=False, broker=None):
term.move(row, 0)
+ _("Available lowest (%): {} ({})").format(
str(MEMORY_AVAILABLE_LOWEST_PERCENTAGE),
MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT.strftime('%Y-%m-%d %H:%M:%S+00:00')
MEMORY_AVAILABLE_LOWEST_PERCENTAGE_AT.strftime(
"%Y-%m-%d %H:%M:%S+00:00"
),
)
)
# for testing
+7 -7
View File
@@ -7,7 +7,7 @@ import sys
class SharedCounter:
""" A synchronized shared counter.
"""A synchronized shared counter.
The locking done by multiprocessing.Value ensures that only a single
process or thread may read or write the in-memory ctypes object. However,
@@ -24,18 +24,18 @@ class SharedCounter:
self.count = multiprocessing.Value("i", n)
def increment(self, n=1):
""" Increment the counter by n (default = 1) """
"""Increment the counter by n (default = 1)"""
with self.count.get_lock():
self.count.value += n
@property
def value(self):
""" Return the value of the counter """
"""Return the value of the counter"""
return self.count.value
class Queue(multiprocessing.queues.Queue):
""" A portable implementation of multiprocessing.Queue.
"""A portable implementation of multiprocessing.Queue.
Because of multithreading / multiprocessing semantics, Queue.qsize() may
raise the NotImplementedError exception on Unix platforms like Mac OS X
@@ -57,7 +57,7 @@ class Queue(multiprocessing.queues.Queue):
self.size = SharedCounter(0)
def __getstate__(self):
return super(Queue, self).__getstate__() + (self.size, )
return super(Queue, self).__getstate__() + (self.size,)
def __setstate__(self, state):
super(Queue, self).__setstate__(state[:-1])
@@ -73,9 +73,9 @@ class Queue(multiprocessing.queues.Queue):
return x
def qsize(self) -> int:
""" Reliable implementation of multiprocessing.Queue.qsize() """
"""Reliable implementation of multiprocessing.Queue.qsize()"""
return self.size.value
def empty(self) -> bool:
""" Reliable implementation of multiprocessing.Queue.empty() """
"""Reliable implementation of multiprocessing.Queue.empty()"""
return not self.qsize() > 0
+2 -1
View File
@@ -1,7 +1,7 @@
import importlib
from django.db.models.signals import post_save
from django.dispatch import receiver, Signal
from django.dispatch import Signal, receiver
from django.utils.translation import gettext_lazy as _
from django_q.conf import logger
@@ -31,6 +31,7 @@ def call_hook(sender, instance, **kwargs):
)
)
# args: task
pre_enqueue = Signal()
+2 -2
View File
@@ -3,9 +3,9 @@ from typing import Union
from django.utils import timezone
from django_q.brokers import get_broker, Broker
from django_q.brokers import Broker, get_broker
from django_q.conf import Conf, logger
from django_q.signing import SignedPackage, BadSignature
from django_q.signing import BadSignature, SignedPackage
class Status:
+3 -3
View File
@@ -1,11 +1,11 @@
"""Provides task functionality."""
# Standard
from multiprocessing import Value
from time import sleep, time
# django
from django.db import IntegrityError
from django.utils import timezone
from multiprocessing import Value
# local
from django_q.brokers import get_broker
@@ -153,7 +153,7 @@ def result(task_id, wait=0, cached=Conf.CACHED):
def result_cached(task_id, wait=0, broker=None):
"""
Return the result from the cache backend
Return the result from the cache backend
"""
if not broker:
broker = get_broker()
@@ -755,7 +755,7 @@ class AsyncTask:
def _sync(pack):
"""Simulate a package travelling through the cluster."""
from django_q.cluster import worker, monitor
from django_q.cluster import monitor, worker
task_queue = Queue()
result_queue = Queue()
+49 -46
View File
@@ -1,4 +1,5 @@
import os
import django
BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
@@ -8,7 +9,7 @@ BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
# See https://docs.djangoproject.com/en/2.2/howto/deployment/checklist/
# SECURITY WARNING: keep the secret key used in production secret!
SECRET_KEY = ')cqmpi+p@n&!u&fu@!m@9h&1bz9mwmstsahe)nf!ms+c$uc=x7'
SECRET_KEY = ")cqmpi+p@n&!u&fu@!m@9h&1bz9mwmstsahe)nf!ms+c$uc=x7"
# SECURITY WARNING: don't run with debug turned on in production!
DEBUG = True
@@ -19,41 +20,41 @@ ALLOWED_HOSTS = []
# Application definition
INSTALLED_APPS = (
'django.contrib.admin',
'django.contrib.auth',
'django.contrib.contenttypes',
'django.contrib.sessions',
'django.contrib.messages',
'django.contrib.staticfiles',
'django_q',
'django_redis'
"django.contrib.admin",
"django.contrib.auth",
"django.contrib.contenttypes",
"django.contrib.sessions",
"django.contrib.messages",
"django.contrib.staticfiles",
"django_q",
"django_redis",
)
MIDDLEWARE_CLASSES = (
'django.contrib.sessions.middleware.SessionMiddleware',
'django.middleware.common.CommonMiddleware',
'django.middleware.csrf.CsrfViewMiddleware',
'django.contrib.auth.middleware.AuthenticationMiddleware',
'django.contrib.messages.middleware.MessageMiddleware',
'django.middleware.clickjacking.XFrameOptionsMiddleware',
"django.contrib.sessions.middleware.SessionMiddleware",
"django.middleware.common.CommonMiddleware",
"django.middleware.csrf.CsrfViewMiddleware",
"django.contrib.auth.middleware.AuthenticationMiddleware",
"django.contrib.messages.middleware.MessageMiddleware",
"django.middleware.clickjacking.XFrameOptionsMiddleware",
)
MIDDLEWARE = MIDDLEWARE_CLASSES
ROOT_URLCONF = 'tests.urls'
ROOT_URLCONF = "tests.urls"
TEMPLATES = [
{
'BACKEND': 'django.template.backends.django.DjangoTemplates',
'DIRS': [],
'APP_DIRS': True,
'OPTIONS': {
'context_processors': [
'django.template.context_processors.debug',
'django.template.context_processors.request',
'django.contrib.auth.context_processors.auth',
'django.contrib.messages.context_processors.messages',
"BACKEND": "django.template.backends.django.DjangoTemplates",
"DIRS": [],
"APP_DIRS": True,
"OPTIONS": {
"context_processors": [
"django.template.context_processors.debug",
"django.template.context_processors.request",
"django.contrib.auth.context_processors.auth",
"django.contrib.messages.context_processors.messages",
],
},
},
@@ -64,9 +65,9 @@ TEMPLATES = [
# https://docs.djangoproject.com/en/2.2/ref/settings/#databases
DATABASES = {
'default': {
'ENGINE': 'django.db.backends.sqlite3',
'NAME': os.path.join(BASE_DIR, 'db.sqlite3'),
"default": {
"ENGINE": "django.db.backends.sqlite3",
"NAME": os.path.join(BASE_DIR, "db.sqlite3"),
}
}
@@ -74,9 +75,9 @@ DATABASES = {
# Internationalization
# https://docs.djangoproject.com/en/2.2/topics/i18n/
LANGUAGE_CODE = 'en-us'
LANGUAGE_CODE = "en-us"
TIME_ZONE = 'UTC'
TIME_ZONE = "UTC"
USE_I18N = True
@@ -85,17 +86,17 @@ USE_L10N = True
USE_TZ = True
LOGGING = {
'version': 1,
'disable_existing_loggers': False,
'handlers': {
'console': {
'class': 'logging.StreamHandler',
"version": 1,
"disable_existing_loggers": False,
"handlers": {
"console": {
"class": "logging.StreamHandler",
},
},
'loggers': {
'django_q': {
'handlers': ['console'],
'level': 'INFO',
"loggers": {
"django_q": {
"handlers": ["console"],
"level": "INFO",
},
},
}
@@ -103,7 +104,7 @@ LOGGING = {
# Static files (CSS, JavaScript, Images)
# https://docs.djangoproject.com/en/2.2/howto/static-files/
STATIC_URL = '/static/'
STATIC_URL = "/static/"
# Django Redis
CACHES = {
@@ -113,13 +114,15 @@ CACHES = {
"OPTIONS": {
"CLIENT_CLASS": "django_redis.client.DefaultClient",
"PARSER_CLASS": "redis.connection.HiredisParser",
}
},
}
}
# Django Q specific
Q_CLUSTER = {'name': 'django_q_test',
'cpu_affinity': 1,
'testing': True,
'log_level': 'DEBUG',
'django_redis': 'default'}
Q_CLUSTER = {
"name": "django_q_test",
"cpu_affinity": 1,
"testing": True,
"log_level": "DEBUG",
"django_redis": "default",
}
+44 -38
View File
@@ -1,80 +1,86 @@
import pytest
from django.urls import reverse
from django.utils import timezone
import pytest
from django_q.tasks import schedule
from django_q.models import Task, Failure, OrmQ
from django_q.humanhash import uuid
from django_q.conf import Conf
from django_q.humanhash import uuid
from django_q.models import Failure, OrmQ, Task
from django_q.signing import SignedPackage
from django_q.tasks import schedule
@pytest.mark.django_db
def test_admin_views(admin_client, monkeypatch):
monkeypatch.setattr(Conf, 'ORM', 'default')
s = schedule('schedule.test')
monkeypatch.setattr(Conf, "ORM", "default")
s = schedule("schedule.test")
tag = uuid()
f = Task.objects.create(
id=tag[1],
name=tag[0],
func='test.fail',
func="test.fail",
started=timezone.now(),
stopped=timezone.now(),
success=False)
success=False,
)
tag = uuid()
t = Task.objects.create(
id=tag[1],
name=tag[0],
func='test.success',
func="test.success",
started=timezone.now(),
stopped=timezone.now(),
success=True)
success=True,
)
q = OrmQ.objects.create(
key='test',
payload=SignedPackage.dumps({'id': 1, 'func': 'test', 'name': 'test'}))
key="test",
payload=SignedPackage.dumps({"id": 1, "func": "test", "name": "test"}),
)
admin_urls = (
# schedule
reverse('admin:django_q_schedule_changelist'),
reverse('admin:django_q_schedule_add'),
reverse('admin:django_q_schedule_change', args=(s.id,)),
reverse('admin:django_q_schedule_history', args=(s.id,)),
reverse('admin:django_q_schedule_delete', args=(s.id,)),
reverse("admin:django_q_schedule_changelist"),
reverse("admin:django_q_schedule_add"),
reverse("admin:django_q_schedule_change", args=(s.id,)),
reverse("admin:django_q_schedule_history", args=(s.id,)),
reverse("admin:django_q_schedule_delete", args=(s.id,)),
# success
reverse('admin:django_q_success_changelist'),
reverse('admin:django_q_success_change', args=(t.id,)),
reverse('admin:django_q_success_history', args=(t.id,)),
reverse('admin:django_q_success_delete', args=(t.id,)),
reverse("admin:django_q_success_changelist"),
reverse("admin:django_q_success_change", args=(t.id,)),
reverse("admin:django_q_success_history", args=(t.id,)),
reverse("admin:django_q_success_delete", args=(t.id,)),
# failure
reverse('admin:django_q_failure_changelist'),
reverse('admin:django_q_failure_change', args=(f.id,)),
reverse('admin:django_q_failure_history', args=(f.id,)),
reverse('admin:django_q_failure_delete', args=(f.id,)),
reverse("admin:django_q_failure_changelist"),
reverse("admin:django_q_failure_change", args=(f.id,)),
reverse("admin:django_q_failure_history", args=(f.id,)),
reverse("admin:django_q_failure_delete", args=(f.id,)),
# orm queue
reverse('admin:django_q_ormq_changelist'),
reverse('admin:django_q_ormq_change', args=(q.id,)),
reverse('admin:django_q_ormq_history', args=(q.id,)),
reverse('admin:django_q_ormq_delete', args=(q.id,)),
reverse("admin:django_q_ormq_changelist"),
reverse("admin:django_q_ormq_change", args=(q.id,)),
reverse("admin:django_q_ormq_history", args=(q.id,)),
reverse("admin:django_q_ormq_delete", args=(q.id,)),
)
for url in admin_urls:
response = admin_client.get(url)
assert response.status_code == 200
# resubmit the failure
url = reverse('admin:django_q_failure_changelist')
data = {'action': 'retry_failed',
'_selected_action': [f.pk]}
url = reverse("admin:django_q_failure_changelist")
data = {"action": "retry_failed", "_selected_action": [f.pk]}
response = admin_client.post(url, data)
assert response.status_code == 302
assert Failure.objects.filter(name=f.id).exists() is False
# change q
url = reverse('admin:django_q_ormq_change', args=(q.id,))
data = {'key': 'default', 'payload': 'test', 'lock_0': '2015-09-17', 'lock_1': '14:31:51', '_save': 'Save'}
url = reverse("admin:django_q_ormq_change", args=(q.id,))
data = {
"key": "default",
"payload": "test",
"lock_0": "2015-09-17",
"lock_1": "14:31:51",
"_save": "Save",
}
response = admin_client.post(url, data)
assert response.status_code == 302
# delete q
url = reverse('admin:django_q_ormq_delete', args=(q.id,))
data = {'post': 'yes'}
url = reverse("admin:django_q_ormq_delete", args=(q.id,))
data = {"post": "yes"}
response = admin_client.post(url, data)
assert response.status_code == 302
+2 -2
View File
@@ -4,7 +4,7 @@ from time import sleep
import pytest
import redis
from django_q.brokers import get_broker, Broker
from django_q.brokers import Broker, get_broker
from django_q.conf import Conf
from django_q.humanhash import uuid
@@ -194,7 +194,7 @@ def canceled_sqs(monkeypatch):
"aws_region": os.getenv("AWS_REGION"),
"aws_access_key_id": os.getenv("AWS_ACCESS_KEY_ID"),
"aws_secret_access_key": os.getenv("AWS_SECRET_ACCESS_KEY"),
"receive_message_wait_time_seconds": 20
"receive_message_wait_time_seconds": 20,
},
)
# check broker
+62 -40
View File
@@ -2,17 +2,30 @@ from multiprocessing import Event, Value
import pytest
from django_q.cluster import pusher, worker, monitor
from django_q.conf import Conf
from django_q.tasks import async_task, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \
async_iter, Chain, async_chain, Iter, AsyncTask
from django_q.brokers import get_broker
from django_q.cluster import monitor, pusher, worker
from django_q.conf import Conf
from django_q.queues import Queue
from django_q.tasks import (
AsyncTask,
Chain,
Iter,
async_chain,
async_iter,
async_task,
count_group,
delete_cached,
delete_group,
fetch,
fetch_group,
result,
result_group,
)
@pytest.fixture
def broker(monkeypatch):
monkeypatch.setattr(Conf, 'DJANGO_REDIS', 'default')
monkeypatch.setattr(Conf, "DJANGO_REDIS", "default")
return get_broker()
@@ -20,16 +33,16 @@ def broker(monkeypatch):
def test_cached(broker):
broker.purge_queue()
broker.cache.clear()
group = 'cache_test'
group = "cache_test"
# queue the tests
task_id = async_task('math.copysign', 1, -1, cached=True, broker=broker)
async_task('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async_task('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async_task('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async_task('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async_task('math.copysign', 1, -1, cached=True, broker=broker, group=group)
async_task('math.popysign', 1, -1, cached=True, broker=broker, group=group)
iter_id = async_iter('math.floor', [i for i in range(10)], cached=True)
task_id = async_task("math.copysign", 1, -1, cached=True, broker=broker)
async_task("math.copysign", 1, -1, cached=True, broker=broker, group=group)
async_task("math.copysign", 1, -1, cached=True, broker=broker, group=group)
async_task("math.copysign", 1, -1, cached=True, broker=broker, group=group)
async_task("math.copysign", 1, -1, cached=True, broker=broker, group=group)
async_task("math.copysign", 1, -1, cached=True, broker=broker, group=group)
async_task("math.popysign", 1, -1, cached=True, broker=broker, group=group)
iter_id = async_iter("math.floor", [i for i in range(10)], cached=True)
# test wait on cache
# test wait timeout
assert result(task_id, wait=10, cached=True) is None
@@ -48,11 +61,11 @@ def test_cached(broker):
pusher(task_queue, stop_event, broker=broker)
assert broker.queue_size() == 0
assert task_queue.qsize() == task_count
task_queue.put('STOP')
task_queue.put("STOP")
result_queue = Queue()
worker(task_queue, result_queue, Value('f', -1))
worker(task_queue, result_queue, Value("f", -1))
assert result_queue.qsize() == task_count
result_queue.put('STOP')
result_queue.put("STOP")
monitor(result_queue)
assert result_queue.qsize() == 0
# assert results
@@ -85,10 +98,10 @@ def test_iter(broker):
it = [i for i in range(10)]
it2 = [(1, -1), (2, -1), (3, -4), (5, 6)]
it3 = (1, 2, 3, 4, 5)
t = async_iter('math.floor', it, sync=True)
t2 = async_iter('math.copysign', it2, sync=True)
t3 = async_iter('math.floor', it3, sync=True)
t4 = async_iter('math.floor', (1,), sync=True)
t = async_iter("math.floor", it, sync=True)
t2 = async_iter("math.copysign", it2, sync=True)
t3 = async_iter("math.floor", it3, sync=True)
t4 = async_iter("math.floor", (1,), sync=True)
result_t = result(t)
assert result_t is not None
task_t = fetch(t)
@@ -97,7 +110,7 @@ def test_iter(broker):
assert result(t3) is not None
assert result(t4)[0] == 1
# test iter class
i = Iter('math.copysign', sync=True, cached=True)
i = Iter("math.copysign", sync=True, cached=True)
i.append(1, -1)
i.append(2, -1)
i.append(3, -4)
@@ -118,9 +131,9 @@ def test_chain(broker):
broker.purge_queue()
broker.cache.clear()
task_chain = Chain(sync=True)
task_chain.append('math.floor', 1)
task_chain.append('math.copysign', 1, -1)
task_chain.append('math.floor', 2)
task_chain.append("math.floor", 1)
task_chain.append("math.copysign", 1, -1)
task_chain.append("math.floor", 2)
assert task_chain.length() == 3
assert task_chain.current() is None
task_chain.run()
@@ -130,7 +143,7 @@ def test_chain(broker):
t = task_chain.fetch()
assert len(t) == task_chain.length()
task_chain.cached = True
task_chain.append('math.floor', 3)
task_chain.append("math.floor", 3)
assert task_chain.length() == 4
task_chain.run()
r = task_chain.result(wait=1000)
@@ -139,16 +152,20 @@ def test_chain(broker):
t = task_chain.fetch()
assert len(t) == task_chain.length()
# test single
rid = async_chain(['django_q.tests.tasks.hello', 'django_q.tests.tasks.hello'], sync=True, cached=True)
assert result_group(rid, cached=True) == ['hello', 'hello']
rid = async_chain(
["django_q.tests.tasks.hello", "django_q.tests.tasks.hello"],
sync=True,
cached=True,
)
assert result_group(rid, cached=True) == ["hello", "hello"]
@pytest.mark.django_db
def test_asynctask_class(broker, monkeypatch):
broker.purge_queue()
broker.cache.clear()
a = AsyncTask('math.copysign')
assert a.func == 'math.copysign'
a = AsyncTask("math.copysign")
assert a.func == "math.copysign"
a.args = (1, -1)
assert a.started is False
a.cached = True
@@ -161,29 +178,34 @@ def test_asynctask_class(broker, monkeypatch):
assert a.result() == -1
assert a.fetch().result == -1
# again with kwargs
a = AsyncTask('math.copysign', 1, -1, cached=True, sync=True, broker=broker)
a = AsyncTask("math.copysign", 1, -1, cached=True, sync=True, broker=broker)
a.run()
assert a.result() == -1
# with q_options
a = AsyncTask('math.copysign', 1, -1, q_options={'cached': True, 'sync': False, 'broker': broker})
a = AsyncTask(
"math.copysign",
1,
-1,
q_options={"cached": True, "sync": False, "broker": broker},
)
assert not a.sync
a.sync = True
assert a.kwargs['q_options']['sync'] is True
assert a.kwargs["q_options"]["sync"] is True
a.run()
assert a.result() == -1
a.group = 'async_class_test'
assert a.group == 'async_class_test'
a.group = "async_class_test"
assert a.group == "async_class_test"
a.save = False
assert not a.save
a.hook = 'djq.tests.tasks.hello'
assert a.hook == 'djq.tests.tasks.hello'
a.hook = "djq.tests.tasks.hello"
assert a.hook == "djq.tests.tasks.hello"
assert a.started is False
a.run()
assert a.result_group() == [-1]
assert a.fetch_group() == [a.fetch()]
# global overrides
monkeypatch.setattr(Conf, 'SYNC', True)
monkeypatch.setattr(Conf, 'CACHED', True)
a = AsyncTask('math.floor', 1.5)
monkeypatch.setattr(Conf, "SYNC", True)
monkeypatch.setattr(Conf, "CACHED", True)
a = AsyncTask("math.floor", 1.5)
a.run()
assert a.result() == 1
+242 -162
View File
@@ -1,25 +1,34 @@
import os
import sys
import threading
import uuid as uuidlib
from multiprocessing import Event, Value
from time import sleep
from django.utils import timezone
import uuid as uuidlib
import os
import pytest
from django.utils import timezone
myPath = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, myPath + '/../')
sys.path.insert(0, myPath + "/../")
from django_q.cluster import Cluster, Sentinel, pusher, worker, monitor, save_task
from django_q.humanhash import DEFAULT_WORDLIST, uuid
from django_q.tasks import fetch, fetch_group, async_task, result, result_group, count_group, delete_group, queue_size
from django_q.models import Task, Success
from django_q.brokers import Broker, get_broker
from django_q.cluster import Cluster, Sentinel, monitor, pusher, save_task, worker
from django_q.conf import Conf
from django_q.status import Stat
from django_q.brokers import get_broker, Broker
from django_q.tests.tasks import multiply, TaskError
from django_q.humanhash import DEFAULT_WORDLIST, uuid
from django_q.models import Success, Task
from django_q.queues import Queue
from django_q.status import Stat
from django_q.tasks import (
async_task,
count_group,
delete_group,
fetch,
fetch_group,
queue_size,
result,
result_group,
)
from django_q.tests.tasks import TaskError, multiply
class WordClass:
@@ -32,7 +41,7 @@ class WordClass:
@pytest.fixture
def broker(monkeypatch):
monkeypatch.setattr(Conf, 'DJANGO_REDIS', 'default')
monkeypatch.setattr(Conf, "DJANGO_REDIS", "default")
return get_broker()
@@ -42,19 +51,21 @@ def test_redis_connection(broker):
@pytest.mark.django_db
def test_sync(broker):
task = async_task('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker, sync=True)
task = async_task(
"django_q.tests.tasks.count_letters", DEFAULT_WORDLIST, broker=broker, sync=True
)
assert result(task) == 1506
@pytest.mark.django_db
def test_sync_raise_exception(broker):
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)
@pytest.mark.django_db
def test_cluster_initial(broker):
broker.list_key = 'initial_test:q'
broker.list_key = "initial_test:q"
broker.delete_queue()
c = Cluster(broker=broker)
assert c.sentinel is None
@@ -80,16 +91,23 @@ def test_sentinel():
stop_event = Event()
stop_event.set()
cluster_id = uuidlib.uuid4()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=get_broker('sentinel_test:q'))
s = Sentinel(
stop_event,
start_event,
cluster_id=cluster_id,
broker=get_broker("sentinel_test:q"),
)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
@pytest.mark.django_db
def test_cluster(broker):
broker.list_key = 'cluster_test:q'
broker.list_key = "cluster_test:q"
broker.delete_queue()
task = async_task('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker)
task = async_task(
"django_q.tests.tasks.count_letters", DEFAULT_WORDLIST, broker=broker
)
assert broker.queue_size() == 1
task_queue = Queue()
assert task_queue.qsize() == 0
@@ -102,12 +120,12 @@ def test_cluster(broker):
assert task_queue.qsize() == 1
assert queue_size(broker=broker) == 0
# Test work
task_queue.put('STOP')
worker(task_queue, result_queue, Value('f', -1))
task_queue.put("STOP")
worker(task_queue, result_queue, Value("f", -1))
assert task_queue.qsize() == 0
assert result_queue.qsize() == 1
# Test monitor
result_queue.put('STOP')
result_queue.put("STOP")
monitor(result_queue)
assert result_queue.qsize() == 0
# check result
@@ -117,33 +135,63 @@ def test_cluster(broker):
@pytest.mark.django_db
def test_enqueue(broker, admin_user):
broker.list_key = 'cluster_test:q'
broker.list_key = "cluster_test:q"
broker.delete_queue()
a = async_task('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result',
broker=broker)
b = async_task('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_cluster.assert_result',
broker=broker)
a = async_task(
"django_q.tests.tasks.count_letters",
DEFAULT_WORDLIST,
hook="django_q.tests.test_cluster.assert_result",
broker=broker,
)
b = async_task(
"django_q.tests.tasks.count_letters2",
WordClass(),
hook="django_q.tests.test_cluster.assert_result",
broker=broker,
)
# unknown argument
c = async_task('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany',
hook='django_q.tests.test_cluster.assert_bad_result', broker=broker)
c = async_task(
"django_q.tests.tasks.count_letters",
DEFAULT_WORDLIST,
"oneargumentoomany",
hook="django_q.tests.test_cluster.assert_bad_result",
broker=broker,
)
# unknown function
d = async_task('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_cluster.assert_bad_result',
broker=broker)
d = async_task(
"django_q.tests.tasks.does_not_exist",
WordClass(),
hook="django_q.tests.test_cluster.assert_bad_result",
broker=broker,
)
# function without result
e = async_task('django_q.tests.tasks.countdown', 100000, broker=broker)
e = async_task("django_q.tests.tasks.countdown", 100000, broker=broker)
# function as instance
f = async_task(multiply, 753, 2, hook=assert_result, broker=broker)
# model as argument
g = async_task('django_q.tests.tasks.get_task_name', Task(name='John'), broker=broker)
g = async_task(
"django_q.tests.tasks.get_task_name", Task(name="John"), broker=broker
)
# args,kwargs, group and broken hook
h = async_task('django_q.tests.tasks.word_multiply', 2, word='django', hook='fail.me', broker=broker)
h = async_task(
"django_q.tests.tasks.word_multiply",
2,
word="django",
hook="fail.me",
broker=broker,
)
# args unpickle test
j = async_task('django_q.tests.tasks.get_user_id', admin_user, broker=broker, group='test_j')
j = async_task(
"django_q.tests.tasks.get_user_id", admin_user, broker=broker, group="test_j"
)
# q_options and save opt_out test
k = async_task('django_q.tests.tasks.get_user_id', admin_user,
q_options={'broker': broker, 'group': 'test_k', 'save': False, 'timeout': 90})
k = async_task(
"django_q.tests.tasks.get_user_id",
admin_user,
q_options={"broker": broker, "group": "test_k", "save": False, "timeout": 90},
)
# test unicode
assert Task(name='Amalia').__str__()=='Amalia'
assert Task(name="Amalia").__str__() == "Amalia"
# check if everything has a task id
assert isinstance(a, str)
assert isinstance(b, str)
@@ -166,19 +214,19 @@ def test_enqueue(broker, admin_user):
pusher(task_queue, stop_event, broker=broker)
assert broker.queue_size() == 0
assert task_queue.qsize() == task_count
task_queue.put('STOP')
task_queue.put("STOP")
# test wait timeout
assert result(j, wait=10) is None
assert fetch(j, wait=10) is None
assert result_group('test_j', wait=10) is None
assert result_group('test_j', count=2, wait=10) is None
assert fetch_group('test_j', wait=10) is None
assert fetch_group('test_j', count=2, wait=10) is None
assert result_group("test_j", wait=10) is None
assert result_group("test_j", count=2, wait=10) is None
assert fetch_group("test_j", wait=10) is None
assert fetch_group("test_j", count=2, wait=10) is None
# let a worker handle them
result_queue = Queue()
worker(task_queue, result_queue, Value('f', -1))
worker(task_queue, result_queue, Value("f", -1))
assert result_queue.qsize() == task_count
result_queue.put('STOP')
result_queue.put("STOP")
# store the results
monitor(result_queue)
assert result_queue.qsize() == 0
@@ -215,7 +263,7 @@ def test_enqueue(broker, admin_user):
result_g = fetch(g)
assert result_g is not None
assert result_g.success is True
assert result(g) == 'John'
assert result(g) == "John"
# task h
result_h = fetch(h)
assert result_h is not None
@@ -230,19 +278,19 @@ def test_enqueue(broker, admin_user):
assert fetch(result_j.name) == result_j
assert result(result_j.name) == result_j.result
# groups
assert result_group('test_j')[0] == result_j.result
assert result_group("test_j")[0] == result_j.result
assert result_j.group_result()[0] == result_j.result
assert result_group('test_j', failures=True)[0] == result_j.result
assert result_group("test_j", failures=True)[0] == result_j.result
assert result_j.group_result(failures=True)[0] == result_j.result
assert fetch_group('test_j')[0].id == [result_j][0].id
assert fetch_group('test_j', failures=False)[0].id == [result_j][0].id
assert count_group('test_j') == 1
assert fetch_group("test_j")[0].id == [result_j][0].id
assert fetch_group("test_j", failures=False)[0].id == [result_j][0].id
assert count_group("test_j") == 1
assert result_j.group_count() == 1
assert count_group('test_j', failures=True) == 0
assert count_group("test_j", failures=True) == 0
assert result_j.group_count(failures=True) == 0
assert delete_group('test_j') == 1
assert delete_group("test_j") == 1
assert result_j.group_delete() == 0
deleted_group = delete_group('test_j', tasks=True)
deleted_group = delete_group("test_j", tasks=True)
assert deleted_group is None or deleted_group[0] == 0 # Django 1.9
deleted_group = result_j.group_delete(tasks=True)
assert deleted_group is None or deleted_group[0] == 0 # Django 1.9
@@ -254,22 +302,31 @@ def test_enqueue(broker, admin_user):
@pytest.mark.django_db
@pytest.mark.parametrize('cluster_config_timeout, async_task_kwargs', (
(1, {}),
(10, {'timeout': 1}),
(None, {'timeout': 1}),
))
@pytest.mark.parametrize(
"cluster_config_timeout, async_task_kwargs",
(
(1, {}),
(10, {"timeout": 1}),
(None, {"timeout": 1}),
),
)
def test_timeout(broker, cluster_config_timeout, async_task_kwargs):
# set up the Sentinel
broker.list_key = 'timeout_test:q'
broker.list_key = "timeout_test:q"
broker.purge_queue()
async_task('time.sleep', 5, broker=broker, **async_task_kwargs)
async_task("time.sleep", 5, broker=broker, **async_task_kwargs)
start_event = Event()
stop_event = Event()
cluster_id = uuidlib.uuid4()
# Set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker, timeout=cluster_config_timeout)
s = Sentinel(
stop_event,
start_event,
cluster_id=cluster_id,
broker=broker,
timeout=cluster_config_timeout,
)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 1
@@ -277,23 +334,32 @@ def test_timeout(broker, cluster_config_timeout, async_task_kwargs):
@pytest.mark.django_db
@pytest.mark.parametrize('cluster_config_timeout, async_task_kwargs', (
(5, {}),
(10, {'timeout': 5}),
(1, {'timeout': 5}),
(None, {'timeout': 5}),
))
@pytest.mark.parametrize(
"cluster_config_timeout, async_task_kwargs",
(
(5, {}),
(10, {"timeout": 5}),
(1, {"timeout": 5}),
(None, {"timeout": 5}),
),
)
def test_timeout_task_finishes(broker, cluster_config_timeout, async_task_kwargs):
# set up the Sentinel
broker.list_key = 'timeout_test:q'
broker.list_key = "timeout_test:q"
broker.purge_queue()
async_task('time.sleep', 3, broker=broker, **async_task_kwargs)
async_task("time.sleep", 3, broker=broker, **async_task_kwargs)
start_event = Event()
stop_event = Event()
cluster_id = uuidlib.uuid4()
# Set a timer to stop the Sentinel
threading.Timer(6, stop_event.set).start()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker, timeout=cluster_config_timeout)
s = Sentinel(
stop_event,
start_event,
cluster_id=cluster_id,
broker=broker,
timeout=cluster_config_timeout,
)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 0
@@ -303,70 +369,71 @@ def test_timeout_task_finishes(broker, cluster_config_timeout, async_task_kwargs
@pytest.mark.django_db
def test_recycle(broker, monkeypatch):
# set up the Sentinel
broker.list_key = 'test_recycle_test:q'
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
broker.list_key = "test_recycle_test:q"
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
start_event = Event()
stop_event = Event()
cluster_id = uuidlib.uuid4()
# override settings
monkeypatch.setattr(Conf, 'RECYCLE', 2)
monkeypatch.setattr(Conf, 'WORKERS', 1)
monkeypatch.setattr(Conf, "RECYCLE", 2)
monkeypatch.setattr(Conf, "WORKERS", 1)
# set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 1
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
task_queue = Queue()
result_queue = Queue()
# push two tasks
pusher(task_queue, stop_event, broker=broker)
pusher(task_queue, stop_event, broker=broker)
# worker should exit on recycle
worker(task_queue, result_queue, Value('f', -1))
worker(task_queue, result_queue, Value("f", -1))
# check if the work has been done
assert result_queue.qsize() == 2
# save_limit test
monkeypatch.setattr(Conf, 'SAVE_LIMIT', 1)
result_queue.put('STOP')
monkeypatch.setattr(Conf, "SAVE_LIMIT", 1)
result_queue.put("STOP")
# run monitor
monitor(result_queue)
assert Success.objects.count() == Conf.SAVE_LIMIT
broker.delete_queue()
@pytest.mark.django_db
def test_max_rss(broker, monkeypatch):
# set up the Sentinel
broker.list_key = 'test_max_rss_test:q'
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
broker.list_key = "test_max_rss_test:q"
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
start_event = Event()
stop_event = Event()
cluster_id = uuidlib.uuid4()
# override settings
monkeypatch.setattr(Conf, 'MAX_RSS', 40000)
monkeypatch.setattr(Conf, 'WORKERS', 1)
monkeypatch.setattr(Conf, "MAX_RSS", 40000)
monkeypatch.setattr(Conf, "WORKERS", 1)
# set a timer to stop the Sentinel
threading.Timer(3, stop_event.set).start()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker)
assert start_event.is_set()
assert s.status() == Conf.STOPPED
assert s.reincarnations == 1
async_task('django_q.tests.tasks.multiply', 2, 2, broker=broker)
async_task("django_q.tests.tasks.multiply", 2, 2, broker=broker)
task_queue = Queue()
result_queue = Queue()
# push the task
pusher(task_queue, stop_event, broker=broker)
# worker should exit on recycle
worker(task_queue, result_queue, Value('f', -1))
worker(task_queue, result_queue, Value("f", -1))
# check if the work has been done
assert result_queue.qsize() == 1
# save_limit test
monkeypatch.setattr(Conf, 'SAVE_LIMIT', 1)
result_queue.put('STOP')
monkeypatch.setattr(Conf, "SAVE_LIMIT", 1)
result_queue.put("STOP")
# run monitor
monitor(result_queue)
assert Success.objects.count() == Conf.SAVE_LIMIT
@@ -375,13 +442,15 @@ def test_max_rss(broker, monkeypatch):
@pytest.mark.django_db
def test_bad_secret(broker, monkeypatch):
broker.list_key = 'test_bad_secret:q'
async_task('math.copysign', 1, -1, broker=broker)
broker.list_key = "test_bad_secret:q"
async_task("math.copysign", 1, -1, broker=broker)
stop_event = Event()
stop_event.set()
start_event = Event()
cluster_id = uuidlib.uuid4()
s = Sentinel(stop_event, start_event, cluster_id=cluster_id, broker=broker, start=False)
s = Sentinel(
stop_event, start_event, cluster_id=cluster_id, broker=broker, start=False
)
Stat(s).save()
# change the SECRET
monkeypatch.setattr(Conf, "SECRET_KEY", "OOPS")
@@ -391,87 +460,94 @@ def test_bad_secret(broker, monkeypatch):
task_queue = Queue()
pusher(task_queue, stop_event, broker=broker)
result_queue = Queue()
task_queue.put('STOP')
worker(task_queue, result_queue, Value('f', -1), )
task_queue.put("STOP")
worker(
task_queue,
result_queue,
Value("f", -1),
)
assert result_queue.qsize() == 0
broker.delete_queue()
@pytest.mark.django_db
def test_attempt_count(broker, monkeypatch):
monkeypatch.setattr(Conf, 'MAX_ATTEMPTS', 3)
monkeypatch.setattr(Conf, "MAX_ATTEMPTS", 3)
tag = uuid()
task = {'id': tag[1],
'name': tag[0],
'func': 'math.copysign',
'args': (1, -1),
'kwargs': {},
'started': timezone.now(),
'stopped': timezone.now(),
'success': False,
'result': None}
task = {
"id": tag[1],
"name": tag[0],
"func": "math.copysign",
"args": (1, -1),
"kwargs": {},
"started": timezone.now(),
"stopped": timezone.now(),
"success": False,
"result": None,
}
# initial save - no success
save_task(task, broker)
assert Task.objects.filter(id=task['id']).exists()
saved_task = Task.objects.get(id=task['id'])
assert Task.objects.filter(id=task["id"]).exists()
saved_task = Task.objects.get(id=task["id"])
assert saved_task.attempt_count == 1
sleep(0.5)
# second save
old_stopped = task['stopped']
task['stopped'] = timezone.now()
old_stopped = task["stopped"]
task["stopped"] = timezone.now()
save_task(task, broker)
saved_task = Task.objects.get(id=task['id'])
saved_task = Task.objects.get(id=task["id"])
assert saved_task.attempt_count == 2
# third save -
task['stopped'] = timezone.now()
task["stopped"] = timezone.now()
save_task(task, broker)
saved_task = Task.objects.get(id=task['id'])
saved_task = Task.objects.get(id=task["id"])
assert saved_task.attempt_count == 3
# task should be removed from queue
assert broker.queue_size() == 0
@pytest.mark.django_db
def test_update_failed(broker):
tag = uuid()
task = {'id': tag[1],
'name': tag[0],
'func': 'math.copysign',
'args': (1, -1),
'kwargs': {},
'started': timezone.now(),
'stopped': timezone.now(),
'success': False,
'result': None}
task = {
"id": tag[1],
"name": tag[0],
"func": "math.copysign",
"args": (1, -1),
"kwargs": {},
"started": timezone.now(),
"stopped": timezone.now(),
"success": False,
"result": None,
}
# initial save - no success
save_task(task, broker)
assert Task.objects.filter(id=task['id']).exists()
saved_task = Task.objects.get(id=task['id'])
assert Task.objects.filter(id=task["id"]).exists()
saved_task = Task.objects.get(id=task["id"])
assert saved_task.success is False
sleep(0.5)
# second save - no success
old_stopped = task['stopped']
task['stopped'] = timezone.now()
old_stopped = task["stopped"]
task["stopped"] = timezone.now()
save_task(task, broker)
saved_task = Task.objects.get(id=task['id'])
saved_task = Task.objects.get(id=task["id"])
assert saved_task.stopped > old_stopped
# third save - success
task['stopped'] = timezone.now()
task['result'] = 'result'
task['success'] = True
task["stopped"] = timezone.now()
task["result"] = "result"
task["success"] = True
save_task(task, broker)
saved_task = Task.objects.get(id=task['id'])
saved_task = Task.objects.get(id=task["id"])
assert saved_task.success is True
# fourth save - no success
task['result'] = None
task['success'] = False
task['stopped'] = old_stopped
task["result"] = None
task["success"] = False
task["stopped"] = old_stopped
save_task(task, broker)
# should not overwrite success
saved_task = Task.objects.get(id=task['id'])
saved_task = Task.objects.get(id=task["id"])
assert saved_task.success is True
assert saved_task.result == 'result'
assert saved_task.result == "result"
@pytest.mark.django_db
@@ -486,47 +562,51 @@ def test_acknowledge_failure_override():
self.acknowledgements[task_id] = count + 1
tag = uuid()
task_fail_ack = {'id': tag[1],
'name': tag[0],
'ack_id': 'test_fail_ack_id',
'ack_failure': True,
'func': 'math.copysign',
'args': (1, -1),
'kwargs': {},
'started': timezone.now(),
'stopped': timezone.now(),
'success': False,
'result': None}
task_fail_ack = {
"id": tag[1],
"name": tag[0],
"ack_id": "test_fail_ack_id",
"ack_failure": True,
"func": "math.copysign",
"args": (1, -1),
"kwargs": {},
"started": timezone.now(),
"stopped": timezone.now(),
"success": False,
"result": None,
}
tag = uuid()
task_fail_no_ack = task_fail_ack.copy()
task_fail_no_ack.update({'id': tag[1],
'name': tag[0],
'ack_id': 'test_fail_no_ack_id'})
del task_fail_no_ack['ack_failure']
task_fail_no_ack.update(
{"id": tag[1], "name": tag[0], "ack_id": "test_fail_no_ack_id"}
)
del task_fail_no_ack["ack_failure"]
tag = uuid()
task_success_ack = task_fail_ack.copy()
task_success_ack.update({
'id': tag[1],
'name': tag[0],
'ack_id': 'test_success_ack_id',
'success': True,
})
del task_success_ack['ack_failure']
task_success_ack.update(
{
"id": tag[1],
"name": tag[0],
"ack_id": "test_success_ack_id",
"success": True,
}
)
del task_success_ack["ack_failure"]
result_queue = Queue()
result_queue.put(task_fail_ack)
result_queue.put(task_fail_no_ack)
result_queue.put(task_success_ack)
result_queue.put('STOP')
broker = VerifyAckMockBroker(list_key='key')
result_queue.put("STOP")
broker = VerifyAckMockBroker(list_key="key")
monitor(result_queue, broker)
assert broker.acknowledgements.get('test_fail_ack_id') == 1
assert broker.acknowledgements.get('test_fail_no_ack_id') is None
assert broker.acknowledgements.get('test_success_ack_id') == 1
assert broker.acknowledgements.get("test_fail_ack_id") == 1
assert broker.acknowledgements.get("test_fail_no_ack_id") is None
assert broker.acknowledgements.get("test_success_ack_id") == 1
@pytest.mark.django_db
+7 -7
View File
@@ -4,22 +4,22 @@ from django.core.management import call_command
@pytest.mark.django_db
def test_qcluster():
call_command('qcluster', run_once=True)
call_command("qcluster", run_once=True)
@pytest.mark.django_db
def test_qmonitor():
call_command('qmonitor', run_once=True)
call_command("qmonitor", run_once=True)
@pytest.mark.django_db
def test_qinfo():
call_command('qinfo')
call_command('qinfo', config=True)
call_command('qinfo', ids=True)
call_command("qinfo")
call_command("qinfo", config=True)
call_command("qinfo", ids=True)
@pytest.mark.django_db
def test_qmemory():
call_command('qmemory', run_once=True)
call_command('qmemory', workers=True, run_once=True)
call_command("qmemory", run_once=True)
call_command("qmemory", workers=True, run_once=True)
+9 -8
View File
@@ -1,12 +1,13 @@
import pytest
import uuid
from django_q.tasks import async_task
import pytest
from django_q.brokers import get_broker
from django_q.cluster import Cluster
from django_q.monitor import monitor, info, get_ids
from django_q.status import Stat
from django_q.conf import Conf
from django_q.monitor import get_ids, info, monitor
from django_q.status import Stat
from django_q.tasks import async_task
@pytest.mark.django_db
@@ -28,9 +29,9 @@ def test_monitor(monkeypatch):
break
assert found_c
# test lock size
monkeypatch.setattr(Conf, 'ORM', 'default')
b = get_broker('monitor_test')
b.enqueue('test')
monkeypatch.setattr(Conf, "ORM", "default")
b = get_broker("monitor_test")
b.enqueue("test")
b.dequeue()
assert b.lock_size() == 1
monitor(run_once=True, broker=b)
@@ -48,4 +49,4 @@ def test_info():
def do_sync():
async_task('django_q.tests.tasks.countdown', 1, sync=True, save=True)
async_task("django_q.tests.tasks.countdown", 1, sync=True, save=True)
+162 -122
View File
@@ -9,89 +9,102 @@ from django.core.exceptions import ValidationError
from django.db import IntegrityError
from django.test import override_settings
from django.utils import timezone
from django.utils.timezone import is_naive
from django_q.brokers import get_broker, Broker
from django_q.cluster import pusher, worker, monitor, scheduler
from django_q.brokers import Broker, get_broker
from django_q.cluster import monitor, pusher, scheduler, worker, localtime
from django_q.conf import Conf
from django_q.queues import Queue
from django_q.tasks import Schedule, fetch, schedule as create_schedule
from django_q.tasks import Schedule, fetch
from django_q.tasks import schedule as create_schedule
from django_q.tests.settings import BASE_DIR
from django_q.tests.testing_utilities.multiple_database_routers import (TestingReplicaDatabaseRouter,
TestingMultipleAppsDatabaseRouter)
from django_q.tests.testing_utilities.multiple_database_routers import (
TestingMultipleAppsDatabaseRouter,
TestingReplicaDatabaseRouter,
)
@pytest.fixture
def broker(monkeypatch) -> Broker:
"""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()
@pytest.fixture
def orm_broker(monkeypatch) -> None:
"""Patches the Conf object setting the ORM attribute to a database named default."""
monkeypatch.setattr(Conf, 'ORM', 'default')
monkeypatch.setattr(Conf, "ORM", "default")
@pytest.fixture
def orm_no_replica_broker(orm_broker, monkeypatch) -> Broker:
"""Generates a Broker with a disabled read replica database configuration."""
monkeypatch.setattr(Conf, 'HAS_REPLICA', False)
return get_broker(list_key='scheduler_test:q')
monkeypatch.setattr(Conf, "HAS_REPLICA", False)
return get_broker(list_key="scheduler_test:q")
@pytest.fixture
def orm_replica_broker(orm_broker, monkeypatch) -> Broker:
"""Generates a Broker with read replica database configuration."""
monkeypatch.setattr(Conf, 'HAS_REPLICA', True)
return get_broker(list_key='scheduler_test:q')
monkeypatch.setattr(Conf, "HAS_REPLICA", True)
return get_broker(list_key="scheduler_test:q")
REPLICA_DATABASE_ROUTERS = [f"{TestingReplicaDatabaseRouter.__module__}.{TestingReplicaDatabaseRouter.__name__}"]
REPLICA_DATABASE_ROUTERS = [
f"{TestingReplicaDatabaseRouter.__module__}.{TestingReplicaDatabaseRouter.__name__}"
]
REPLICA_DATABASES = {
'default': {
'ENGINE': 'django.db.backends.sqlite3',
'NAME': os.path.join(BASE_DIR, 'db.sqlite3'),
"default": {
"ENGINE": "django.db.backends.sqlite3",
"NAME": os.path.join(BASE_DIR, "db.sqlite3"),
},
'replica': {
'ENGINE': 'django.db.backends.sqlite3',
'NAME': os.path.join(BASE_DIR, 'db.sqlite3'),
"replica": {
"ENGINE": "django.db.backends.sqlite3",
"NAME": os.path.join(BASE_DIR, "db.sqlite3"),
},
}
MULTIPLE_APPS_DATABASE_ROUTERS = [
f"{TestingMultipleAppsDatabaseRouter.__module__}.{TestingMultipleAppsDatabaseRouter.__name__}"]
f"{TestingMultipleAppsDatabaseRouter.__module__}.{TestingMultipleAppsDatabaseRouter.__name__}"
]
MULTIPLE_APPS_DATABASES = {
'default': {
'ENGINE': 'django.db.backends.sqlite3',
'NAME': os.path.join(BASE_DIR, 'db.sqlite3'),
"default": {
"ENGINE": "django.db.backends.sqlite3",
"NAME": os.path.join(BASE_DIR, "db.sqlite3"),
},
'admin': {
'ENGINE': 'django.db.backends.sqlite3',
'NAME': os.path.join(BASE_DIR, 'db.sqlite3'),
"admin": {
"ENGINE": "django.db.backends.sqlite3",
"NAME": os.path.join(BASE_DIR, "db.sqlite3"),
},
}
@pytest.mark.django_db
def test_scheduler(broker, monkeypatch):
broker.list_key = 'scheduler_test:q'
broker.list_key = "scheduler_test:q"
broker.delete_queue()
schedule = create_schedule('math.copysign',
1, -1,
name='test math',
hook='django_q.tests.tasks.result',
schedule_type=Schedule.HOURLY,
repeats=1)
schedule = create_schedule(
"math.copysign",
1,
-1,
name="test math",
hook="django_q.tests.tasks.result",
schedule_type=Schedule.HOURLY,
repeats=1,
)
assert schedule.last_run() is None
# check duplicate constraint
with pytest.raises(IntegrityError):
schedule = create_schedule('math.copysign',
1, -1,
name='test math',
hook='django_q.tests.tasks.result',
schedule_type=Schedule.HOURLY,
repeats=1)
schedule = create_schedule(
"math.copysign",
1,
-1,
name="test math",
hook="django_q.tests.tasks.result",
schedule_type=Schedule.HOURLY,
repeats=1,
)
# run scheduler
scheduler(broker=broker)
# set up the workflow
@@ -102,12 +115,12 @@ def test_scheduler(broker, monkeypatch):
pusher(task_queue, stop_event, broker=broker)
assert task_queue.qsize() == 1
assert broker.queue_size() == 0
task_queue.put('STOP')
task_queue.put("STOP")
# let a worker handle them
result_queue = Queue()
worker(task_queue, result_queue, Value('b', -1))
worker(task_queue, result_queue, Value("b", -1))
assert result_queue.qsize() == 1
result_queue.put('STOP')
result_queue.put("STOP")
# store the results
monitor(result_queue)
assert result_queue.qsize() == 0
@@ -121,99 +134,113 @@ def test_scheduler(broker, monkeypatch):
assert task.success is True
assert task.result < 0
# Once schedule with delete
once_schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='django',
schedule_type=Schedule.ONCE,
repeats=-1,
hook='django_q.tests.tasks.result'
)
assert hasattr(once_schedule, 'pk') is True
once_schedule = create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="django",
schedule_type=Schedule.ONCE,
repeats=-1,
hook="django_q.tests.tasks.result",
)
assert hasattr(once_schedule, "pk") is True
# negative repeats
always_schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='django',
schedule_type=Schedule.DAILY,
repeats=-1,
hook='django_q.tests.tasks.result'
)
assert hasattr(always_schedule, 'pk') is True
always_schedule = create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="django",
schedule_type=Schedule.DAILY,
repeats=-1,
hook="django_q.tests.tasks.result",
)
assert hasattr(always_schedule, "pk") is True
# Minute schedule
minute_schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='django',
schedule_type=Schedule.MINUTES,
minutes=10)
assert hasattr(minute_schedule, 'pk') is True
minute_schedule = create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="django",
schedule_type=Schedule.MINUTES,
minutes=10,
)
assert hasattr(minute_schedule, "pk") is True
# Cron schedule
cron_schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='django',
schedule_type=Schedule.CRON,
cron="0 22 * * 1-5")
assert hasattr(cron_schedule, 'pk') is True
cron_schedule = create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="django",
schedule_type=Schedule.CRON,
cron="0 22 * * 1-5",
)
assert hasattr(cron_schedule, "pk") is True
assert cron_schedule.full_clean() is None
assert cron_schedule.__str__() == 'django_q.tests.tasks.word_multiply'
assert cron_schedule.__str__() == "django_q.tests.tasks.word_multiply"
with pytest.raises(ValidationError):
create_schedule('django_q.tests.tasks.word_multiply',
2,
word='django',
schedule_type=Schedule.CRON,
cron="0 22 * * 1-12")
create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="django",
schedule_type=Schedule.CRON,
cron="0 22 * * 1-12",
)
# All other types
for t in Schedule.TYPE:
if t[0] == Schedule.CRON:
continue
schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='django',
schedule_type=t[0],
repeats=1,
hook='django_q.tests.tasks.result'
)
schedule = create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="django",
schedule_type=t[0],
repeats=1,
hook="django_q.tests.tasks.result",
)
assert schedule is not None
assert schedule.last_run() is None
scheduler(broker=broker)
# via model
Schedule.objects.create(func='django_q.tests.tasks.word_multiply',
args='2',
kwargs='word="django"',
schedule_type=Schedule.DAILY
)
Schedule.objects.create(
func="django_q.tests.tasks.word_multiply",
args="2",
kwargs='word="django"',
schedule_type=Schedule.DAILY,
)
# scheduler
scheduler(broker=broker)
# ONCE schedule should be deleted
assert Schedule.objects.filter(pk=once_schedule.pk).exists() is False
# Catch up On
monkeypatch.setattr(Conf, 'CATCH_UP', True)
monkeypatch.setattr(Conf, "CATCH_UP", True)
now = timezone.now()
schedule = create_schedule('django_q.tests.tasks.word_multiply',
2,
word='catch_up',
schedule_type=Schedule.HOURLY,
next_run=timezone.now() - timedelta(hours=12),
repeats=-1
)
schedule = create_schedule(
"django_q.tests.tasks.word_multiply",
2,
word="catch_up",
schedule_type=Schedule.HOURLY,
next_run=timezone.now() - timedelta(hours=12),
repeats=-1,
)
scheduler(broker=broker)
schedule = Schedule.objects.get(pk=schedule.pk)
assert schedule.next_run < now
# Catch up off
monkeypatch.setattr(Conf, 'CATCH_UP', False)
monkeypatch.setattr(Conf, "CATCH_UP", False)
scheduler(broker=broker)
schedule = Schedule.objects.get(pk=schedule.pk)
assert schedule.next_run > now
# Done
broker.delete_queue()
monkeypatch.setattr(Conf, 'PREFIX', 'some_cluster_name')
monkeypatch.setattr(Conf, "PREFIX", "some_cluster_name")
# create a schedule on another cluster
schedule = create_schedule('math.copysign',
1, -1,
name='test schedule on a another cluster',
hook='django_q.tests.tasks.result',
schedule_type=Schedule.HOURLY,
cluster="some_other_cluster_name",
repeats=1)
schedule = create_schedule(
"math.copysign",
1,
-1,
name="test schedule on a another cluster",
hook="django_q.tests.tasks.result",
schedule_type=Schedule.HOURLY,
cluster="some_other_cluster_name",
repeats=1,
)
# run scheduler
scheduler(broker=broker)
# set up the workflow
@@ -226,15 +253,18 @@ def test_scheduler(broker, monkeypatch):
# queue must be empty
assert task_queue.qsize() == 0
monkeypatch.setattr(Conf, 'PREFIX', 'default')
monkeypatch.setattr(Conf, "PREFIX", "default")
# create a schedule on the same cluster
schedule = create_schedule('math.copysign',
1, -1,
name='test schedule with no cluster',
hook='django_q.tests.tasks.result',
schedule_type=Schedule.HOURLY,
cluster="default",
repeats=1)
schedule = create_schedule(
"math.copysign",
1,
-1,
name="test schedule with no cluster",
hook="django_q.tests.tasks.result",
schedule_type=Schedule.HOURLY,
cluster="default",
repeats=1,
)
# run scheduler
scheduler(broker=broker)
# set up the workflow
@@ -249,11 +279,12 @@ def test_scheduler(broker, monkeypatch):
@override_settings(
DATABASE_ROUTERS=REPLICA_DATABASE_ROUTERS,
DATABASES=REPLICA_DATABASES
DATABASE_ROUTERS=REPLICA_DATABASE_ROUTERS, DATABASES=REPLICA_DATABASES
)
@pytest.mark.django_db
def test_scheduler_atomic_transaction_must_specify_a_database_when_no_replicas_are_used(orm_no_replica_broker: Broker):
def test_scheduler_atomic_transaction_must_specify_a_database_when_no_replicas_are_used(
orm_no_replica_broker: Broker,
):
"""
GIVEN a environment without a read replica database
WHEN the scheduler is called
@@ -267,12 +298,12 @@ def test_scheduler_atomic_transaction_must_specify_a_database_when_no_replicas_a
@override_settings(
DATABASE_ROUTERS=REPLICA_DATABASE_ROUTERS,
DATABASES=REPLICA_DATABASES
DATABASE_ROUTERS=REPLICA_DATABASE_ROUTERS, DATABASES=REPLICA_DATABASES
)
@pytest.mark.django_db
def test_scheduler_atomic_transaction_must_specify_no_database_when_read_write_replicas_are_used(
orm_replica_broker: Broker):
orm_replica_broker: Broker,
):
"""
GIVEN a environment with a read/write configured replica database
WHEN the scheduler is called
@@ -285,12 +316,12 @@ def test_scheduler_atomic_transaction_must_specify_no_database_when_read_write_r
@override_settings(
DATABASE_ROUTERS=MULTIPLE_APPS_DATABASE_ROUTERS,
DATABASES=MULTIPLE_APPS_DATABASES
DATABASE_ROUTERS=MULTIPLE_APPS_DATABASE_ROUTERS, DATABASES=MULTIPLE_APPS_DATABASES
)
@pytest.mark.django_db
def test_scheduler_atomic_transaction_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
WHEN the scheduler is called
@@ -300,5 +331,14 @@ def test_scheduler_atomic_transaction_must_specify_the_database_based_on_router_
with mock.patch("django_q.cluster.db") as mocked_db:
scheduler(broker=broker)
# The router should correctly set the database to use!
assert broker.connection.db == 'default'
assert broker.connection.db == "default"
mocked_db.transaction.atomic.assert_called_with(using=broker.connection.db)
def test_localtime():
assert not is_naive(localtime())
@override_settings(USE_TZ=False)
def test_naive_localtime():
assert is_naive(localtime())
@@ -25,14 +25,14 @@ class TestingMultipleAppsDatabaseRouter:
@staticmethod
def is_admin(model):
return model._meta.app_label in ['admin']
return model._meta.app_label in ["admin"]
def db_for_read(self, model, **hints):
if self.is_admin(model):
return 'admin'
return 'default'
return "admin"
return "default"
def db_for_write(self, model, **hints):
if self.is_admin(model):
return 'admin'
return 'default'
return "admin"
return "default"
+2 -2
View File
@@ -1,6 +1,6 @@
from django.urls import re_path
from django.contrib import admin
from django.urls import re_path
urlpatterns = [
re_path(r'^admin/', admin.site.urls),
re_path(r"^admin/", admin.site.urls),
]