Prepare retried jobs before requeue

This commit is contained in:
Codex
2026-08-01 17:51:59 +02:00
parent edb565d855
commit f417a96811
6 changed files with 153 additions and 19 deletions
+16 -16
View File
@@ -1258,12 +1258,15 @@ class DownloadQueue:
self._save_job_locked(job)
return job
def _prepare_job(self, job):
if self.job_prepare_fn is None:
return job
prepared = self.job_prepare_fn(dict(job))
return prepared if prepared is not None else job
def add(self, payload):
job = build_job(payload, self.config_getter())
if self.job_prepare_fn is not None:
prepared = self.job_prepare_fn(dict(job))
if prepared is not None:
job = prepared
job = self._prepare_job(job)
with self.lock:
reusable = self._find_reusable_failed_job_locked(job)
if reusable is not None:
@@ -1347,24 +1350,21 @@ class DownloadQueue:
raise ValueError("Running jobs cannot be retried")
if job["status"] not in {"failed", "canceled"}:
raise ValueError("Only failed or canceled jobs can be retried")
job = self._prepare_job(job)
with self.lock:
job = self._reset_retryable_job_locked(job)
self.wakeup.set()
return job
def retry_all_failed(self):
count = 0
with self.lock, self._connect() as conn:
rows = conn.execute("SELECT * FROM jobs WHERE status = 'failed'").fetchall()
for row in rows:
job = self._job_from_row(row)
job["status"] = "pending"
job["exit_code"] = None
job["pid"] = None
job["started_at"] = None
job["finished_at"] = None
job["updated_at"] = now_iso()
job["log"] = []
self._upsert_job_conn(conn, job)
with self.lock:
with self._connect() as conn:
rows = conn.execute("SELECT * FROM jobs WHERE status = 'failed'").fetchall()
for row in rows:
job = self._prepare_job(self._job_from_row(row))
with self.lock:
self._reset_retryable_job_locked(job)
count += 1
if count:
self.wakeup.set()