Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions mrq/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,13 @@ def add_parser_args(parser, config_type):
type=str,
help='Adds random latency to the network calls, zero to N seconds. Can be a range (1-2)')

parser.add_argument(
'--default_job_ttl',
default=180 * 24 * 3600,
action='store',
type=float,
help='Seconds the tasks are kept in MongoDB when statuses are not success, abort, cancel and started')

parser.add_argument(
'--default_job_result_ttl',
default=7 * 24 * 3600,
Expand Down
14 changes: 11 additions & 3 deletions mrq/job.py
Original file line number Diff line number Diff line change
Expand Up @@ -102,11 +102,10 @@ def fetch(self, start=False, full_data=True):
"path": 1,
"params": 1,
"status": 1,
"retry_count": 1
"retry_count": 1,
}

if start:

self.datestarted = datetime.datetime.utcnow()
self.set_data(self.collection.find_and_modify(
{
Expand All @@ -117,6 +116,9 @@ def fetch(self, start=False, full_data=True):
"status": "started",
"datestarted": self.datestarted,
"worker": self.worker.id
},
"$unset": {
"dateexpires": 1 # we don't want started jobs to expire unexpectedly
}},
projection=fields)
)
Expand Down Expand Up @@ -153,7 +155,8 @@ def set_data(self, data):
task_def = self.get_task_config()

self.timeout = task_def.get("timeout", cfg["default_job_timeout"])
self.result_ttl = task_def.get("result_ttl", cfg["default_job_result_ttl"])
self.default_ttl = task_def.get("default_ttl", cfg["default_job_ttl"])
self.result_ttl = task_def.get("result_ttl", cfg["default_job_result_ttl"]) # success ttl
self.abort_ttl = task_def.get("abort_ttl", cfg["default_job_abort_ttl"])
self.cancel_ttl = task_def.get("cancel_ttl", cfg["default_job_cancel_ttl"])
self.max_retries = task_def.get("max_retries", cfg["default_job_max_retries"])
Expand Down Expand Up @@ -448,6 +451,11 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
"status": status,
"dateupdated": now
}

# we don't want started jobs to expire unexpectedly
if status not in ["started", "success", "abort", "cancel"] and self.default_ttl is not None:
db_updates["dateexpires"] = (self.data.get("datequeued") or now) + datetime.timedelta(days=self.default_ttl)

db_updates.update(updates or {})

if self.datestarted:
Expand Down
5 changes: 4 additions & 1 deletion mrq/queue_regular.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,8 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
"status": "started",
"datestarted": datetime.datetime.utcnow(),
"worker": worker.id if worker else None
}, "$unset": {
"dateexpires": 1 # we don't want started jobs to expire unexpectedly
}},
sort=sort_order,
return_document=ReturnDocument.AFTER,
Expand All @@ -101,7 +103,8 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
"params": 1,
"status": 1,
"retry_count": 1,
"queue": 1
"queue": 1,
"datequeued": 1
}
)

Expand Down
10 changes: 10 additions & 0 deletions tests/fixtures/config-scheduler8.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
NAME = "testworker"

TASKS = {
"tests.tasks.general.Retry": {
"default_ttl": 2,
"queue": "tests"
}
}

QUEUES = ["high", "default", "low", "tests"]
2 changes: 2 additions & 0 deletions tests/test_retry.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@


def test_retry(worker):
worker.start(flags="--config tests/fixtures/config-scheduler8.py")

job_id = worker.send_task(
"tests.tasks.general.Retry", {"queue": "noexec", "delay": 60}, block=False)
Expand All @@ -15,6 +16,7 @@ def test_retry(worker):
assert job_data["queue"] == "noexec"
assert job_data["status"] == "retry"
assert job_data["dateretry"] > datetime.datetime.utcnow()
assert datetime.datetime.utcnow() + datetime.timedelta(days=1) < job_data["dateexpires"] < datetime.datetime.utcnow() + datetime.timedelta(days=3)
assert job_data.get("result") is None


Expand Down