Skip to content

Commit 1ed40d1

Browse files
author
Florian Perucki
committed
pause subqueues and the root queue is paused
1 parent 1638a89 commit 1ed40d1

3 files changed

Lines changed: 81 additions & 6 deletions

File tree

mrq/queue.py

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -158,21 +158,36 @@ def unserialize_job_ids(self, job_ids):
158158
else:
159159
return [x.encode('hex') for x in job_ids]
160160

161+
def _get_pausable_id(self):
162+
"""
163+
Get the queue id (either id or root_id) that should be used to pause/unpause the current queue
164+
TODO: handle subqueues with more than one level, e.g. "queue/subqueue/"
165+
"""
166+
queue = self.id
167+
delimiter = context.get_current_config().get("subqueues_delimiter")
168+
if delimiter is not None and self.id.endswith(delimiter):
169+
queue = self.root_id
170+
return queue
171+
161172
def pause(self):
162173
""" Adds this queue to the set of paused queues """
163-
context.connections.redis.sadd(Queue.redis_key_paused_queues(), self.id)
174+
context.connections.redis.sadd(Queue.redis_key_paused_queues(), self._get_pausable_id())
164175

165176
def is_paused(self):
166177
"""
167178
Returns wether the queue is paused or not.
168179
Warning: this does NOT ensure that the queue was effectively added to
169-
the list of paused queues. See the 'paused_queues_refresh_interval' option.
180+
the set of paused queues. See the 'paused_queues_refresh_interval' option.
170181
"""
171-
return context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id)
182+
root_is_paused = False
183+
if self.root_id != self.id:
184+
root_is_paused = context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.root_id)
185+
186+
return root_is_paused or context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id)
172187

173188
def resume(self):
174189
""" Resumes a paused queue """
175-
context.connections.redis.srem(Queue.redis_key_paused_queues(), self.id)
190+
context.connections.redis.srem(Queue.redis_key_paused_queues(), self._get_pausable_id())
176191

177192
def size(self):
178193
""" Returns the total number of jobs on the queue """

mrq/worker.py

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -231,7 +231,7 @@ def greenlet_paused_queues(self):
231231

232232
while True:
233233

234-
# Update the process-local list of known queues
234+
# Update the process-local list of paused queues
235235
Queue.paused_queues = Queue.redis_paused_queues()
236236
time.sleep(self.config["paused_queues_refresh_interval"])
237237

@@ -483,7 +483,12 @@ def work_loop(self, max_jobs=None):
483483
gevent.sleep(0.01)
484484

485485
jobs = []
486-
available_queues = [queue for queue in self.queues if queue.id not in Queue.paused_queues]
486+
487+
available_queues = [
488+
queue for queue in self.queues
489+
if queue.root_id not in Queue.paused_queues and
490+
queue.id not in Queue.paused_queues
491+
]
487492

488493
for queue_i in xrange(len(available_queues)):
489494

tests/test_pause.py

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
from mrq.job import Job
2+
import pytest
23
from mrq.queue import Queue, send_task
34
import time
5+
from mrq.context import set_current_config, get_config
46

57

68
def test_pause_resume(worker):
@@ -73,3 +75,56 @@ def test_pause_refresh_interval(worker):
7375
assert job1["result"] == {"a": 41}
7476

7577
worker.stop()
78+
79+
80+
def test_pause_subqueue(worker):
81+
82+
# set config in current context in order to have a subqueue delimiter
83+
set_current_config(get_config(config_type="worker"))
84+
85+
worker.start(queues="high high/", flags="--subqueues_refresh_interval=1 --paused_queues_refresh_interval=1")
86+
87+
Queue("high").pause()
88+
89+
assert Queue("high/").is_paused()
90+
91+
# wait for the paused_queues list to be refreshed
92+
time.sleep(2)
93+
94+
job_id1 = send_task(
95+
"tests.tasks.general.MongoInsert", {"a": 41},
96+
queue="high")
97+
98+
job_id2 = send_task(
99+
"tests.tasks.general.MongoInsert", {"a": 43},
100+
queue="high/subqueue")
101+
102+
# wait a bit to make sure the jobs status will still be queued
103+
time.sleep(5)
104+
105+
job1 = Job(job_id1).fetch().data
106+
job2 = Job(job_id2).fetch().data
107+
108+
assert job1["status"] == "queued"
109+
assert job2["status"] == "queued"
110+
111+
assert worker.mongodb_jobs.tests_inserts.count() == 0
112+
113+
Queue("high/").resume()
114+
115+
Job(job_id1).wait(poll_interval=0.01)
116+
117+
Job(job_id2).wait(poll_interval=0.01)
118+
119+
job1 = Job(job_id1).fetch().data
120+
job2 = Job(job_id2).fetch().data
121+
122+
assert job1["status"] == "success"
123+
assert job1["result"] == {"a": 41}
124+
125+
assert job2["status"] == "success"
126+
assert job2["result"] == {"a": 43}
127+
128+
assert worker.mongodb_jobs.tests_inserts.count() == 2
129+
130+
worker.stop()

0 commit comments

Comments
 (0)