stop and continue for task

This commit is contained in:
root
2026-09-27 15:45:51 +08:00
parent b6160430a0
commit 43644c211b
7 changed files with 312 additions and 31 deletions
+5
View File
@@ -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 processing does not occupy download slots. Runtime logs are available in the
Log tab and remain in memory only. 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: Run tests from the repository root:
```bash ```bash
@@ -91,11 +91,29 @@ function renderTask(task) {
const actions = fragment.querySelector(".task-actions"); const actions = fragment.querySelector(".task-actions");
if (["queued", "running"].includes(task.status)) { 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"); const cancel = make("button", "secondary-action", "Cancel");
cancel.type = "button"; cancel.type = "button";
cancel.addEventListener("click", () => mutate(`/api/tasks/${task.id}/cancel`, "Cancel this task?")); cancel.addEventListener("click", () => mutate(`/api/tasks/${task.id}/cancel`, "Cancel this task?"));
actions.append(cancel); 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)) { if (["partial", "failed", "cancelled"].includes(task.status)) {
const retryProcessing = make("button", "secondary-action", "Retry processing"); const retryProcessing = make("button", "secondary-action", "Retry processing");
retryProcessing.type = "button"; retryProcessing.type = "button";
@@ -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="completed"] { border-left-color: #4a8f55; }
.task-row[data-status="partial"], .task-row[data-status="failed"] { border-left-color: #e36b3e; } .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="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-row-header { display: flex; justify-content: space-between; gap: 16px; align-items: flex-start; }
.task-main { min-width: 0; } .task-main { min-width: 0; }
.task-title { font-weight: 720; overflow-wrap: anywhere; } .task-title { font-weight: 720; overflow-wrap: anywhere; }
+95 -14
View File
@@ -4,6 +4,7 @@ from __future__ import annotations
import copy import copy
import threading import threading
from collections import deque
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from typing import Any, Iterable 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, "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._tasks: dict[int, dict[str, Any]] = {}
self._download_queue: deque[tuple[int, int]] = deque()
self._logs: list[dict[str, Any]] = [] self._logs: list[dict[str, Any]] = []
self._next_task_id = 1 self._next_task_id = 1
self._next_item_id = 1 self._next_item_id = 1
@@ -86,6 +88,7 @@ class TaskState:
"proxy": proxy, "proxy": proxy,
"status": "queued", "status": "queued",
"cancel_requested": False, "cancel_requested": False,
"pause_requested": False,
"created_at": _now(), "created_at": _now(),
"started_at": None, "started_at": None,
"completed_at": None, "completed_at": None,
@@ -117,24 +120,29 @@ class TaskState:
"warning": None, "warning": None,
"output_path": None, "output_path": None,
"cover_path": None, "cover_path": None,
"resume_phase": "download",
"queue_generation": 0,
"started_at": None, "started_at": None,
"completed_at": None, "completed_at": None,
} }
) )
self._event_unlocked(task, None, "info", "Task queued") self._event_unlocked(task, None, "info", "Task queued")
self._tasks[task_id] = task self._tasks[task_id] = task
self._download_queue.extend((int(item["id"]), 0) for item in task["items"])
return task_id return task_id
def claim_next_item(self) -> tuple[dict[str, Any], dict[str, Any]] | None: def claim_next_item(self) -> tuple[dict[str, Any], dict[str, Any]] | None:
with self._lock: with self._lock:
for task in self._tasks.values(): while self._download_queue:
if task["status"] == "queued": item_id, generation = self._download_queue.popleft()
task.update(status="running", started_at=_now(), cancel_requested=False) try:
self._event_unlocked(task, None, "info", "Task started") task, item = self._find_task_and_item_unlocked(item_id)
if task["status"] != "running" or task["cancel_requested"]: except KeyError:
continue continue
for item in task["items"]: if task["status"] == "queued":
if item["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 continue
item.update(status="running", stage="Preparing", started_at=_now(), message="Preparing download") item.update(status="running", stage="Preparing", started_at=_now(), message="Preparing download")
self._event_unlocked(task, int(item["id"]), "info", "Video started") self._event_unlocked(task, int(item["id"]), "info", "Video started")
@@ -166,8 +174,11 @@ class TaskState:
return False return False
now = _now() now = _now()
task["cancel_requested"] = True task["cancel_requested"] = True
task["pause_requested"] = False
if task["status"] in {"paused", "pausing"}:
task["status"] = "running"
for item in task["items"]: 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) item.update(status="cancelled", stage="Cancelled", completed_at=now)
self._event_unlocked(task, None, "info", "Cancellation requested") self._event_unlocked(task, None, "info", "Cancellation requested")
self._finalize_task_unlocked(task) self._finalize_task_unlocked(task)
@@ -177,6 +188,57 @@ class TaskState:
with self._lock: with self._lock:
return bool(self._tasks[task_id]["cancel_requested"]) 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: def finalize_task_if_ready(self, task_id: int) -> None:
with self._lock: with self._lock:
task = self._tasks[task_id] task = self._tasks[task_id]
@@ -192,7 +254,9 @@ class TaskState:
return [] return []
for item in retryable: for item in retryable:
self._reset_for_video_retry_unlocked(item) 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") self._event_unlocked(task, None, "info", "Task retry queued")
return [int(item["id"]) for item in retryable] return [int(item["id"]) for item in retryable]
@@ -202,14 +266,15 @@ class TaskState:
if task["status"] == "running" or item["status"] == "completed": if task["status"] == "running" or item["status"] == "completed":
return False return False
self._reset_for_video_retry_unlocked(item) 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") self._event_unlocked(task, item_id, "info", "Video retry queued")
return True return True
def queue_processing(self, item_id: int, *, stage: str = "Waiting for ffmpeg") -> tuple[dict[str, Any], dict[str, Any]]: def queue_processing(self, item_id: int, *, stage: str = "Waiting for ffmpeg") -> tuple[dict[str, Any], dict[str, Any]]:
with self._lock: with self._lock:
task, item = self._find_task_and_item_unlocked(item_id) 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") raise RuntimeError("Task is cancelled")
item.update( item.update(
status="processing_queued", status="processing_queued",
@@ -228,7 +293,7 @@ class TaskState:
task, item = self._find_task_and_item_unlocked(item_id) task, item = self._find_task_and_item_unlocked(item_id)
if task["status"] == "running" or item["status"] != "failed": if task["status"] == "running" or item["status"] != "failed":
return None 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()) item.update(warning=None, started_at=item["started_at"] or _now())
self._event_unlocked(task, item_id, "info", "Processing retry queued") self._event_unlocked(task, item_id, "info", "Processing retry queued")
return self.queue_processing(item_id) 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"}] retryable = [item for item in task["items"] if item["status"] in {"failed", "completed_warning"}]
if not retryable: if not retryable:
return [] 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]]] = [] claimed: list[tuple[dict[str, Any], dict[str, Any]]] = []
for item in retryable: for item in retryable:
item.update(warning=None, started_at=item["started_at"] or _now()) 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) task, item = self._find_task_and_item_unlocked(item_id)
if task["status"] == "running" or item["status"] != "completed_warning": if task["status"] == "running" or item["status"] != "completed_warning":
return None 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()) item.update(warning=None, started_at=item["started_at"] or _now())
self._event_unlocked(task, item_id, "info", "Cover retry queued") self._event_unlocked(task, item_id, "info", "Cover retry queued")
return self.queue_processing(item_id, stage="Waiting for ffmpeg") return self.queue_processing(item_id, stage="Waiting for ffmpeg")
@@ -277,10 +342,15 @@ class TaskState:
warning=None, warning=None,
output_path=None, output_path=None,
cover_path=None, cover_path=None,
resume_phase="download",
started_at=None, started_at=None,
completed_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]: def _find_item_unlocked(self, item_id: int) -> dict[str, Any]:
return self._find_task_and_item_unlocked(item_id)[1] return self._find_task_and_item_unlocked(item_id)[1]
@@ -334,6 +404,9 @@ class TaskState:
self._next_log_id += 1 self._next_log_id += 1
def _finalize_task_unlocked(self, task: dict[str, Any]) -> None: 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"}: if task["status"] not in {"queued", "running"}:
return return
statuses = {item["status"] for item in task["items"]} statuses = {item["status"] for item in task["items"]}
@@ -351,3 +424,11 @@ class TaskState:
status = "failed" status = "failed"
task.update(status=status, completed_at=_now()) task.update(status=status, completed_at=_now())
self._event_unlocked(task, None, "info" if status == "completed" else "warning", f"Task {status}") 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")
+28 -1
View File
@@ -101,6 +101,15 @@ class DownloadService:
def retry_item_cover(self, item_id: int) -> bool: def retry_item_cover(self, item_id: int) -> bool:
return self.runner.retry_item_cover(item_id) 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]: def logs(self, level: str | None) -> dict[str, Any]:
try: try:
return {"logs": self.store.logs(level)} return {"logs": self.store.logs(level)}
@@ -207,13 +216,31 @@ class ServiceRequestHandler(BaseHTTPRequestHandler):
return return
if self.path.endswith("/cancel") and self.path.startswith("/api/tasks/"): if self.path.endswith("/cancel") and self.path.startswith("/api/tasks/"):
task_id = _id_from_path(self.path[: -len("/cancel")], "/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") self._error(HTTPStatus.CONFLICT, "Task cannot be cancelled")
else: else:
task = self.server.service.get_task(task_id) task = self.server.service.get_task(task_id)
assert task is not None assert task is not None
self._json(HTTPStatus.OK, task) self._json(HTTPStatus.OK, task)
return 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/"): if self.path.endswith("/retry-video") and self.path.startswith("/api/tasks/"):
task_id = _id_from_path(self.path[: -len("/retry-video")], "/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): if task_id is None or not self.server.service.retry_task_video(task_id):
+104 -12
View File
@@ -43,6 +43,11 @@ class PostprocessJob:
output_path: Path output_path: Path
result: SegmentDownloadResult | None = None result: SegmentDownloadResult | None = None
cover_only: bool = False cover_only: bool = False
generation: int = 0
class TaskPausedError(Exception):
"""Raised when a task reaches a pause checkpoint."""
class TaskRunner: 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._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._ffmpeg_executor = concurrent.futures.ThreadPoolExecutor(max_workers=self.worker_capacity, thread_name_prefix="m3u8downloaderd-ffmpeg")
self._jobs: dict[int, PostprocessJob] = {} self._jobs: dict[int, PostprocessJob] = {}
self._pending: deque[int] = deque() self._pending: deque[tuple[int, int]] = deque()
self._jobs_lock = threading.RLock() self._jobs_lock = threading.RLock()
self._processes: set[subprocess.Popen[str]] = set() self._processes: dict[subprocess.Popen[str], int] = {}
self._processes_lock = threading.Lock() self._processes_lock = threading.Lock()
def start(self) -> None: def start(self) -> None:
@@ -88,18 +93,57 @@ class TaskRunner:
for job in jobs: for job in jobs:
self._cleanup_job(job) 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: def retry_item_video(self, item_id: int) -> bool:
task, _ = self._task_item(item_id)
if not self.store.retry_item_video(item_id): if not self.store.retry_item_video(item_id):
return False return False
self._discard_job(item_id) self._discard_job(item_id)
self._cleanup_workspace(task, item_id)
return True return True
def retry_task_video(self, task_id: int) -> bool: def retry_task_video(self, task_id: int) -> bool:
item_ids = self.store.retry_task_video(task_id) item_ids = self.store.retry_task_video(task_id)
if not item_ids: if not item_ids:
return False return False
task = self.store.get_task(task_id)
for item_id in item_ids: for item_id in item_ids:
self._discard_job(item_id) self._discard_job(item_id)
if task is not None:
self._cleanup_workspace(task, item_id)
return True return True
def retry_item_processing(self, item_id: int) -> bool: def retry_item_processing(self, item_id: int) -> bool:
@@ -141,7 +185,7 @@ class TaskRunner:
self._jobs[item_id] = job self._jobs[item_id] = job
else: else:
job.cover_only = True job.cover_only = True
self._pending.append(item_id) self._enqueue_existing(item_id)
return True return True
def _dispatch_downloads(self) -> None: def _dispatch_downloads(self) -> None:
@@ -191,7 +235,7 @@ class TaskRunner:
progress: dict[str, tuple[int, int]] = {} progress: dict[str, tuple[int, int]] = {}
def on_event(event: object) -> bool | None: 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 return False
kind = getattr(event, "kind", None) kind = getattr(event, "kind", None)
track_id = getattr(event, "track_id", None) or "unknown" 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") 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) 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}" 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)) 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) 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): if self._stop.is_set() or self.store.cancel_requested(task_id):
raise DownloadCancelledError("Download cancelled by user") 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"])) 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)) self._remember_job(PostprocessJob(task_id, item_id, output_dir, output_path, result=result))
except DownloadCancelledError: except DownloadCancelledError:
if self.store.pause_requested(task_id):
self.store.mark_item_paused(item_id, "download")
return
if workspace is not None: if workspace is not None:
shutil.rmtree(workspace, ignore_errors=True) shutil.rmtree(workspace, ignore_errors=True)
self._cancel_item(task_id, item_id) self._cancel_item(task_id, item_id)
@@ -239,6 +288,10 @@ class TaskRunner:
self._cancel_item(job.task_id, job.item_id) self._cancel_item(job.task_id, job.item_id)
self._discard_job(item_id) self._discard_job(item_id)
return 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: try:
if job.cover_only: if job.cover_only:
self._retry_cover_job(job, task) self._retry_cover_job(job, task)
@@ -247,6 +300,9 @@ class TaskRunner:
except DownloadCancelledError: except DownloadCancelledError:
self._cancel_item(job.task_id, job.item_id) self._cancel_item(job.task_id, job.item_id)
self._discard_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: except Exception as error:
self._mark_failed(job.task_id, job.item_id, str(error), "Processing failed") self._mark_failed(job.task_id, job.item_id, str(error), "Processing failed")
finally: finally:
@@ -364,7 +420,7 @@ class TaskRunner:
command = [*command[:-1], "-progress", "pipe:1", "-nostats", command[-1]] 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) process = subprocess.Popen(command, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, start_new_session=True)
with self._processes_lock: with self._processes_lock:
self._processes.add(process) self._processes[process] = job.task_id
output: list[str] = [] output: list[str] = []
try: try:
assert process.stdout is not None 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): if self._stop.is_set() or self.store.cancel_requested(job.task_id):
self._terminate_process(process) self._terminate_process(process)
raise DownloadCancelledError("Processing cancelled") 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("=") key, separator, value = line.strip().partition("=")
if separator and key in {"out_time_us", "out_time_ms"} and duration > 0: if separator and key in {"out_time_us", "out_time_ms"} and duration > 0:
try: try:
progress(max(0.0, min(1.0, float(value) / 1_000_000 / duration))) progress(max(0.0, min(1.0, float(value) / 1_000_000 / duration)))
except ValueError: except ValueError:
pass pass
if self.store.pause_requested(job.task_id):
raise TaskPausedError("Processing paused")
if process.wait(): if process.wait():
raise RuntimeError(f"{summary}: {''.join(output).strip()[-1500:] or summary}") raise RuntimeError(f"{summary}: {''.join(output).strip()[-1500:] or summary}")
progress(1.0) progress(1.0)
finally: finally:
with self._processes_lock: 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: 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) 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: def _remember_job(self, job: PostprocessJob) -> None:
with self._jobs_lock: with self._jobs_lock:
self._jobs[job.item_id] = job self._jobs[job.item_id] = job
self._pending.append(job.item_id)
self.store.queue_processing(job.item_id) self.store.queue_processing(job.item_id)
self._enqueue_existing(job.item_id)
def _enqueue_existing(self, item_id: int) -> None: def _enqueue_existing(self, item_id: int) -> None:
with self._jobs_lock: with self._jobs_lock:
if item_id in self._jobs: job = self._jobs.get(item_id)
self._pending.append(item_id) if job is not None:
job.generation += 1
self._pending.append((item_id, job.generation))
def _next_postprocess_item(self) -> int | None: def _next_postprocess_item(self) -> int | None:
with self._jobs_lock: with self._jobs_lock:
while self._pending: while self._pending:
item_id = self._pending.popleft() item_id, generation = self._pending.popleft()
job = self._jobs.get(item_id) 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 continue
if self.store.cancel_requested(job.task_id): if self.store.cancel_requested(job.task_id):
self._jobs.pop(item_id, None) self._jobs.pop(item_id, None)
@@ -480,6 +545,33 @@ class TaskRunner:
for process in processes: for process in processes:
self._terminate_process(process) 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 @staticmethod
def _terminate_process(process: subprocess.Popen[str]) -> None: def _terminate_process(process: subprocess.Popen[str]) -> None:
if process.poll() is not None: if process.poll() is not None:
+57
View File
@@ -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]["proxy"] == "http://127.0.0.1:7890"
assert task[1]["output_dir"] == str(config.download_root / "batch") assert task[1]["output_dir"] == str(config.download_root / "batch")
assert "m3u8_referer" not in task[1]["items"][0] 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") service.store.log("error", "Expected log entry")
logs = _get(base, "/api/logs?level=error") logs = _get(base, "/api/logs?level=error")
assert logs[0] == 200 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 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: def test_runner_runs_multiple_videos_from_one_task_concurrently(tmp_path: Path) -> None:
store = TaskState(tmp_path, max_active_downloads=2) store = TaskState(tmp_path, max_active_downloads=2)
runner = TaskRunner(store, worker_capacity=2, ffmpeg="ffmpeg") runner = TaskRunner(store, worker_capacity=2, ffmpeg="ffmpeg")