2012-01-28 06:58:40 +00:00
|
|
|
from tests import RQTestCase
|
2012-01-30 18:41:13 +00:00
|
|
|
from tests import testjob, failing_job
|
2012-01-28 06:58:40 +00:00
|
|
|
from rq import Queue, Worker
|
2012-01-30 18:41:13 +00:00
|
|
|
from rq.job import Job
|
2012-01-28 06:58:40 +00:00
|
|
|
|
|
|
|
|
|
|
|
class TestWorker(RQTestCase):
|
|
|
|
def test_create_worker(self):
|
|
|
|
"""Worker creation."""
|
|
|
|
fooq, barq = Queue('foo'), Queue('bar')
|
|
|
|
w = Worker([fooq, barq])
|
|
|
|
self.assertEquals(w.queues, [fooq, barq])
|
|
|
|
|
|
|
|
def test_work_and_quit(self):
|
|
|
|
"""Worker processes work, then quits."""
|
|
|
|
fooq, barq = Queue('foo'), Queue('bar')
|
|
|
|
w = Worker([fooq, barq])
|
|
|
|
self.assertEquals(w.work(burst=True), False, 'Did not expect any work on the queue.')
|
|
|
|
|
|
|
|
fooq.enqueue(testjob, name='Frank')
|
|
|
|
self.assertEquals(w.work(burst=True), True, 'Expected at least some work done.')
|
|
|
|
|
2012-01-30 18:41:13 +00:00
|
|
|
def test_work_is_unreadable(self):
|
2012-02-08 14:18:24 +00:00
|
|
|
"""Unreadable jobs are put on the failure queue."""
|
2012-01-30 18:41:13 +00:00
|
|
|
q = Queue()
|
|
|
|
failure_q = Queue('failure')
|
|
|
|
|
|
|
|
self.assertEquals(failure_q.count, 0)
|
|
|
|
self.assertEquals(q.count, 0)
|
|
|
|
|
|
|
|
# NOTE: We have to fake this enqueueing for this test case.
|
|
|
|
# What we're simulating here is a call to a function that is not
|
|
|
|
# importable from the worker process.
|
2012-02-07 23:40:43 +00:00
|
|
|
job = Job.for_call(failing_job, 3)
|
2012-02-08 13:18:17 +00:00
|
|
|
job.save()
|
|
|
|
data = self.testconn.hget(job.key, 'data')
|
|
|
|
invalid_data = data.replace('failing_job', 'nonexisting_job')
|
|
|
|
self.testconn.hset(job.key, 'data', invalid_data)
|
2012-01-30 18:41:13 +00:00
|
|
|
|
|
|
|
# We use the low-level internal function to enqueue any data (bypassing
|
|
|
|
# validity checks)
|
2012-02-08 14:18:24 +00:00
|
|
|
q.push_job_id(job.id)
|
2012-01-30 18:41:13 +00:00
|
|
|
|
|
|
|
self.assertEquals(q.count, 1)
|
|
|
|
|
|
|
|
# All set, we're going to process it
|
|
|
|
w = Worker([q])
|
|
|
|
w.work(burst=True) # should silently pass
|
|
|
|
self.assertEquals(q.count, 0)
|
|
|
|
|
|
|
|
self.assertEquals(failure_q.count, 1)
|
|
|
|
|
2012-01-30 19:55:52 +00:00
|
|
|
def test_work_fails(self):
|
2012-02-08 14:18:24 +00:00
|
|
|
"""Failing jobs are put on the failure queue."""
|
2012-01-30 19:55:52 +00:00
|
|
|
q = Queue()
|
|
|
|
failure_q = Queue('failure')
|
|
|
|
|
|
|
|
self.assertEquals(failure_q.count, 0)
|
|
|
|
self.assertEquals(q.count, 0)
|
|
|
|
|
|
|
|
q.enqueue(failing_job)
|
|
|
|
self.assertEquals(q.count, 1)
|
2012-02-08 14:18:24 +00:00
|
|
|
|
2012-01-30 19:55:52 +00:00
|
|
|
w = Worker([q])
|
|
|
|
w.work(burst=True) # should silently pass
|
|
|
|
self.assertEquals(q.count, 0)
|
|
|
|
self.assertEquals(failure_q.count, 1)
|
|
|
|
|
2012-01-28 06:58:40 +00:00
|
|
|
|