From 5718e5fe91bdef092ebd8bf2091e3322c3d80e4f Mon Sep 17 00:00:00 2001 From: Dymas Date: Tue, 1 Sep 2026 13:39:27 +0200 Subject: [PATCH] Resume failed episode batch retries --- CHANGELOG.md | 5 +++ README.md | 2 +- VERSION | 2 +- app_support.py | 2 +- queue_jobs.py | 107 ++++++++++++++++++++++++++++++++++++++++++++++++- test_app.py | 48 ++++++++++++++++++++++ 6 files changed, 161 insertions(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9c0bdbc..0250576 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,10 @@ # Changelog +## 0.52.20 - 2026-09-01 + +- Fixed retrying a failed multi-episode provider download so already finalized episodes are skipped and only the remaining episodes are requested again. +- Preserved the original episode range in queue display and watchlist sync while tracking a narrowed retry episode range internally. + ## 0.52.19 - 2026-09-01 - Changed the provider update check to report newer local provider patches separately from upstream updates, avoiding false update warnings for the patched Anikoto module. diff --git a/README.md b/README.md index e90636e..7d2f562 100644 --- a/README.md +++ b/README.md @@ -100,7 +100,7 @@ Keep `./.kaizoku` mounted for production instances. That directory contains the ## Download Flow -Kaizoku stores provider-backed show IDs as `provider:id`, for example `anikoto:some-show-slug`. Queue jobs resolve the episode source through `providers/bridge.js`, then `provider_downloader.py` downloads the media with `ffmpeg` into a staging directory. If a provider returns a master HLS playlist, Kaizoku selects the highest-bandwidth variant before starting `ffmpeg`. When a media playlist uses extensionless CDN segments, Kaizoku skips direct ffmpeg and uses a StrawVerse-style segment downloader: it fetches the media playlist, downloads and concatenates segments itself, strips short PNG wrappers when present, retries temporary HTTP failures such as `429 Too Many Requests`, then remuxes the local transport stream to MP4. Direct ffmpeg attempts also have a timeout guard so stalled HLS inputs can fall back cleanly. Each episode is written as a temporary `.mp4.part` file and moved into the downloads library after that episode succeeds, so already-finished episodes from a larger batch survive if a later episode fails. Finalization preserves episode numbers from staged `SxxEyy` or `Episode yy` filenames before applying configured season/episode offsets. +Kaizoku stores provider-backed show IDs as `provider:id`, for example `anikoto:some-show-slug`. Queue jobs resolve the episode source through `providers/bridge.js`, then `provider_downloader.py` downloads the media with `ffmpeg` into a staging directory. If a provider returns a master HLS playlist, Kaizoku selects the highest-bandwidth variant before starting `ffmpeg`. When a media playlist uses extensionless CDN segments, Kaizoku skips direct ffmpeg and uses a StrawVerse-style segment downloader: it fetches the media playlist, downloads and concatenates segments itself, strips short PNG wrappers when present, retries temporary HTTP failures such as `429 Too Many Requests`, then remuxes the local transport stream to MP4. Direct ffmpeg attempts also have a timeout guard so stalled HLS inputs can fall back cleanly. Each episode is written as a temporary `.mp4.part` file and moved into the downloads library after that episode succeeds, so already-finished episodes from a larger batch survive if a later episode fails. Retrying that failed queue job requests only the remaining episodes while keeping the original episode range for display and watchlist sync. Finalization preserves episode numbers from staged `SxxEyy` or `Episode yy` filenames before applying configured season/episode offsets. If the primary provider cannot list, resolve, or download a requested episode, Kaizoku searches the same title on the remaining providers and tries the matching episode there. Existing finalization code moves staged files into the configured library layout, preserving data already present in production download folders. diff --git a/VERSION b/VERSION index f0f4383..e3ff91f 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.52.19 +0.52.20 diff --git a/app_support.py b/app_support.py index ccdab65..bbab2ec 100644 --- a/app_support.py +++ b/app_support.py @@ -1201,7 +1201,7 @@ def command_for_job(job, backend="kaizoku", download_path=None): "--title", str(job.get("anime_name") or job.get("title") or job.get("query") or "Anime"), "--episodes", - str(job.get("episodes") or "1"), + str(job.get("retry_episodes") or job.get("episodes") or "1"), "--mode", str(job.get("mode") or "sub"), "--quality", diff --git a/queue_jobs.py b/queue_jobs.py index 524bafd..952959c 100644 --- a/queue_jobs.py +++ b/queue_jobs.py @@ -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): diff --git a/test_app.py b/test_app.py index 8ba3a2b..517e662 100644 --- a/test_app.py +++ b/test_app.py @@ -613,6 +613,30 @@ class QueueApiTests(unittest.TestCase): self.assertEqual(retried["status"], "pending") self.assertEqual(retried["result_index"], 4) + def test_covered_episodes_includes_completed_and_remaining_retry_parts(self): + created = self.queue.add( + { + "show_id": "show-42", + "query": "Retry Coverage", + "title": "Retry Coverage", + "anime_name": "Retry Coverage", + "mode": "sub", + "quality": "best", + "episodes": "1-3", + "download_dir": "/tmp/example", + "season": "1", + } + ) + with self.queue.lock: + job = self.queue._find(created["id"]) + job["completed_episodes"] = ["1"] + job["retry_episodes"] = "2-3" + self.queue._save_job_locked(job) + + covered = self.queue.covered_episodes("show-42", "sub", ["1", "2", "3", "4"]) + + self.assertEqual(covered, {"1", "2", "3"}) + def test_add_reuses_matching_failed_job_instead_of_creating_duplicate(self): payload = { "show_id": "show-42", @@ -1099,6 +1123,16 @@ class DownloadQueueWorkerFailureTests(unittest.TestCase): self.assertFalse(staged_episode.exists()) self.assertFalse(staging_dir.exists()) self.assertIn(f"Saved: {final_episode}", stored["log"]) + self.assertEqual(stored["completed_episodes"], ["1"]) + self.assertEqual(stored["retry_episodes"], "2") + + retried = queue.retry(job["id"]) + command = APP.app_support.command_for_job(retried) + + self.assertEqual(retried["status"], "pending") + self.assertEqual(retried["episodes"], "1-2") + self.assertEqual(retried["retry_episodes"], "2") + self.assertEqual(command[command.index("--episodes") + 1], "2") def test_successful_detached_download_logs_sync_skipped(self): queue = APP.DownloadQueue( @@ -2158,6 +2192,20 @@ class WatchlistCompletionTests(unittest.TestCase): self.assertNotIn("-S", command) + def test_command_for_job_uses_retry_episode_spec_when_present(self): + command = APP.app_support.command_for_job( + { + "query": "Queue Show", + "quality": "best", + "episodes": "1-12", + "retry_episodes": "3-12", + "mode": "dub", + }, + download_path="/tmp/kaizoku-staging", + ) + + self.assertEqual(command[command.index("--episodes") + 1], "3-12") + def test_download_watchlist_item_queues_full_episode_range_for_selected_mode(self): item = { "show_id": "show-42",