Skip to content

Commit ea9c9a2

Browse files
committed
New method of checking timeouted jobs, add maxconcurrency status in dashboard + test fixes
1 parent 9c4f913 commit ea9c9a2

10 files changed

Lines changed: 80 additions & 31 deletions

File tree

Makefile

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,7 @@ pep8:
6767
autopep8:
6868
autopep8 --max-line-length 99 -aaaaaaaa --in-place --recursive mrq
6969

70-
pypi: linterrors
70+
pypi: linterrors linterrors3
7171
python setup.py sdist upload
7272

7373
build_docs:

mrq/dashboard/static/js/views/jobs.js

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -296,6 +296,7 @@ define(["jquery", "underscore", "views/generic/datatablepage", "models"],functio
296296
'timeout': "label-danger",
297297
'failed': "label-danger",
298298
'maxretries': "label-danger",
299+
'maxconcurrency': "label-danger",
299300
'interrupt': "label-danger",
300301
'cancel': "label-warning",
301302
'abort': "label-warning",

mrq/dashboard/templates/index.html

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -397,13 +397,16 @@ <h4 class="modal-title"></h4>
397397
"failed": "failed",
398398
"retry": "retry",
399399
"maxretries": "maxretries",
400+
"maxconcurrency": "maxconcurrency",
400401
"timeout": "timeout",
401402
"interrupt": "interrupt",
402403
"cancel": "cancel",
403404
"abort": "abort",
404405
"queued-started-success": "{OK}",
405-
"failed-retry-maxretries-timeout-interrupt-cancel-abort": "{NOT OK}",
406-
"failed-maxretries-timeout-abort": "{ERROR}"
406+
"failed-retry-maxretries-maxconcurrency-timeout-interrupt-cancel-abort": "{NOT OK}",
407+
"failed-maxretries-maxconcurrency-timeout": "{ERROR}",
408+
"failed-retry-maxretries-maxconcurrency-interrupt-timeout": "{ERROR+RETRY}",
409+
407410
}, function(k, i) { %>
408411
<option <%= filters.status==i?"selected='selected'":"" %> value="<%= i %>"><%= k %></option>
409412
<% }) %>

mrq/job.py

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,10 @@
2323
from . import context
2424

2525

26+
FINAL_STATUSES = {"timeout", "abort", "failed", "success", "interrupt", "retry", "maxretries", "maxconcurrency"}
27+
TRANSIENT_STATUSES = {"cancel", "queued", "started"}
28+
29+
2630
class Job(object):
2731

2832
timeout = None
@@ -438,6 +442,11 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
438442
if self.id is None:
439443
return
440444

445+
# Forbid some status transitions
446+
if self.data and self.data.get("status") in FINAL_STATUSES and status not in TRANSIENT_STATUSES:
447+
context.log.error("Can't go from status %s to %s" % (self.data["status"], status))
448+
return
449+
441450
context.metric("jobs.status.%s" % status)
442451

443452
if self.stored is False and self.statuses_no_storage is not None and status in self.statuses_no_storage:
@@ -476,7 +485,8 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
476485
db_updates["traceback"] = trace
477486
db_updates["exceptiontype"] = exc.__name__
478487

479-
self._save_traceback_history(status, trace, exc)
488+
if self.data:
489+
self.data.update(db_updates)
480490

481491
# In the most common case, we allow an optimization on Mongo writes
482492
if status == "success":
@@ -500,8 +510,9 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
500510
"_id": self.id
501511
}, {"$set": db_updates}, w=w, j=j, manipulate=False)
502512

503-
if self.data:
504-
self.data.update(db_updates)
513+
if exception:
514+
self._save_traceback_history(status, trace, exc)
515+
505516

506517
def set_current_io(self, io_data):
507518

mrq/worker.py

Lines changed: 24 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -356,6 +356,28 @@ def report_worker(self, w=0):
356356
except Exception as e: # pylint: disable=broad-except
357357
self.log.debug("Worker report failed: %s" % e)
358358

