-
Notifications
You must be signed in to change notification settings - Fork 115
Expand file tree
/
Copy pathtest_subqueues.py
More file actions
76 lines (50 loc) · 2.25 KB
/
Copy pathtest_subqueues.py
File metadata and controls
76 lines (50 loc) · 2.25 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
import time
import pytest
from mrq.job import Job
from mrq.queue import Queue
@pytest.mark.parametrize(["queues", "enqueue_on"], [
[["main/", "second/"], ["main/", "main/sub", "main/sub/nested", "second/x"]],
[["prefix/main/"], ["prefix/main/", "prefix/main/sub", "prefix/main/sub/nested"]],
])
def test_matchable_subqueues(worker, queues, enqueue_on):
worker.start(queues=" ".join(queues), flags="--subqueues_refresh_interval=0.1")
job_ids = []
for subqueue in enqueue_on:
job_id = worker.send_task("tests.tasks.general.GetTime", {}, queue=subqueue, block=False)
job_ids.append(job_id)
assert all([Job(j).wait(poll_interval=0.01, timeout=3) for j in job_ids])
worker.stop()
@pytest.mark.parametrize(["queue", "enqueue_on"], [
["main/", ["/main", "main_", "/", "main", "other"]],
["prefix/main/", ["prefix", "prefix/other", "prefix/main"]],
])
def test_unmatchable_subqueues(worker, queue, enqueue_on):
worker.start(queues=queue, flags="--subqueues_refresh_interval=0.1")
job_ids = []
for subqueue in enqueue_on:
job_id = worker.send_task("tests.tasks.general.GetTime", {}, queue=subqueue, block=False)
job_ids.append(job_id)
time.sleep(2)
results = [Job(j).fetch().data.get("status") for j in job_ids]
# ensure tasks are not consumed by a worker
assert results == ["queued"] * len(results)
worker.stop()
@pytest.mark.parametrize(["delimiter"], ["/", ".", "-"])
def test_custom_delimiters(worker, delimiter):
queue = "main" + delimiter
subqueue = queue + "subqueue"
worker.start(queues=queue, flags="--subqueues_refresh_interval=0.1 --subqueues_delimiter=%s" % delimiter)
job_id = worker.send_task("tests.tasks.general.GetTime", {}, queue=subqueue, block=False)
Job(job_id).wait(poll_interval=0.01)
worker.stop()
def test_refresh_interval(worker):
""" Tests that a refresh interval of 0 disables the subqueue detection """
worker.start(queues="test/", flags="--subqueues_refresh_interval=0")
time.sleep(2)
job_id1 = worker.send_task(
"tests.tasks.general.GetTime", {"a": 41},
queue="test/subqueue", block=False)
time.sleep(5)
job1 = Job(job_id1).fetch().data
assert job1["status"] == "queued"
worker.stop()