-
Notifications
You must be signed in to change notification settings - Fork 115
Expand file tree
/
Copy pathtest_jobaction.py
More file actions
167 lines (117 loc) · 5.49 KB
/
Copy pathtest_jobaction.py
File metadata and controls
167 lines (117 loc) · 5.49 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
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
from mrq.job import Job
from mrq.queue import Queue
import time
import pytest
OPTS = []
for p_query in [
# Query, number_matching
({"path": "tests.tasks.general.MongoInsert"}, 3),
({"queue": "q1"}, 2),
({"params": "{\"a\": 44}"}, 1)
]:
OPTS.append([p_query])
def test_cancel_by_worker(worker):
from bson import ObjectId
job_id = worker.send_task("tests.tasks.general.Add", {"a": 41, "b": 1}, queue="default", block=False)
job = Job(job_id)
job.wait(poll_interval=0.01)
job_data = job.fetch().data
worker.send_task("mrq.basetasks.utils.JobAction", {"action": "cancel", "worker": str(job_data["worker"])},
block=True)
job_data = job.fetch().data
assert job_data["status"] == "cancel"
@pytest.mark.parametrize(["p_query"], OPTS)
def test_cancel_by_path(worker, p_query):
expected_action_jobs = p_query[1]
# Start the worker with only one greenlet so that tasks execute
# sequentially
worker.start(flags="--greenlets 1", queues="default q1 q2")
job_ids = []
job_ids.append(worker.send_task("tests.tasks.general.Add", {
"a": 41, "b": 1, "sleep": 4}, queue="default", block=False))
params = {
"action": "cancel",
"status": "queued"
}
params.update(p_query[0])
job_ids.append(worker.send_task(
"tests.tasks.general.MongoInsert", {"a": 42}, queue="q1", block=False))
job_ids.append(worker.send_task(
"tests.tasks.general.MongoInsert", {"a": 42}, queue="q2", block=False))
job_ids.append(worker.send_task(
"tests.tasks.general.MongoInsert", {"a": 43}, queue="q2", block=False))
job_ids.append(worker.send_task(
"tests.tasks.general.MongoInsert2", {"a": 44}, queue="q1", block=False))
requeue_job = worker.send_task(
"mrq.basetasks.utils.JobAction", params, block=False)
Job(job_ids[-1]).wait(poll_interval=0.01)
# Leave some time to unqueue job_id4 without executing.
time.sleep(1)
worker.stop(deps=False)
jobs = [Job(job_id).fetch().data for job_id in job_ids]
assert jobs[0]["status"] == "success"
assert jobs[0]["result"] == 42
assert Job(requeue_job).fetch().data["result"]["cancelled"] == expected_action_jobs
# Check that the right number of jobs ran.
assert worker.mongodb_jobs.tests_inserts.count() == len(
job_ids) - 1 - expected_action_jobs
action_jobs = list(worker.mongodb_jobs.mrq_jobs.find({"status": "cancel"}))
assert len(action_jobs) == expected_action_jobs
assert set([x.get("result") for x in action_jobs]) == set([None])
assert Queue("default").size() == 0
assert Queue("q1").size() == 0
assert Queue("q2").size() == 0
worker.mongodb_jobs.tests_inserts.remove({})
# Then requeue the same jobs
params = {
"action": "requeue"
}
params.update(p_query[0])
worker.start(flags="--gevent 1", start_deps=False, queues="default", flush=False)
ret = worker.send_task("mrq.basetasks.utils.JobAction", params, block=True)
assert ret["requeued"] == expected_action_jobs
worker.stop(deps=False)
assert worker.mongodb_jobs.mrq_jobs.find(
{"status": "queued"}).count() == expected_action_jobs
assert Queue("default").size() + Queue("q1").size() + \
Queue("q2").size() == expected_action_jobs
worker.stop_deps()
def test_cleaning_jobs(worker):
def test(qname, action, val1, val2):
worker.start()
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("mrq.basetasks.utils.JobAction", {
"path": "tests.tasks.general.MongoInsert",
"status": "queued",
"action": action,
"queue": qname
}, block=False, queue="testMrq")
time.sleep(1)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
worker.send_task("tests.tasks.general.MongoInsert", {"a": 43}, block=False, queue=qname)
assert worker.mongodb_jobs.mrq_jobs.count({"status": "queued"}) == val1
worker.stop(deps=False)
worker.start(queues="testMrq", deps=False)
worker.wait_for_idle()
assert worker.mongodb_jobs.mrq_jobs.count({"status": "queued"}) == val2
worker.stop(deps=True)
# Test action: cancel
test("test_cancel", "cancel", 12, 6)
# Test action requeued
test("test_requeued", "requeue", 12, 11)
# Test action: requeue_retry
test("test_requeue_retry", "requeue_retry", 12, 11)
# Test run task with no job queued
worker.start()
worker.send_task("mrq.basetasks.utils.JobAction", {}, block=False)
worker.wait_for_idle()
assert worker.mongodb_jobs.mrq_jobs.count({"status": "queued"}) == 0
worker.stop(deps=True)