Skip to content

Commit 9c1550e

Browse files
committed
More test fixes
1 parent 4744e2b commit 9c1550e

6 files changed

Lines changed: 73 additions & 57 deletions

File tree

mrq/monkey.py

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -99,9 +99,9 @@ def mrq_monkey_patched(self, *args, **kwargs):
9999
ret = base_method(self, *args, **kwargs)
100100
finally:
101101
stop_time = time.time()
102-
102+
103103
job = None
104-
104+
105105
if config["trace_io"]:
106106
job = get_current_job()
107107
if job:
@@ -388,6 +388,10 @@ class mrq_patched_pymongo_cursor(Cursor):
388388

389389
# Some dark magic is needed here to cope with python's name mangling for private variables.
390390
def _Cursor__send_message(self, *args, **kwargs):
391+
392+
if len(args) and args[0].name == "find":
393+
return Cursor._Cursor__send_message(self, *args, **kwargs)
394+
391395
# print self.__dict__
392396
job = get_current_job()
393397

tests/conftest.py

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -76,11 +76,11 @@ def start(self, cmdline=None, env=None, expected_children=0):
7676
if expected_children > 0:
7777
psutil_process = psutil.Process(self.process.pid)
7878

79-
# print "Expecting %s children, got %s" % (expected_children,
80-
# psutil_process.get_children(recursive=False))
8179
while True:
82-
self.process_children = psutil_process.children(
83-
recursive=True)
80+
self.process_children = psutil_process.children(recursive=True)
81+
# print("Expecting %s children of pid %s, got %s" % (
82+
# expected_children, self.process.pid, len(self.process_children))
83+
# )
8484
if len(self.process_children) >= expected_children:
8585
break
8686
time.sleep(0.1)
@@ -247,6 +247,11 @@ def get_report(self, with_memory=False):
247247
return data
248248

249249
def get_wait_for_idle(self):
250+
251+
if "--processes" in self.cmdline:
252+
print("Warning: get_wait_for_idle() doesn't support multiprocess workers yet")
253+
return False
254+
250255
try:
251256
wait_for_net_service("127.0.0.1", 20020, poll_interval=0.01)
252257
f = urllib.request.urlopen("http://127.0.0.1:20020/wait_for_idle")

tests/test_context.py

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -68,17 +68,18 @@ def test_context_metric_queue(worker):
6868

6969
worker.start(flags=" --config tests/fixtures/config-metric.py")
7070

71+
# Will send 1 task inside!
7172
worker.send_task("tests.tasks.general.SendTask", {
7273
"path": "tests.tasks.general.Add", "params": {"a": 41, "b": 1}})
7374

7475
metrics = json.loads(
7576
worker.send_task("tests.tasks.general.GetMetrics", {}))
7677

7778
# GetMetrics is also a task!
78-
assert metrics.get("queues.default.dequeued") == 2
79-
assert metrics.get("queues.all.dequeued") == 2
80-
assert metrics.get("jobs.status.started") == 2
81-
assert metrics.get("jobs.status.success") == 1 # At the time it is run, GetMetrics isn't success yet.
79+
assert metrics.get("queues.default.dequeued") == 3
80+
assert metrics.get("queues.all.dequeued") == 3
81+
assert metrics.get("jobs.status.started") == 3
82+
assert metrics.get("jobs.status.success") == 2 # At the time it is run, GetMetrics isn't success yet.
8283

8384
TEST_LOCAL_METRICS.get("queues.default.enqueued") == 2
8485
TEST_LOCAL_METRICS.get("queues.all.enqueued") == 2

tests/test_io_hooks.py

Lines changed: 21 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -74,11 +74,13 @@ def test_io_hooks_mongodb(worker):
7474

7575
worker.start(flags=" --config tests/fixtures/config-io-hooks.py")
7676

