diff --git a/mrq/config.py b/mrq/config.py index c0401852..fc7f6444 100644 --- a/mrq/config.py +++ b/mrq/config.py @@ -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, diff --git a/mrq/job.py b/mrq/job.py index c25d6da7..849a5149 100644 --- a/mrq/job.py +++ b/mrq/job.py @@ -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( { @@ -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) ) @@ -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"]) @@ -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: diff --git a/mrq/queue_regular.py b/mrq/queue_regular.py index a0e4a3f6..5fe6b085 100644 --- a/mrq/queue_regular.py +++ b/mrq/queue_regular.py @@ -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, @@ -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 } ) diff --git a/tests/fixtures/config-scheduler8.py b/tests/fixtures/config-scheduler8.py new file mode 100644 index 00000000..acb28334 --- /dev/null +++ b/tests/fixtures/config-scheduler8.py @@ -0,0 +1,10 @@ +NAME = "testworker" + +TASKS = { + "tests.tasks.general.Retry": { + "default_ttl": 2, + "queue": "tests" + } +} + +QUEUES = ["high", "default", "low", "tests"] diff --git a/tests/test_retry.py b/tests/test_retry.py index ddbbe686..b945c8b7 100644 --- a/tests/test_retry.py +++ b/tests/test_retry.py @@ -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) @@ -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