mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-05 20:28:12 +08:00
Merge pull request #447 from Djailla/cleanup
[cleanup] Few cleanup commit for linting and migrations
This commit is contained in:
@@ -5,7 +5,7 @@ from django.core.cache import caches, InvalidCacheBackendError
|
|||||||
from django_q.conf import Conf
|
from django_q.conf import Conf
|
||||||
|
|
||||||
|
|
||||||
class Broker(object):
|
class Broker:
|
||||||
def __init__(self, list_key=Conf.PREFIX):
|
def __init__(self, list_key=Conf.PREFIX):
|
||||||
self.connection = self.get_connection(list_key)
|
self.connection = self.get_connection(list_key)
|
||||||
self.list_key = list_key
|
self.list_key = list_key
|
||||||
|
|||||||
@@ -57,7 +57,6 @@ class Sqs(Broker):
|
|||||||
del(config['aws_region'])
|
del(config['aws_region'])
|
||||||
return Session(**config)
|
return Session(**config)
|
||||||
|
|
||||||
|
|
||||||
def get_queue(self):
|
def get_queue(self):
|
||||||
self.sqs = self.connection.resource('sqs')
|
self.sqs = self.connection.resource('sqs')
|
||||||
return self.sqs.create_queue(QueueName=self.list_key)
|
return self.sqs.create_queue(QueueName=self.list_key)
|
||||||
|
|||||||
@@ -9,10 +9,10 @@ class IronMQBroker(Broker):
|
|||||||
return self.connection.post(task)['ids'][0]
|
return self.connection.post(task)['ids'][0]
|
||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
timeout = Conf.RETRY or None
|
timeout = Conf.RETRY or None
|
||||||
tasks = self.connection.get(timeout=timeout, wait=1, max=Conf.BULK)['messages']
|
tasks = self.connection.get(timeout=timeout, wait=1, max=Conf.BULK)['messages']
|
||||||
if tasks:
|
if tasks:
|
||||||
return [(t['id'], t['body']) for t in tasks]
|
return [(t['id'], t['body']) for t in tasks]
|
||||||
|
|
||||||
def ping(self):
|
def ping(self):
|
||||||
return self.connection.name == self.list_key
|
return self.connection.name == self.list_key
|
||||||
|
|||||||
@@ -60,7 +60,7 @@ class ORM(Broker):
|
|||||||
|
|
||||||
def dequeue(self):
|
def dequeue(self):
|
||||||
tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[
|
tasks = self.get_connection().filter(key=self.list_key, lock__lt=_timeout())[
|
||||||
0 : Conf.BULK
|
0: Conf.BULK
|
||||||
]
|
]
|
||||||
if tasks:
|
if tasks:
|
||||||
task_list = []
|
task_list = []
|
||||||
|
|||||||
+2
-8
@@ -1,9 +1,3 @@
|
|||||||
# Future
|
|
||||||
from __future__ import absolute_import
|
|
||||||
from __future__ import division
|
|
||||||
from __future__ import print_function
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
import ast
|
import ast
|
||||||
|
|
||||||
# Standard
|
# Standard
|
||||||
@@ -36,7 +30,7 @@ from django_q.signing import SignedPackage, BadSignature
|
|||||||
from django_q.status import Stat, Status
|
from django_q.status import Stat, Status
|
||||||
|
|
||||||
|
|
||||||
class Cluster(object):
|
class Cluster:
|
||||||
def __init__(self, broker=None):
|
def __init__(self, broker=None):
|
||||||
self.broker = broker or get_broker()
|
self.broker = broker or get_broker()
|
||||||
self.sentinel = None
|
self.sentinel = None
|
||||||
@@ -120,7 +114,7 @@ class Cluster(object):
|
|||||||
return self.start_event is None and self.stop_event is None and self.sentinel
|
return self.start_event is None and self.stop_event is None and self.sentinel
|
||||||
|
|
||||||
|
|
||||||
class Sentinel(object):
|
class Sentinel:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
stop_event,
|
stop_event,
|
||||||
|
|||||||
+2
-2
@@ -22,7 +22,7 @@ except ImportError:
|
|||||||
psutil = None
|
psutil = None
|
||||||
|
|
||||||
|
|
||||||
class Conf(object):
|
class Conf:
|
||||||
"""
|
"""
|
||||||
Configuration class
|
Configuration class
|
||||||
"""
|
"""
|
||||||
@@ -195,7 +195,7 @@ if not logger.handlers:
|
|||||||
|
|
||||||
|
|
||||||
# Error Reporting Interface
|
# Error Reporting Interface
|
||||||
class ErrorReporter(object):
|
class ErrorReporter:
|
||||||
|
|
||||||
# initialize with iterator of reporters (better name, targets?)
|
# initialize with iterator of reporters (better name, targets?)
|
||||||
def __init__(self, reporters):
|
def __init__(self, reporters):
|
||||||
|
|||||||
@@ -50,7 +50,7 @@ DEFAULT_WORDLIST = (
|
|||||||
'zulu')
|
'zulu')
|
||||||
|
|
||||||
|
|
||||||
class HumanHasher(object):
|
class HumanHasher:
|
||||||
|
|
||||||
"""
|
"""
|
||||||
Transforms hex digests to human-readable strings.
|
Transforms hex digests to human-readable strings.
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
import picklefield.fields
|
import picklefield.fields
|
||||||
import django.utils.timezone
|
import django.utils.timezone
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import models, migrations
|
from django.db import models, migrations
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import migrations, models
|
from django.db import migrations, models
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,3 @@
|
|||||||
# -*- coding: utf-8 -*-
|
|
||||||
from __future__ import unicode_literals
|
|
||||||
|
|
||||||
from django.db import migrations, models
|
from django.db import migrations, models
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -7,7 +7,7 @@ import multiprocessing
|
|||||||
import multiprocessing.queues
|
import multiprocessing.queues
|
||||||
|
|
||||||
|
|
||||||
class SharedCounter(object):
|
class SharedCounter:
|
||||||
""" A synchronized shared counter.
|
""" A synchronized shared counter.
|
||||||
|
|
||||||
The locking done by multiprocessing.Value ensures that only a single
|
The locking done by multiprocessing.Value ensures that only a single
|
||||||
|
|||||||
+2
-2
@@ -11,7 +11,7 @@ from django_q.conf import Conf
|
|||||||
BadSignature = signing.BadSignature
|
BadSignature = signing.BadSignature
|
||||||
|
|
||||||
|
|
||||||
class SignedPackage(object):
|
class SignedPackage:
|
||||||
|
|
||||||
"""Wraps Django's signing module with custom Pickle serializer."""
|
"""Wraps Django's signing module with custom Pickle serializer."""
|
||||||
|
|
||||||
@@ -31,7 +31,7 @@ class SignedPackage(object):
|
|||||||
serializer=PickleSerializer)
|
serializer=PickleSerializer)
|
||||||
|
|
||||||
|
|
||||||
class PickleSerializer(object):
|
class PickleSerializer:
|
||||||
|
|
||||||
"""Simple wrapper around Pickle for signing.dumps and signing.loads."""
|
"""Simple wrapper around Pickle for signing.dumps and signing.loads."""
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -5,7 +5,7 @@ from django_q.conf import Conf, logger
|
|||||||
from django_q.signing import SignedPackage, BadSignature
|
from django_q.signing import SignedPackage, BadSignature
|
||||||
|
|
||||||
|
|
||||||
class Status(object):
|
class Status:
|
||||||
"""Cluster status base class."""
|
"""Cluster status base class."""
|
||||||
|
|
||||||
def __init__(self, pid, cluster_id):
|
def __init__(self, pid, cluster_id):
|
||||||
|
|||||||
+4
-5
@@ -7,9 +7,8 @@ from django.db import IntegrityError
|
|||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
from multiprocessing import Value
|
from multiprocessing import Value
|
||||||
|
|
||||||
from django_q.brokers import get_broker
|
|
||||||
|
|
||||||
# local
|
# local
|
||||||
|
from django_q.brokers import get_broker
|
||||||
from django_q.conf import Conf, logger
|
from django_q.conf import Conf, logger
|
||||||
from django_q.humanhash import uuid
|
from django_q.humanhash import uuid
|
||||||
from django_q.models import Schedule, Task
|
from django_q.models import Schedule, Task
|
||||||
@@ -480,7 +479,7 @@ def async_chain(chain, group=None, cached=Conf.CACHED, sync=Conf.SYNC, broker=No
|
|||||||
return group
|
return group
|
||||||
|
|
||||||
|
|
||||||
class Iter(object):
|
class Iter:
|
||||||
"""
|
"""
|
||||||
An async task with iterable arguments
|
An async task with iterable arguments
|
||||||
"""
|
"""
|
||||||
@@ -550,7 +549,7 @@ class Iter(object):
|
|||||||
return len(self.args)
|
return len(self.args)
|
||||||
|
|
||||||
|
|
||||||
class Chain(object):
|
class Chain:
|
||||||
"""
|
"""
|
||||||
A sequential chain of tasks
|
A sequential chain of tasks
|
||||||
"""
|
"""
|
||||||
@@ -634,7 +633,7 @@ class Chain(object):
|
|||||||
return len(self.chain)
|
return len(self.chain)
|
||||||
|
|
||||||
|
|
||||||
class AsyncTask(object):
|
class AsyncTask:
|
||||||
"""
|
"""
|
||||||
an async task
|
an async task
|
||||||
"""
|
"""
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ from django_q.tests.tasks import multiply, TaskError
|
|||||||
from django_q.queues import Queue
|
from django_q.queues import Queue
|
||||||
|
|
||||||
|
|
||||||
class WordClass(object):
|
class WordClass:
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.word_list = DEFAULT_WORDLIST
|
self.word_list = DEFAULT_WORDLIST
|
||||||
|
|
||||||
@@ -45,6 +45,7 @@ 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
|
assert result(task) == 1506
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
def test_sync_raise_exception(broker):
|
def test_sync_raise_exception(broker):
|
||||||
with pytest.raises(TaskError):
|
with pytest.raises(TaskError):
|
||||||
@@ -401,6 +402,7 @@ def test_update_failed(broker):
|
|||||||
assert saved_task.success is True
|
assert saved_task.success is True
|
||||||
assert saved_task.result == 'result'
|
assert saved_task.result == 'result'
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
def test_acknowledge_failure_override():
|
def test_acknowledge_failure_override():
|
||||||
class VerifyAckMockBroker(Broker):
|
class VerifyAckMockBroker(Broker):
|
||||||
@@ -434,10 +436,12 @@ def test_acknowledge_failure_override():
|
|||||||
|
|
||||||
tag = uuid()
|
tag = uuid()
|
||||||
task_success_ack = task_fail_ack.copy()
|
task_success_ack = task_fail_ack.copy()
|
||||||
task_success_ack.update({'id': tag[1],
|
task_success_ack.update({
|
||||||
'name': tag[0],
|
'id': tag[1],
|
||||||
'ack_id': 'test_success_ack_id',
|
'name': tag[0],
|
||||||
'success': True,})
|
'ack_id': 'test_success_ack_id',
|
||||||
|
'success': True,
|
||||||
|
})
|
||||||
del task_success_ack['ack_failure']
|
del task_success_ack['ack_failure']
|
||||||
|
|
||||||
result_queue = Queue()
|
result_queue = Queue()
|
||||||
@@ -453,6 +457,7 @@ def test_acknowledge_failure_override():
|
|||||||
assert broker.acknowledgements.get('test_fail_no_ack_id') is None
|
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_success_ack_id') == 1
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.django_db
|
@pytest.mark.django_db
|
||||||
def assert_result(task):
|
def assert_result(task):
|
||||||
assert task is not None
|
assert task is not None
|
||||||
|
|||||||
+1
-1
@@ -3384,7 +3384,7 @@ import sys
|
|||||||
import base64
|
import base64
|
||||||
import zlib
|
import zlib
|
||||||
|
|
||||||
class DictImporter(object):
|
class DictImporter:
|
||||||
def __init__(self, sources):
|
def __init__(self, sources):
|
||||||
self.sources = sources
|
self.sources = sources
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user