mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-04 05:38:11 +08:00
Replaces async occurrences with alternatives
* async is now a reserved word in python3.7 * Rename async function to enqueue * Rename all async_ functions to enqueue_ * Rename Async class to AsyncTask * Updates the docs.
This commit is contained in:
@@ -12,7 +12,7 @@ default_app_config = 'django_q.apps.DjangoQConfig'
|
||||
# root imports will slowly be deprecated.
|
||||
# please import from the relevant sub modules
|
||||
if django.VERSION[:2] < (1, 9):
|
||||
from .tasks import async, schedule, result, result_group, fetch, fetch_group, count_group, delete_group, queue_size
|
||||
from .tasks import enqueue, schedule, result, result_group, fetch, fetch_group, count_group, delete_group, queue_size
|
||||
from .models import Task, Schedule, Success, Failure
|
||||
from .cluster import Cluster
|
||||
from .status import Stat
|
||||
|
||||
+2
-2
@@ -2,7 +2,7 @@
|
||||
from django.contrib import admin
|
||||
from django.utils.translation import ugettext_lazy as _
|
||||
|
||||
from django_q.tasks import async
|
||||
from django_q.tasks import enqueue
|
||||
from django_q.models import Success, Failure, Schedule, OrmQ
|
||||
from django_q.conf import Conf
|
||||
|
||||
@@ -41,7 +41,7 @@ class TaskAdmin(admin.ModelAdmin):
|
||||
def retry_failed(FailAdmin, request, queryset):
|
||||
"""Submit selected tasks back to the queue."""
|
||||
for task in queryset:
|
||||
async(task.func, *task.args or (), hook=task.hook, **task.kwargs or {})
|
||||
enqueue(task.func, *task.args or (), hook=task.hook, **task.kwargs or {})
|
||||
task.delete()
|
||||
|
||||
|
||||
|
||||
+5
-5
@@ -405,9 +405,9 @@ def save_task(task, broker):
|
||||
# SAVE LIMIT < 0 : Don't save success
|
||||
if not task.get('save', Conf.SAVE_LIMIT >= 0) and task['success']:
|
||||
return
|
||||
# async next in a chain
|
||||
# enqueues next in a chain
|
||||
if task.get('chain', None):
|
||||
tasks.async_chain(task['chain'], group=task['group'], cached=task['cached'], sync=task['sync'], broker=broker)
|
||||
tasks.enqueue_chain(task['chain'], group=task['group'], cached=task['cached'], sync=task['sync'], broker=broker)
|
||||
# SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning
|
||||
db.close_old_connections()
|
||||
try:
|
||||
@@ -473,9 +473,9 @@ def save_cached(task, broker):
|
||||
# save the group list
|
||||
group_list.append(task_key)
|
||||
broker.cache.set(group_key, group_list, timeout)
|
||||
# async next in a chain
|
||||
# enqueue next in a chain
|
||||
if task.get('chain', None):
|
||||
tasks.async_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker)
|
||||
tasks.enqueue_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker)
|
||||
# save the task
|
||||
broker.cache.set(task_key,
|
||||
SignedPackage.dumps(task),
|
||||
@@ -536,7 +536,7 @@ def scheduler(broker=None):
|
||||
q_options['broker'] = broker
|
||||
q_options['group'] = q_options.get('group', s.name or s.id)
|
||||
kwargs['q_options'] = q_options
|
||||
s.task = tasks.async(s.func, *args, **kwargs)
|
||||
s.task = tasks.enqueue(s.func, *args, **kwargs)
|
||||
# log it
|
||||
if not s.task:
|
||||
logger.error(
|
||||
|
||||
+13
-13
@@ -17,7 +17,7 @@ from django_q.signals import pre_enqueue
|
||||
from django_q.queues import Queue
|
||||
|
||||
|
||||
def async(func, *args, **kwargs):
|
||||
def enqueue(func, *args, **kwargs):
|
||||
"""Queue a task for the cluster."""
|
||||
keywords = kwargs.copy()
|
||||
opt_keys = ('hook', 'group', 'save', 'sync', 'cached', 'ack_failure', 'iter_count', 'iter_cached', 'chain', 'broker')
|
||||
@@ -390,9 +390,9 @@ def queue_size(broker=None):
|
||||
return broker.queue_size()
|
||||
|
||||
|
||||
def async_iter(func, args_iter, **kwargs):
|
||||
def enqueue_iter(func, args_iter, **kwargs):
|
||||
"""
|
||||
async a function with iterable arguments
|
||||
enqueues a function with iterable arguments
|
||||
"""
|
||||
iter_count = len(args_iter)
|
||||
iter_group = uuid()[1]
|
||||
@@ -411,13 +411,13 @@ def async_iter(func, args_iter, **kwargs):
|
||||
for args in args_iter:
|
||||
if type(args) is not tuple:
|
||||
args = (args,)
|
||||
async(func, *args, **options)
|
||||
enqueue(func, *args, **options)
|
||||
return iter_group
|
||||
|
||||
|
||||
def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None):
|
||||
def enqueue_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=None):
|
||||
"""
|
||||
async a chain of tasks
|
||||
enqueues a chain of tasks
|
||||
the chain must be in the format [(func,(args),{kwargs}),(func,(args),{kwargs})]
|
||||
"""
|
||||
if not group:
|
||||
@@ -436,7 +436,7 @@ def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=No
|
||||
kwargs['cached'] = cached
|
||||
kwargs['sync'] = sync
|
||||
kwargs['broker'] = broker or get_broker()
|
||||
async(task[0], *args, **kwargs)
|
||||
enqueue(task[0], *args, **kwargs)
|
||||
return group
|
||||
|
||||
|
||||
@@ -472,7 +472,7 @@ class Iter(object):
|
||||
self.kwargs['cached'] = self.cached
|
||||
self.kwargs['sync'] = self.sync
|
||||
self.kwargs['broker'] = self.broker
|
||||
self.id = async_iter(self.func, self.args, **self.kwargs)
|
||||
self.id = enqueue_iter(self.func, self.args, **self.kwargs)
|
||||
self.started = True
|
||||
return self.id
|
||||
|
||||
@@ -518,7 +518,7 @@ class Chain(object):
|
||||
def append(self, func, *args, **kwargs):
|
||||
"""
|
||||
add a task to the chain
|
||||
takes the same parameters as async()
|
||||
takes the same parameters as enqueue()
|
||||
"""
|
||||
self.chain.append((func, args, kwargs))
|
||||
# remove existing results
|
||||
@@ -532,8 +532,8 @@ class Chain(object):
|
||||
Start queueing the chain to the worker cluster
|
||||
:return: the chain's group id
|
||||
"""
|
||||
self.group = async_chain(chain=self.chain[:], group=self.group, cached=self.cached, sync=self.sync,
|
||||
broker=self.broker)
|
||||
self.group = enqueue_chain(chain=self.chain[:], group=self.group, cached=self.cached, sync=self.sync,
|
||||
broker=self.broker)
|
||||
self.started = True
|
||||
return self.group
|
||||
|
||||
@@ -573,7 +573,7 @@ class Chain(object):
|
||||
return len(self.chain)
|
||||
|
||||
|
||||
class Async(object):
|
||||
class AsyncTask(object):
|
||||
"""
|
||||
an async task
|
||||
"""
|
||||
@@ -647,7 +647,7 @@ class Async(object):
|
||||
return self.kwargs.get(key, default)
|
||||
|
||||
def run(self):
|
||||
self.id = async(self.func, *self.args, **self.kwargs)
|
||||
self.id = enqueue(self.func, *self.args, **self.kwargs)
|
||||
self.started = True
|
||||
return self.id
|
||||
|
||||
|
||||
@@ -5,8 +5,8 @@ import pytest
|
||||
from django_q.cluster import pusher, worker, monitor
|
||||
from django_q.compat import range
|
||||
from django_q.conf import Conf
|
||||
from django_q.tasks import async, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \
|
||||
async_iter, Chain, async_chain, Iter, Async
|
||||
from django_q.tasks import enqueue, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \
|
||||
enqueue_iter, Chain, enqueue_chain, Iter, AsyncTask
|
||||
from django_q.brokers import get_broker
|
||||
from django_q.queues import Queue
|
||||
|
||||
@@ -23,14 +23,14 @@ def test_cached(broker):
|
||||
broker.cache.clear()
|
||||
group = 'cache_test'
|
||||
# queue the tests
|
||||
task_id = async('math.copysign', 1, -1, cached=True, broker=broker)
|
||||
async('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
async('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
async('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
async('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
async('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
async('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 = enqueue('math.copysign', 1, -1, cached=True, broker=broker)
|
||||
enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
enqueue('math.copysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
enqueue('math.popysign', 1, -1, cached=True, broker=broker, group=group)
|
||||
iter_id = enqueue_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
|
||||
@@ -86,10 +86,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 = enqueue_iter('math.floor', it, sync=True)
|
||||
t2 = enqueue_iter('math.copysign', it2, sync=True)
|
||||
t3 = enqueue_iter('math.floor', it3, sync=True)
|
||||
t4 = enqueue_iter('math.floor', (1,), sync=True)
|
||||
result_t = result(t)
|
||||
assert result_t is not None
|
||||
task_t = fetch(t)
|
||||
@@ -140,15 +140,15 @@ 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)
|
||||
rid = enqueue_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_async_class(broker, monkeypatch):
|
||||
def test_asynctask_class(broker, monkeypatch):
|
||||
broker.purge_queue()
|
||||
broker.cache.clear()
|
||||
a = Async('math.copysign')
|
||||
a = AsyncTask('math.copysign')
|
||||
assert a.func == 'math.copysign'
|
||||
a.args = (1, -1)
|
||||
assert a.started is False
|
||||
@@ -162,11 +162,11 @@ def test_async_class(broker, monkeypatch):
|
||||
assert a.result() == -1
|
||||
assert a.fetch().result == -1
|
||||
# again with kwargs
|
||||
a = Async('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 = Async('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 a.sync is False
|
||||
a.sync = True
|
||||
assert a.kwargs['q_options']['sync'] is True
|
||||
@@ -185,6 +185,6 @@ def test_async_class(broker, monkeypatch):
|
||||
# global overrides
|
||||
monkeypatch.setattr(Conf, 'SYNC', True)
|
||||
monkeypatch.setattr(Conf, 'CACHED', True)
|
||||
a = Async('math.floor', 1.5)
|
||||
a = AsyncTask('math.floor', 1.5)
|
||||
a.run()
|
||||
assert a.result() == 1
|
||||
|
||||
@@ -13,7 +13,7 @@ sys.path.insert(0, myPath + '/../')
|
||||
from django_q.cluster import Cluster, Sentinel, pusher, worker, monitor, save_task
|
||||
from django_q.compat import range
|
||||
from django_q.humanhash import DEFAULT_WORDLIST, uuid
|
||||
from django_q.tasks import fetch, fetch_group, async, result, result_group, count_group, delete_group, queue_size
|
||||
from django_q.tasks import fetch, fetch_group, enqueue, result, result_group, count_group, delete_group, queue_size
|
||||
from django_q.models import Task, Success
|
||||
from django_q.conf import Conf
|
||||
from django_q.status import Stat
|
||||
@@ -42,7 +42,7 @@ def test_redis_connection(broker):
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_sync(broker):
|
||||
task = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker, sync=True)
|
||||
task = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker, sync=True)
|
||||
assert result(task) == 1506
|
||||
|
||||
|
||||
@@ -82,7 +82,7 @@ def test_sentinel():
|
||||
def test_cluster(broker):
|
||||
broker.list_key = 'cluster_test:q'
|
||||
broker.delete_queue()
|
||||
task = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker)
|
||||
task = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker)
|
||||
assert broker.queue_size() == 1
|
||||
task_queue = Queue()
|
||||
assert task_queue.qsize() == 0
|
||||
@@ -109,32 +109,32 @@ def test_cluster(broker):
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
def test_async(broker, admin_user):
|
||||
def test_enqueue(broker, admin_user):
|
||||
broker.list_key = 'cluster_test:q'
|
||||
broker.delete_queue()
|
||||
a = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result',
|
||||
broker=broker)
|
||||
b = async('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_cluster.assert_result',
|
||||
broker=broker)
|
||||
a = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, hook='django_q.tests.test_cluster.assert_result',
|
||||
broker=broker)
|
||||
b = enqueue('django_q.tests.tasks.count_letters2', WordClass(), hook='django_q.tests.test_cluster.assert_result',
|
||||
broker=broker)
|
||||
# unknown argument
|
||||
c = async('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany',
|
||||
hook='django_q.tests.test_cluster.assert_bad_result', broker=broker)
|
||||
c = enqueue('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, 'oneargumentoomany',
|
||||
hook='django_q.tests.test_cluster.assert_bad_result', broker=broker)
|
||||
# unknown function
|
||||
d = async('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_cluster.assert_bad_result',
|
||||
broker=broker)
|
||||
d = enqueue('django_q.tests.tasks.does_not_exist', WordClass(), hook='django_q.tests.test_cluster.assert_bad_result',
|
||||
broker=broker)
|
||||
# function without result
|
||||
e = async('django_q.tests.tasks.countdown', 100000, broker=broker)
|
||||
e = enqueue('django_q.tests.tasks.countdown', 100000, broker=broker)
|
||||
# function as instance
|
||||
f = async(multiply, 753, 2, hook=assert_result, broker=broker)
|
||||
f = enqueue(multiply, 753, 2, hook=assert_result, broker=broker)
|
||||
# model as argument
|
||||
g = async('django_q.tests.tasks.get_task_name', Task(name='John'), broker=broker)
|
||||
g = enqueue('django_q.tests.tasks.get_task_name', Task(name='John'), broker=broker)
|
||||
# args,kwargs, group and broken hook
|
||||
h = async('django_q.tests.tasks.word_multiply', 2, word='django', hook='fail.me', broker=broker)
|
||||
h = enqueue('django_q.tests.tasks.word_multiply', 2, word='django', hook='fail.me', broker=broker)
|
||||
# args unpickle test
|
||||
j = async('django_q.tests.tasks.get_user_id', admin_user, broker=broker, group='test_j')
|
||||
j = enqueue('django_q.tests.tasks.get_user_id', admin_user, broker=broker, group='test_j')
|
||||
# q_options and save opt_out test
|
||||
k = async('django_q.tests.tasks.get_user_id', admin_user,
|
||||
q_options={'broker': broker, 'group': 'test_k', 'save': False, 'timeout': 90})
|
||||
k = enqueue('django_q.tests.tasks.get_user_id', admin_user,
|
||||
q_options={'broker': broker, 'group': 'test_k', 'save': False, 'timeout': 90})
|
||||
# check if everything has a task id
|
||||
assert isinstance(a, str)
|
||||
assert isinstance(b, str)
|
||||
@@ -249,7 +249,7 @@ def test_timeout(broker):
|
||||
# set up the Sentinel
|
||||
broker.list_key = 'timeout_test:q'
|
||||
broker.purge_queue()
|
||||
async('django_q.tests.tasks.count_forever', broker=broker)
|
||||
enqueue('django_q.tests.tasks.count_forever', broker=broker)
|
||||
start_event = Event()
|
||||
stop_event = Event()
|
||||
# Set a timer to stop the Sentinel
|
||||
@@ -265,7 +265,7 @@ def test_timeout(broker):
|
||||
def test_timeout_override(broker):
|
||||
# set up the Sentinel
|
||||
broker.list_key = 'timeout_override_test:q'
|
||||
async('django_q.tests.tasks.count_forever', broker=broker, timeout=1)
|
||||
enqueue('django_q.tests.tasks.count_forever', broker=broker, timeout=1)
|
||||
start_event = Event()
|
||||
stop_event = Event()
|
||||
# Set a timer to stop the Sentinel
|
||||
@@ -281,9 +281,9 @@ def test_timeout_override(broker):
|
||||
def test_recycle(broker, monkeypatch):
|
||||
# set up the Sentinel
|
||||
broker.list_key = 'test_recycle_test:q'
|
||||
async('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
async('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
async('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
start_event = Event()
|
||||
stop_event = Event()
|
||||
# override settings
|
||||
@@ -295,8 +295,8 @@ def test_recycle(broker, monkeypatch):
|
||||
assert start_event.is_set()
|
||||
assert s.status() == Conf.STOPPED
|
||||
assert s.reincarnations == 1
|
||||
async('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
async('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
enqueue('django_q.tests.tasks.multiply', 2, 2, broker=broker)
|
||||
task_queue = Queue()
|
||||
result_queue = Queue()
|
||||
# push two tasks
|
||||
@@ -318,7 +318,7 @@ def test_recycle(broker, monkeypatch):
|
||||
@pytest.mark.django_db
|
||||
def test_bad_secret(broker, monkeypatch):
|
||||
broker.list_key = 'test_bad_secret:q'
|
||||
async('math.copysign', 1, -1, broker=broker)
|
||||
enqueue('math.copysign', 1, -1, broker=broker)
|
||||
stop_event = Event()
|
||||
stop_event.set()
|
||||
start_event = Event()
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import pytest
|
||||
|
||||
from django_q.tasks import async
|
||||
from django_q.tasks import enqueue
|
||||
from django_q.brokers import get_broker
|
||||
from django_q.cluster import Cluster
|
||||
from django_q.compat import range
|
||||
@@ -46,4 +46,4 @@ def test_info():
|
||||
|
||||
|
||||
def do_sync():
|
||||
async('django_q.tests.tasks.countdown', 1, sync=True, save=True)
|
||||
enqueue('django_q.tests.tasks.countdown', 1, sync=True, save=True)
|
||||
|
||||
Reference in New Issue
Block a user