Skip to content

Commit 8a75885

Browse files
Introducing ttl for aborted jobs (default is 24h), configurable through mrq-config
1 parent 6132345 commit 8a75885

4 files changed

Lines changed: 26 additions & 5 deletions

File tree

mrq/config.py

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -156,7 +156,14 @@ def add_parser_args(parser, config_type):
156156
default=7 * 24 * 3600,
157157
action='store',
158158
type=int,
159-
help='Seconds the results are kept in MongoDB when status in (success, cancel, abort)')
159+
help='Seconds the results are kept in MongoDB when status in (success)')
160+
161+
parser.add_argument(
162+
'--default_job_aborted_or_canceled_ttl',
163+
default=24 * 3600,
164+
action='store',
165+
type=int,
166+
help='Seconds the tasks are kept in MongoDB when status in (cancel, abort)')
160167

161168
parser.add_argument(
162169
'--default_job_timeout',

mrq/job.py

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ class Job(object):
2323

2424
timeout = None
2525
result_ttl = None
26+
aborted_or_canceled_ttl = None
2627
max_retries = None
2728
retry_delay = None
2829

@@ -127,6 +128,7 @@ def set_data(self, data):
127128

128129
self.timeout = task_def.get("timeout", cfg["default_job_timeout"])
129130
self.result_ttl = task_def.get("result_ttl", cfg["default_job_result_ttl"])
131+
self.aborted_or_canceled_ttl = task_def.get("aborted_or_canceled_ttl", cfg["default_job_aborted_or_canceled_ttl"])
130132
self.max_retries = task_def.get("max_retries", cfg["default_job_max_retries"])
131133
self.retry_delay = task_def.get("retry_delay", cfg["default_job_retry_delay"])
132134

@@ -340,13 +342,21 @@ def save_success(self, result=None):
340342

341343
def save_cancel(self):
342344

343-
dateexpires = datetime.datetime.utcnow() + datetime.timedelta(seconds=self.result_ttl)
345+
dateexpires = datetime.datetime.utcnow() + datetime.timedelta(seconds=self.aborted_or_canceled_ttl)
344346
updates = {
345347
"dateexpires": dateexpires
346348
}
347349

348350
self._save_status("cancel", updates)
349351

352+
def save_abort(self):
353+
dateexpires = datetime.datetime.utcnow() + datetime.timedelta(seconds=self.aborted_or_canceled_ttl)
354+
updates = {
355+
"dateexpires": dateexpires
356+
}
357+
358+
self._save_status("abort", updates)
359+
350360
def _save_status(self, status, updates=None, exception=False, w=1):
351361

352362
if self.id is None:

mrq/worker.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -538,7 +538,7 @@ def perform_job(self, job):
538538

539539
except AbortInterrupt:
540540
self.log.error("Caught abort")
541-
job._save_status("abort", exception=True)
541+
job.save_abort()
542542

543543
except TimeoutInterrupt:
544544
self.log.error("Job timeouted after %s seconds" % job.timeout)

tests/test_abort.py

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
from mrq.job import Job
22
from mrq.queue import Queue
3-
import time
3+
from datetime import datetime
4+
from datetime import timedelta
45

56

67
def test_abort(worker):
@@ -14,4 +15,7 @@ def test_abort(worker):
1415
db_jobs = list(worker.mongodb_jobs.mrq_jobs.find())
1516
assert len(db_jobs) == 1
1617

17-
assert db_jobs[0]["status"] == "abort"
18+
job = db_jobs[0]
19+
assert job["status"] == "abort"
20+
assert job.get("dateexpires") is not None
21+
assert job["dateexpires"] < datetime.utcnow() + timedelta(hours=24)

0 commit comments

Comments
 (0)