mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 16:48:11 +08:00
Merge pull request #279 from Eagllus/django_2_compat
Django 2 compatiblity
This commit is contained in:
+2
-2
@@ -10,7 +10,7 @@ import signal
|
|||||||
import socket
|
import socket
|
||||||
import ast
|
import ast
|
||||||
from time import sleep
|
from time import sleep
|
||||||
from multiprocessing import Queue, Event, Process, Value, current_process
|
from multiprocessing import Event, Process, Value, current_process
|
||||||
|
|
||||||
# external
|
# external
|
||||||
import arrow
|
import arrow
|
||||||
@@ -30,7 +30,7 @@ from django_q.models import Task, Success, Schedule
|
|||||||
from django_q.status import Stat, Status
|
from django_q.status import Stat, Status
|
||||||
from django_q.brokers import get_broker
|
from django_q.brokers import get_broker
|
||||||
from django_q.signals import pre_execute
|
from django_q.signals import pre_execute
|
||||||
|
from django_q.queues import Queue
|
||||||
|
|
||||||
|
|
||||||
class Cluster(object):
|
class Cluster(object):
|
||||||
|
|||||||
+4
-1
@@ -1,7 +1,7 @@
|
|||||||
import logging
|
import logging
|
||||||
from copy import deepcopy
|
from copy import deepcopy
|
||||||
from signal import signal
|
from signal import signal
|
||||||
from multiprocessing import cpu_count, Queue
|
from multiprocessing import cpu_count
|
||||||
|
|
||||||
# django
|
# django
|
||||||
from django.utils.translation import ugettext_lazy as _
|
from django.utils.translation import ugettext_lazy as _
|
||||||
@@ -11,6 +11,9 @@ from django.conf import settings
|
|||||||
import os
|
import os
|
||||||
import pkg_resources
|
import pkg_resources
|
||||||
|
|
||||||
|
# local
|
||||||
|
from django_q.queues import Queue
|
||||||
|
|
||||||
# optional
|
# optional
|
||||||
try:
|
try:
|
||||||
import psutil
|
import psutil
|
||||||
|
|||||||
@@ -0,0 +1,76 @@
|
|||||||
|
from __future__ import unicode_literals
|
||||||
|
|
||||||
|
import datetime
|
||||||
|
import time
|
||||||
|
import zlib
|
||||||
|
|
||||||
|
from django.utils import baseconv
|
||||||
|
from django.utils.crypto import constant_time_compare
|
||||||
|
from django.utils.encoding import force_bytes, force_str, force_text
|
||||||
|
from django.core.signing import BadSignature, SignatureExpired, b64_decode, JSONSerializer, \
|
||||||
|
Signer as Sgnr, TimestampSigner as TsS, dumps
|
||||||
|
|
||||||
|
dumps = dumps
|
||||||
|
|
||||||
|
|
||||||
|
"""
|
||||||
|
The loads function is the same as the `django.core.signing.loads` function
|
||||||
|
The difference is that `this` loads function calls `TimestampSigner` and `Signer`
|
||||||
|
"""
|
||||||
|
def loads(s, key=None, salt='django.core.signing', serializer=JSONSerializer, max_age=None):
|
||||||
|
"""
|
||||||
|
Reverse of dumps(), raise BadSignature if signature fails.
|
||||||
|
|
||||||
|
The serializer is expected to accept a bytestring.
|
||||||
|
"""
|
||||||
|
# TimestampSigner.unsign() returns str but base64 and zlib compression
|
||||||
|
# operate on bytes.
|
||||||
|
base64d = force_bytes(TimestampSigner(key, salt=salt).unsign(s, max_age=max_age))
|
||||||
|
decompress = False
|
||||||
|
if base64d[:1] == b'.':
|
||||||
|
# It's compressed; uncompress it first
|
||||||
|
base64d = base64d[1:]
|
||||||
|
decompress = True
|
||||||
|
data = b64_decode(base64d)
|
||||||
|
if decompress:
|
||||||
|
data = zlib.decompress(data)
|
||||||
|
return serializer().loads(data)
|
||||||
|
|
||||||
|
|
||||||
|
class Signer(Sgnr):
|
||||||
|
|
||||||
|
def unsign(self, signed_value):
|
||||||
|
# force_str is removed in Django 2.0
|
||||||
|
signed_value = force_str(signed_value)
|
||||||
|
if self.sep not in signed_value:
|
||||||
|
raise BadSignature('No "%s" found in value' % self.sep)
|
||||||
|
value, sig = signed_value.rsplit(self.sep, 1)
|
||||||
|
if constant_time_compare(sig, self.signature(value)):
|
||||||
|
# force_text is removed in Django 2.0
|
||||||
|
return force_text(value)
|
||||||
|
raise BadSignature('Signature "%s" does not match' % sig)
|
||||||
|
|
||||||
|
|
||||||
|
"""
|
||||||
|
TimestampSigner is also the same as `django.core.signing.TimestampSigner` but is
|
||||||
|
calling `this` Signer.
|
||||||
|
"""
|
||||||
|
class TimestampSigner(Signer, TsS):
|
||||||
|
|
||||||
|
def unsign(self, value, max_age=None):
|
||||||
|
"""
|
||||||
|
Retrieve original value and check it wasn't signed more
|
||||||
|
than max_age seconds ago.
|
||||||
|
"""
|
||||||
|
result = super().unsign(value)
|
||||||
|
value, timestamp = result.rsplit(self.sep, 1)
|
||||||
|
timestamp = baseconv.base62.decode(timestamp)
|
||||||
|
if max_age is not None:
|
||||||
|
if isinstance(max_age, datetime.timedelta):
|
||||||
|
max_age = max_age.total_seconds()
|
||||||
|
# Check timestamp is not older than max_age
|
||||||
|
age = time.time() - timestamp
|
||||||
|
if age > max_age:
|
||||||
|
raise SignatureExpired(
|
||||||
|
'Signature age %s > %s seconds' % (age, max_age))
|
||||||
|
return value
|
||||||
@@ -0,0 +1,69 @@
|
|||||||
|
"""
|
||||||
|
The code is derived from https://github.com/althonos/pronto/commit/3384010dfb4fc7c66a219f59276adef3288a886b
|
||||||
|
"""
|
||||||
|
|
||||||
|
import multiprocessing
|
||||||
|
import multiprocessing.queues
|
||||||
|
|
||||||
|
|
||||||
|
class SharedCounter(object):
|
||||||
|
""" A synchronized shared counter.
|
||||||
|
|
||||||
|
The locking done by multiprocessing.Value ensures that only a single
|
||||||
|
process or thread may read or write the in-memory ctypes object. However,
|
||||||
|
in order to do n += 1, Python performs a read followed by a write, so a
|
||||||
|
second process may read the old value before the new one is written by
|
||||||
|
the first process. The solution is to use a multiprocessing.Lock to
|
||||||
|
guarantee the atomicity of the modifications to Value.
|
||||||
|
|
||||||
|
This class comes almost entirely from Eli Bendersky's blog:
|
||||||
|
http://eli.thegreenplace.net/2012/01/04/shared-counter-with-pythons-multiprocessing/
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, n=0):
|
||||||
|
self.count = multiprocessing.Value('i', n)
|
||||||
|
|
||||||
|
def increment(self, n=1):
|
||||||
|
""" Increment the counter by n (default = 1) """
|
||||||
|
with self.count.get_lock():
|
||||||
|
self.count.value += n
|
||||||
|
|
||||||
|
@property
|
||||||
|
def value(self):
|
||||||
|
""" Return the value of the counter """
|
||||||
|
return self.count.value
|
||||||
|
|
||||||
|
|
||||||
|
class Queue(multiprocessing.queues.Queue):
|
||||||
|
""" A portable implementation of multiprocessing.Queue.
|
||||||
|
|
||||||
|
Because of multithreading / multiprocessing semantics, Queue.qsize() may
|
||||||
|
raise the NotImplementedError exception on Unix platforms like Mac OS X
|
||||||
|
where sem_getvalue() is not implemented. This subclass addresses this
|
||||||
|
problem by using a synchronized shared counter (initialized to zero) and
|
||||||
|
increasing / decreasing its value every time the put() and get() methods
|
||||||
|
are called, respectively. This not only prevents NotImplementedError from
|
||||||
|
being raised, but also allows us to implement a reliable version of both
|
||||||
|
qsize() and empty().
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, *args, **kwargs):
|
||||||
|
super(Queue, self).__init__(*args, ctx=multiprocessing.get_context(), **kwargs)
|
||||||
|
self.size = SharedCounter(0)
|
||||||
|
|
||||||
|
def put(self, *args, **kwargs):
|
||||||
|
super(Queue, self).put(*args, **kwargs)
|
||||||
|
self.size.increment(1)
|
||||||
|
|
||||||
|
def get(self, *args, **kwargs):
|
||||||
|
x = super(Queue, self).get(*args, **kwargs)
|
||||||
|
self.size.increment(-1)
|
||||||
|
return x
|
||||||
|
|
||||||
|
def qsize(self):
|
||||||
|
""" Reliable implementation of multiprocessing.Queue.qsize() """
|
||||||
|
return self.size.value
|
||||||
|
|
||||||
|
def empty(self):
|
||||||
|
""" Reliable implementation of multiprocessing.Queue.empty() """
|
||||||
|
return not self.qsize() > 0
|
||||||
+1
-1
@@ -4,7 +4,7 @@ try:
|
|||||||
except ImportError:
|
except ImportError:
|
||||||
import pickle
|
import pickle
|
||||||
|
|
||||||
from django.core import signing
|
from django_q import core_signing as signing
|
||||||
|
|
||||||
from django_q.conf import Conf
|
from django_q.conf import Conf
|
||||||
|
|
||||||
|
|||||||
+2
-1
@@ -1,5 +1,5 @@
|
|||||||
"""Provides task functionality."""
|
"""Provides task functionality."""
|
||||||
from multiprocessing import Queue, Value
|
from multiprocessing import Value
|
||||||
|
|
||||||
# django
|
# django
|
||||||
from django.db import IntegrityError
|
from django.db import IntegrityError
|
||||||
@@ -14,6 +14,7 @@ from django_q.models import Schedule, Task
|
|||||||
from django_q.humanhash import uuid
|
from django_q.humanhash import uuid
|
||||||
from django_q.brokers import get_broker
|
from django_q.brokers import get_broker
|
||||||
from django_q.signals import pre_enqueue
|
from django_q.signals import pre_enqueue
|
||||||
|
from django_q.queues import Queue
|
||||||
|
|
||||||
|
|
||||||
def async(func, *args, **kwargs):
|
def async(func, *args, **kwargs):
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import os
|
import os
|
||||||
|
import django
|
||||||
|
|
||||||
BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||||
|
|
||||||
@@ -28,6 +29,7 @@ INSTALLED_APPS = (
|
|||||||
'django_redis'
|
'django_redis'
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
MIDDLEWARE_CLASSES = (
|
MIDDLEWARE_CLASSES = (
|
||||||
'django.contrib.sessions.middleware.SessionMiddleware',
|
'django.contrib.sessions.middleware.SessionMiddleware',
|
||||||
'django.middleware.common.CommonMiddleware',
|
'django.middleware.common.CommonMiddleware',
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
from multiprocessing import Event, Queue, Value
|
from multiprocessing import Event, Value
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
@@ -8,6 +8,7 @@ from django_q.conf import Conf
|
|||||||
from django_q.tasks import async, result, fetch, count_group, result_group, fetch_group, delete_group, delete_cached, \
|
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
|
async_iter, Chain, async_chain, Iter, Async
|
||||||
from django_q.brokers import get_broker
|
from django_q.brokers import get_broker
|
||||||
|
from django_q.queues import Queue
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
import sys
|
import sys
|
||||||
import threading
|
import threading
|
||||||
from multiprocessing import Queue, Event, Value
|
from multiprocessing import Event, Value
|
||||||
from time import sleep
|
from time import sleep
|
||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
|
|
||||||
@@ -19,6 +19,7 @@ from django_q.conf import Conf
|
|||||||
from django_q.status import Stat
|
from django_q.status import Stat
|
||||||
from django_q.brokers import get_broker
|
from django_q.brokers import get_broker
|
||||||
from django_q.tests.tasks import multiply
|
from django_q.tests.tasks import multiply
|
||||||
|
from django_q.queues import Queue
|
||||||
|
|
||||||
|
|
||||||
class WordClass(object):
|
class WordClass(object):
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
from datetime import timedelta
|
from datetime import timedelta
|
||||||
from multiprocessing import Queue, Event, Value
|
from multiprocessing import Event, Value
|
||||||
|
|
||||||
import arrow
|
import arrow
|
||||||
import pytest
|
import pytest
|
||||||
@@ -10,6 +10,7 @@ from django_q.brokers import get_broker
|
|||||||
from django_q.cluster import pusher, worker, monitor, scheduler
|
from django_q.cluster import pusher, worker, monitor, scheduler
|
||||||
from django_q.conf import Conf
|
from django_q.conf import Conf
|
||||||
from django_q.tasks import Schedule, fetch, schedule as create_schedule
|
from django_q.tasks import Schedule, fetch, schedule as create_schedule
|
||||||
|
from django_q.queues import Queue
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
|
|||||||
Reference in New Issue
Block a user