|
|
@ -3,9 +3,10 @@ from __future__ import (absolute_import, division, print_function,
|
|
|
|
unicode_literals)
|
|
|
|
unicode_literals)
|
|
|
|
|
|
|
|
|
|
|
|
from datetime import datetime
|
|
|
|
from datetime import datetime
|
|
|
|
import time
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
import time
|
|
|
|
import sys
|
|
|
|
import sys
|
|
|
|
|
|
|
|
import zlib
|
|
|
|
|
|
|
|
|
|
|
|
is_py2 = sys.version[0] == '2'
|
|
|
|
is_py2 = sys.version[0] == '2'
|
|
|
|
if is_py2:
|
|
|
|
if is_py2:
|
|
|
@ -15,7 +16,7 @@ else:
|
|
|
|
|
|
|
|
|
|
|
|
from tests import fixtures, RQTestCase
|
|
|
|
from tests import fixtures, RQTestCase
|
|
|
|
|
|
|
|
|
|
|
|
from rq.compat import PY2
|
|
|
|
from rq.compat import PY2, as_text
|
|
|
|
from rq.exceptions import NoSuchJobError, UnpickleError
|
|
|
|
from rq.exceptions import NoSuchJobError, UnpickleError
|
|
|
|
from rq.job import Job, get_current_job, JobStatus, cancel_job, requeue_job
|
|
|
|
from rq.job import Job, get_current_job, JobStatus, cancel_job, requeue_job
|
|
|
|
from rq.queue import Queue, get_failed_queue
|
|
|
|
from rq.queue import Queue, get_failed_queue
|
|
|
@ -160,7 +161,7 @@ class TestJob(RQTestCase):
|
|
|
|
self.assertEqual(self.testconn.type(job.key), b'hash')
|
|
|
|
self.assertEqual(self.testconn.type(job.key), b'hash')
|
|
|
|
|
|
|
|
|
|
|
|
# Saving writes pickled job data
|
|
|
|
# Saving writes pickled job data
|
|
|
|
unpickled_data = loads(self.testconn.hget(job.key, 'data'))
|
|
|
|
unpickled_data = loads(zlib.decompress(self.testconn.hget(job.key, 'data')))
|
|
|
|
self.assertEqual(unpickled_data[0], 'tests.fixtures.some_calculation')
|
|
|
|
self.assertEqual(unpickled_data[0], 'tests.fixtures.some_calculation')
|
|
|
|
|
|
|
|
|
|
|
|
def test_fetch(self):
|
|
|
|
def test_fetch(self):
|
|
|
@ -236,7 +237,8 @@ class TestJob(RQTestCase):
|
|
|
|
def test_fetching_unreadable_data(self):
|
|
|
|
def test_fetching_unreadable_data(self):
|
|
|
|
"""Fetching succeeds on unreadable data, but lazy props fail."""
|
|
|
|
"""Fetching succeeds on unreadable data, but lazy props fail."""
|
|
|
|
# Set up
|
|
|
|
# Set up
|
|
|
|
job = Job.create(func=fixtures.some_calculation, args=(3, 4), kwargs=dict(z=2))
|
|
|
|
job = Job.create(func=fixtures.some_calculation, args=(3, 4),
|
|
|
|
|
|
|
|
kwargs=dict(z=2))
|
|
|
|
job.save()
|
|
|
|
job.save()
|
|
|
|
|
|
|
|
|
|
|
|
# Just replace the data hkey with some random noise
|
|
|
|
# Just replace the data hkey with some random noise
|
|
|
@ -255,14 +257,57 @@ class TestJob(RQTestCase):
|
|
|
|
# Now slightly modify the job to make it unimportable (this is
|
|
|
|
# Now slightly modify the job to make it unimportable (this is
|
|
|
|
# equivalent to a worker not having the most up-to-date source code
|
|
|
|
# equivalent to a worker not having the most up-to-date source code
|
|
|
|
# and unable to import the function)
|
|
|
|
# and unable to import the function)
|
|
|
|
data = self.testconn.hget(job.key, 'data')
|
|
|
|
job_data = job.data
|
|
|
|
unimportable_data = data.replace(b'say_hello', b'nay_hello')
|
|
|
|
unimportable_data = job_data.replace(b'say_hello', b'nay_hello')
|
|
|
|
self.testconn.hset(job.key, 'data', unimportable_data)
|
|
|
|
|
|
|
|
|
|
|
|
self.testconn.hset(job.key, 'data', zlib.compress(unimportable_data))
|
|
|
|
|
|
|
|
|
|
|
|
job.refresh()
|
|
|
|
job.refresh()
|
|
|
|
with self.assertRaises(AttributeError):
|
|
|
|
with self.assertRaises(AttributeError):
|
|
|
|
job.func # accessing the func property should fail
|
|
|
|
job.func # accessing the func property should fail
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_compressed_exc_info_handling(self):
|
|
|
|
|
|
|
|
"""Jobs handle both compressed and uncompressed exc_info"""
|
|
|
|
|
|
|
|
exception_string = 'Some exception'
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
job = Job.create(func=fixtures.say_hello, args=('Lionel',))
|
|
|
|
|
|
|
|
job.exc_info = exception_string
|
|
|
|
|
|
|
|
job.save()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# exc_info is stored in compressed format
|
|
|
|
|
|
|
|
exc_info = self.testconn.hget(job.key, 'exc_info')
|
|
|
|
|
|
|
|
self.assertEqual(
|
|
|
|
|
|
|
|
as_text(zlib.decompress(exc_info)),
|
|
|
|
|
|
|
|
exception_string
|
|
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
job.refresh()
|
|
|
|
|
|
|
|
self.assertEqual(job.exc_info, exception_string)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# Uncompressed exc_info is also handled
|
|
|
|
|
|
|
|
self.testconn.hset(job.key, 'exc_info', exception_string)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
job.refresh()
|
|
|
|
|
|
|
|
self.assertEqual(job.exc_info, exception_string)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_compressed_job_data_handling(self):
|
|
|
|
|
|
|
|
"""Jobs handle both compressed and uncompressed data"""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
job = Job.create(func=fixtures.say_hello, args=('Lionel',))
|
|
|
|
|
|
|
|
job.save()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# Job data is stored in compressed format
|
|
|
|
|
|
|
|
job_data = job.data
|
|
|
|
|
|
|
|
self.assertEqual(
|
|
|
|
|
|
|
|
zlib.compress(job_data),
|
|
|
|
|
|
|
|
self.testconn.hget(job.key, 'data')
|
|
|
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
self.testconn.hset(job.key, 'data', job_data)
|
|
|
|
|
|
|
|
job.refresh()
|
|
|
|
|
|
|
|
self.assertEqual(job.data, job_data)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_custom_meta_is_persisted(self):
|
|
|
|
def test_custom_meta_is_persisted(self):
|
|
|
|
"""Additional meta data on jobs are stored persisted correctly."""
|
|
|
|
"""Additional meta data on jobs are stored persisted correctly."""
|
|
|
|
job = Job.create(func=fixtures.say_hello, args=('Lionel',))
|
|
|
|
job = Job.create(func=fixtures.say_hello, args=('Lionel',))
|
|
|
@ -457,7 +502,7 @@ class TestJob(RQTestCase):
|
|
|
|
"""test if a job created with ttl expires [issue502]"""
|
|
|
|
"""test if a job created with ttl expires [issue502]"""
|
|
|
|
queue = Queue(connection=self.testconn)
|
|
|
|
queue = Queue(connection=self.testconn)
|
|
|
|
queue.enqueue(fixtures.say_hello, job_id="1234", ttl=1)
|
|
|
|
queue.enqueue(fixtures.say_hello, job_id="1234", ttl=1)
|
|
|
|
time.sleep(1)
|
|
|
|
time.sleep(1.1)
|
|
|
|
self.assertEqual(0, len(queue.get_jobs()))
|
|
|
|
self.assertEqual(0, len(queue.get_jobs()))
|
|
|
|
|
|
|
|
|
|
|
|
def test_create_and_cancel_job(self):
|
|
|
|
def test_create_and_cancel_job(self):
|
|
|
|