Skip to content

Commit 5ee1335

Browse files
author
ismael
committed
Merge branch '0.9.x' into task_expiry
* 0.9.x: update known queues when requeuing jobs fix makefile add minified js
2 parents 92bd2ea + b8d7112 commit 5ee1335

6 files changed

Lines changed: 48 additions & 5 deletions

File tree

Makefile

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,9 @@ clean:
4444
find . -path ./venv -prune -o -name "*.pyc" -exec rm {} \;
4545
find . -name __pycache__ | xargs rm -r
4646

47+
build_dashboard:
48+
cd mrq/dashboard/static && npm install && mkdir -p bin && npm run build
49+
4750
dashboard:
4851
python mrq/dashboard/app.py
4952

mrq/basetasks/utils.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,12 @@ def perform_action(self, action, query, destination_queue):
134134
"_id": {"$in": jobs_by_queue[queue]}
135135
}, {"$set": updates}, multi=True)
136136

137+
if destination_queue is None:
138+
Queue.ensure_known_queues(jobs_by_queue.iterkeys())
139+
140+
if destination_queue is not None:
141+
Queue.ensure_known_queues([destination_queue])
142+
137143
print(stats)
138144

139145
return stats

mrq/dashboard/static/bin/0.bundle.js

Lines changed: 3 additions & 3 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

mrq/dashboard/static/package.json

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,7 @@
22
"name": "mrq-dashboard",
33
"version": "0.9.1",
44
"description": "PricingAssistant MRQ dashboard",
5-
"dependencies": {
6-
},
5+
"dependencies": {},
76
"engines": {
87
"node": "^7.5.0"
98
},

mrq/queue.py

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010

1111
import sys
1212
from future import standard_library
13+
from itertools import chain
1314

1415
PY3 = sys.version_info > (3,)
1516
standard_library.install_aliases()
@@ -119,6 +120,13 @@ def get_retry_queue(self):
119120
""" Return the name of the queue where retried jobs will be queued """
120121
return self.id
121122

123+
@classmethod
124+
def ensure_known_queues(cls, queues):
125+
""" List all previously known queues """
126+
now = time.time()
127+
params = chain.from_iterable((now, queue) for queue in queues)
128+
context.connections.redis.zadd(Queue.redis_key_known_queues(), *params)
129+
122130
def add_to_known_queues(self, timestamp=None):
123131
""" Adds this queue to the shared list of known queues """
124132
now = timestamp or time.time()

tests/test_general.py

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -179,6 +179,33 @@ def test_known_queues_lifecycle(worker):
179179
# Still not there.
180180
assert set(Queue.redis_known_queues().keys()) == set(["default", "xtest", "test_timed_set"])
181181

182+
# Now we're going to test that the known queues are correctly updated when requeuing a job
183+
184+
# Queue the job again
185+
send_task("tests.tasks.general.Add", {"a": 41, "b": 1, "sleep": 1}, queue="x")
186+
187+
worker.send_task("mrq.basetasks.cleaning.CleanKnownQueues", {}, block=True)
188+
# Requeue it in a different queue
189+
params = {
190+
"action": "requeue",
191+
"destination_queue": "x2"
192+
}
193+
worker.send_task("mrq.basetasks.utils.JobAction", params, block=True)
194+
195+
assert set(Queue.redis_known_queues().keys()) == set(["x", "x2", "default", "xtest", "test_timed_set"])
196+
197+
Queue("x2").empty()
198+
assert set(Queue.redis_known_queues().keys()) == set(["default", "x", "xtest", "test_timed_set"])
199+
200+
# Requeue it in the same queue
201+
params = {
202+
"action": "requeue"
203+
}
204+
worker.send_task("mrq.basetasks.utils.JobAction", params, block=True)
205+
206+
# The queue should be back
207+
assert set(Queue.redis_known_queues().keys()) == set(["default", "x", "x2", "xtest", "test_timed_set"])
208+
182209

183210
def test_general_exception_status(worker):
184211

0 commit comments

Comments
 (0)