diff --git a/.travis.yml b/.travis.yml index 9c438ba..7d6cb02 100644 --- a/.travis.yml +++ b/.travis.yml @@ -5,20 +5,20 @@ services: - mongodb python: - - "2.7" - "3.6" + - "3.7" env: - - DJANGO=2.0 + - DJANGO=2.1 - DJANGO=1.11.11 - - DJANGO=1.8.19 matrix: exclude: - - python: "2.7" - env: DJANGO=2.0 + - python: "3.7" + env: DJANGO=1.11.11 -sudo: false +sudo: true +dist: xenial addons: apt: diff --git a/CHANGELOG.md b/CHANGELOG.md index 1ca6da1..b114e0a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -600,4 +600,4 @@ ## [v0.1.0](https://github.com/koed00/django-q/tree/v0.1.0) (2015-06-28) -\* *This Change Log was automatically generated by [github_changelog_generator](https://github.com/skywinder/Github-Changelog-Generator)* \ No newline at end of file +\* *This Change Log was automatically generated by [github_changelog_generator](https://github.com/skywinder/Github-Changelog-Generator)* diff --git a/README.rst b/README.rst index 74bf8ee..6414f5d 100644 --- a/README.rst +++ b/README.rst @@ -26,12 +26,12 @@ Features Requirements ~~~~~~~~~~~~ -- `Django `__ > = 1.8 +- `Django `__ > = 1.11.11 - `Django-picklefield `__ - `Arrow `__ - `Blessed `__ -Tested with: Python 2.7 & 3.6. Django 1.8.19, 1.11.11 and 2.0.x +Tested with: Python 3.6. 3.7 Django 1.11.11 and 2.0.x Brokers ~~~~~~~ @@ -110,19 +110,19 @@ Check overall statistics with:: Creating Tasks ~~~~~~~~~~~~~~ -Use `async` from your code to quickly offload tasks: +Use `async_task` from your code to quickly offload tasks: .. code:: python - from django_q.tasks import async, result + from django_q.tasks import async_task, result # create the task - async('math.copysign', 2, -2) + async_task('math.copysign', 2, -2) # or with a reference import math.copysign - task_id = async(copysign, 2, -2) + task_id = async_task(copysign, 2, -2) # get the result task_result = result(task_id) @@ -133,7 +133,7 @@ Use `async` from your code to quickly offload tasks: # but in most cases you will want to use a hook: - async('math.modf', 2.5, hook='hooks.print_result') + async_task('math.modf', 2.5, hook='hooks.print_result') # hooks.py def print_result(task): diff --git a/django_q/__init__.py b/django_q/__init__.py index 5225644..a8dd924 100644 --- a/django_q/__init__.py +++ b/django_q/__init__.py @@ -1,21 +1,6 @@ -# import os -# import sys -import django - -# myPath = os.path.dirname(os.path.abspath(__file__)) -# sys.path.insert(0, myPath) - -VERSION = (0, 9, 4) +VERSION = (1, 0, 0) 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 .models import Task, Schedule, Success, Failure - from .cluster import Cluster - from .status import Stat - from .brokers import get_broker __all__ = ['conf', 'cluster', 'models', 'tasks'] diff --git a/django_q/admin.py b/django_q/admin.py index f8511ad..5b9a99a 100644 --- a/django_q/admin.py +++ b/django_q/admin.py @@ -2,9 +2,9 @@ from django.contrib import admin from django.utils.translation import ugettext_lazy as _ -from django_q.tasks import async -from django_q.models import Success, Failure, Schedule, OrmQ from django_q.conf import Conf +from django_q.models import Success, Failure, Schedule, OrmQ +from django_q.tasks import async_task class TaskAdmin(admin.ModelAdmin): @@ -19,7 +19,7 @@ class TaskAdmin(admin.ModelAdmin): 'group' ) - def has_add_permission(self, request, obj=None): + def has_add_permission(self, request): """Don't allow adds.""" return False @@ -34,14 +34,13 @@ class TaskAdmin(admin.ModelAdmin): def get_readonly_fields(self, request, obj=None): """Set all fields readonly.""" - return list(self.readonly_fields) + \ - [field.name for field in obj._meta.fields] + return list(self.readonly_fields) + [field.name for field in obj._meta.fields] 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 {}) + async_task(task.func, *task.args or (), hook=task.hook, **task.kwargs or {}) task.delete() @@ -56,10 +55,10 @@ class FailAdmin(admin.ModelAdmin): 'func', 'started', 'stopped', - 'result' + 'short_result' ) - def has_add_permission(self, request, obj=None): + def has_add_permission(self, request): """Don't allow adds.""" return False @@ -70,8 +69,7 @@ class FailAdmin(admin.ModelAdmin): def get_readonly_fields(self, request, obj=None): """Set all fields readonly.""" - return list(self.readonly_fields) + \ - [field.name for field in obj._meta.fields] + return list(self.readonly_fields) + [field.name for field in obj._meta.fields] class ScheduleAdmin(admin.ModelAdmin): @@ -113,10 +111,11 @@ class QueueAdmin(admin.ModelAdmin): def get_queryset(self, request): return super(QueueAdmin, self).get_queryset(request).using(Conf.ORM) - def has_add_permission(self, request, obj=None): + def has_add_permission(self, request): """Don't allow adds.""" return False + admin.site.register(Schedule, ScheduleAdmin) admin.site.register(Success, TaskAdmin) admin.site.register(Failure, FailAdmin) diff --git a/django_q/cluster.py b/django_q/cluster.py index 3959e29..dc1b46c 100644 --- a/django_q/cluster.py +++ b/django_q/cluster.py @@ -4,32 +4,31 @@ from __future__ import division from __future__ import print_function from __future__ import unicode_literals +from time import sleep + +# external +import arrow +import ast # Standard import importlib import signal import socket -import ast -from time import sleep -from multiprocessing import Event, Process, Value, current_process - -# external -import arrow - +import traceback # Django +from django import db from django.utils import timezone from django.utils.translation import ugettext_lazy as _ -from django import db +from multiprocessing import Event, Process, Value, current_process # Local from django_q import tasks -from django_q.compat import range -from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter, rollbar +from django_q.brokers import get_broker +from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter from django_q.models import Task, Success, Schedule +from django_q.queues import Queue +from django_q.signals import pre_execute from django_q.signing import SignedPackage, BadSignature from django_q.status import Stat, Status -from django_q.brokers import get_broker -from django_q.signals import pre_execute -from django_q.queues import Queue class Cluster(object): @@ -287,7 +286,7 @@ def pusher(task_queue, event, broker=None): try: task_set = broker.dequeue() except Exception as e: - logger.error(e) + logger.error(e, traceback.format_exc()) # broker probably crashed. Let the sentinel handle it. sleep(10) break @@ -298,7 +297,7 @@ def pusher(task_queue, event, broker=None): try: task = SignedPackage.loads(task[1]) except (TypeError, BadSignature) as e: - logger.error(e) + logger.error(e, traceback.format_exc()) broker.fail(ack_id) continue task['ack_id'] = ack_id @@ -366,8 +365,6 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT): result = (e, False) if error_reporter: error_reporter.report() - if rollbar: - rollbar.report_exc_info() # We're still going if not result: db.close_old_connections() @@ -380,11 +377,9 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT): res = f(*task['args'], **task['kwargs']) result = (res, True) except Exception as e: - result = ('{}'.format(e), False) + result = ('{} : {}'.format(e, traceback.format_exc()), False) if error_reporter: error_reporter.report() - if rollbar: - rollbar.report_exc_info() # Process result task['result'] = result[0] task['success'] = result[1] @@ -405,7 +400,7 @@ 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) # SAVE LIMIT > 0: Prune database, SAVE_LIMIT 0: No pruning @@ -473,7 +468,7 @@ 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 + # async_task next in a chain if task.get('chain', None): tasks.async_chain(task['chain'], group=group, cached=task['cached'], sync=task['sync'], broker=broker) # save the task @@ -536,15 +531,15 @@ 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.async_task(s.func, *args, **kwargs) # log it if not s.task: logger.error( - _('{} failed to create a task from schedule [{}]').format(current_process().name, - s.name or s.id)) + _('{} failed to create a task from schedule [{}]').format(current_process().name, + s.name or s.id)) else: logger.info( - _('{} created a task from schedule [{}]').format(current_process().name, s.name or s.id)) + _('{} created a task from schedule [{}]').format(current_process().name, s.name or s.id)) # default behavior is to delete a ONCE schedule if s.schedule_type == s.ONCE: if s.repeats < 0: diff --git a/django_q/compat.py b/django_q/compat.py deleted file mode 100644 index 0ba80c4..0000000 --- a/django_q/compat.py +++ /dev/null @@ -1,14 +0,0 @@ -from __future__ import absolute_import -""" -Compatibility layer. - -Intentionally replaces use of python-future -""" - -# https://github.com/Koed00/django-q/issues/4 - -try: - range = xrange -except NameError: - range = range - diff --git a/django_q/conf.py b/django_q/conf.py index 2fd2e6e..6db6126 100644 --- a/django_q/conf.py +++ b/django_q/conf.py @@ -150,9 +150,6 @@ class Conf(object): # The redis stats key Q_STAT = 'django_q:{}:cluster'.format(PREFIX) - # Optional rollbar key - ROLLBAR = conf.get('rollbar', {}) - # Optional error reporting setup ERROR_REPORTER = conf.get('error_reporter', {}) @@ -190,19 +187,6 @@ if not logger.handlers: logger.addHandler(handler) -# rollbar -if Conf.ROLLBAR: - rollbar_conf = deepcopy(Conf.ROLLBAR) - try: - import rollbar - rollbar.init(rollbar_conf.pop('access_token'), environment=rollbar_conf.pop('environment'), **rollbar_conf) - except ImportError: - rollbar = None - -else: - rollbar = None - - # Error Reporting Interface class ErrorReporter(object): diff --git a/django_q/models.py b/django_q/models.py index b43f81c..4666f36 100644 --- a/django_q/models.py +++ b/django_q/models.py @@ -1,8 +1,7 @@ from django import get_version -try: - from django.urls import reverse -except ImportError: # Django < 1.10 - from django.core.urlresolvers import reverse +from django.template.defaultfilters import truncatechars + +from django.urls import reverse from django.utils.html import format_html from django.utils.translation import ugettext_lazy as _ from django.db import models @@ -82,6 +81,10 @@ class Task(models.Model): def time_taken(self): return (self.stopped - self.started).total_seconds() + @property + def short_result(self): + return truncatechars(self.result, 100) + def __unicode__(self): return u'{}'.format(self.name or self.id) diff --git a/django_q/tasks.py b/django_q/tasks.py index 04bbb2f..4769c3d 100644 --- a/django_q/tasks.py +++ b/django_q/tasks.py @@ -1,26 +1,28 @@ """Provides task functionality.""" # Standard from time import sleep, time -from multiprocessing import Value # django from django.db import IntegrityError from django.utils import timezone +from multiprocessing import Value -# local -from django_q.signing import SignedPackage -from django_q.conf import Conf, logger -from django_q.models import Schedule, Task -from django_q.humanhash import uuid from django_q.brokers import get_broker -from django_q.signals import pre_enqueue +# local +from django_q.cluster import worker, monitor +from django_q.conf import Conf, logger +from django_q.humanhash import uuid +from django_q.models import Schedule, Task from django_q.queues import Queue +from django_q.signals import pre_enqueue +from django_q.signing import SignedPackage -def async(func, *args, **kwargs): +def async_task(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') + opt_keys = ( + 'hook', 'group', 'save', 'sync', 'cached', 'ack_failure', 'iter_count', 'iter_cached', 'chain', 'broker') q_options = keywords.pop('q_options', {}) # get an id tag = uuid() @@ -392,7 +394,7 @@ def queue_size(broker=None): def async_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] @@ -409,15 +411,15 @@ def async_iter(func, args_iter, **kwargs): broker = options['broker'] broker.cache.set('{}:{}:args'.format(broker.list_key, iter_group), SignedPackage.dumps(args_iter)) for args in args_iter: - if type(args) is not tuple: + if not isinstance(args, tuple): args = (args,) - async(func, *args, **options) + async_task(func, *args, **options) return iter_group def async_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 +438,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) + async_task(task[0], *args, **kwargs) return group @@ -518,7 +520,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 async_task() """ self.chain.append((func, args, kwargs)) # remove existing results @@ -573,7 +575,7 @@ class Chain(object): return len(self.chain) -class Async(object): +class AsyncTask(object): """ an async task """ @@ -647,7 +649,7 @@ class Async(object): return self.kwargs.get(key, default) def run(self): - self.id = async(self.func, *self.args, **self.kwargs) + self.id = async_task(self.func, *self.args, **self.kwargs) self.started = True return self.id @@ -673,10 +675,6 @@ class Async(object): def _sync(pack): - # Python 2.6 is unable to handle this import on top of the file - # because it creates a circular dependency between tasks and cluster - from django_q.cluster import worker, monitor - """Simulate a package travelling through the cluster.""" task_queue = Queue() result_queue = Queue() @@ -686,4 +684,8 @@ def _sync(pack): worker(task_queue, result_queue, Value('f', -1)) result_queue.put('STOP') monitor(result_queue) + task_queue.close() + task_queue.join_thread() + result_queue.close() + result_queue.join_thread() return task['id'] diff --git a/django_q/tests/test_admin.py b/django_q/tests/test_admin.py index 7c2acef..52e8d41 100644 --- a/django_q/tests/test_admin.py +++ b/django_q/tests/test_admin.py @@ -1,7 +1,4 @@ -try: - from django.urls import reverse -except ImportError: # Django < 1.10 - from django.core.urlresolvers import reverse +from django.urls import reverse from django.utils import timezone import pytest diff --git a/django_q/tests/test_brokers.py b/django_q/tests/test_brokers.py index 6ef7cd4..9d40b38 100644 --- a/django_q/tests/test_brokers.py +++ b/django_q/tests/test_brokers.py @@ -5,7 +5,6 @@ import pytest import redis from django_q.brokers import get_broker, Broker -from django_q.compat import range from django_q.conf import Conf from django_q.humanhash import uuid @@ -63,7 +62,7 @@ def test_disque(monkeypatch): assert broker.info() is not None # clear before we start broker.delete_queue() - # enqueue + # async_task broker.enqueue('test') assert broker.queue_size() == 1 # dequeue @@ -127,7 +126,7 @@ def test_ironmq(monkeypatch): # clear before we start broker.purge_queue() assert broker.queue_size() == 0 - # enqueue + # async_task broker.enqueue('test') # dequeue task = broker.dequeue()[0] @@ -136,7 +135,7 @@ def test_ironmq(monkeypatch): assert broker.dequeue() is None # Retry test # monkeypatch.setattr(Conf, 'RETRY', 1) - # broker.enqueue('test') + # broker.async_task('test') # assert broker.dequeue() is not None # sleep(3) # assert broker.dequeue() is not None @@ -180,7 +179,7 @@ def canceled_sqs(monkeypatch): assert broker.ping() is True assert broker.info() is not None assert broker.queue_size() == 0 - # enqueue + # async_task broker.enqueue('test') # dequeue task = broker.dequeue()[0] @@ -240,7 +239,7 @@ def test_orm(monkeypatch): assert broker.info() is not None # clear before we start broker.delete_queue() - # enqueue + # async_task broker.enqueue('test') assert broker.queue_size() == 1 # dequeue @@ -297,7 +296,7 @@ def test_mongo(monkeypatch): assert broker.info() is not None # clear before we start broker.delete_queue() - # enqueue + # async_task broker.enqueue('test') assert broker.queue_size() == 1 # dequeue diff --git a/django_q/tests/test_cached.py b/django_q/tests/test_cached.py index e32c60c..8c618c9 100644 --- a/django_q/tests/test_cached.py +++ b/django_q/tests/test_cached.py @@ -3,10 +3,9 @@ from multiprocessing import Event, Value 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 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.queues import Queue @@ -23,13 +22,13 @@ 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) + 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 @@ -145,10 +144,10 @@ def test_chain(broker): @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 +161,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 +184,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 diff --git a/django_q/tests/test_cluster.py b/django_q/tests/test_cluster.py index 6c4b8c3..3355052 100644 --- a/django_q/tests/test_cluster.py +++ b/django_q/tests/test_cluster.py @@ -11,9 +11,8 @@ myPath = os.path.dirname(os.path.abspath(__file__)) 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, async_task, 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 +41,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 = async_task('django_q.tests.tasks.count_letters', DEFAULT_WORDLIST, broker=broker, sync=True) assert result(task) == 1506 @@ -82,7 +81,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 = 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 @@ -109,32 +108,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 = 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('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('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('django_q.tests.tasks.countdown', 100000, broker=broker) + e = async_task('django_q.tests.tasks.countdown', 100000, broker=broker) # function as instance - f = async(multiply, 753, 2, hook=assert_result, broker=broker) + f = async_task(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 = async_task('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 = async_task('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 = 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('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}) # check if everything has a task id assert isinstance(a, str) assert isinstance(b, str) @@ -249,7 +248,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) + async_task('django_q.tests.tasks.count_forever', broker=broker) start_event = Event() stop_event = Event() # Set a timer to stop the Sentinel @@ -265,7 +264,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) + async_task('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 +280,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) + 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() # override settings @@ -295,8 +294,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) + 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 @@ -318,7 +317,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) + async_task('math.copysign', 1, -1, broker=broker) stop_event = Event() stop_event.set() start_event = Event() diff --git a/django_q/tests/test_monitor.py b/django_q/tests/test_monitor.py index 0d43b7f..5719973 100644 --- a/django_q/tests/test_monitor.py +++ b/django_q/tests/test_monitor.py @@ -1,9 +1,8 @@ import pytest -from django_q.tasks import async +from django_q.tasks import async_task from django_q.brokers import get_broker from django_q.cluster import Cluster -from django_q.compat import range from django_q.monitor import monitor, info from django_q.status import Stat from django_q.conf import Conf @@ -46,4 +45,4 @@ def test_info(): def do_sync(): - async('django_q.tests.tasks.countdown', 1, sync=True, save=True) + async_task('django_q.tests.tasks.countdown', 1, sync=True, save=True) diff --git a/docs/brokers.rst b/docs/brokers.rst index ed06b42..109a653 100644 --- a/docs/brokers.rst +++ b/docs/brokers.rst @@ -140,7 +140,7 @@ You can override this class if you want to contribute and support your own broke .. py:class:: Broker - .. py:method:: enqueue(task) + .. py:method:: async_task(task) Sends a task package to the broker queue and returns a tracking id if available. diff --git a/docs/chain.rst b/docs/chain.rst index 11298bc..b372dca 100644 --- a/docs/chain.rst +++ b/docs/chain.rst @@ -6,13 +6,13 @@ Sometimes you want to run tasks sequentially. For that you can use the :func:`as .. code-block:: python - # Async a chain of tasks + # async a chain of tasks from django_q.tasks import async_chain, result_group # the chain must be in the format # [(func,(args),{kwargs}),(func,(args),{kwargs}),..] group_id = async_chain([('math.copysign', (1, -1)), - ('math.floor', (1,))]) + ('math.floor', (1,))]) # get group result result_group(group_id, count=2) @@ -63,7 +63,7 @@ Reference .. py:method:: append(func, *args, **kwargs) - Append a task to the chain. Takes the same arguments as :func:`async` + Append a task to the chain. Takes the same arguments as :func:`async_task` :return: the current number of tasks in the chain :rtype: int @@ -102,4 +102,4 @@ Reference get the length of the chain - :return int: length of the chain \ No newline at end of file + :return int: length of the chain diff --git a/docs/conf.py b/docs/conf.py index fc4de62..5b37c95 100644 --- a/docs/conf.py +++ b/docs/conf.py @@ -71,9 +71,9 @@ author = 'Ilan Steemers' # built documents. # # The short X.Y version. -version = '0.9' +version = '1.0' # The full version, including alpha/beta/rc tags. -release = '0.9.4' +release = '1.0.0' # The language for content autogenerated by Sphinx. Refer to documentation # for a list of supported languages. diff --git a/docs/configure.rst b/docs/configure.rst index 638de05..824ef83 100644 --- a/docs/configure.rst +++ b/docs/configure.rst @@ -64,7 +64,7 @@ Set this to something that makes sense for your project. Can be overridden for i ack_failures ~~~~~~~~~~~~ -When set to ``True``, also acknowledge unsuccessful tasks. This causes failed tasks to be considered as successful deliveries, thereby removing them from the task queue. Can also be set per-task by passing the ``ack_failure`` option to :func:`async`. Defaults to ``False``. +When set to ``True``, also acknowledge unsuccessful tasks. This causes failed tasks to be considered as successful deliveries, thereby removing them from the task queue. Can also be set per-task by passing the ``ack_failure`` option to :func:`async_task`. Defaults to ``False``. .. _retry: @@ -101,7 +101,7 @@ Guard loop sleep in seconds, must be greater than 0 and less than 60. sync ~~~~ -When set to ``True`` this configuration option forces all :func:`async` calls to be run with ``sync=True``. +When set to ``True`` this configuration option forces all :func:`async_task` calls to be run with ``sync=True``. Effectively making everything synchronous. Useful for testing. Defaults to ``False``. .. _queue_limit: @@ -380,25 +380,6 @@ To enable installed error reporters, you must provide the configuration settings For more information on error reporters and developing error reporting plugins for Django Q, see :doc:`errors`. -rollbar -~~~~~~~ -You can redirect worker exceptions directly to your `Rollbar `__ dashboard by installing the python notifier with ``pip install rollbar`` and adding this configuration dictionary to your config:: - - # rollbar config - Q_CLUSTER = { - 'rollbar': { - 'access_token': '32we33a92a5224jiww8982', - 'environment': 'Django-Q' - } - } - -Please check the Pyrollbar `configuration reference `__ for more options. -Note that you will need a `Rollbar `__ account and access token to use this feature. - - -.. note:: - The ``rollbar`` setting is included for backwards compatibility, for those who utilized rollbar configuration before the ``error_reporter`` interface was introduced. Note that Rollbar support can be configured either via the ``rollbar`` setting, or via the ``django-q-rollbar`` package and enabled via the ``error_reporter`` setting above. - cpu_affinity ~~~~~~~~~~~~ diff --git a/docs/examples.rst b/docs/examples.rst index bc6158f..7d7b5a1 100644 --- a/docs/examples.rst +++ b/docs/examples.rst @@ -12,18 +12,18 @@ Sending an email can take a while so why not queue it: # Welcome mail with follow up example from datetime import timedelta from django.utils import timezone - from django_q.tasks import async, schedule + from django_q.tasks import async_task, schedule from django_q.models import Schedule def welcome_mail(user): msg = 'Welcome to our website' # send this message right away - async('django.core.mail.send_mail', - 'Welcome', - msg, - 'from@example.com', - [user.email]) + async_task('django.core.mail.send_mail', + 'Welcome', + msg, + 'from@example.com', + [user.email]) # and this follow up email in one hour msg = 'Here are some tips to get you started...' schedule('django.core.mail.send_mail', @@ -51,7 +51,7 @@ A good place to use async tasks are Django's model signals. You don't want to de from django.contrib.auth.models import User from django.db.models.signals import pre_save from django.dispatch import receiver - from django_q.tasks import async + from django_q.tasks import async_task # set up the pre_save signal for our user @receiver(pre_save, sender=User) @@ -64,7 +64,7 @@ A good place to use async tasks are Django's model signals. You don't want to de # has his email changed? if not user.email == instance.email: # tell everyone - async('tasks.inform_everyone', instance) + async_task('tasks.inform_everyone', instance) The task will send a message to everyone else informing them that the users email address has changed. Note that this adds almost no overhead to the save action: @@ -87,8 +87,8 @@ The task will send a message to everyone else informing them that the users emai for u in User.objects.exclude(pk=user.pk): msg = 'Dear {}, {} has a new email address: {}' msg = msg.format(u.username, user.username, user.email) - async('django.core.mail.send_mail', - 'New email', msg, 'from@example.com', [u.email]) + async_task('django.core.mail.send_mail', + 'New email', msg, 'from@example.com', [u.email]) Of course you can do other things beside sending emails. These are just generic examples. You can use signals with async to update fields in other objects too. @@ -104,19 +104,19 @@ In this example the user requests a report and we let the cluster do the generat .. code-block:: python # Report generation with hook example - from django_q.tasks import async + from django_q.tasks import async_task # views.py # user requests a report. def create_report(request): - async('tasks.create_html_report', - request.user, - hook='tasks.email_report') + async_task('tasks.create_html_report', + request.user, + hook='tasks.email_report') .. code-block:: python # tasks.py - from django_q.tasks import async + from django_q.tasks import async_task # report generator def create_html_report(user): @@ -127,16 +127,16 @@ In this example the user requests a report and we let the cluster do the generat def email_report(task): if task.success: # Email the report - async('django.core.mail.send_mail', - 'The report you requested', - task.result, - 'from@example.com', - task.args[0].email) + async_task('django.core.mail.send_mail', + 'The report you requested', + task.result, + 'from@example.com', + task.args[0].email) else: # Tell the admins something went wrong - async('django.core.mail.mail_admins', - 'Report generation failed', - task.result) + async_task('django.core.mail.mail_admins', + 'Report generation failed', + task.result) The hook is practical here, because it allows us to detach the sending task from the report generation function and to report on possible failures. @@ -152,12 +152,12 @@ here's an example of how you can have Django Q take care of your indexes in real from .models import Document from django.db.models.signals import post_save from django.dispatch import receiver - from django_q.tasks import async + from django_q.tasks import async_task # hook up the post save handler @receiver(post_save, sender=Document) def document_changed(sender, instance, **kwargs): - async('tasks.index_object', sender, instance, save=False) + async_task('tasks.index_object', sender, instance, save=False) # turn off result saving to not flood your database .. code-block:: python @@ -177,7 +177,7 @@ here's an example of how you can have Django Q take care of your indexes in real index.update_object(instance, using=backend) Now every time a Document is saved, your indexes will be updated without causing a delay in your save action. -You could expand this to dealing with deletes, by adding a ``post_delete`` signal and calling ``index.remove_object`` in the async function. +You could expand this to dealing with deletes, by adding a ``post_delete`` signal and calling ``index.remove_object`` in the async_task function. .. _shell: @@ -187,13 +187,13 @@ You can execute or schedule shell commands using Pythons :mod:`subprocess` modul .. code-block:: python - from django_q.tasks import async, result + from django_q.tasks import async_task, result # make a backup copy of setup.py - async('subprocess.call', ['cp', 'setup.py', 'setup.py.bak']) + async_task('subprocess.call', ['cp', 'setup.py', 'setup.py.bak']) # call ls -l and dump the output - task_id=async('subprocess.check_output', ['ls', '-l']) + task_id=async_task('subprocess.check_output', ['ls', '-l']) # get the result dir_list = result(task_id) @@ -202,10 +202,10 @@ In Python 3.5 the subprocess module has changed quite a bit and returns a :class .. code-block:: python - from django_q.tasks import async, result + from django_q.tasks import async_task, result # make a backup copy of setup.py - tid = async('subprocess.run', ['cp', 'setup.py', 'setup.py.bak']) + tid = async_task('subprocess.run', ['cp', 'setup.py', 'setup.py.bak']) # get the result r=result(tid, 500) @@ -220,22 +220,22 @@ In Python 3.5 the subprocess module has changed quite a bit and returns a :class from subprocess import PIPE # call ls -l and pipe the output - tid = async('subprocess.run', ['ls', '-l'], stdout=PIPE) + tid = async_task('subprocess.run', ['ls', '-l'], stdout=PIPE) # get the result res = result(tid, 500) # print the output print(res.stdout) -Instead of :func:`async` you can of course also use :func:`schedule` to schedule commands. +Instead of :func:`async_task` you can of course also use :func:`schedule` to schedule commands. For regular Django management commands, it is easier to call them directly: .. code-block:: python - from django_q.tasks import async, schedule + from django_q.tasks import async_task, schedule - async('django.core.management.call_command','clearsessions') + async_task('django.core.management.call_command','clearsessions') # or clear those sessions every hour @@ -255,7 +255,7 @@ Adapted from `Sebastian Raschka's blog `__ Django Q aims to use as much of Django's standard offerings as possible - The code is tested against Django versions `1.8.19 LTS`, `1.11.11` and `2.0.x`. + The code is tested against Django versions `1.11.11 LTS` and `2.0.x`. + Please note that Django versions below 2.0 do not support Python 3.7 - `Django-picklefield `__ @@ -122,14 +123,14 @@ Other known issues are: Python ~~~~~~ -The code is always tested against the latest version of Python 2 and Python 3 and we try to stay compatible with the last two versions of each. -Current tests are performed with Python 2.7.14 and 3.6.3 +The code is always tested against the latest version Python 3 and we try to stay compatible with the last two versions of each. +Current tests are performed with 3.6 and 3.7 If you do encounter any regressions with earlier versions, please submit an issue on `github `__ .. note:: - Django 1.7.10 or earlier is not compatible with Python 3.5 Django releases before 1.11 are not officially supported on Python 3.6 + Django releases before 2.0 are not supported on Python 3.7 Open-source packages ~~~~~~~~~~~~~~~~~~~~ @@ -139,9 +140,10 @@ You can reference the `requirements =1.8', 'django-picklefield', 'blessed', 'arrow'], + install_requires=['django>=1.11', 'django-picklefield', 'blessed', 'arrow'], test_requires=['pytest', 'pytest-django', ], cmdclass={'test': PyTest}, classifiers=[ - 'Development Status :: 4 - Beta', + 'Development Status :: 5 - Production/Stable', 'Environment :: Web Environment', 'Framework :: Django', 'Intended Audience :: Developers', @@ -48,13 +48,13 @@ setup( 'Operating System :: POSIX', 'Operating System :: MacOS', 'Programming Language :: Python', - 'Programming Language :: Python :: 2', - 'Programming Language :: Python :: 2.7', 'Programming Language :: Python :: 3', 'Programming Language :: Python :: 3.4', 'Programming Language :: Python :: 3.5', 'Programming Language :: Python :: 3.6', + 'Programming Language :: Python :: 3.7', 'Topic :: Internet :: WWW/HTTP', + 'Topic :: System :: Distributed Computing', 'Topic :: Software Development :: Libraries :: Python Modules', ], entry_points={