From 43644c211b96453c6cd778942bf4e3d483413dcb Mon Sep 17 00:00:00 2001 From: root Date: Sun, 27 Sep 2026 15:45:51 +0800 Subject: [PATCH] stop and continue for task --- m3u8downloaderd/README.md | 5 + .../src/m3u8downloaderd/static/app.js | 18 +++ .../src/m3u8downloaderd/static/styles.css | 1 + m3u8downloaderd/src/m3u8downloaderd/tasks.py | 117 +++++++++++++++--- m3u8downloaderd/src/m3u8downloaderd/web.py | 29 ++++- m3u8downloaderd/src/m3u8downloaderd/worker.py | 116 +++++++++++++++-- m3u8downloaderd/tests/test_service.py | 57 +++++++++ 7 files changed, 312 insertions(+), 31 deletions(-) diff --git a/m3u8downloaderd/README.md b/m3u8downloaderd/README.md index 0166663..087cba9 100644 --- a/m3u8downloaderd/README.md +++ b/m3u8downloaderd/README.md @@ -43,6 +43,11 @@ of 1 to 32. Finished segment downloads enter a FIFO ffmpeg queue so slow media processing does not occupy download slots. Runtime logs are available in the Log tab and remain in memory only. +Each queued or active task can be paused from the task list. A pause waits for +an active segment request to finish, retains downloaded segments, and stops +active ffmpeg work. Resuming returns segment downloads to the download queue; +ffmpeg work restarts its per-video post-processing from the beginning. + Run tests from the repository root: ```bash diff --git a/m3u8downloaderd/src/m3u8downloaderd/static/app.js b/m3u8downloaderd/src/m3u8downloaderd/static/app.js index f95b105..e121ba0 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/static/app.js +++ b/m3u8downloaderd/src/m3u8downloaderd/static/app.js @@ -91,11 +91,29 @@ function renderTask(task) { const actions = fragment.querySelector(".task-actions"); if (["queued", "running"].includes(task.status)) { + const pause = make("button", "secondary-action", "Pause"); + pause.type = "button"; + pause.addEventListener("click", () => mutate(`/api/tasks/${task.id}/pause`)); + const cancel = make("button", "secondary-action", "Cancel"); + cancel.type = "button"; + cancel.addEventListener("click", () => mutate(`/api/tasks/${task.id}/cancel`, "Cancel this task?")); + actions.append(pause, cancel); + } + if (task.status === "pausing") { const cancel = make("button", "secondary-action", "Cancel"); cancel.type = "button"; cancel.addEventListener("click", () => mutate(`/api/tasks/${task.id}/cancel`, "Cancel this task?")); actions.append(cancel); } + if (task.status === "paused") { + const resume = make("button", "primary-action", "Resume"); + resume.type = "button"; + resume.addEventListener("click", () => mutate(`/api/tasks/${task.id}/resume`)); + const cancel = make("button", "secondary-action", "Cancel"); + cancel.type = "button"; + cancel.addEventListener("click", () => mutate(`/api/tasks/${task.id}/cancel`, "Cancel this task?")); + actions.append(resume, cancel); + } if (["partial", "failed", "cancelled"].includes(task.status)) { const retryProcessing = make("button", "secondary-action", "Retry processing"); retryProcessing.type = "button"; diff --git a/m3u8downloaderd/src/m3u8downloaderd/static/styles.css b/m3u8downloaderd/src/m3u8downloaderd/static/styles.css index e5f1ab2..904a580 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/static/styles.css +++ b/m3u8downloaderd/src/m3u8downloaderd/static/styles.css @@ -61,6 +61,7 @@ h2 { margin-bottom: 0; font-size: 1.15rem; } .task-row[data-status="completed"] { border-left-color: #4a8f55; } .task-row[data-status="partial"], .task-row[data-status="failed"] { border-left-color: #e36b3e; } .task-row[data-status="cancelled"] { border-left-color: #7d689f; } +.task-row[data-status="pausing"], .task-row[data-status="paused"] { border-left-color: #b45288; } .task-row-header { display: flex; justify-content: space-between; gap: 16px; align-items: flex-start; } .task-main { min-width: 0; } .task-title { font-weight: 720; overflow-wrap: anywhere; } diff --git a/m3u8downloaderd/src/m3u8downloaderd/tasks.py b/m3u8downloaderd/src/m3u8downloaderd/tasks.py index 63f41f1..d3e0b2f 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/tasks.py +++ b/m3u8downloaderd/src/m3u8downloaderd/tasks.py @@ -4,6 +4,7 @@ from __future__ import annotations import copy import threading +from collections import deque from datetime import datetime, timezone from pathlib import Path from typing import Any, Iterable @@ -32,6 +33,7 @@ class TaskState: "max_active_ffmpeg": max_active_ffmpeg if max_active_ffmpeg is not None else max_active_downloads, } self._tasks: dict[int, dict[str, Any]] = {} + self._download_queue: deque[tuple[int, int]] = deque() self._logs: list[dict[str, Any]] = [] self._next_task_id = 1 self._next_item_id = 1 @@ -86,6 +88,7 @@ class TaskState: "proxy": proxy, "status": "queued", "cancel_requested": False, + "pause_requested": False, "created_at": _now(), "started_at": None, "completed_at": None, @@ -117,28 +120,33 @@ class TaskState: "warning": None, "output_path": None, "cover_path": None, + "resume_phase": "download", + "queue_generation": 0, "started_at": None, "completed_at": None, } ) self._event_unlocked(task, None, "info", "Task queued") self._tasks[task_id] = task + self._download_queue.extend((int(item["id"]), 0) for item in task["items"]) return task_id def claim_next_item(self) -> tuple[dict[str, Any], dict[str, Any]] | None: with self._lock: - for task in self._tasks.values(): - if task["status"] == "queued": - task.update(status="running", started_at=_now(), cancel_requested=False) - self._event_unlocked(task, None, "info", "Task started") - if task["status"] != "running" or task["cancel_requested"]: + while self._download_queue: + item_id, generation = self._download_queue.popleft() + try: + task, item = self._find_task_and_item_unlocked(item_id) + except KeyError: continue - for item in task["items"]: - if item["status"] != "queued": - continue - item.update(status="running", stage="Preparing", started_at=_now(), message="Preparing download") - self._event_unlocked(task, int(item["id"]), "info", "Video started") - return copy.deepcopy(task), copy.deepcopy(item) + if task["status"] == "queued": + task.update(status="running", started_at=_now(), cancel_requested=False, pause_requested=False) + self._event_unlocked(task, None, "info", "Task started") + if task["status"] != "running" or task["cancel_requested"] or task["pause_requested"] or item["status"] != "queued" or item["queue_generation"] != generation: + continue + item.update(status="running", stage="Preparing", started_at=_now(), message="Preparing download") + self._event_unlocked(task, int(item["id"]), "info", "Video started") + return copy.deepcopy(task), copy.deepcopy(item) return None def get_task(self, task_id: int) -> dict[str, Any] | None: @@ -166,8 +174,11 @@ class TaskState: return False now = _now() task["cancel_requested"] = True + task["pause_requested"] = False + if task["status"] in {"paused", "pausing"}: + task["status"] = "running" for item in task["items"]: - if item["status"] in {"queued", "processing_queued"}: + if item["status"] in {"queued", "processing_queued", "paused"}: item.update(status="cancelled", stage="Cancelled", completed_at=now) self._event_unlocked(task, None, "info", "Cancellation requested") self._finalize_task_unlocked(task) @@ -177,6 +188,57 @@ class TaskState: with self._lock: return bool(self._tasks[task_id]["cancel_requested"]) + def pause_requested(self, task_id: int) -> bool: + with self._lock: + return bool(self._tasks[task_id]["pause_requested"]) + + def pause_task(self, task_id: int) -> bool: + with self._lock: + task = self._tasks.get(task_id) + if task is None or task["status"] not in {"queued", "running"}: + return False + task.update(status="pausing", pause_requested=True) + for item in task["items"]: + if item["status"] == "queued": + item.update(status="paused", stage="Paused", resume_phase="download", message="Paused") + elif item["status"] == "processing_queued": + item.update(status="paused", stage="Paused", resume_phase="processing", message="Paused") + elif item["status"] == "running": + item.update(stage="Pausing", message="Pausing") + self._event_unlocked(task, None, "info", "Pause requested") + self._settle_pause_unlocked(task) + return True + + def mark_item_paused(self, item_id: int, resume_phase: str) -> None: + with self._lock: + task, item = self._find_task_and_item_unlocked(item_id) + if not task["pause_requested"]: + return + item.update(status="paused", stage="Paused", resume_phase=resume_phase, message="Paused", completed_at=None) + self._settle_pause_unlocked(task) + + def resume_task(self, task_id: int) -> tuple[list[int], list[int]] | None: + with self._lock: + task = self._tasks.get(task_id) + if task is None or task["status"] != "paused": + return None + download_ids: list[int] = [] + processing_ids: list[int] = [] + for item in task["items"]: + if item["status"] != "paused": + continue + if item["resume_phase"] == "processing": + item.update(status="processing_queued", stage="Waiting for ffmpeg", progress_phase="processing", message="Waiting for ffmpeg") + processing_ids.append(int(item["id"])) + else: + item.update(status="queued", stage="Queued", progress_phase="download", message="Queued") + download_ids.append(int(item["id"])) + task.update(status="running", pause_requested=False, cancel_requested=False, completed_at=None) + for item_id in download_ids: + self._enqueue_download_unlocked(self._find_item_unlocked(item_id)) + self._event_unlocked(task, None, "info", "Task resumed") + return download_ids, processing_ids + def finalize_task_if_ready(self, task_id: int) -> None: with self._lock: task = self._tasks[task_id] @@ -192,7 +254,9 @@ class TaskState: return [] for item in retryable: self._reset_for_video_retry_unlocked(item) - task.update(status="queued", cancel_requested=False, started_at=None, completed_at=None, error=None) + task.update(status="queued", cancel_requested=False, pause_requested=False, started_at=None, completed_at=None, error=None) + for item in retryable: + self._enqueue_download_unlocked(item) self._event_unlocked(task, None, "info", "Task retry queued") return [int(item["id"]) for item in retryable] @@ -202,14 +266,15 @@ class TaskState: if task["status"] == "running" or item["status"] == "completed": return False self._reset_for_video_retry_unlocked(item) - task.update(status="queued", cancel_requested=False, started_at=None, completed_at=None, error=None) + task.update(status="queued", cancel_requested=False, pause_requested=False, started_at=None, completed_at=None, error=None) + self._enqueue_download_unlocked(item) self._event_unlocked(task, item_id, "info", "Video retry queued") return True def queue_processing(self, item_id: int, *, stage: str = "Waiting for ffmpeg") -> tuple[dict[str, Any], dict[str, Any]]: with self._lock: task, item = self._find_task_and_item_unlocked(item_id) - if task["cancel_requested"]: + if task["cancel_requested"] or task["pause_requested"]: raise RuntimeError("Task is cancelled") item.update( status="processing_queued", @@ -228,7 +293,7 @@ class TaskState: task, item = self._find_task_and_item_unlocked(item_id) if task["status"] == "running" or item["status"] != "failed": return None - task.update(status="running", cancel_requested=False, started_at=task["started_at"] or _now(), completed_at=None, error=None) + task.update(status="running", cancel_requested=False, pause_requested=False, started_at=task["started_at"] or _now(), completed_at=None, error=None) item.update(warning=None, started_at=item["started_at"] or _now()) self._event_unlocked(task, item_id, "info", "Processing retry queued") return self.queue_processing(item_id) @@ -241,7 +306,7 @@ class TaskState: retryable = [item for item in task["items"] if item["status"] in {"failed", "completed_warning"}] if not retryable: return [] - task.update(status="running", cancel_requested=False, started_at=task["started_at"] or _now(), completed_at=None, error=None) + task.update(status="running", cancel_requested=False, pause_requested=False, started_at=task["started_at"] or _now(), completed_at=None, error=None) claimed: list[tuple[dict[str, Any], dict[str, Any]]] = [] for item in retryable: item.update(warning=None, started_at=item["started_at"] or _now()) @@ -258,7 +323,7 @@ class TaskState: task, item = self._find_task_and_item_unlocked(item_id) if task["status"] == "running" or item["status"] != "completed_warning": return None - task.update(status="running", cancel_requested=False, started_at=task["started_at"] or _now(), completed_at=None, error=None) + task.update(status="running", cancel_requested=False, pause_requested=False, started_at=task["started_at"] or _now(), completed_at=None, error=None) item.update(warning=None, started_at=item["started_at"] or _now()) self._event_unlocked(task, item_id, "info", "Cover retry queued") return self.queue_processing(item_id, stage="Waiting for ffmpeg") @@ -277,10 +342,15 @@ class TaskState: warning=None, output_path=None, cover_path=None, + resume_phase="download", started_at=None, completed_at=None, ) + def _enqueue_download_unlocked(self, item: dict[str, Any]) -> None: + item["queue_generation"] = int(item.get("queue_generation", 0)) + 1 + self._download_queue.append((int(item["id"]), int(item["queue_generation"]))) + def _find_item_unlocked(self, item_id: int) -> dict[str, Any]: return self._find_task_and_item_unlocked(item_id)[1] @@ -334,6 +404,9 @@ class TaskState: self._next_log_id += 1 def _finalize_task_unlocked(self, task: dict[str, Any]) -> None: + if task["status"] == "pausing": + self._settle_pause_unlocked(task) + return if task["status"] not in {"queued", "running"}: return statuses = {item["status"] for item in task["items"]} @@ -351,3 +424,11 @@ class TaskState: status = "failed" task.update(status=status, completed_at=_now()) self._event_unlocked(task, None, "info" if status == "completed" else "warning", f"Task {status}") + + def _settle_pause_unlocked(self, task: dict[str, Any]) -> None: + if task["status"] != "pausing": + return + if any(item["status"] == "running" for item in task["items"]): + return + task.update(status="paused", completed_at=None) + self._event_unlocked(task, None, "info", "Task paused") diff --git a/m3u8downloaderd/src/m3u8downloaderd/web.py b/m3u8downloaderd/src/m3u8downloaderd/web.py index 63a620a..38d9ed6 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/web.py +++ b/m3u8downloaderd/src/m3u8downloaderd/web.py @@ -101,6 +101,15 @@ class DownloadService: def retry_item_cover(self, item_id: int) -> bool: return self.runner.retry_item_cover(item_id) + def pause_task(self, task_id: int) -> bool: + return self.runner.pause_task(task_id) + + def resume_task(self, task_id: int) -> bool: + return self.runner.resume_task(task_id) + + def cancel_task(self, task_id: int) -> bool: + return self.runner.cancel_task(task_id) + def logs(self, level: str | None) -> dict[str, Any]: try: return {"logs": self.store.logs(level)} @@ -207,13 +216,31 @@ class ServiceRequestHandler(BaseHTTPRequestHandler): return if self.path.endswith("/cancel") and self.path.startswith("/api/tasks/"): task_id = _id_from_path(self.path[: -len("/cancel")], "/api/tasks/") - if task_id is None or not self.server.service.store.cancel_task(task_id): + if task_id is None or not self.server.service.cancel_task(task_id): self._error(HTTPStatus.CONFLICT, "Task cannot be cancelled") else: task = self.server.service.get_task(task_id) assert task is not None self._json(HTTPStatus.OK, task) return + if self.path.endswith("/pause") and self.path.startswith("/api/tasks/"): + task_id = _id_from_path(self.path[: -len("/pause")], "/api/tasks/") + if task_id is None or not self.server.service.pause_task(task_id): + self._error(HTTPStatus.CONFLICT, "Task cannot be paused") + else: + task = self.server.service.get_task(task_id) + assert task is not None + self._json(HTTPStatus.OK, task) + return + if self.path.endswith("/resume") and self.path.startswith("/api/tasks/"): + task_id = _id_from_path(self.path[: -len("/resume")], "/api/tasks/") + if task_id is None or not self.server.service.resume_task(task_id): + self._error(HTTPStatus.CONFLICT, "Task cannot be resumed") + else: + task = self.server.service.get_task(task_id) + assert task is not None + self._json(HTTPStatus.OK, task) + return if self.path.endswith("/retry-video") and self.path.startswith("/api/tasks/"): task_id = _id_from_path(self.path[: -len("/retry-video")], "/api/tasks/") if task_id is None or not self.server.service.retry_task_video(task_id): diff --git a/m3u8downloaderd/src/m3u8downloaderd/worker.py b/m3u8downloaderd/src/m3u8downloaderd/worker.py index 226b23a..5e36629 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/worker.py +++ b/m3u8downloaderd/src/m3u8downloaderd/worker.py @@ -43,6 +43,11 @@ class PostprocessJob: output_path: Path result: SegmentDownloadResult | None = None cover_only: bool = False + generation: int = 0 + + +class TaskPausedError(Exception): + """Raised when a task reaches a pause checkpoint.""" class TaskRunner: @@ -60,9 +65,9 @@ class TaskRunner: self._download_executor = concurrent.futures.ThreadPoolExecutor(max_workers=self.worker_capacity, thread_name_prefix="m3u8downloaderd-download") self._ffmpeg_executor = concurrent.futures.ThreadPoolExecutor(max_workers=self.worker_capacity, thread_name_prefix="m3u8downloaderd-ffmpeg") self._jobs: dict[int, PostprocessJob] = {} - self._pending: deque[int] = deque() + self._pending: deque[tuple[int, int]] = deque() self._jobs_lock = threading.RLock() - self._processes: set[subprocess.Popen[str]] = set() + self._processes: dict[subprocess.Popen[str], int] = {} self._processes_lock = threading.Lock() def start(self) -> None: @@ -88,18 +93,57 @@ class TaskRunner: for job in jobs: self._cleanup_job(job) + def pause_task(self, task_id: int) -> bool: + if not self.store.pause_task(task_id): + return False + self._terminate_task_ffmpeg(task_id) + return True + + def resume_task(self, task_id: int) -> bool: + resumed = self.store.resume_task(task_id) + if resumed is None: + return False + _, processing_ids = resumed + for item_id in processing_ids: + if self._has_job(item_id): + self._enqueue_existing(item_id) + else: + task = self.store.get_task(task_id) + if task is not None: + item = next((candidate for candidate in task["items"] if int(candidate["id"]) == item_id), None) + if item is not None: + self._mark_missing_intermediate(task, item) + return True + + def cancel_task(self, task_id: int) -> bool: + task = self.store.get_task(task_id) + if task is None or not self.store.cancel_task(task_id): + return False + self._terminate_task_ffmpeg(task_id) + for item in task["items"]: + item_id = int(item["id"]) + if item["status"] != "running": + self._discard_job(item_id) + self._cleanup_workspace(task, item_id) + return True + def retry_item_video(self, item_id: int) -> bool: + task, _ = self._task_item(item_id) if not self.store.retry_item_video(item_id): return False self._discard_job(item_id) + self._cleanup_workspace(task, item_id) return True def retry_task_video(self, task_id: int) -> bool: item_ids = self.store.retry_task_video(task_id) if not item_ids: return False + task = self.store.get_task(task_id) for item_id in item_ids: self._discard_job(item_id) + if task is not None: + self._cleanup_workspace(task, item_id) return True def retry_item_processing(self, item_id: int) -> bool: @@ -141,7 +185,7 @@ class TaskRunner: self._jobs[item_id] = job else: job.cover_only = True - self._pending.append(item_id) + self._enqueue_existing(item_id) return True def _dispatch_downloads(self) -> None: @@ -191,7 +235,7 @@ class TaskRunner: progress: dict[str, tuple[int, int]] = {} def on_event(event: object) -> bool | None: - if self._stop.is_set() or self.store.cancel_requested(task_id): + if self._stop.is_set() or self.store.cancel_requested(task_id) or self.store.pause_requested(task_id): return False kind = getattr(event, "kind", None) track_id = getattr(event, "track_id", None) or "unknown" @@ -212,14 +256,19 @@ class TaskRunner: self.store.update_item(item_id, stage="Downloading", progress_phase="download", message="Inspecting media playlist") options = RequestOptions(headers={"Referer": str(item["m3u8_referer"])}, proxy=str(task["proxy"]) or None, use_system_proxy=False) workspace = output_dir / f".m3u8downloaderd-{task_id}-{item_id}" - shutil.rmtree(workspace, ignore_errors=True) request = DownloadRequest(output_dir=output_dir, temporary_dir=workspace, file_name=str(item["title"]), thread_count=max(1, os.cpu_count() or 1)) result = N_m3u8DL(options).download_segments_url(str(item["m3u8_url"]), request, on_event=on_event) if self._stop.is_set() or self.store.cancel_requested(task_id): raise DownloadCancelledError("Download cancelled by user") + if self.store.pause_requested(task_id): + self.store.mark_item_paused(item_id, "download") + return output_path = self._reserve_output_path(output_dir, str(item["title"])) self._remember_job(PostprocessJob(task_id, item_id, output_dir, output_path, result=result)) except DownloadCancelledError: + if self.store.pause_requested(task_id): + self.store.mark_item_paused(item_id, "download") + return if workspace is not None: shutil.rmtree(workspace, ignore_errors=True) self._cancel_item(task_id, item_id) @@ -239,6 +288,10 @@ class TaskRunner: self._cancel_item(job.task_id, job.item_id) self._discard_job(item_id) return + if self.store.pause_requested(job.task_id): + self._cleanup_processing_outputs(job) + self.store.mark_item_paused(job.item_id, "processing") + return try: if job.cover_only: self._retry_cover_job(job, task) @@ -247,6 +300,9 @@ class TaskRunner: except DownloadCancelledError: self._cancel_item(job.task_id, job.item_id) self._discard_job(item_id) + except TaskPausedError: + self._cleanup_processing_outputs(job) + self.store.mark_item_paused(job.item_id, "processing") except Exception as error: self._mark_failed(job.task_id, job.item_id, str(error), "Processing failed") finally: @@ -364,7 +420,7 @@ class TaskRunner: command = [*command[:-1], "-progress", "pipe:1", "-nostats", command[-1]] process = subprocess.Popen(command, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, start_new_session=True) with self._processes_lock: - self._processes.add(process) + self._processes[process] = job.task_id output: list[str] = [] try: assert process.stdout is not None @@ -373,18 +429,23 @@ class TaskRunner: if self._stop.is_set() or self.store.cancel_requested(job.task_id): self._terminate_process(process) raise DownloadCancelledError("Processing cancelled") + if self.store.pause_requested(job.task_id): + self._terminate_process(process) + raise TaskPausedError("Processing paused") key, separator, value = line.strip().partition("=") if separator and key in {"out_time_us", "out_time_ms"} and duration > 0: try: progress(max(0.0, min(1.0, float(value) / 1_000_000 / duration))) except ValueError: pass + if self.store.pause_requested(job.task_id): + raise TaskPausedError("Processing paused") if process.wait(): raise RuntimeError(f"{summary}: {''.join(output).strip()[-1500:] or summary}") progress(1.0) finally: with self._processes_lock: - self._processes.discard(process) + self._processes.pop(process, None) def _update_processing(self, job: PostprocessJob, stage: str, completed: float, total: float) -> None: self.store.update_item(job.item_id, status="running", stage=stage, progress_phase="processing", processing_completed=max(0.0, min(completed, total)), processing_total=max(total, 1.0), message=stage) @@ -412,20 +473,24 @@ class TaskRunner: def _remember_job(self, job: PostprocessJob) -> None: with self._jobs_lock: self._jobs[job.item_id] = job - self._pending.append(job.item_id) self.store.queue_processing(job.item_id) + self._enqueue_existing(job.item_id) def _enqueue_existing(self, item_id: int) -> None: with self._jobs_lock: - if item_id in self._jobs: - self._pending.append(item_id) + job = self._jobs.get(item_id) + if job is not None: + job.generation += 1 + self._pending.append((item_id, job.generation)) def _next_postprocess_item(self) -> int | None: with self._jobs_lock: while self._pending: - item_id = self._pending.popleft() + item_id, generation = self._pending.popleft() job = self._jobs.get(item_id) - if job is None: + if job is None or job.generation != generation: + continue + if self.store.pause_requested(job.task_id): continue if self.store.cancel_requested(job.task_id): self._jobs.pop(item_id, None) @@ -480,6 +545,33 @@ class TaskRunner: for process in processes: self._terminate_process(process) + def _terminate_task_ffmpeg(self, task_id: int) -> None: + with self._processes_lock: + processes = [process for process, owner in self._processes.items() if owner == task_id] + for process in processes: + self._terminate_process(process) + + def _cleanup_processing_outputs(self, job: PostprocessJob) -> None: + if job.result is not None: + for path in job.result.temporary_dir.glob("*.merged.*"): + path.unlink(missing_ok=True) + for path in job.result.temporary_dir.glob("*.concat.txt"): + path.unlink(missing_ok=True) + (job.output_dir / f".{job.output_path.stem}.{job.item_id}.mux.tmp.mp4").unlink(missing_ok=True) + (job.output_dir / f".{job.output_path.stem}.{job.item_id}.cover.tmp.mp4").unlink(missing_ok=True) + + def _task_item(self, item_id: int) -> tuple[dict[str, object], dict[str, object]]: + for task in self.store.list_tasks(): + for item in task["items"]: + if int(item["id"]) == item_id: + return task, item + raise KeyError(item_id) + + @staticmethod + def _cleanup_workspace(task: dict[str, object], item_id: int) -> None: + output_dir = Path(str(task["output_dir"])) + shutil.rmtree(output_dir / f".m3u8downloaderd-{int(task['id'])}-{item_id}", ignore_errors=True) + @staticmethod def _terminate_process(process: subprocess.Popen[str]) -> None: if process.poll() is not None: diff --git a/m3u8downloaderd/tests/test_service.py b/m3u8downloaderd/tests/test_service.py index 2a61b2c..60dd7df 100644 --- a/m3u8downloaderd/tests/test_service.py +++ b/m3u8downloaderd/tests/test_service.py @@ -67,6 +67,12 @@ def test_http_api_validates_and_snapshots_settings(tmp_path: Path) -> None: assert task[1]["proxy"] == "http://127.0.0.1:7890" assert task[1]["output_dir"] == str(config.download_root / "batch") assert "m3u8_referer" not in task[1]["items"][0] + paused = _request(base, f"/api/tasks/{task[1]['id']}/pause", {}) + assert paused[0] == 200 + assert paused[1]["status"] == "paused" + resumed = _request(base, f"/api/tasks/{task[1]['id']}/resume", {}) + assert resumed[0] == 200 + assert resumed[1]["status"] == "running" service.store.log("error", "Expected log entry") logs = _get(base, "/api/logs?level=error") assert logs[0] == 200 @@ -125,6 +131,57 @@ def test_tasks_and_settings_are_not_retained_after_restart(tmp_path: Path) -> No assert restarted.bootstrap()["settings"]["max_active_downloads"] == 1 +def test_pause_waits_for_running_item_then_resumes_at_download_queue_tail(tmp_path: Path) -> None: + store = TaskState(tmp_path) + payload = [_item("http://127.0.0.1"), {**_item("http://127.0.0.1"), "code": "episode-2", "title": "episode-2"}] + task_id = store.create_task( + title="batch", + base_dir=tmp_path, + output_dir=tmp_path / "batch", + proxy="", + items=decode_download_items(json.dumps(payload)), + ) + + claimed = store.claim_next_item() + assert claimed is not None + _, first = claimed + assert store.pause_task(task_id) + pausing = store.get_task(task_id) + assert pausing is not None + assert pausing["status"] == "pausing" + assert pausing["items"][1]["status"] == "paused" + + store.mark_item_paused(int(first["id"]), "download") + paused = store.get_task(task_id) + assert paused is not None + assert paused["status"] == "paused" + + resumed = store.resume_task(task_id) + assert resumed == ([int(first["id"]), int(paused["items"][1]["id"])], []) + next_claimed = store.claim_next_item() + assert next_claimed is not None + _, next_item = next_claimed + assert next_item["id"] == first["id"] + + +def test_cancelling_a_paused_task_reaches_cancelled(tmp_path: Path) -> None: + store = TaskState(tmp_path) + task_id = store.create_task( + title="batch", + base_dir=tmp_path, + output_dir=tmp_path / "batch", + proxy="", + items=decode_download_items(json.dumps([_item("http://127.0.0.1")])), + ) + + assert store.pause_task(task_id) + assert store.get_task(task_id)["status"] == "paused" # type: ignore[index] + assert store.cancel_task(task_id) + cancelled = store.get_task(task_id) + assert cancelled is not None + assert cancelled["status"] == "cancelled" + + def test_runner_runs_multiple_videos_from_one_task_concurrently(tmp_path: Path) -> None: store = TaskState(tmp_path, max_active_downloads=2) runner = TaskRunner(store, worker_capacity=2, ffmpeg="ffmpeg")