Add queue download progress

This commit is contained in:
Codex
2026-08-09 12:55:06 +02:00
parent 499271c3b8
commit 11f43bba05
7 changed files with 271 additions and 14 deletions
+6
View File
@@ -1,5 +1,11 @@
# Changelog # Changelog
## 0.51.8 - 2026-08-09
- Added structured Kaizoku downloader progress events for episode starts, provider resolution, segment downloads, remuxing, and completion.
- Added per-job Queue progress bars with percent, episode position, and segment counts for active downloads.
- Made downloader progress output flush immediately so the Queue page updates while segments are downloading.
## 0.51.7 - 2026-08-09 ## 0.51.7 - 2026-08-09
- Changed provider HLS downloads to preflight media playlists and skip direct ffmpeg when extensionless CDN segment URLs are detected. - Changed provider HLS downloads to preflight media playlists and skip direct ffmpeg when extensionless CDN segment URLs are detected.
+1
View File
@@ -8,6 +8,7 @@ Kaizoku is a local web app for searching, tracking, and downloading anime from A
- Choose the active search provider directly from the Search page. - Choose the active search provider directly from the Search page.
- Prefer provider-supplied search artwork, with a local title-based thumbnail fallback when provider artwork is missing or broken. - Prefer provider-supplied search artwork, with a local title-based thumbnail fallback when provider artwork is missing or broken.
- Queue single episodes or batches in subbed or dubbed mode. - Queue single episodes or batches in subbed or dubbed mode.
- Monitor active Queue jobs with per-episode progress, segment counts, and live downloader output.
- Open result actions in a floating window for watchlist, media type, library name, season, episode range, and download folder choices. - Open result actions in a floating window for watchlist, media type, library name, season, episode range, and download folder choices.
- Fall back across the other configured providers when an episode or stream cannot be resolved on the primary provider. - Fall back across the other configured providers when an episode or stream cannot be resolved on the primary provider.
- Save files with Jellyfin-friendly layout: `TV/Series Name/Season 01/Series Name - S01E01.mp4`. - Save files with Jellyfin-friendly layout: `TV/Series Name/Season 01/Series Name - S01E01.mp4`.
+1 -1
View File
@@ -1 +1 @@
0.51.7 0.51.8
+119 -13
View File
@@ -63,6 +63,11 @@ def bridge(command, provider, *args):
raise RuntimeError(f"Provider bridge returned invalid JSON: {exc}") from exc raise RuntimeError(f"Provider bridge returned invalid JSON: {exc}") from exc
def emit_progress(**payload):
clean = {key: value for key, value in payload.items() if value is not None}
print(f"KAIZOKU_PROGRESS {json.dumps(clean, sort_keys=True)}", flush=True)
def provider_order(primary): def provider_order(primary):
ordered = [] ordered = []
primary = str(primary or "").strip().lower() primary = str(primary or "").strip().lower()
@@ -361,7 +366,7 @@ def decrypt_aes128_segment(data, key, iv_value):
return proc.stdout return proc.stdout
def download_hls_segments(stream, input_url, target, partial): def download_hls_segments(stream, input_url, target, partial, episode_number=None, episode_index=None, episode_total=None):
input_url, playlist = hls_media_playlist(stream, input_url) input_url, playlist = hls_media_playlist(stream, input_url)
segments = parse_hls_segments(input_url, playlist) segments = parse_hls_segments(input_url, playlist)
if not segments: if not segments:
@@ -373,7 +378,19 @@ def download_hls_segments(stream, input_url, target, partial):
segment_dir.mkdir(parents=True, exist_ok=True) segment_dir.mkdir(parents=True, exist_ok=True)
key_cache = {} key_cache = {}
try: try:
print(f"Downloading {len(segments)} HLS playlist segments directly...") message = f"Downloading {len(segments)} HLS playlist segments directly..."
print(message, flush=True)
emit_progress(
phase="segments",
message=message,
episode=episode_number,
episode_index=episode_index,
episode_total=episode_total,
segment=0,
segment_total=len(segments),
percent=0,
)
last_percent = -1
with ts_file.open("wb") as joined: with ts_file.open("wb") as joined:
for index, segment in enumerate(segments, start=1): for index, segment in enumerate(segments, start=1):
data = fetch_bytes(segment["url"], headers=stream.get("headers") or {}, timeout=60) data = fetch_bytes(segment["url"], headers=stream.get("headers") or {}, timeout=60)
@@ -386,8 +403,31 @@ def download_hls_segments(stream, input_url, target, partial):
if not data: if not data:
raise RuntimeError(f"Segment {index} was empty.") raise RuntimeError(f"Segment {index} was empty.")
joined.write(data) joined.write(data)
if index == 1 or index == len(segments) or index % 25 == 0: percent = int(index * 100 / len(segments))
print(f"Downloaded segment {index}/{len(segments)}") if index == 1 or index == len(segments) or percent >= last_percent + 2:
last_percent = percent
message = f"Episode {episode_number}: downloaded segment {index}/{len(segments)} ({percent}%)."
print(message, flush=True)
emit_progress(
phase="segments",
message=message,
episode=episode_number,
episode_index=episode_index,
episode_total=episode_total,
segment=index,
segment_total=len(segments),
percent=percent,
)
emit_progress(
phase="remux",
message=f"Episode {episode_number}: remuxing downloaded segments.",
episode=episode_number,
episode_index=episode_index,
episode_total=episode_total,
segment=len(segments),
segment_total=len(segments),
percent=99,
)
cmd = [ cmd = [
"ffmpeg", "ffmpeg",
"-hide_banner", "-hide_banner",
@@ -446,7 +486,7 @@ def download_subtitles(subtitles, output_base):
return saved return saved
def download_episode(stream, target): def download_episode(stream, target, episode_number=None, episode_index=None, episode_total=None):
input_url = resolve_hls_input_url(stream) input_url = resolve_hls_input_url(stream)
partial = target.with_suffix(target.suffix + ".part") partial = target.with_suffix(target.suffix + ".part")
for path in (target, partial): for path in (target, partial):
@@ -459,8 +499,25 @@ def download_episode(stream, target):
media_url, playlist = hls_media_playlist(stream, input_url) media_url, playlist = hls_media_playlist(stream, input_url)
segments = parse_hls_segments(media_url, playlist) segments = parse_hls_segments(media_url, playlist)
if hls_segments_need_native_download(segments): if hls_segments_need_native_download(segments):
print("HLS playlist uses extensionless provider segments; skipping direct ffmpeg.") message = "HLS playlist uses extensionless provider segments; skipping direct ffmpeg."
code = download_hls_segments(stream, media_url, target, partial) print(message, flush=True)
emit_progress(
phase="preflight",
message=message,
episode=episode_number,
episode_index=episode_index,
episode_total=episode_total,
percent=0,
)
code = download_hls_segments(
stream,
media_url,
target,
partial,
episode_number=episode_number,
episode_index=episode_index,
episode_total=episode_total,
)
if code == 0: if code == 0:
return 0 return 0
return code return code
@@ -510,7 +567,15 @@ def download_episode(stream, target):
return 0 return 0
if stream.get("isM3U8") or ".m3u8" in str(input_url): if stream.get("isM3U8") or ".m3u8" in str(input_url):
try: try:
code = download_hls_segments(stream, input_url, target, partial) code = download_hls_segments(
stream,
input_url,
target,
partial,
episode_number=episode_number,
episode_index=episode_index,
episode_total=episode_total,
)
if code == 0: if code == 0:
return 0 return 0
except Exception as exc: except Exception as exc:
@@ -567,7 +632,7 @@ def main():
print(f"Title: {series_title}") print(f"Title: {series_title}")
print(f"Episodes: {', '.join(str(ep.get('number')) for ep in wanted)}") print(f"Episodes: {', '.join(str(ep.get('number')) for ep in wanted)}")
exit_code = 0 exit_code = 0
for ep in wanted: for episode_index, ep in enumerate(wanted, start=1):
number = ep.get("number") number = ep.get("number")
ep_id = ep.get("id") ep_id = ep.get("id")
padded = f"{int(float(number)):02d}" if str(number).replace(".", "", 1).isdigit() else str(number) padded = f"{int(float(number)):02d}" if str(number).replace(".", "", 1).isdigit() else str(number)
@@ -575,13 +640,46 @@ def main():
target = output_dir / f"{basename}.mp4" target = output_dir / f"{basename}.mp4"
last_error = None last_error = None
downloaded = False downloaded = False
emit_progress(
phase="episode",
message=f"Starting episode {number} ({episode_index}/{len(wanted)}).",
episode=number,
episode_index=episode_index,
episode_total=len(wanted),
percent=0,
)
for active_provider, _active_show_id, active_ep in provider_episode_candidates(provider, provider_show_id, args.title or info.get("title"), ep): for active_provider, _active_show_id, active_ep in provider_episode_candidates(provider, provider_show_id, args.title or info.get("title"), ep):
try: try:
active_ep_id = active_ep.get("id") active_ep_id = active_ep.get("id")
print(f"Resolving episode {number} on {active_provider} ({args.mode}, {args.quality})...") emit_progress(
phase="resolve",
message=f"Resolving episode {number} on {active_provider} ({args.mode}, {args.quality}).",
episode=number,
episode_index=episode_index,
episode_total=len(wanted),
provider=active_provider,
percent=0,
)
print(f"Resolving episode {number} on {active_provider} ({args.mode}, {args.quality})...", flush=True)
stream = bridge("resolve", active_provider, active_ep_id, args.mode, args.quality) stream = bridge("resolve", active_provider, active_ep_id, args.mode, args.quality)
print(f"Downloading episode {number} from {active_provider} {stream.get('quality') or 'auto'}...") emit_progress(
code = download_episode(stream, target) phase="download",
message=f"Downloading episode {number} from {active_provider} {stream.get('quality') or 'auto'}.",
episode=number,
episode_index=episode_index,
episode_total=len(wanted),
provider=active_provider,
quality=stream.get("quality") or "auto",
percent=0,
)
print(f"Downloading episode {number} from {active_provider} {stream.get('quality') or 'auto'}...", flush=True)
code = download_episode(
stream,
target,
episode_number=number,
episode_index=episode_index,
episode_total=len(wanted),
)
if code != 0: if code != 0:
last_error = RuntimeError(f"ffmpeg failed on {active_provider} with exit code {code}") last_error = RuntimeError(f"ffmpeg failed on {active_provider} with exit code {code}")
print(str(last_error), file=sys.stderr) print(str(last_error), file=sys.stderr)
@@ -591,7 +689,15 @@ def main():
print(f"Saved subtitle: {subtitle_path.name}") print(f"Saved subtitle: {subtitle_path.name}")
except Exception as exc: except Exception as exc:
print(f"Subtitle download skipped: {exc}") print(f"Subtitle download skipped: {exc}")
print(f"Saved: {target.name}") emit_progress(
phase="done",
message=f"Episode {number} saved.",
episode=number,
episode_index=episode_index,
episode_total=len(wanted),
percent=100,
)
print(f"Saved: {target.name}", flush=True)
downloaded = True downloaded = True
break break
except Exception as exc: except Exception as exc:
+45
View File
@@ -1254,6 +1254,7 @@ class DownloadQueue:
job["finished_at"] = None job["finished_at"] = None
job["updated_at"] = now_iso() job["updated_at"] = now_iso()
job["cancel_requested"] = False job["cancel_requested"] = False
job.pop("progress", None)
job["log"] = [] job["log"] = []
self._save_job_locked(job) self._save_job_locked(job)
return job return job
@@ -1421,6 +1422,29 @@ class DownloadQueue:
return self._job_from_row(row) return self._job_from_row(row)
raise KeyError("Job not found") raise KeyError("Job not found")
def _progress_from_line(self, text):
if not str(text or "").startswith("KAIZOKU_PROGRESS "):
return None
payload = str(text).split(" ", 1)[1]
try:
progress = json.loads(payload)
except json.JSONDecodeError:
return None
if not isinstance(progress, dict):
return None
try:
progress["percent"] = max(0, min(100, int(float(progress.get("percent") or 0))))
except (TypeError, ValueError):
progress["percent"] = 0
for key in ("episode_index", "episode_total", "segment", "segment_total"):
if progress.get(key) in (None, ""):
continue
try:
progress[key] = int(float(progress[key]))
except (TypeError, ValueError):
progress.pop(key, None)
return progress
def _next_pending(self): def _next_pending(self):
with self.lock: with self.lock:
with self._connect() as conn: with self._connect() as conn:
@@ -1435,8 +1459,15 @@ class DownloadQueue:
text = strip_control(line).strip() text = strip_control(line).strip()
if not text: if not text:
return return
progress = self._progress_from_line(text)
if progress is not None:
text = str(progress.get("message") or "").strip()
if not text:
text = f"Download progress: {progress.get('percent', 0)}%."
self._stdout_log(job, text) self._stdout_log(job, text)
with self.lock: with self.lock:
if progress is not None:
job["progress"] = progress
job.setdefault("log", []).append(text) job.setdefault("log", []).append(text)
job["log"] = job["log"][-MAX_LOG_LINES:] job["log"] = job["log"][-MAX_LOG_LINES:]
job["updated_at"] = now_iso() job["updated_at"] = now_iso()
@@ -1590,6 +1621,11 @@ class DownloadQueue:
job["target_dir"] = str(target_dir) job["target_dir"] = str(target_dir)
job["download_backend"] = attempts[0]["backend"] job["download_backend"] = attempts[0]["backend"]
job["fallback_used"] = False job["fallback_used"] = False
job["progress"] = {
"phase": "queued",
"message": "Download queued.",
"percent": 0,
}
job["log"] = [ job["log"] = [
f"Starting: {job['command']}", f"Starting: {job['command']}",
f"Staging in: {job['staging_dir']}", f"Staging in: {job['staging_dir']}",
@@ -1710,9 +1746,11 @@ class DownloadQueue:
job["updated_at"] = job["finished_at"] job["updated_at"] = job["finished_at"]
if canceled: if canceled:
job["status"] = "canceled" job["status"] = "canceled"
job.setdefault("progress", {})["message"] = "Download canceled."
job.setdefault("log", []).append("Canceled.") job.setdefault("log", []).append("Canceled.")
elif finalize_error: elif finalize_error:
job["status"] = "failed" job["status"] = "failed"
job.setdefault("progress", {})["message"] = "Download failed during library finalization."
message = str(finalize_error).strip() message = str(finalize_error).strip()
if "No downloaded files were found in the staging folder" in message: if "No downloaded files were found in the staging folder" in message:
job.setdefault("log", []).append( job.setdefault("log", []).append(
@@ -1723,6 +1761,12 @@ class DownloadQueue:
cleanup_staging = True cleanup_staging = True
elif exit_code == 0: elif exit_code == 0:
job["status"] = "done" job["status"] = "done"
job["progress"] = {
**(job.get("progress") or {}),
"phase": "done",
"message": "Download completed.",
"percent": 100,
}
if successful_backend and successful_backend != attempts[0]["backend"]: if successful_backend and successful_backend != attempts[0]["backend"]:
job.setdefault("log", []).append(f"Fallback download completed with {successful_backend}.") job.setdefault("log", []).append(f"Fallback download completed with {successful_backend}.")
for path in moved_files: for path in moved_files:
@@ -1737,6 +1781,7 @@ class DownloadQueue:
job.setdefault("log", []).append("Watchlist download state synced.") job.setdefault("log", []).append("Watchlist download state synced.")
else: else:
job["status"] = "failed" job["status"] = "failed"
job.setdefault("progress", {})["message"] = "Download failed."
job.setdefault("log", []).append( job.setdefault("log", []).append(
f"{job.get('download_backend') or 'Download'} failed with exit code {exit_code}." f"{job.get('download_backend') or 'Download'} failed with exit code {exit_code}."
) )
+75
View File
@@ -269,6 +269,36 @@ QUEUE_HTML = r"""<!doctype html>
.status.done { color: #082312; background: var(--good); } .status.done { color: #082312; background: var(--good); }
.status.failed { color: #2d0910; background: var(--bad); } .status.failed { color: #2d0910; background: var(--bad); }
.status.canceled { color: #2d1c00; background: var(--warn); } .status.canceled { color: #2d1c00; background: var(--warn); }
.progress-wrap {
display: grid;
gap: 8px;
min-width: 0;
}
.progress-meta {
display: flex;
align-items: center;
justify-content: space-between;
gap: 10px;
color: var(--muted);
font-size: 12px;
}
.progress-label {
overflow-wrap: anywhere;
}
.progress-track {
height: 10px;
border-radius: 999px;
overflow: hidden;
background: rgba(255, 255, 255, 0.08);
border: 1px solid rgba(255, 255, 255, 0.06);
}
.progress-fill {
width: 0%;
height: 100%;
border-radius: inherit;
background: linear-gradient(90deg, var(--accent), #75c9ff);
transition: width 0.28s ease;
}
pre { pre {
margin: 0; margin: 0;
max-height: 220px; max-height: 220px;
@@ -408,6 +438,38 @@ QUEUE_HTML = r"""<!doctype html>
return [...(job.log || [])].slice(-28).reverse().join("\n"); return [...(job.log || [])].slice(-28).reverse().join("\n");
} }
function progressPercent(job) {
const progress = job.progress || {};
const value = Number(progress.percent || 0);
if (!Number.isFinite(value)) return job.status === "done" ? 100 : 0;
return Math.max(0, Math.min(100, Math.round(value)));
}
function progressLabel(job) {
const progress = job.progress || {};
if (progress.message) return progress.message;
if (job.status === "running") return "Download running.";
if (job.status === "pending") return "Waiting for worker.";
if (job.status === "done") return "Download completed.";
if (job.status === "failed") return "Download failed.";
if (job.status === "canceled") return "Download canceled.";
return "Waiting.";
}
function progressDetail(job) {
const progress = job.progress || {};
const parts = [];
if (progress.episode_index && progress.episode_total) {
parts.push(`episode ${progress.episode_index}/${progress.episode_total}`);
} else if (progress.episode) {
parts.push(`episode ${progress.episode}`);
}
if (progress.segment && progress.segment_total) {
parts.push(`segment ${progress.segment}/${progress.segment_total}`);
}
return parts.join(" · ");
}
function renderJobActions(job, actions) { function renderJobActions(job, actions) {
actions.innerHTML = ""; actions.innerHTML = "";
if (job.job_type === "jellyfin_handoff") { if (job.job_type === "jellyfin_handoff") {
@@ -436,6 +498,11 @@ QUEUE_HTML = r"""<!doctype html>
status.className = "status"; status.className = "status";
status.classList.add(job.status); status.classList.add(job.status);
item.querySelector("pre").textContent = buildJobLog(job) || "Waiting..."; item.querySelector("pre").textContent = buildJobLog(job) || "Waiting...";
const percent = progressPercent(job);
item.querySelector(".progress-fill").style.width = `${percent}%`;
item.querySelector(".progress-percent").textContent = `${percent}%`;
item.querySelector(".progress-label").textContent = progressLabel(job);
item.querySelector(".progress-detail").textContent = progressDetail(job);
item.querySelector(".toolbar .muted").textContent = job.exit_code === null || job.exit_code === undefined ? "" : `exit ${job.exit_code}`; item.querySelector(".toolbar .muted").textContent = job.exit_code === null || job.exit_code === undefined ? "" : `exit ${job.exit_code}`;
renderJobActions(job, item.querySelector(".row")); renderJobActions(job, item.querySelector(".row"));
state.jobCards[job.id] = item; state.jobCards[job.id] = item;
@@ -456,6 +523,14 @@ QUEUE_HTML = r"""<!doctype html>
<span class="status"></span> <span class="status"></span>
</div> </div>
<div class="job-body"> <div class="job-body">
<div class="progress-wrap">
<div class="progress-meta">
<span class="progress-label"></span>
<span class="progress-percent"></span>
</div>
<div class="progress-track" aria-hidden="true"><div class="progress-fill"></div></div>
<div class="progress-meta"><span class="progress-detail"></span></div>
</div>
<pre></pre> <pre></pre>
<div class="toolbar"> <div class="toolbar">
<span class="muted"></span> <span class="muted"></span>
+24
View File
@@ -403,6 +403,23 @@ class DownloadQueueCancelTests(unittest.TestCase):
) )
self.assertEqual(job["log"], ["Downloading episode 1"]) self.assertEqual(job["log"], ["Downloading episode 1"])
def test_append_log_parses_kaizoku_progress_lines(self):
queue = object.__new__(APP.DownloadQueue)
queue.lock = threading.RLock()
queue._save_job_locked = lambda _job: None
job = {"id": "job-1", "download_backend": "kaizoku", "log": [], "updated_at": APP.now_iso()}
APP.DownloadQueue._append_log(
queue,
job,
'KAIZOKU_PROGRESS {"episode":1,"episode_index":1,"episode_total":12,"message":"Episode 1: downloaded segment 25/100 (25%).","percent":25,"segment":25,"segment_total":100}',
)
self.assertEqual(job["progress"]["percent"], 25)
self.assertEqual(job["progress"]["episode_total"], 12)
self.assertEqual(job["progress"]["segment_total"], 100)
self.assertEqual(job["log"], ["Episode 1: downloaded segment 25/100 (25%)."])
class QueueApiTests(unittest.TestCase): class QueueApiTests(unittest.TestCase):
def setUp(self): def setUp(self):
@@ -4123,6 +4140,13 @@ class TemplateHelperTests(unittest.TestCase):
self.assertIn('const progress = `${job.completed || 0}/${job.total || 0} entries`;', APP.QUEUE_HTML) self.assertIn('const progress = `${job.completed || 0}/${job.total || 0} entries`;', APP.QUEUE_HTML)
self.assertIn('const moved = `${job.moved || 0} moved`;', APP.QUEUE_HTML) self.assertIn('const moved = `${job.moved || 0} moved`;', APP.QUEUE_HTML)
def test_queue_page_renders_download_progress_bar(self):
self.assertIn('class="progress-wrap"', APP.QUEUE_HTML)
self.assertIn('class="progress-fill"', APP.QUEUE_HTML)
self.assertIn("function progressPercent(job)", APP.QUEUE_HTML)
self.assertIn('item.querySelector(".progress-fill").style.width = `${percent}%`;', APP.QUEUE_HTML)
self.assertIn("segment ${progress.segment}/${progress.segment_total}", APP.QUEUE_HTML)
def test_search_page_queues_selected_title_with_result_index(self): def test_search_page_queues_selected_title_with_result_index(self):
self.assertIn('query: state.selected.title,', APP.INDEX_HTML) self.assertIn('query: state.selected.title,', APP.INDEX_HTML)
self.assertIn('result_index: state.selected.index,', APP.INDEX_HTML) self.assertIn('result_index: state.selected.index,', APP.INDEX_HTML)