mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-05 19:58:12 +08:00
minor linting
This commit is contained in:
+6
-7
@@ -86,9 +86,8 @@ def monitor(run_once=False):
|
|||||||
|
|
||||||
|
|
||||||
class Status(object):
|
class Status(object):
|
||||||
"""
|
|
||||||
Cluster status base class
|
"""Cluster status base class."""
|
||||||
"""
|
|
||||||
|
|
||||||
def __init__(self, pid):
|
def __init__(self, pid):
|
||||||
self.workers = []
|
self.workers = []
|
||||||
@@ -106,9 +105,8 @@ class Status(object):
|
|||||||
|
|
||||||
|
|
||||||
class Stat(Status):
|
class Stat(Status):
|
||||||
"""
|
|
||||||
Status object for Cluster monitoring
|
"""Status object for Cluster monitoring."""
|
||||||
"""
|
|
||||||
|
|
||||||
def __init__(self, sentinel):
|
def __init__(self, sentinel):
|
||||||
super(Stat, self).__init__(sentinel.parent_pid or sentinel.pid)
|
super(Stat, self).__init__(sentinel.parent_pid or sentinel.pid)
|
||||||
@@ -171,7 +169,8 @@ class Stat(Status):
|
|||||||
@staticmethod
|
@staticmethod
|
||||||
def get_all(r=redis_client):
|
def get_all(r=redis_client):
|
||||||
"""
|
"""
|
||||||
Gets status for all currently running clusters with the same prefix and secret key
|
Get the status for all currently running clusters with the same prefix
|
||||||
|
and secret key.
|
||||||
:return: list of type Stat
|
:return: list of type Stat
|
||||||
"""
|
"""
|
||||||
stats = []
|
stats = []
|
||||||
|
|||||||
+5
-7
@@ -1,3 +1,4 @@
|
|||||||
|
"""Package signing."""
|
||||||
try:
|
try:
|
||||||
import cPickle as pickle
|
import cPickle as pickle
|
||||||
except ImportError:
|
except ImportError:
|
||||||
@@ -11,9 +12,8 @@ BadSignature = signing.BadSignature
|
|||||||
|
|
||||||
|
|
||||||
class SignedPackage(object):
|
class SignedPackage(object):
|
||||||
"""
|
|
||||||
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=Conf.COMPRESSED):
|
||||||
@@ -32,10 +32,8 @@ class SignedPackage(object):
|
|||||||
|
|
||||||
|
|
||||||
class PickleSerializer(object):
|
class PickleSerializer(object):
|
||||||
"""
|
|
||||||
Simple wrapper around Pickle for signing.dumps and
|
"""Simple wrapper around Pickle for signing.dumps and signing.loads."""
|
||||||
signing.loads.
|
|
||||||
"""
|
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def dumps(obj):
|
def dumps(obj):
|
||||||
|
|||||||
+29
-22
@@ -1,3 +1,4 @@
|
|||||||
|
"""Provides task functionalities."""
|
||||||
from multiprocessing import Queue, Value
|
from multiprocessing import Queue, Value
|
||||||
|
|
||||||
# django
|
# django
|
||||||
@@ -12,9 +13,7 @@ from django_q.humanhash import uuid
|
|||||||
|
|
||||||
|
|
||||||
def async(func, *args, **kwargs):
|
def async(func, *args, **kwargs):
|
||||||
"""
|
"""Send a task to the cluster."""
|
||||||
Sends a task to the cluster
|
|
||||||
"""
|
|
||||||
# get options from q_options dict or direct from kwargs
|
# get options from q_options dict or direct from kwargs
|
||||||
options = kwargs.pop('q_options', kwargs)
|
options = kwargs.pop('q_options', kwargs)
|
||||||
hook = options.pop('hook', None)
|
hook = options.pop('hook', None)
|
||||||
@@ -26,7 +25,10 @@ def async(func, *args, **kwargs):
|
|||||||
# get an id
|
# get an id
|
||||||
tag = uuid()
|
tag = uuid()
|
||||||
# build the task package
|
# build the task package
|
||||||
task = {'id': tag[1], 'name': tag[0], 'func': func, 'args': args, 'kwargs': kwargs,
|
task = {'id': tag[1], 'name': tag[0],
|
||||||
|
'func': func,
|
||||||
|
'args': args,
|
||||||
|
'kwargs': kwargs,
|
||||||
'started': timezone.now()}
|
'started': timezone.now()}
|
||||||
# add optionals
|
# add optionals
|
||||||
if hook:
|
if hook:
|
||||||
@@ -47,19 +49,20 @@ def async(func, *args, **kwargs):
|
|||||||
|
|
||||||
def schedule(func, *args, **kwargs):
|
def schedule(func, *args, **kwargs):
|
||||||
"""
|
"""
|
||||||
:param func: function to schedule
|
Create a schedule.
|
||||||
:param args: function arguments
|
|
||||||
:param name: optional name for the schedule
|
:param func: function to schedule.
|
||||||
:param hook: optional result hook function
|
:param args: function arguments.
|
||||||
|
:param name: optional name for the schedule.
|
||||||
|
:param hook: optional result hook function.
|
||||||
:type schedule_type: Schedule.TYPE
|
:type schedule_type: Schedule.TYPE
|
||||||
:param repeats: how many times to repeat. 0=never, -1=always
|
:param repeats: how many times to repeat. 0=never, -1=always.
|
||||||
:param next_run: Next scheduled run
|
:param next_run: Next scheduled run.
|
||||||
:type next_run: datetime.datetime
|
:type next_run: datetime.datetime
|
||||||
:param kwargs: function keyword arguments
|
:param kwargs: function keyword arguments.
|
||||||
:return: the schedule object
|
:return: the schedule object.
|
||||||
:rtype: Schedule
|
:rtype: Schedule
|
||||||
"""
|
"""
|
||||||
|
|
||||||
name = kwargs.pop('name', None)
|
name = kwargs.pop('name', None)
|
||||||
hook = kwargs.pop('hook', None)
|
hook = kwargs.pop('hook', None)
|
||||||
schedule_type = kwargs.pop('schedule_type', Schedule.ONCE)
|
schedule_type = kwargs.pop('schedule_type', Schedule.ONCE)
|
||||||
@@ -79,7 +82,8 @@ def schedule(func, *args, **kwargs):
|
|||||||
|
|
||||||
def result(task_id):
|
def result(task_id):
|
||||||
"""
|
"""
|
||||||
Returns the result of the named task
|
Return the result of the named task.
|
||||||
|
|
||||||
:type task_id: str or uuid
|
:type task_id: str or uuid
|
||||||
:param task_id: the task name or uuid
|
:param task_id: the task name or uuid
|
||||||
:return: the result object of this task
|
:return: the result object of this task
|
||||||
@@ -90,7 +94,8 @@ def result(task_id):
|
|||||||
|
|
||||||
def result_group(group_id, failures=False):
|
def result_group(group_id, failures=False):
|
||||||
"""
|
"""
|
||||||
returns a list of results for a task group
|
Return a list of results for a task group.
|
||||||
|
|
||||||
:param str group_id: the group id
|
:param str group_id: the group id
|
||||||
:param bool failures: set to True to include failures
|
:param bool failures: set to True to include failures
|
||||||
:return: list or results
|
:return: list or results
|
||||||
@@ -100,7 +105,8 @@ def result_group(group_id, failures=False):
|
|||||||
|
|
||||||
def fetch(task_id):
|
def fetch(task_id):
|
||||||
"""
|
"""
|
||||||
Returns the processed task
|
Return the processed task.
|
||||||
|
|
||||||
:param task_id: the task name or uuid
|
:param task_id: the task name or uuid
|
||||||
:type task_id: str or uuid
|
:type task_id: str or uuid
|
||||||
:return: the full task object
|
:return: the full task object
|
||||||
@@ -111,17 +117,19 @@ def fetch(task_id):
|
|||||||
|
|
||||||
def fetch_group(group_id, failures=True):
|
def fetch_group(group_id, failures=True):
|
||||||
"""
|
"""
|
||||||
Returns a list of Tasks for a task group
|
Return a list of Tasks for a task group.
|
||||||
|
|
||||||
:param str group_id: the group id
|
:param str group_id: the group id
|
||||||
:param bool failures: set to False to exclude failures
|
:param bool failures: set to False to exclude failures
|
||||||
:return: list of Tasks
|
:return: list of Tasks
|
||||||
"""
|
"""
|
||||||
|
|
||||||
return Task.get_task_group(group_id, failures)
|
return Task.get_task_group(group_id, failures)
|
||||||
|
|
||||||
|
|
||||||
def count_group(group_id, failures=False):
|
def count_group(group_id, failures=False):
|
||||||
"""
|
"""
|
||||||
|
Count the results in a group.
|
||||||
|
|
||||||
:param str group_id: the group id
|
:param str group_id: the group id
|
||||||
:param bool failures: Returns failure count if True
|
:param bool failures: Returns failure count if True
|
||||||
:return: the number of tasks/results in a group
|
:return: the number of tasks/results in a group
|
||||||
@@ -132,6 +140,8 @@ def count_group(group_id, failures=False):
|
|||||||
|
|
||||||
def delete_group(group_id, tasks=False):
|
def delete_group(group_id, tasks=False):
|
||||||
"""
|
"""
|
||||||
|
Delete a group.
|
||||||
|
|
||||||
:param str group_id: the group id
|
:param str group_id: the group id
|
||||||
:param bool tasks: If set to True this will also delete the group tasks.
|
:param bool tasks: If set to True this will also delete the group tasks.
|
||||||
Otherwise just the group label is removed.
|
Otherwise just the group label is removed.
|
||||||
@@ -141,10 +151,7 @@ def delete_group(group_id, tasks=False):
|
|||||||
|
|
||||||
|
|
||||||
def _sync(task_id, pack):
|
def _sync(task_id, pack):
|
||||||
"""
|
"""Simulate a package travelling through the cluster."""
|
||||||
Simulates a package travelling through the cluster.
|
|
||||||
|
|
||||||
"""
|
|
||||||
task_queue = Queue()
|
task_queue = Queue()
|
||||||
result_queue = Queue()
|
result_queue = Queue()
|
||||||
task_queue.put(pack)
|
task_queue.put(pack)
|
||||||
|
|||||||
Reference in New Issue
Block a user