2012-01-28 06:58:40 +00:00
|
|
|
from tests import RQTestCase
|
|
|
|
from tests import testjob
|
|
|
|
from pickle import dumps
|
|
|
|
from rq import Queue
|
2012-01-30 15:49:15 +00:00
|
|
|
from rq.exceptions import UnpickleError
|
2011-11-15 07:13:16 +00:00
|
|
|
|
|
|
|
|
2011-11-14 13:18:21 +00:00
|
|
|
class TestQueue(RQTestCase):
|
|
|
|
def test_create_queue(self):
|
|
|
|
"""Creating queues."""
|
|
|
|
q = Queue('my-queue')
|
|
|
|
self.assertEquals(q.name, 'my-queue')
|
|
|
|
|
2011-11-15 20:15:51 +00:00
|
|
|
def test_create_default_queue(self):
|
|
|
|
"""Instantiating the default queue."""
|
|
|
|
q = Queue()
|
|
|
|
self.assertEquals(q.name, 'default')
|
|
|
|
|
2011-11-16 11:45:16 +00:00
|
|
|
|
|
|
|
def test_equality(self):
|
|
|
|
"""Mathematical equality of queues."""
|
|
|
|
q1 = Queue('foo')
|
|
|
|
q2 = Queue('foo')
|
|
|
|
q3 = Queue('bar')
|
|
|
|
|
|
|
|
self.assertEquals(q1, q2)
|
|
|
|
self.assertEquals(q2, q1)
|
|
|
|
self.assertNotEquals(q1, q3)
|
|
|
|
self.assertNotEquals(q2, q3)
|
|
|
|
|
|
|
|
|
2011-11-14 13:18:21 +00:00
|
|
|
def test_queue_empty(self):
|
|
|
|
"""Detecting empty queues."""
|
|
|
|
q = Queue('my-queue')
|
|
|
|
self.assertEquals(q.empty, True)
|
|
|
|
|
2012-01-28 06:58:40 +00:00
|
|
|
self.testconn.rpush('rq:queue:my-queue', 'some val')
|
2011-11-14 13:18:21 +00:00
|
|
|
self.assertEquals(q.empty, False)
|
2011-11-14 11:10:59 +00:00
|
|
|
|
2011-11-15 08:36:32 +00:00
|
|
|
|
|
|
|
def test_enqueue(self):
|
2011-11-15 20:15:51 +00:00
|
|
|
"""Putting work on queues."""
|
2011-11-15 07:13:16 +00:00
|
|
|
q = Queue('my-queue')
|
|
|
|
self.assertEquals(q.empty, True)
|
|
|
|
|
2011-11-15 07:43:06 +00:00
|
|
|
# testjob spec holds which queue this is sent to
|
2011-11-15 20:15:51 +00:00
|
|
|
q.enqueue(testjob, 'Nick', foo='bar')
|
2011-11-15 07:43:06 +00:00
|
|
|
self.assertEquals(q.empty, False)
|
|
|
|
self.assertQueueContains(q, testjob)
|
|
|
|
|
2011-11-15 08:36:29 +00:00
|
|
|
def test_dequeue(self):
|
|
|
|
"""Fetching work from specific queue."""
|
|
|
|
q = Queue('foo')
|
2011-11-15 20:15:51 +00:00
|
|
|
q.enqueue(testjob, 'Rick', foo='bar')
|
2011-11-15 08:36:29 +00:00
|
|
|
|
|
|
|
# Pull it off the queue (normally, a worker would do this)
|
2011-11-16 11:44:33 +00:00
|
|
|
job = q.dequeue()
|
|
|
|
self.assertEquals(job.func, testjob)
|
|
|
|
self.assertEquals(job.origin, q)
|
|
|
|
self.assertEquals(job.args[0], 'Rick')
|
|
|
|
self.assertEquals(job.kwargs['foo'], 'bar')
|
|
|
|
|
|
|
|
|
|
|
|
def test_dequeue_any(self):
|
|
|
|
"""Fetching work from any given queue."""
|
|
|
|
fooq = Queue('foo')
|
|
|
|
barq = Queue('bar')
|
|
|
|
|
|
|
|
self.assertEquals(Queue.dequeue_any([fooq, barq], False), None)
|
|
|
|
|
|
|
|
# Enqueue a single item
|
|
|
|
barq.enqueue(testjob)
|
|
|
|
job = Queue.dequeue_any([fooq, barq], False)
|
|
|
|
self.assertEquals(job.func, testjob)
|
|
|
|
|
|
|
|
# Enqueue items on both queues
|
|
|
|
barq.enqueue(testjob, 'for Bar')
|
|
|
|
fooq.enqueue(testjob, 'for Foo')
|
|
|
|
|
|
|
|
job = Queue.dequeue_any([fooq, barq], False)
|
|
|
|
self.assertEquals(job.func, testjob)
|
|
|
|
self.assertEquals(job.origin, fooq)
|
|
|
|
self.assertEquals(job.args[0], 'for Foo', 'Foo should be dequeued first.')
|
|
|
|
|
|
|
|
job = Queue.dequeue_any([fooq, barq], False)
|
|
|
|
self.assertEquals(job.func, testjob)
|
|
|
|
self.assertEquals(job.origin, barq)
|
|
|
|
self.assertEquals(job.args[0], 'for Bar', 'Bar should be dequeued second.')
|
|
|
|
|
|
|
|
|
2012-01-27 14:15:13 +00:00
|
|
|
def test_dequeue_unpicklable_data(self):
|
|
|
|
"""Error handling of invalid pickle data."""
|
|
|
|
|
|
|
|
# Push non-pickle data on the queue
|
|
|
|
q = Queue('foo')
|
|
|
|
blob = 'this is nothing like pickled data'
|
|
|
|
self.testconn.rpush(q._key, blob)
|
|
|
|
|
2012-01-30 15:49:15 +00:00
|
|
|
with self.assertRaises(UnpickleError):
|
2012-01-27 14:15:13 +00:00
|
|
|
q.dequeue() # error occurs when perform()'ing
|
|
|
|
|
|
|
|
# Push value pickle data, but not representing a job tuple
|
|
|
|
q = Queue('foo')
|
2012-01-30 15:49:15 +00:00
|
|
|
blob = dumps('this is pickled, but not a job tuple')
|
2012-01-27 14:15:13 +00:00
|
|
|
self.testconn.rpush(q._key, blob)
|
|
|
|
|
2012-01-30 15:49:15 +00:00
|
|
|
with self.assertRaises(UnpickleError):
|
2012-01-27 14:15:13 +00:00
|
|
|
q.dequeue() # error occurs when perform()'ing
|
|
|
|
|
|
|
|
# Push slightly incorrect pickled data onto the queue (simulate
|
|
|
|
# a function that can't be imported from the worker)
|
|
|
|
q = Queue('foo')
|
|
|
|
|
|
|
|
job_tuple = dumps((testjob, [], dict(name='Frank'), 'unused'))
|
|
|
|
blob = job_tuple.replace('testjob', 'fooobar')
|
|
|
|
self.testconn.rpush(q._key, blob)
|
|
|
|
|
2012-01-30 15:49:15 +00:00
|
|
|
with self.assertRaises(UnpickleError):
|
2012-01-27 14:15:13 +00:00
|
|
|
q.dequeue() # error occurs when dequeue()'ing
|
|
|
|
|