mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-06 11:28:11 +08:00
Added full package signing and compression through Django signing module and JsonPickle
This commit is contained in:
+22
-11
@@ -1,17 +1,17 @@
|
|||||||
import importlib
|
import importlib
|
||||||
import logging
|
import logging
|
||||||
import pickle
|
|
||||||
import signal
|
import signal
|
||||||
from multiprocessing import cpu_count, Queue, Event, Process, current_process
|
from multiprocessing import cpu_count, Queue, Event, Process, current_process
|
||||||
import sys
|
import sys
|
||||||
from time import sleep
|
from time import sleep
|
||||||
|
|
||||||
|
import jsonpickle
|
||||||
import coloredlogs
|
import coloredlogs
|
||||||
|
from django.core.signing import BadSignature
|
||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
import redis
|
import redis
|
||||||
|
|
||||||
from django.core.signing import Signer, BadSignature
|
from django.core import signing
|
||||||
|
|
||||||
from django_q.apps import LOG_LEVEL, SECRET_KEY, SAVE_LIMIT
|
from django_q.apps import LOG_LEVEL, SECRET_KEY, SAVE_LIMIT
|
||||||
from django_q.humanhash import uuid
|
from django_q.humanhash import uuid
|
||||||
@@ -19,7 +19,6 @@ from django_q.models import Task, Success
|
|||||||
|
|
||||||
prefix = 'django_q'
|
prefix = 'django_q'
|
||||||
q_list = '{}:q'.format(prefix)
|
q_list = '{}:q'.format(prefix)
|
||||||
signer = Signer(SECRET_KEY)
|
|
||||||
logger = logging.getLogger('django-q')
|
logger = logging.getLogger('django-q')
|
||||||
coloredlogs.install(level=getattr(logging, LOG_LEVEL))
|
coloredlogs.install(level=getattr(logging, LOG_LEVEL))
|
||||||
|
|
||||||
@@ -120,13 +119,10 @@ class Cluster(object):
|
|||||||
for pack in iter(task_queue.get, 'STOP'):
|
for pack in iter(task_queue.get, 'STOP'):
|
||||||
# unpickle the task
|
# unpickle the task
|
||||||
try:
|
try:
|
||||||
task = pickle.loads(pack)
|
task = signing.loads(pack, key=SECRET_KEY, salt='django_q.q', serializer=JSONPickleSerializer)
|
||||||
except TypeError as e:
|
except TypeError as e:
|
||||||
logger.error(e)
|
logger.error(e)
|
||||||
continue
|
continue
|
||||||
# check signature
|
|
||||||
try:
|
|
||||||
task[0] = signer.unsign(task[0])
|
|
||||||
except BadSignature as e:
|
except BadSignature as e:
|
||||||
task[0] = task[0].rsplit(":", 1)[0]
|
task[0] = task[0].rsplit(":", 1)[0]
|
||||||
task.append(timezone.now())
|
task.append(timezone.now())
|
||||||
@@ -197,9 +193,24 @@ class Cluster(object):
|
|||||||
|
|
||||||
|
|
||||||
def async(func, *args, **kwargs):
|
def async(func, *args, **kwargs):
|
||||||
# [name, func, args, kwargs, started, finished, result, success]
|
"""
|
||||||
name = signer.sign(uuid()[0])
|
Schedules a module function
|
||||||
pack = pickle.dumps([name, func, args, kwargs, timezone.now()])
|
[name, func, args, kwargs, started, finished, result, success]
|
||||||
|
"""
|
||||||
|
name = uuid()[0]
|
||||||
|
pack = signing.dumps([name, func, args, kwargs, timezone.now()], key=SECRET_KEY, salt='django_q.q', compress=True,
|
||||||
|
serializer=JSONPickleSerializer)
|
||||||
r.rpush(q_list, pack)
|
r.rpush(q_list, pack)
|
||||||
logger.debug('Pushed {}'.format(pack))
|
logger.debug('Pushed {}'.format(pack))
|
||||||
return name
|
return name
|
||||||
|
|
||||||
|
class JSONPickleSerializer(object):
|
||||||
|
"""
|
||||||
|
Simple wrapper around jsonpickle to be used in signing.dumps and
|
||||||
|
signing.loads.
|
||||||
|
"""
|
||||||
|
def dumps(self, obj):
|
||||||
|
return jsonpickle.dumps(obj).encode('latin-1')
|
||||||
|
|
||||||
|
def loads(self, data):
|
||||||
|
return jsonpickle.loads(data.decode('latin-1'))
|
||||||
|
|||||||
@@ -3,5 +3,6 @@ django-picklefield==0.3.1
|
|||||||
Django==1.8.2
|
Django==1.8.2
|
||||||
hiredis==0.2.0
|
hiredis==0.2.0
|
||||||
humanfriendly==1.27
|
humanfriendly==1.27
|
||||||
|
jsonpickle==0.9.2
|
||||||
logutils==0.3.3
|
logutils==0.3.3
|
||||||
redis==2.10.3
|
redis==2.10.3
|
||||||
|
|||||||
Reference in New Issue
Block a user