mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-04 00:18:11 +08:00
Hints and lints
This commit is contained in:
+27
-19
@@ -1,4 +1,5 @@
|
|||||||
import ast
|
import ast
|
||||||
|
|
||||||
# Standard
|
# Standard
|
||||||
import importlib
|
import importlib
|
||||||
import signal
|
import signal
|
||||||
@@ -10,6 +11,7 @@ from time import sleep
|
|||||||
|
|
||||||
# external
|
# external
|
||||||
import arrow
|
import arrow
|
||||||
|
|
||||||
# Django
|
# Django
|
||||||
from django import db
|
from django import db
|
||||||
from django.conf import settings
|
from django.conf import settings
|
||||||
@@ -101,10 +103,10 @@ class Cluster:
|
|||||||
@property
|
@property
|
||||||
def is_stopping(self) -> bool:
|
def is_stopping(self) -> bool:
|
||||||
return (
|
return (
|
||||||
self.stop_event
|
self.stop_event
|
||||||
and self.start_event
|
and self.start_event
|
||||||
and self.start_event.is_set()
|
and self.start_event.is_set()
|
||||||
and self.stop_event.is_set()
|
and self.stop_event.is_set()
|
||||||
)
|
)
|
||||||
|
|
||||||
@property
|
@property
|
||||||
@@ -114,13 +116,13 @@ class Cluster:
|
|||||||
|
|
||||||
class Sentinel:
|
class Sentinel:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
stop_event,
|
stop_event,
|
||||||
start_event,
|
start_event,
|
||||||
cluster_id,
|
cluster_id,
|
||||||
broker=None,
|
broker=None,
|
||||||
timeout=Conf.TIMEOUT,
|
timeout=Conf.TIMEOUT,
|
||||||
start=True,
|
start=True,
|
||||||
):
|
):
|
||||||
# Make sure we catch signals for the pool
|
# Make sure we catch signals for the pool
|
||||||
signal.signal(signal.SIGINT, signal.SIG_IGN)
|
signal.signal(signal.SIGINT, signal.SIG_IGN)
|
||||||
@@ -344,9 +346,10 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
|
|||||||
logger.info(_(f"{current_process().name} stopped pushing tasks"))
|
logger.info(_(f"{current_process().name} stopped pushing tasks"))
|
||||||
|
|
||||||
|
|
||||||
def monitor(result_queue, broker=None):
|
def monitor(result_queue: Queue, broker: Broker = None):
|
||||||
"""
|
"""
|
||||||
Gets finished tasks from the result queue and saves them to Django
|
Gets finished tasks from the result queue and saves them to Django
|
||||||
|
:type broker: brokers.Broker
|
||||||
:type result_queue: multiprocessing.Queue
|
:type result_queue: multiprocessing.Queue
|
||||||
"""
|
"""
|
||||||
if not broker:
|
if not broker:
|
||||||
@@ -373,9 +376,12 @@ def monitor(result_queue, broker=None):
|
|||||||
logger.info(_(f"{name} stopped monitoring results"))
|
logger.info(_(f"{name} stopped monitoring results"))
|
||||||
|
|
||||||
|
|
||||||
def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
def worker(
|
||||||
|
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.
|
||||||
: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
|
||||||
@@ -434,9 +440,11 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
|||||||
logger.info(_(f"{name} stopped doing work"))
|
logger.info(_(f"{name} stopped doing work"))
|
||||||
|
|
||||||
|
|
||||||
def save_task(task, broker):
|
def save_task(task, broker: Broker):
|
||||||
"""
|
"""
|
||||||
Saves the task package to Django or the cache
|
Saves the task package to Django or the cache
|
||||||
|
:param task: the task package
|
||||||
|
:type broker: brokers.Broker
|
||||||
"""
|
"""
|
||||||
# SAVE LIMIT < 0 : Don't save success
|
# SAVE LIMIT < 0 : Don't save success
|
||||||
if not task.get("save", Conf.SAVE_LIMIT >= 0) and task["success"]:
|
if not task.get("save", Conf.SAVE_LIMIT >= 0) and task["success"]:
|
||||||
@@ -482,7 +490,7 @@ def save_task(task, broker):
|
|||||||
logger.error(e)
|
logger.error(e)
|
||||||
|
|
||||||
|
|
||||||
def save_cached(task, broker):
|
def save_cached(task, broker: Broker):
|
||||||
task_key = f'{broker.list_key}:{task["id"]}'
|
task_key = f'{broker.list_key}:{task["id"]}'
|
||||||
timeout = task["cached"]
|
timeout = task["cached"]
|
||||||
if timeout is True:
|
if timeout is True:
|
||||||
@@ -534,7 +542,7 @@ def save_cached(task, broker):
|
|||||||
logger.error(e)
|
logger.error(e)
|
||||||
|
|
||||||
|
|
||||||
def scheduler(broker=None):
|
def scheduler(broker: Broker = None):
|
||||||
"""
|
"""
|
||||||
Creates a task from a schedule at the scheduled time and schedules next run
|
Creates a task from a schedule at the scheduled time and schedules next run
|
||||||
"""
|
"""
|
||||||
@@ -544,9 +552,9 @@ def scheduler(broker=None):
|
|||||||
try:
|
try:
|
||||||
with db.transaction.atomic(using=Schedule.objects.db):
|
with db.transaction.atomic(using=Schedule.objects.db):
|
||||||
for s in (
|
for s in (
|
||||||
Schedule.objects.select_for_update()
|
Schedule.objects.select_for_update()
|
||||||
.exclude(repeats=0)
|
.exclude(repeats=0)
|
||||||
.filter(next_run__lt=timezone.now())
|
.filter(next_run__lt=timezone.now())
|
||||||
):
|
):
|
||||||
args = ()
|
args = ()
|
||||||
kwargs = {}
|
kwargs = {}
|
||||||
|
|||||||
+1
-6
@@ -1,6 +1,4 @@
|
|||||||
import logging
|
import logging
|
||||||
|
|
||||||
# external
|
|
||||||
import os
|
import os
|
||||||
from copy import deepcopy
|
from copy import deepcopy
|
||||||
from multiprocessing import cpu_count
|
from multiprocessing import cpu_count
|
||||||
@@ -8,11 +6,8 @@ from signal import signal
|
|||||||
|
|
||||||
import pkg_resources
|
import pkg_resources
|
||||||
from django.conf import settings
|
from django.conf import settings
|
||||||
|
|
||||||
# django
|
|
||||||
from django.utils.translation import gettext_lazy as _
|
from django.utils.translation import gettext_lazy as _
|
||||||
|
|
||||||
# local
|
|
||||||
from django_q.queues import Queue
|
from django_q.queues import Queue
|
||||||
|
|
||||||
# optional
|
# optional
|
||||||
@@ -216,7 +211,7 @@ if Conf.ERROR_REPORTER:
|
|||||||
# and instantiate an ErrorReporter using the provided config
|
# and instantiate an ErrorReporter using the provided config
|
||||||
for name, conf in error_conf.items():
|
for name, conf in error_conf.items():
|
||||||
for entry in pkg_resources.iter_entry_points(
|
for entry in pkg_resources.iter_entry_points(
|
||||||
"djangoq.errorreporters", name
|
"djangoq.errorreporters", name
|
||||||
):
|
):
|
||||||
Reporter = entry.load()
|
Reporter = entry.load()
|
||||||
reporters.append(Reporter(**conf))
|
reporters.append(Reporter(**conf))
|
||||||
|
|||||||
+2
-2
@@ -11,7 +11,7 @@ class SignedPackage:
|
|||||||
"""Wraps Django's signing module with custom Pickle serializer."""
|
"""Wraps Django's signing module with custom Pickle serializer."""
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def dumps(obj, compressed=Conf.COMPRESSED):
|
def dumps(obj, compressed: bool = Conf.COMPRESSED) -> str:
|
||||||
return signing.dumps(
|
return signing.dumps(
|
||||||
obj,
|
obj,
|
||||||
key=Conf.SECRET_KEY,
|
key=Conf.SECRET_KEY,
|
||||||
@@ -21,7 +21,7 @@ class SignedPackage:
|
|||||||
)
|
)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def loads(obj):
|
def loads(obj) -> any:
|
||||||
return signing.loads(
|
return signing.loads(
|
||||||
obj, key=Conf.SECRET_KEY, salt=Conf.PREFIX, serializer=PickleSerializer
|
obj, key=Conf.SECRET_KEY, salt=Conf.PREFIX, serializer=PickleSerializer
|
||||||
)
|
)
|
||||||
|
|||||||
+1
-1
@@ -97,7 +97,7 @@ class Stat(Status):
|
|||||||
return Status(pid=pid, cluster_id=cluster_id)
|
return Status(pid=pid, cluster_id=cluster_id)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def get_all(broker: Broker = None)->list:
|
def get_all(broker: Broker = None) -> list:
|
||||||
"""
|
"""
|
||||||
Get the status for all currently running clusters with the same prefix
|
Get the status for all currently running clusters with the same prefix
|
||||||
and secret key.
|
and secret key.
|
||||||
|
|||||||
Reference in New Issue
Block a user