Skip to content

Commit d08e0e2

Browse files
committed
Test for --max_latency (#67)
1 parent 96f586a commit d08e0e2

2 files changed

Lines changed: 47 additions & 8 deletions

File tree

tests/tasks/general.py

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,9 @@
1-
from time import sleep
21
from mrq.task import Task
32
from mrq.context import (log, retry_current_job, connections, get_current_config, get_current_job,
43
set_job_progress, subpool_map, queue_job, abort_current_job)
54
import urllib2
65
import json
6+
import time
77

88

99
class Add(Task):
@@ -15,7 +15,7 @@ def run(self, params):
1515

1616
if params.get("sleep", 0):
1717
log.info("sleeping", params.get("sleep", 0))
18-
sleep(params.get("sleep", 0))
18+
time.sleep(params.get("sleep", 0))
1919

2020
return res
2121

@@ -24,6 +24,13 @@ class TimeoutFromConfig(Add):
2424
pass
2525

2626

27+
class GetTime(Task):
28+
29+
def run(self, params):
30+
31+
return time.time()
32+
33+
2734
class Fetch(Task):
2835

2936
def run(self, params):
@@ -53,7 +60,7 @@ def run(self, params):
5360
LEAKS.append(["1" for _ in range(params.get("size", 0))])
5461

5562
if params.get("sleep", 0) > 0:
56-
sleep(params.get("sleep", 0))
63+
time.sleep(params.get("sleep", 0))
5764

5865
return params.get("return")
5966

@@ -80,7 +87,7 @@ class RaiseException(Task):
8087

8188
def run(self, params):
8289

83-
sleep(params.get("sleep", 0))
90+
time.sleep(params.get("sleep", 0))
8491

8592
raise Exception(params.get("message", ""))
8693

@@ -96,7 +103,7 @@ class ReturnParams(Task):
96103

97104
def run(self, params):
98105

99-
sleep(params.get("sleep", 0))
106+
time.sleep(params.get("sleep", 0))
100107

101108
return params
102109

@@ -107,7 +114,7 @@ def run(self, params):
107114

108115
for i in range(1, 100):
109116
set_job_progress(0.01 * i, save=params["save"])
110-
sleep(0.1)
117+
time.sleep(0.1)
111118

112119

113120
class MongoInsert(Task):
@@ -118,7 +125,7 @@ def run(self, params):
118125
{"params": params}, manipulate=False)
119126

120127
if params.get("sleep", 0) > 0:
121-
sleep(params.get("sleep", 0))
128+
time.sleep(params.get("sleep", 0))
122129

123130
return params
124131

@@ -141,7 +148,7 @@ def inner(self, x):
141148
if x == "exception":
142149
raise Exception(x)
143150

144-
sleep(x)
151+
time.sleep(x)
145152

146153
return x
147154

tests/test_performance.py

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,38 @@
22
from mrq.queue import Queue
33
import pytest
44

5+
@pytest.mark.parametrize(["p_max_latency", "p_min_observed_latency", "p_max_observed_latency"], [
6+
[1, 0.03, 1],
7+
[0.01, -1, 0.02]
8+
])
9+
def test_job_max_latency(worker, p_max_latency, p_min_observed_latency, p_max_observed_latency):
10+
11+
worker.start(flags=" --greenlets=1 --max_latency=%s" % (p_max_latency), trace=False)
12+
13+
def get_latency():
14+
t = time.time()
15+
return worker.send_task("tests.tasks.general.GetTime", {}) - t
16+
17+
# Warm up the worker
18+
get_latency()
19+
20+
# This is the latency induced by our test system & general task work
21+
# We're on the same machine so even in different processes time.time() should be pretty reliable
22+
base_latency = get_latency()
23+
print "Base latency: %ss" % base_latency
24+
25+
min_latency = min([get_latency() for _ in range(0, 5)])
26+
print "FYI, min latency = %ss" % min_latency
27+
28+
# Sleep a while with an idle worker to make the poll interval go up
29+
time.sleep(30)
30+
31+
latency = get_latency() - base_latency
32+
33+
print "Observed latency: %ss" % latency
34+
35+
assert p_min_observed_latency <= latency < p_max_observed_latency
36+
537

638
@pytest.mark.parametrize(["p_latency", "p_min", "p_max"], [
739
[0, 0, 3],

0 commit comments

Comments
 (0)