77-
worker.send_task(
77+
ret = worker.send_task(
7878
"tests.tasks.io.TestIo",
7979
{"test": "mongodb-full-getmore"}
8080
)
8181

82+
print(ret)
83+
8284
events = json.loads(
8385
worker.send_task("tests.tasks.general.GetIoHookEvents", {}))
8486

@@ -87,8 +89,6 @@ def test_io_hooks_mongodb(worker):
8789
for evt in job_events:
8890
print(evt)
8991

90-
assert len(job_events) == 4 * 2
91-
9292
# First, insert
9393
assert job_events[0]["hook"] == "mongodb_pre"
9494
assert job_events[1]["hook"] == "mongodb_post"
@@ -119,13 +119,25 @@ def test_io_hooks_mongodb(worker):
119119
assert job_events[4]["method"] == "cursor"
120120
assert job_events[5]["method"] == "cursor"
121121

122-
# Result MongoDB update
123-
122+
# Then getmore query
124123
assert job_events[6]["hook"] == "mongodb_pre"
125124
assert job_events[7]["hook"] == "mongodb_post"
126125

127-
assert job_events[6]["method"] == "update"
128-
assert job_events[7]["method"] == "update"
126+
assert job_events[6]["collection"] == "mrq.tests_inserts"
127+
assert job_events[7]["collection"] == "mrq.tests_inserts"
128+
129+
assert job_events[6]["method"] == "cursor"
130+
assert job_events[7]["method"] == "cursor"
131+
132+
# Result MongoDB update
133+
134+
assert job_events[8]["hook"] == "mongodb_pre"
135+
assert job_events[9]["hook"] == "mongodb_post"
136+
137+
assert job_events[8]["method"] == "update"
138+
assert job_events[9]["method"] == "update"
139+
140+
assert job_events[8]["collection"] == "mrq.mrq_jobs"
141+
assert job_events[9]["collection"] == "mrq.mrq_jobs"
129142

130-
assert job_events[6]["collection"] == "mrq.mrq_jobs"
131-
assert job_events[7]["collection"] == "mrq.mrq_jobs"
143+
assert len(job_events) == 5 * 2

tests/test_jobinspect.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ def test_current_job_inspect(worker):
1717
job_id = worker.send_task(
1818
"tests.tasks.general.MongoInsert", {"a": 41, "b": 1, "sleep": 3}, block=False)
1919

20-
time.sleep(1)
20+
time.sleep(2)
2121

2222
# Test the HTTP admin API
2323
admin_worker = json.loads(urllib.request.urlopen("http://localhost:20020").read().decode('utf-8'))

tests/test_parallel.py

Lines changed: 31 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,8 @@ def test_parallel_100sleeps(worker, p_flags):
1212

1313
worker.start(flags=p_flags)
1414

15+
print("Worker started. Queueing sleeps")
16+
1517
start_time = time.time()
1618

1719
# This will sleep a total of 100 seconds
@@ -27,56 +29,48 @@ def test_parallel_100sleeps(worker, p_flags):
2729
assert result == list(range(100))
2830

2931

30-
@pytest.mark.parametrize(["p_greenlets"], [
31-
[1],
32-
[2]
32+
@pytest.mark.parametrize(["p_greenlets", "p_strategy"], [
33+
[g, s]
34+
for g in [1, 2]
35+
for s in ["", "parallel", "burst"]
3336
])
34-
def test_dequeue_strategy(worker, p_greenlets):
37+
def test_dequeue_strategy(worker, p_greenlets, p_strategy):
3538

3639
worker.start_deps(flush=True)
3740

3841
worker.send_task(
39-
"tests.tasks.general.MongoInsert", {"a": 41, "sleep": 2}, queue="q1", block=False, start=False)
40-
worker.send_task(
41-
"tests.tasks.general.MongoInsert", {"a": 42, "sleep": 2}, queue="q2", block=False, start=False)
42+
"tests.tasks.general.MongoInsert", {"a": 41, "sleep": 1}, queue="q1", block=False, start=False)
4243
worker.send_task(
43-
"tests.tasks.general.MongoInsert", {"a": 41, "sleep": 2}, queue="q1", block=False, start=False)
44+
"tests.tasks.general.MongoInsert", {"a": 42, "sleep": 1}, queue="q2", block=False, start=False)
4445
worker.send_task(
45-
"tests.tasks.general.MongoInsert", {"a": 42, "sleep": 2}, queue="q2", block=False, start=False)
46+
"tests.tasks.general.MongoInsert", {"a": 43, "sleep": 1}, queue="q1", block=False, start=False)
4647
worker.send_task(
47-
"tests.tasks.general.MongoInsert", {"a": 43, "sleep": 2}, queue="q3", block=False, start=False)
48+
"tests.tasks.general.MongoInsert", {"a": 44, "sleep": 1}, queue="q2", block=False, start=False)
4849
worker.send_task(
49-
"tests.tasks.general.MongoInsert", {"a": 43, "sleep": 2}, queue="q3", block=False, start=False)
50-
51-
time.sleep(0.1)
52-
53-
worker.start(flags="--dequeue_strategy parallel --greenlets %s" % p_greenlets, queues="q1 q2", deps=False, start=False)
54-
55-
if p_greenlets == 1:
56-
time.sleep(1 + 2)
57-
else:
58-
time.sleep(1)
50+
"tests.tasks.general.MongoInsert", {"a": 45, "sleep": 1}, queue="q3", block=False, start=False)
5951

60-
# Should be dequeued in parallel
61-
assert connections.mongodb_jobs.tests_inserts.count({"params.a": 41}) == 1
62-
assert connections.mongodb_jobs.tests_inserts.count({"params.a": 42}) == 1
63-
assert connections.mongodb_jobs.tests_inserts.count() == 2
52+
time.sleep(0.5)
6453

65-
worker.stop(deps=False, sig=9)
66-
time.sleep(1)
54+
flags = "--greenlets %s" % p_greenlets
55+
if p_strategy:
56+
flags += " --dequeue_strategy %s" % p_strategy
6757

68-
worker.start(flags="--dequeue_strategy burst --greenlets 2", queues="q3", deps=False)
58+
print("Worker has flags %s" % flags)
59+
worker.start(flags=flags, queues="q1 q2", deps=False, block=False)
6960

70-
time.sleep(3)
61+
gotit = worker.get_wait_for_idle()
7162

72-
assert connections.mongodb_jobs.tests_inserts.count({"params.a": 43}) == 2
73-
74-
# Worker should be stopped now so even if we queue nothing will happen.
75-
worker.send_task(
76-
"tests.tasks.general.MongoInsert", {"a": 43, "sleep": 2}, queue="q3", block=False, start=False)
77-
78-
time.sleep(2)
63+
if p_strategy == "burst":
64+
assert not gotit # because worker should be stopped already
65+
else:
66+
assert gotit
7967

80-
assert connections.mongodb_jobs.tests_inserts.count({"params.a": 43}) == 2
68+
inserts = list(connections.mongodb_jobs.tests_inserts.find(sort=[("_id", 1)]))
69+
order = [row["params"]["a"] for row in inserts]
8170

82-
worker.stop()
71+
if p_strategy == "parallel":
72+
assert set(order[0:2]) == set([41, 42])
73+
assert set(order[2:4]) == set([43, 44])
74+
else:
75+
assert set(order[0:2]) == set([41, 43])
76+
assert set(order[2:4]) == set([42, 44])

0 commit comments

Comments
 (0)