Skip to content

Commit b153975

Browse files
committed
Merge branch 'master' of github.com:pricingassistant/mrq
2 parents 198dd6a + ab939ed commit b153975

4 files changed

Lines changed: 62 additions & 22 deletions

File tree

mrq/scheduler.py

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ def sync_tasks(self, tasks):
3636
# The date part will be discarded in check()
3737
task["dailytime"] = datetime.datetime.combine(datetime.datetime.utcnow(), task["dailytime"])
3838
task["interval"] = 3600 * 24
39+
3940
self.collection.insert(task)
4041
log.debug("Scheduler: added %s" % task["hash"])
4142

@@ -55,19 +56,32 @@ def check(self):
5556

5657
dailytime = task.get("dailytime").time()
5758

58-
if task.get("datelastqueued") and task.get("datelastqueued").time().isoformat()[0:8] != dailytime.isoformat()[0:8]:
59-
log.debug("Adjusting the time of scheduled task %s to %s" % (task["_id"], dailytime))
60-
61-
self.collection.update({"_id": task["_id"]}, {"$set": {
62-
"datelastqueued": datetime.datetime.combine(task.get("datelastqueued").date() - datetime.timedelta(days=1), dailytime)
59+
time_datelastqueued = task.get("datelastqueued").time().isoformat()[0:8]
60+
time_dailytime = dailytime.isoformat()[0:8]
61+
if task.get("datelastqueued") and time_datelastqueued != time_dailytime:
62+
log.debug("Adjusting the time of scheduled task %s from %s to %s" % (task["_id"], time_datelastqueued, time_dailytime))
63+
64+
# Make sure we don't queue the task in a loop by adjusting the time
65+
if time_datelastqueued < time_dailytime:
66+
adjusted_datelastqueued = datetime.datetime.combine(task.get("datelastqueued").date() - datetime.timedelta(days=1), dailytime)
67+
else:
68+
adjusted_datelastqueued = datetime.datetime.combine(task.get("datelastqueued").date(), dailytime)
69+
70+
# We do find_and_modify and not update() because several check() may be happening
71+
# at the same time.
72+
self.collection.find_and_modify({
73+
"_id": task["_id"],
74+
"datelastqueued": task.get("datelastqueued")
75+
}, {"$set": {
76+
"datelastqueued": adjusted_datelastqueued
6377
}})
6478
self.refresh()
6579

6680
task_data = self.collection.find_and_modify({
6781
"_id": task["_id"],
6882
"datelastqueued": {"$lt": last_time}
6983
}, {"$set": {
70-
"datelastqueued": datetime.datetime.utcnow()
84+
"datelastqueued": now
7185
}})
7286

7387
if task_data:

mrq/worker.py

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -110,8 +110,8 @@ def connect(self, force=False):
110110

111111
def ensure_indexes(self):
112112

113-
self.mongodb_logs.mrq_logs.ensure_index([("job", 1)], background=True)
114-
self.mongodb_logs.mrq_logs.ensure_index([("worker", 1)], background=True, sparse=True)
113+
self.mongodb_logs.mrq_logs.ensure_index([("job", 1)], background=False)
114+
self.mongodb_logs.mrq_logs.ensure_index([("worker", 1)], background=False, sparse=True)
115115

116116
if self.config["mongodb_logs_size"] > 0:
117117

@@ -120,14 +120,16 @@ def ensure_indexes(self):
120120
except:
121121
pass
122122

123-
self.mongodb_logs.mrq_workers.ensure_index([("status", 1)], background=True)
124-
self.mongodb_logs.mrq_workers.ensure_index([("datereported", 1)], background=True, expireAfterSeconds=3600)
123+
self.mongodb_logs.mrq_workers.ensure_index([("status", 1)], background=False)
124+
self.mongodb_logs.mrq_workers.ensure_index([("datereported", 1)], background=False, expireAfterSeconds=3600)
125125

126-
self.mongodb_jobs.mrq_jobs.ensure_index([("status", 1)], background=True)
127-
self.mongodb_jobs.mrq_jobs.ensure_index([("path", 1), ("status", 1)], background=True)
128-
self.mongodb_jobs.mrq_jobs.ensure_index([("worker", 1), ("status", 1)], background=True, sparse=True)
129-
self.mongodb_jobs.mrq_jobs.ensure_index([("queue", 1), ("status", 1)], background=True)
130-
self.mongodb_jobs.mrq_jobs.ensure_index([("dateexpires", 1)], sparse=True, background=True, expireAfterSeconds=0)
126+
self.mongodb_jobs.mrq_jobs.ensure_index([("status", 1)], background=False)
127+
self.mongodb_jobs.mrq_jobs.ensure_index([("path", 1), ("status", 1)], background=False)
128+
self.mongodb_jobs.mrq_jobs.ensure_index([("worker", 1), ("status", 1)], background=False, sparse=True)
129+
self.mongodb_jobs.mrq_jobs.ensure_index([("queue", 1), ("status", 1)], background=False)
130+
self.mongodb_jobs.mrq_jobs.ensure_index([("dateexpires", 1)], sparse=True, background=False, expireAfterSeconds=0)
131+
132+
self.mongodb_jobs.mrq_scheduled_jobs.ensure_index([("hash", 1)], unique=True, background=False, drop_dups=True)
131133

132134
try:
133135
# This will be default in MongoDB 2.6

tests/fixtures/config-scheduler3.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,13 @@
11
import datetime
2+
import os
23

34
SCHEDULER_TASKS = [
45
{
56
"path": "mrq.basetasks.tests.general.MongoInsert",
67
"params": {
78
"a": 1
89
},
9-
"dailytime": (datetime.datetime.utcnow() + datetime.timedelta(seconds=3)).time()
10+
"dailytime": datetime.datetime.fromtimestamp(float(os.environ.get("MRQ_TEST_SCHEDULER_TIME"))).time()
1011
},
1112
]
1213

tests/test_scheduler.py

Lines changed: 29 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,21 @@
11
from bson import ObjectId
22
import urllib2
33
import time
4+
import pytest
5+
import datetime
46

57

6-
def test_scheduler_simple(worker):
8+
# We want to test that launching the scheduler several times queues tasks only once.
9+
PROCESS_CONFIGS = [
10+
["--gevent 1"],
11+
["--gevent 1 --processes 5"]
12+
]
713

8-
worker.start(flags="--scheduler --config tests/fixtures/config-scheduler1.py")
14+
15+
@pytest.mark.parametrize(["p_flags"], PROCESS_CONFIGS)
16+
def test_scheduler_simple(worker, p_flags):
17+
18+
worker.start(flags="--scheduler --config tests/fixtures/config-scheduler1.py %s" % p_flags)
919

1020
collection = worker.mongodb_logs.tests_inserts
1121
scheduled_jobs = worker.mongodb_jobs.mrq_scheduled_jobs
@@ -30,7 +40,7 @@ def test_scheduler_simple(worker):
3040
collection.remove({})
3141

3242
# Start with new config
33-
worker.start(deps=False, flags="--scheduler --config tests/fixtures/config-scheduler2.py")
43+
worker.start(deps=False, flags="--scheduler --config tests/fixtures/config-scheduler2.py %s" % p_flags)
3444

3545
time.sleep(2)
3646

@@ -43,14 +53,21 @@ def test_scheduler_simple(worker):
4353
assert len(inserts) == 3
4454

4555

46-
def test_scheduler_dailytime(worker):
56+
@pytest.mark.parametrize(["p_flags"], PROCESS_CONFIGS)
57+
def test_scheduler_dailytime(worker, p_flags):
4758

4859
# Task is scheduled in 3 seconds
49-
worker.start(flags="--scheduler --config tests/fixtures/config-scheduler3.py")
60+
worker.start(
61+
flags="--scheduler --config tests/fixtures/config-scheduler3.py %s" % p_flags,
62+
env={
63+
64+
# We need to pass this in the environment so that each worker has the exact same hash
65+
"MRQ_TEST_SCHEDULER_TIME": str(time.time() + 4)
66+
})
5067

5168
# It will be done a first time immediately
5269

53-
time.sleep(1)
70+
time.sleep(3)
5471

5572
collection = worker.mongodb_logs.tests_inserts
5673

@@ -60,3 +77,9 @@ def test_scheduler_dailytime(worker):
6077
time.sleep(4)
6178

6279
assert collection.find().count() == 2
80+
81+
# Nothing more should happen today
82+
time.sleep(4)
83+
84+
assert collection.find().count() == 2
85+

0 commit comments

Comments
 (0)