Skip to content

Commit 7c1e229

Browse files
committed
Fix agent shutdown
1 parent 639f950 commit 7c1e229

4 files changed

Lines changed: 16 additions & 11 deletions

File tree

mrq/agent.py

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -31,9 +31,11 @@ def work(self):
3131

3232
self.pool.start()
3333

34-
self.pool.wait()
35-
36-
connections.mongodb_jobs.mrq_agents.delete_one({"_id": self.id})
34+
try:
35+
self.pool.wait()
36+
finally:
37+
self.shutdown_now()
38+
connections.mongodb_jobs.mrq_agents.delete_one({"_id": self.id})
3739

3840
def shutdown_now(self):
3941
self.pool.terminate()
@@ -81,7 +83,8 @@ def get_agent_report(self):
8183
"available_cpu": get_current_config()["available_cpu"],
8284
"available_memory": get_current_config()["available_memory"],
8385
"worker_group": self.worker_group,
84-
86+
"datereported": datetime.datetime.utcnow(),
87+
"dateexpires": datetime.datetime.utcnow() + datetime.timedelta(seconds=(self.config["report_interval"] * 3) + 5)
8588
}
8689
return report
8790

@@ -184,8 +187,7 @@ def orchestrate(self):
184187
connections.mongodb_jobs.mrq_agents.update_one({"_id": agent["_id"]}, {"$set": {
185188
"desired_workers": agent["new_desired_workers"],
186189
"free_cpu": agent["free_cpu"],
187-
"free_memory": agent["free_memory"],
188-
"datereported": datetime.datetime.utcnow()
190+
"free_memory": agent["free_memory"]
189191
}})
190192

191193
def get_desired_workers_for_group(self, group):
@@ -206,4 +208,3 @@ def fetch_worker_group_agents(self):
206208

207209
def fetch_worker_group_definition(self):
208210
return connections.mongodb_jobs.mrq_workergroups.find_one({"_id": self.worker_group})
209-

mrq/worker.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -151,7 +151,9 @@ def ensure_indexes(self):
151151
[("hash", 1)], unique=True, background=False)
152152

153153
self.mongodb_jobs.mrq_agents.ensure_index(
154-
[("datereported", 1)], background=False, expireAfterSeconds=300)
154+
[("datereported", 1)], background=False)
155+
self.mongodb_jobs.mrq_agents.ensure_index(
156+
[("dateexpires", 1)], background=False, expireAfterSeconds=0)
155157
self.mongodb_jobs.mrq_agents.ensure_index(
156158
[("worker_group", 1)], background=False)
157159

tests/conftest.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,7 @@ def stop(self, force=False, timeout=None, block=True, sig=15):
9797

9898
if self.process is not None:
9999

100-
print("Test sending signal %s to %s" % (sig, self.process.pid))
100+
print("Sending signal %s to pid %s" % (sig, self.process.pid))
101101
os.kill(self.process.pid, sig)
102102

103103
# When sending a sigkill to the process, we also want to kill the

tests/test_agent.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -153,14 +153,14 @@ def test_agent_process(worker):
153153

154154
connections.mongodb_jobs.mrq_workergroups.insert_one({"_id": "xxx", "profiles": [
155155
{
156-
"command": "mrq-worker a",
156+
"command": "mrq-worker a --report_interval=1",
157157
"memory": 100,
158158
"cpu": 100,
159159
"min_count": 1
160160
}
161161
]})
162162

163-
time.sleep(3)
163+
time.sleep(7)
164164

165165
assert connections.mongodb_jobs.mrq_workers.count() == 1
166166
w = connections.mongodb_jobs.mrq_workers.find_one()
@@ -179,6 +179,8 @@ def test_agent_process(worker):
179179

180180
worker.stop(deps=False)
181181

182+
time.sleep(2)
183+
182184
assert connections.mongodb_jobs.mrq_agents.count() == 0
183185

184186
worker.stop_deps()

0 commit comments

Comments
 (0)