359+
def greenlet_timeouts(self):
360+
""" This greenlet kills jobs in other greenlets if they timeout.
361+
"""
362+
363+
while True:
364+
now = datetime.datetime.utcnow()
365+
for greenlet in list(self.gevent_pool):
366+
job = get_current_job(id(greenlet))
367+
if job and job.timeout and job.datestarted:
368+
expires = job.datestarted + datetime.timedelta(seconds=job.timeout)
369+
if now > expires:
370+
greenlet.kill(block=False)
371+
if job.data["status"] != "timeout":
372+
updates = {
373+
"exceptiontype": "TimeoutInterrupt",
374+
"traceback": "".join(traceback.format_stack(greenlet.gr_frame))
375+
}
376+
job._save_status("timeout", updates=updates, exception=False)
377+
378+
time.sleep(1)
379+
380+
359381
def greenlet_admin(self):
360382
""" This greenlet is used to get status information about the worker
361383
when --admin_port was given
@@ -456,6 +478,8 @@ def work_init(self):
456478
if self.config["admin_port"]:
457479
self.greenlets["admin"] = gevent.spawn(self.greenlet_admin)
458480

481+
self.greenlets["timeouts"] = gevent.spawn(self.greenlet_timeouts)
482+
459483
if self.config["scheduler"] and self.config["scheduler_interval"] > 0:
460484

461485
from .scheduler import Scheduler
@@ -645,19 +669,6 @@ def perform_job(self, job):
645669

646670
set_current_job(job)
647671

648-
gevent_timeout = None
649-
if job.timeout:
650-
651-
gevent_timeout = gevent.Timeout(
652-
job.timeout,
653-
TimeoutInterrupt(
654-
'Job exceeded maximum timeout value in greenlet (%d seconds).' %
655-
job.timeout
656-
)
657-
)
658-
659-
gevent_timeout.start()
660-
661672
try:
662673
job.perform()
663674

@@ -691,9 +702,6 @@ def perform_job(self, job):
691702

692703
finally:
693704

694-
if gevent_timeout:
695-
gevent_timeout.cancel()
696-
697705
set_current_job(None)
698706

699707
self.done_jobs += 1

tests/fixtures/config-scheduler6.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,9 @@
55
{
66
"path": "tests.tasks.general.MongoInsert",
77
"params": {
8-
"monthday": i
8+
"monthday": i + 1
99
},
10-
"monthday": i,
10+
"monthday": i + 1,
1111
"dailytime": datetime.datetime.fromtimestamp(float(os.environ.get("MRQ_TEST_SCHEDULER_TIME"))).time()
1212
} for i in range(31)
1313
]

tests/tasks/general.py

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,16 @@
1717
class Add(Task):
1818

1919
def run(self, params):
20+
if params.get("broadexcept"):
21+
try:
22+
return self._add(params)
23+
except BaseException, e: # Will catch greenlet.GreenletExit & all others
24+
print("Got base exception %s" % type(e))
25+
return
26+
else:
27+
return self._add(params)
28+
29+
def _add(self, params):
2030
log.info("adding", params)
2131
res = params.get("a", 0) + params.get("b", 0)
2232

@@ -296,4 +306,14 @@ def run(self, params):
296306

297307
class QueueAllKnown(Task):
298308
def run(self, params):
299-
return list(Queue.all_known())
309+
return list(Queue.all_known())
310+
311+
312+
class Uninterruptable(Task):
313+
314+
def run(self, params):
315+
while True:
316+
try:
317+
time.sleep(1)
318+
except:
319+
pass

tests/test_general.py

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -189,12 +189,14 @@ def test_known_queues_lifecycle(worker):
189189
time.sleep(1)
190190

191191
all_known_plus_sub = worker.send_task("tests.tasks.general.QueueAllKnown", {}, queue="default")
192-
assert set(all_known_plus_sub).difference(set(all_known)) == set(["test_raw/sub"])
192+
assert set(all_known_plus_sub) == set(all_known).union(set(["test_raw/sub"]))
193193

194-
Queue("test_raw/sub").remove_raw_jobs(["a", "b", "c"])
194+
# This behavious was removed in https://github.com/pricingassistant/mrq/commit/dcb7c954c998d0d8f32b799da0fb0aa11524e5b9
195+
# We might restore it if we find a way to improve performance
196+
# Queue("test_raw/sub").remove_raw_jobs(["a", "b", "c"])
195197

196-
all_known_plus_sub = worker.send_task("tests.tasks.general.QueueAllKnown", {}, queue="default")
197-
assert set(all_known_plus_sub).difference(set(all_known)) == set()
198+
# all_known_plus_sub = worker.send_task("tests.tasks.general.QueueAllKnown", {}, queue="default")
199+
# assert set(all_known_plus_sub) == set(all_known)
198200

199201

200202
def test_general_exception_status(worker):

tests/test_scheduler.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -166,7 +166,7 @@ def test_scheduler_weekday_dailytime(worker):
166166

167167

168168
def test_scheduler_monthday(worker):
169-
# Task is scheduled in 3 seconds
169+
# Task is scheduled in 10 seconds
170170
worker.start(
171171
flags="--scheduler --config tests/fixtures/config-scheduler6.py",
172172
env={
@@ -183,7 +183,7 @@ def test_scheduler_monthday(worker):
183183
assert len(inserts) == 0
184184

185185
# And then only once
186-
time.sleep(6)
186+
time.sleep(10)
187187
inserts = list(collection.find())
188188
assert len(inserts) == 1
189189
assert collection.find({"params.monthday": datetime.datetime.utcnow().day}).count() == 1

tests/test_timeout.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,10 @@ def test_timeout_normal(worker):
1212
"a": 1, "b": 2, "sleep": 1000}, block=True, accept_statuses=["timeout"])
1313
assert r != 3
1414

15+
r = worker.send_task("tests.tasks.general.TimeoutFromConfig", {
16+
"a": 1, "b": 2, "sleep": 1000, "broadexcept": True}, block=True, accept_statuses=["timeout"])
17+
assert r != 3
18+
1519

1620
def test_timeout_global_config(worker):
1721

0 commit comments

Comments
 (0)