Skip to content

Commit 57800f8

Browse files
committed
Agent fixes
1 parent 5d6cc22 commit 57800f8

7 files changed

Lines changed: 88 additions & 21 deletions

File tree

.gitignore

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ develop-eggs
1717
lib
1818
lib64
1919
__pycache__
20+
.cache
2021

2122
# Installer logs
2223
pip-log.txt

mrq/agent.py

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ def __init__(self, worker_group=None):
1717
self.id = ObjectId()
1818
self.worker_group = worker_group or get_current_config()["worker_group"]
1919
self.pool = ProcessPool()
20+
self.config = get_current_config()
2021

2122
def work(self):
2223

@@ -50,7 +51,7 @@ def greenlet_manage(self):
5051
try:
5152
self.manage()
5253
except Exception as e: # pylint: disable=broad-except
53-
self.log.error("When reporting: %s" % e)
54+
log.error("When reporting: %s" % e)
5455
finally:
5556
time.sleep(self.config["report_interval"])
5657

@@ -59,11 +60,13 @@ def manage(self):
5960
report = self.get_agent_report()
6061

6162
try:
62-
db = self.mongodb_jobs.mrq_agents.find_and_modify({
63+
db = connections.mongodb_jobs.mrq_agents.find_and_modify({
6364
"_id": ObjectId(self.id)
6465
}, {"$set": report}, upsert=True)
66+
if not db:
67+
return
6568
except Exception as e: # pylint: disable=broad-except
66-
self.log.debug("Agent report failed: %s" % e)
69+
log.debug("Agent report failed: %s" % e)
6770
return
6871

6972
# If the desired_workers was changed by an orchestrator, apply the changes locally
@@ -74,16 +77,17 @@ def get_agent_report(self):
7477
report = {
7578
"current_workers": [p["command"] for p in self.pool.processes],
7679
"available_cpu": get_current_config()["available_cpu"],
77-
"available_memory": get_current_config()["available_memory"]
80+
"available_memory": get_current_config()["available_memory"],
81+
"worker_group": self.worker_group
7882
}
7983
return report
8084

8185
def greenlet_orchestrate(self):
8286

8387
while True:
84-
with connections.redis.lock(self.redis_agent_orchestrator_key, timeout=60):
88+
with connections.redis.lock(self.redis_agent_orchestrator_key, timeout=self.config["orchestrate_interval"] + 10):
8589
self.orchestrate()
86-
time.sleep(30)
90+
time.sleep(self.config["orchestrate_interval"])
8791

8892
@property
8993
def redis_agent_orchestrator_key(self):

mrq/bin/mrq_agent.py

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -13,16 +13,18 @@
1313
import sys
1414
import argparse
1515

16-
from .config import get_config
17-
from .agent import Agent
18-
from .context import set_current_config
16+
sys.path.insert(0, os.getcwd())
17+
18+
from mrq import config
19+
from mrq.agent import Agent
20+
from mrq.context import set_current_config
1921

2022

2123
def main():
2224

2325
parser = argparse.ArgumentParser(description='Start a MRQ agent')
2426

25-
cfg = get_config(parser=parser, config_type="agent", sources=("file", "env", "args"))
27+
cfg = config.get_config(parser=parser, config_type="agent", sources=("file", "env", "args"))
2628

2729
set_current_config(cfg)
2830

mrq/config.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -280,6 +280,20 @@ def add_parser_args(parser, config_type):
280280
type=int,
281281
help="How much CPU units this agent's workers can use. We recommend using 1024 per CPU.")
282282

283+
parser.add_argument(
284+
'--orchestrate_interval',
285+
default=30,
286+
action="store",
287+
type=float,
288+
help="How much seconds to wait between orchestration runs.")
289+
290+
parser.add_argument(
291+
'--report_interval',
292+
default=10,
293+
action='store',
294+
type=float,
295+
help='Seconds between agent reports to MongoDB')
296+
283297
# Worker-specific args
284298
elif config_type == "worker":
285299

mrq/worker.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -422,6 +422,8 @@ def work(self):
422422

423423
self.work_loop(max_jobs=self.max_jobs)
424424

425+
self.work_stop()
426+
425427
def work_init(self):
426428

427429
self.connect()

tests/conftest.py

Lines changed: 17 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -136,24 +136,30 @@ def __init__(self, request, **kwargs):
136136

137137
self.started = False
138138

139-
def start(self, flush=True, deps=True, trace=True, **kwargs):
139+
def start(self, flush=True, deps=True, trace=True, agent=False, **kwargs):
140140

141141
self.started = True
142142

143143
if deps:
144144
self.start_deps(flush=flush)
145145

146146
processes = 0
147-
m = re.search(r"--processes (\d+)", kwargs.get("flags", ""))
148-
if m:
149-
processes = int(m.group(1))
150-
151-
cmdline = "python mrq/bin/mrq_worker.py --mongodb_logs_size 0 %s %s %s %s" % (
152-
"--admin_port 20020" if (processes <= 1) else "",
153-
"--trace_io --trace_greenlets" if trace else "",
154-
kwargs.get("flags", ""),
155-
kwargs.get("queues", "high default low")
156-
)
147+
148+
if agent:
149+
cmdline = "python mrq/bin/mrq_agent.py %s" % kwargs.get("flags", "")
150+
151+
else:
152+
153+
m = re.search(r"--processes (\d+)", kwargs.get("flags", ""))
154+
if m:
155+
processes = int(m.group(1))
156+
157+
cmdline = "python mrq/bin/mrq_worker.py --mongodb_logs_size 0 %s %s %s %s" % (
158+
"--admin_port 20020" if (processes <= 1) else "",
159+
"--trace_io --trace_greenlets" if trace else "",
160+
kwargs.get("flags", ""),
161+
kwargs.get("queues", "high default low")
162+
)
157163

158164
print(cmdline)
159165
ProcessFixture.start(self, cmdline=cmdline, env=kwargs.get("env"), expected_children=processes)

tests/test_agent.py

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
from mrq.agent import Agent
22
from mrq.context import connections
3+
import time
34

45

56
def scenario(profiles, agents):
@@ -136,3 +137,40 @@ def test_orchestration_scenarios(worker):
136137
}
137138

138139
worker.stop()
140+
141+
142+
def test_agent_process(worker):
143+
144+
worker.start(agent=True, flags="--worker_group xxx --available_memory=500 --available_cpu=500 --orchestrate_interval=1 --report_interval=1")
145+
146+
time.sleep(3)
147+
148+
agents = list(connections.mongodb_jobs.mrq_agents.find())
149+
150+
assert len(agents) == 1
151+
152+
assert connections.mongodb_jobs.mrq_workers.count() == 0
153+
154+
connections.mongodb_jobs.mrq_workergroups.insert_one({"_id": "xxx", "profiles": [
155+
{
156+
"command": "mrq-worker a",
157+
"memory": 100,
158+
"cpu": 100,
159+
"min_count": 1
160+
}
161+
]})
162+
163+
time.sleep(3)
164+
165+
assert connections.mongodb_jobs.mrq_workers.count() == 1
166+
worker = connections.mongodb_jobs.mrq_workers.find_one()
167+
assert worker["status"] == "wait"
168+
169+
connections.mongodb_jobs.mrq_workergroups.update_one({"_id": "xxx"}, {"$set": {"profiles": [
170+
171+
]}})
172+
173+
time.sleep(4)
174+
175+
worker = connections.mongodb_jobs.mrq_workers.find_one()
176+
assert worker["status"] == "stop"

0 commit comments

Comments
 (0)