Resume failed episode batch retries
This commit is contained in:
+105
-2
@@ -1288,6 +1288,98 @@ class DownloadQueue:
|
||||
self._save_job_locked(job)
|
||||
return job
|
||||
|
||||
def _episode_key(self, value):
|
||||
text = str(value or "").strip()
|
||||
if not text:
|
||||
return ""
|
||||
try:
|
||||
number = float(text)
|
||||
except ValueError:
|
||||
return text
|
||||
if number.is_integer():
|
||||
return str(int(number))
|
||||
return str(number).rstrip("0").rstrip(".")
|
||||
|
||||
def _episode_spec_values(self, episode_spec):
|
||||
values = []
|
||||
seen = set()
|
||||
for part in re.split(r"\s+", str(episode_spec or "").strip()):
|
||||
if not part:
|
||||
continue
|
||||
if "-" in part:
|
||||
start, end = [piece.strip() for piece in part.split("-", 1)]
|
||||
try:
|
||||
left = int(float(start))
|
||||
right = int(float(end))
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
step = 1 if right >= left else -1
|
||||
for number in range(left, right + step, step):
|
||||
key = self._episode_key(number)
|
||||
if key not in seen:
|
||||
seen.add(key)
|
||||
values.append(str(number))
|
||||
continue
|
||||
key = self._episode_key(part)
|
||||
if key and key not in seen:
|
||||
seen.add(key)
|
||||
values.append(str(part))
|
||||
return values
|
||||
|
||||
def _episode_spec_from_values(self, values):
|
||||
normalized = [str(value).strip() for value in values or [] if str(value).strip()]
|
||||
if not normalized:
|
||||
return ""
|
||||
parts = []
|
||||
index = 0
|
||||
while index < len(normalized):
|
||||
current = normalized[index]
|
||||
try:
|
||||
parsed = float(current)
|
||||
except (TypeError, ValueError):
|
||||
parts.append(current)
|
||||
index += 1
|
||||
continue
|
||||
if not parsed.is_integer() or current != str(int(parsed)):
|
||||
parts.append(current)
|
||||
index += 1
|
||||
continue
|
||||
start = int(parsed)
|
||||
end = start
|
||||
next_index = index + 1
|
||||
while next_index < len(normalized):
|
||||
try:
|
||||
parsed_candidate = float(normalized[next_index])
|
||||
except (TypeError, ValueError):
|
||||
break
|
||||
if not parsed_candidate.is_integer() or normalized[next_index] != str(int(parsed_candidate)):
|
||||
break
|
||||
candidate = int(parsed_candidate)
|
||||
if candidate != end + 1:
|
||||
break
|
||||
end = candidate
|
||||
next_index += 1
|
||||
parts.append(str(start) if end == start else f"{start}-{end}")
|
||||
index = next_index
|
||||
return " ".join(parts)
|
||||
|
||||
def _mark_episode_completed_for_retry(self, job, episode):
|
||||
episode_key = self._episode_key(episode)
|
||||
if not episode_key:
|
||||
return
|
||||
completed = [str(value).strip() for value in job.get("completed_episodes") or [] if str(value).strip()]
|
||||
completed_keys = {self._episode_key(value) for value in completed}
|
||||
if episode_key not in completed_keys:
|
||||
completed.append(str(episode).strip())
|
||||
completed_keys.add(episode_key)
|
||||
pending_values = [
|
||||
value
|
||||
for value in self._episode_spec_values(job.get("episodes"))
|
||||
if self._episode_key(value) not in completed_keys
|
||||
]
|
||||
job["completed_episodes"] = completed
|
||||
job["retry_episodes"] = self._episode_spec_from_values(pending_values)
|
||||
|
||||
def _prepare_job(self, job):
|
||||
if self.job_prepare_fn is None:
|
||||
return job
|
||||
@@ -1341,7 +1433,7 @@ class DownloadQueue:
|
||||
with self.lock, self._connect() as conn:
|
||||
rows = conn.execute(
|
||||
"""
|
||||
SELECT episodes
|
||||
SELECT *
|
||||
FROM jobs
|
||||
WHERE show_id = ? AND mode = ? AND status IN ('pending', 'running', 'done')
|
||||
ORDER BY created_at ASC
|
||||
@@ -1350,7 +1442,17 @@ class DownloadQueue:
|
||||
).fetchall()
|
||||
covered = set()
|
||||
for row in rows:
|
||||
covered.update(self._expand_episode_spec((row["episodes"] if isinstance(row, sqlite3.Row) else row[0]), available_episodes))
|
||||
job = self._job_from_row(row) if isinstance(row, sqlite3.Row) else {"episodes": row[0]}
|
||||
episode_spec = job.get("episodes")
|
||||
if job.get("status") in {"pending", "running"} and "retry_episodes" in job:
|
||||
episode_spec = job.get("retry_episodes")
|
||||
covered.update(
|
||||
value
|
||||
for value in available_episodes
|
||||
if self._episode_key(value)
|
||||
in {self._episode_key(episode) for episode in job.get("completed_episodes") or []}
|
||||
)
|
||||
covered.update(self._expand_episode_spec(episode_spec, available_episodes))
|
||||
return covered
|
||||
|
||||
def detach_show(self, show_id):
|
||||
@@ -1545,6 +1647,7 @@ class DownloadQueue:
|
||||
if moved:
|
||||
finalized.append(episode_key)
|
||||
job.setdefault("_moved_files", []).extend(moved)
|
||||
self._mark_episode_completed_for_retry(job, episode)
|
||||
return moved
|
||||
|
||||
def _cleanup_staging_dir(self, job):
|
||||
|
||||
Reference in New Issue
Block a user