Skip to content

Commit 55e8310

Browse files
committed
Add missing index + better worker.stop in tests
1 parent 4895e66 commit 55e8310

11 files changed

Lines changed: 29 additions & 30 deletions

mrq/basetasks/indexes.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,9 @@ def run(self, params):
3030
connections.mongodb_jobs.mrq_jobs.ensure_index(
3131
[("dateretry", 1)], sparse=True, background=True)
3232
connections.mongodb_jobs.mrq_jobs.ensure_index(
33-
[("datequeued", 1)], sparse=True, background=True)
33+
[("datequeued", 1)], background=True)
34+
connections.mongodb_jobs.mrq_jobs.ensure_index(
35+
[("queue", 1), ("status", 1), ("datequeued", 1), ("_id", 1)], background=True)
3436

3537
connections.mongodb_jobs.mrq_scheduled_jobs.ensure_index(
3638
[("hash", 1)], unique=True, background=False)

mrq/config.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -337,6 +337,12 @@ def add_parser_args(parser, config_type):
337337
action='store_true',
338338
help='Run the scheduler')
339339

340+
parser.add_argument(
341+
'--ensure_indexes',
342+
default=False,
343+
action='store_true',
344+
help='Ensures the internal MongoDB indexes of MRQ are built, or does so in the background')
345+
340346
parser.add_argument(
341347
'--scheduler_interval',
342348
default=60,

mrq/worker.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
from .exceptions import (TimeoutInterrupt, StopRequested, JobInterrupt, AbortInterrupt,
2323
RetryInterrupt, MaxRetriesInterrupt, MaxConcurrencyInterrupt)
2424
from .context import (set_current_worker, set_current_job, get_current_job, get_current_config,
25-
connections, enable_greenlet_tracing)
25+
connections, enable_greenlet_tracing, run_task)
2626
from .queue import Queue
2727
from .utils import MongoJSONEncoder, MovingAverage
2828
from .processes import Process
@@ -101,6 +101,9 @@ def __init__(self):
101101
"total": 0
102102
}
103103

104+
if self.config["ensure_indexes"]:
105+
run_task("mrq.basetasks.indexes.EnsureIndexes", {})
106+
104107
@property
105108
def config(self):
106109
return get_current_config()

tests/conftest.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -100,6 +100,7 @@ def stop(self, force=False, timeout=None, block=True, sig=15):
100100
# Call this only one time.
101101
if self.stopped and not force:
102102
return
103+
103104
self.stopped = True
104105
self.started = False
105106

tests/test_agent.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -144,8 +144,6 @@ def test_orchestration_scenarios(worker):
144144
"worker2": ["MRQ_WORKER_PROFILE=a mrq-worker a"]
145145
}
146146

147-
worker.stop()
148-
149147

150148
def test_agent_process(worker):
151149

tests/test_cli.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,5 +27,3 @@ def test_cli_run_nonblocking(worker):
2727

2828
assert job1.data["status"] == "success"
2929
assert job1.data["result"] == 42
30-
31-
worker.stop()

tests/test_interrupts.py

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import pytest
88
from mrq.context import connections
99
import json
10+
import os
1011

1112

1213
PROCESS_CONFIGS = [
@@ -32,13 +33,13 @@ def test_interrupt_worker_gracefully(worker, p_flags):
3233
assert job["status"] == "started"
3334

3435
# Stop the worker gracefully. first job should still finish!
35-
worker.stop(block=False, deps=False)
36+
os.kill(worker.process.pid, 2)
3637

3738
time.sleep(1)
3839

3940
# Should not be accepting new jobs!
4041
job_id2 = worker.send_task(
41-
"tests.tasks.general.Add", {"a": 42, "b": 1, "sleep": 4}, block=False)
42+
"tests.tasks.general.Add", {"a": 42, "b": 1, "sleep": 4}, block=False, start=False)
4243

4344
time.sleep(1)
4445

@@ -54,8 +55,6 @@ def test_interrupt_worker_gracefully(worker, p_flags):
5455
job = Job(job_id2).fetch().data
5556
assert job.get("status") == "queued"
5657

57-
worker.stop_deps()
58-
5958

6059
@pytest.mark.parametrize(["p_flags"], PROCESS_CONFIGS)
6160
def test_interrupt_worker_double_sigint(worker, p_flags):
@@ -75,13 +74,13 @@ def test_interrupt_worker_double_sigint(worker, p_flags):
7574
assert job["status"] == "started"
7675

7776
# Stop the worker gracefully. first job should still finish!
78-
worker.stop(block=False, deps=False)
77+
os.kill(worker.process.pid, 2)
7978

8079
time.sleep(1)
8180

8281
# Should not be accepting new jobs!
8382
job_id2 = worker.send_task(
84-
"tests.tasks.general.Add", {"a": 42, "b": 1, "sleep": 20}, block=False)
83+
"tests.tasks.general.Add", {"a": 42, "b": 1, "sleep": 20}, block=False, start=False)
8584

8685
time.sleep(1)
8786

@@ -92,7 +91,7 @@ def test_interrupt_worker_double_sigint(worker, p_flags):
9291
assert job["status"] == "started"
9392

9493
# Sending a second kill -2 should make it stop
95-
worker.stop(block=True, deps=False, force=True)
94+
os.kill(worker.process.pid, 2)
9695

9796
while Job(job_id).fetch().data["status"] == "started":
9897
time.sleep(0.1)
@@ -356,8 +355,10 @@ def test_interrupt_maxconcurrency(worker):
356355
assert set(job_statuses) == set(["success", "maxconcurrency"])
357356

358357
# the job concurrency key must be equal to 0
359-
last_job_id = worker.send_task("tests.tasks.concurrency.LockedAdd",
360-
{"a": 1, "b": 1, "sleep": 2}, block=False
358+
last_job_id = worker.send_task(
359+
"tests.tasks.concurrency.LockedAdd",
360+
{"a": 1, "b": 1, "sleep": 2},
361+
block=False
361362
)
362363

363364
last_job = Job(last_job_id).wait(poll_interval=0.01)

tests/test_memoryleaks.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -29,8 +29,6 @@ def test_max_memory_restart(worker):
2929
# We must have been restarted at least once.
3030
assert worker.mongodb_jobs.mrq_workers.find().count() > 1
3131

32-
worker.stop()
33-
3432

3533
def get_diff_after_jobs(worker, n_tasks, leak, sleep=0):
3634

tests/test_pause.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -47,8 +47,6 @@ def test_pause_resume(worker):
4747

4848
assert worker.mongodb_jobs.tests_inserts.count() == 2
4949

50-
worker.stop()
51-
5250

5351
def test_pause_refresh_interval(worker):
5452

@@ -74,8 +72,6 @@ def test_pause_refresh_interval(worker):
7472
assert job1["status"] == "success"
7573
assert job1["result"] == {"a": 41}
7674

77-
worker.stop()
78-
7975

8076
def test_pause_subqueue(worker):
8177

@@ -126,5 +122,3 @@ def test_pause_subqueue(worker):
126122
assert job2["result"] == {"a": 43}
127123

128124
assert worker.mongodb_jobs.tests_inserts.count() == 2
129-
130-
worker.stop()

tests/test_performance.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
])
2222
def test_job_max_latency(worker, p_max_latency, p_min_observed_latency, p_max_observed_latency):
2323

24-
worker.start(flags=" --greenlets=1 --max_latency=%s" % (p_max_latency), trace=False)
24+
worker.start(flags=" --ensure_indexes --greenlets=1 --max_latency=%s" % (p_max_latency), trace=False)
2525

2626
def get_latency():
2727
t = time.time()
@@ -76,7 +76,7 @@ def test_network_latency(worker, p_latency, p_min, p_max):
7676

7777
def benchmark_task(worker, taskpath, taskparams, tasks=1000, greenlets=50, processes=0, max_seconds=10, profile=False, quiet=True, raw=False, queues="default", config=None):
7878

79-
worker.start(flags="--processes %s --greenlets %s%s%s%s" % (
79+
worker.start(flags="--ensure_indexes --processes %s --greenlets %s%s%s%s" % (
8080
processes,
8181
greenlets,
8282
" --profile" if profile else "",
@@ -321,9 +321,9 @@ def test_performance_queue_cancel_requeue(worker):
321321
def test_worker_efficiency(worker, p_queue_type, p_greenlets, p_min_efficiency):
322322

323323
if p_queue_type == "regular":
324-
worker.start(trace=False, flags="--greenlets %s" % p_greenlets, queues="default")
324+
worker.start(trace=False, flags="--ensure_indexes --greenlets %s" % p_greenlets, queues="default")
325325
elif p_queue_type == "raw":
326-
worker.start(trace=False, flags="--greenlets %s --config tests/fixtures/config-raw1.py" % p_greenlets,
326+
worker.start(trace=False, flags="--ensure_indexes --greenlets %s --config tests/fixtures/config-raw1.py" % p_greenlets,
327327
queues="testperformance_efficiency_raw")
328328
elif p_queue_type == "raw_nostorage":
329329
worker.start(trace=False, flags="--greenlets %s --config tests/fixtures/config-raw1.py" % p_greenlets,
@@ -360,7 +360,7 @@ def test_worker_efficiency(worker, p_queue_type, p_greenlets, p_min_efficiency):
360360
))
361361

362362
# We can't be perfectly efficient!
363-
assert perfect_time < total_time
363+
assert (perfect_time - 1) < total_time
364364

365365
# But we should be at least 80% efficient
366366
assert perfect_time > total_time * p_min_efficiency

0 commit comments

Comments
 (0)