diff --git a/N_m3u8DL_py/README.md b/N_m3u8DL_py/README.md index b6749d5..0ec2f59 100644 --- a/N_m3u8DL_py/README.md +++ b/N_m3u8DL_py/README.md @@ -49,6 +49,15 @@ result = client.download_url( DownloadRequest(output_dir=Path("downloads")), ) +# Download ordered segment files without merging tracks. The caller owns and +# removes result.temporary_dir after its own post-processing is complete. +segments = client.download_segments( + media, + DownloadRequest(output_dir=Path("downloads"), temporary_dir=Path("work")), +) +for track in segments.tracks: + print(track.track_id, track.initialization, track.segments) + video = next(track for track in media.tracks if track.resolution) video_only = client.download( media, diff --git a/N_m3u8DL_py/src/n_m3u8dl_py/__init__.py b/N_m3u8DL_py/src/n_m3u8dl_py/__init__.py index 6744a26..ea4b2cd 100644 --- a/N_m3u8DL_py/src/n_m3u8dl_py/__init__.py +++ b/N_m3u8DL_py/src/n_m3u8dl_py/__init__.py @@ -3,6 +3,7 @@ from .api import ( DownloadRequest, DownloadResult, + DownloadedTrackSegments, DownloadedFile, MediaInfo, MediaTrack, @@ -11,6 +12,7 @@ from .api import ( N_m3u8DL, RequestOptions, TrackSelection, + SegmentDownloadResult, ) from .errors import DownloadCancelledError, DownloadError, ManifestError, N_m3u8DLError, SelectionError from .events import DownloadEvent, DownloadEventKind @@ -25,6 +27,7 @@ __all__ = [ "DownloadEventKind", "DownloadRequest", "DownloadResult", + "DownloadedTrackSegments", "DownloadedFile", "EncryptMethod", "ManifestError", @@ -37,5 +40,6 @@ __all__ = [ "N_m3u8DLError", "RequestOptions", "SelectionError", + "SegmentDownloadResult", "TrackSelection", ] diff --git a/N_m3u8DL_py/src/n_m3u8dl_py/api.py b/N_m3u8DL_py/src/n_m3u8dl_py/api.py index 3df0817..3022a75 100644 --- a/N_m3u8DL_py/src/n_m3u8dl_py/api.py +++ b/N_m3u8DL_py/src/n_m3u8dl_py/api.py @@ -176,6 +176,32 @@ class DownloadResult: temporary_dir: Path +@dataclass(frozen=True) +class DownloadedTrackSegments: + """Downloaded segment files for one selected track, in playback order.""" + + track_id: str + media_type: MediaType + duration: float + directory: Path + initialization: Path | None + segments: tuple[Path, ...] + + +@dataclass(frozen=True) +class SegmentDownloadResult: + """Result of :meth:`N_m3u8DL.download_segments`. + + The caller owns ``temporary_dir`` and must remove it after consuming the + returned segment paths. + """ + + media_info: MediaInfo + selected_track_ids: tuple[str, ...] + tracks: tuple[DownloadedTrackSegments, ...] + temporary_dir: Path + + @dataclass(frozen=True) class MuxResult: """Output generated by :meth:`N_m3u8DL.mux`.""" @@ -242,6 +268,100 @@ class N_m3u8DL: media_info = self.inspect(input_value, options=options, on_event=on_event) return self.download(media_info, request, on_event=on_event) + def download_segments_url( + self, + input_value: str | Path, + request: DownloadRequest | None = None, + *, + options: RequestOptions | None = None, + on_event: EventCallback | None = None, + ) -> SegmentDownloadResult: + """Inspect and download segment files without merging media tracks.""" + media_info = self.inspect(input_value, options=options, on_event=on_event) + return self.download_segments(media_info, request, on_event=on_event) + + def download_segments( + self, + media_info: MediaInfo, + request: DownloadRequest | None = None, + *, + on_event: EventCallback | None = None, + ) -> SegmentDownloadResult: + """Download selected tracks into ordered segment directories. + + This method never invokes ffmpeg and deliberately leaves the temporary + directory in place for the caller to process and remove. + """ + request = request or DownloadRequest() + try: + self._emit(on_event, DownloadEventKind.DOWNLOAD_STARTED, message=f"Preparing {media_info.protocol} segment download") + selected_ids, selected_streams = self._select_streams(media_info, request.selection) + media_info._source.fetch_playlists(selected_streams) + self._emit(on_event, DownloadEventKind.PLAYLIST_LOADED, message=f"Loaded {len(selected_streams)} selected playlists") + apply_custom_range(selected_streams, request.custom_range) + clean_ads(selected_streams, list(request.ad_keywords)) + if not all(stream.playlist and stream.playlist.segments for stream in selected_streams): + raise SelectionError("One or more selected tracks have no downloadable media segments") + output_dir = request.output_dir + save_name = valid_filename(request.file_name or inferred_name(media_info.input_value), 180) + temporary_root = (request.temporary_dir or output_dir / ".n_m3u8dl") / save_name + options = DownloadOptions( + tmp_dir=temporary_root, + save_dir=output_dir, + save_name=save_name, + save_pattern=request.save_pattern, + thread_count=request.thread_count, + retries=request.retry_count, + binary_merge=False, + skip_merge=True, + del_after_done=False, + check_segments_count=request.check_segments_count, + max_speed=request.max_speed, + subtitle_format=request.subtitle_format, + auto_subtitle_fix=request.auto_subtitle_fix, + ffmpeg=None, + use_ffmpeg_concat_demuxer=False, + decryption_engine=request.decryption_engine, + decryption_binary=None, + keys=[], + ) + track_ids = {id(stream): track_id for track_id, stream in media_info._streams_by_id.items()} + manager = DownloadManager( + HttpClient( + dict(media_info._request_options.headers), + media_info._request_options.timeout, + media_info._request_options.proxy, + media_info._request_options.use_system_proxy, + ), + options, + event_callback=on_event, + track_ids=track_ids, + ) + downloaded = manager.download_segments(selected_streams, concurrent_tracks=False) + tracks = tuple( + DownloadedTrackSegments( + track_id=track_id, + media_type=self._media_type(stream), + duration=stream.playlist.total_duration if stream.playlist else 0.0, + directory=result.directory, + initialization=result.initialization, + segments=result.segments, + ) + for track_id, stream, result in zip(selected_ids, selected_streams, downloaded, strict=True) + ) + except DownloadCancelledError: + self._emit_cancelled(on_event) + raise + except N_m3u8DLError as error: + self._emit(on_event, DownloadEventKind.FAILED, message=str(error)) + raise + except Exception as error: + self._emit(on_event, DownloadEventKind.FAILED, message=str(error)) + raise DownloadError(str(error)) from error + result = SegmentDownloadResult(media_info, tuple(selected_ids), tracks, temporary_root) + self._emit(on_event, DownloadEventKind.DOWNLOAD_COMPLETED, message=f"Downloaded segments for {len(tracks)} track(s)") + return result + def download(self, media_info: MediaInfo, request: DownloadRequest | None = None, *, on_event: EventCallback | None = None) -> DownloadResult: request = request or DownloadRequest() try: diff --git a/N_m3u8DL_py/src/n_m3u8dl_py/downloader.py b/N_m3u8DL_py/src/n_m3u8dl_py/downloader.py index 2e28ac2..12b1bba 100644 --- a/N_m3u8DL_py/src/n_m3u8dl_py/downloader.py +++ b/N_m3u8DL_py/src/n_m3u8dl_py/downloader.py @@ -39,6 +39,15 @@ class DownloadOptions: keys: list[str] +@dataclass(frozen=True) +class DownloadedSegments: + """Downloaded, ordered segment files for one media stream.""" + + directory: Path + initialization: Path | None + segments: tuple[Path, ...] + + class DownloadManager: def __init__( self, @@ -67,7 +76,55 @@ class DownloadManager: return [future.result() for future in futures] return [self.download_stream(stream, number) for number, stream in enumerate(streams, 1)] + def download_segments(self, streams: list[StreamSpec], concurrent_tracks: bool = False) -> list[DownloadedSegments]: + """Download streams without merging them or deleting their work directories.""" + self.options.tmp_dir.mkdir(parents=True, exist_ok=True) + self.options.save_dir.mkdir(parents=True, exist_ok=True) + if concurrent_tracks and len(streams) > 1: + workers = min(len(streams), max(1, self.options.thread_count)) + with concurrent.futures.ThreadPoolExecutor(max_workers=workers) as executor: + futures = [executor.submit(self.download_stream_segments, stream, number) for number, stream in enumerate(streams, 1)] + return [future.result() for future in futures] + return [self.download_stream_segments(stream, number) for number, stream in enumerate(streams, 1)] + def download_stream(self, stream: StreamSpec, number: int) -> Path: + downloaded = self.download_stream_segments(stream, number) + work_dir = downloaded.directory + init_path = downloaded.initialization + files = list(downloaded.segments) + if self.options.skip_merge: + self._emit( + DownloadEventKind.TRACK_COMPLETED, + stream, + f"Downloaded segments to {work_dir}", + completed_segments=len(files), + total_segments=len(files), + output_path=work_dir, + ) + return work_dir + output = self._output_path(stream, number) + if stream.media_type is MediaType.SUBTITLES and self.options.auto_subtitle_fix: + self._merge_subtitles(files, output) + else: + merge_files = ([init_path] if init_path else []) + files + force_binary = init_path is not None or self.options.binary_merge + output = self._merge_media([item for item in merge_files if item], output, force_binary, stream) + if any(segment.encrypt_info.method is EncryptMethod.CENC for segment in stream.playlist.segments): + self._decrypt_cenc(output, stream) + self._report(f"Saved {output}") + self._emit( + DownloadEventKind.TRACK_COMPLETED, + stream, + f"Saved {output}", + completed_segments=len(files), + total_segments=len(files), + output_path=output, + ) + if self.options.del_after_done: + shutil.rmtree(work_dir, ignore_errors=True) + return output + + def download_stream_segments(self, stream: StreamSpec, number: int) -> DownloadedSegments: self._check_cancelled() if not stream.playlist or not stream.playlist.segments: raise ValueError(f"No media segments for stream: {stream.describe()}") @@ -79,7 +136,6 @@ class DownloadManager: self._report(f"Downloading {stream.describe()}") self._emit(DownloadEventKind.TRACK_STARTED, stream, f"Downloading {stream.describe()}", total_segments=len(segments)) self._check_cancelled() - downloaded: list[Path] = [] init_path: Path | None = None if stream.playlist.media_init: init_path = work_dir / "00000_init.bin" @@ -112,37 +168,7 @@ class DownloadManager: downloaded = [paths[index] for index in range(1, len(segments) + 1)] if self.options.check_segments_count and len(downloaded) != len(segments): raise RuntimeError(f"Segment count mismatch for {stream.short_name()}") - if self.options.skip_merge: - self._emit( - DownloadEventKind.TRACK_COMPLETED, - stream, - f"Downloaded segments to {work_dir}", - completed_segments=len(segments), - total_segments=len(segments), - output_path=work_dir, - ) - return work_dir - output = self._output_path(stream, number) - if stream.media_type is MediaType.SUBTITLES and self.options.auto_subtitle_fix: - self._merge_subtitles(downloaded, output) - else: - files = ([init_path] if init_path else []) + downloaded - force_binary = init_path is not None or self.options.binary_merge - output = self._merge_media([item for item in files if item], output, force_binary, stream) - if any(segment.encrypt_info.method is EncryptMethod.CENC for segment in segments): - self._decrypt_cenc(output, stream) - self._report(f"Saved {output}") - self._emit( - DownloadEventKind.TRACK_COMPLETED, - stream, - f"Saved {output}", - completed_segments=len(segments), - total_segments=len(segments), - output_path=output, - ) - if self.options.del_after_done: - shutil.rmtree(work_dir, ignore_errors=True) - return output + return DownloadedSegments(work_dir, init_path, tuple(downloaded)) def _download_segment(self, segment: MediaSegment, destination: Path) -> Path: self._check_cancelled() diff --git a/N_m3u8DL_py/tests/test_api.py b/N_m3u8DL_py/tests/test_api.py index c8a9e92..79ac11b 100644 --- a/N_m3u8DL_py/tests/test_api.py +++ b/N_m3u8DL_py/tests/test_api.py @@ -1,5 +1,6 @@ from __future__ import annotations +import shutil from pathlib import Path import pytest @@ -76,6 +77,30 @@ def test_api_download_url_is_one_call_convenience(tmp_path: Path) -> None: assert result.files[0].path.read_bytes() == b"one" +def test_api_downloads_segments_without_merging_or_ffmpeg(tmp_path: Path) -> None: + (tmp_path / "one.ts").write_bytes(b"one") + (tmp_path / "two.ts").write_bytes(b"two") + manifest = tmp_path / "movie.m3u8" + manifest.write_text("#EXTM3U\n#EXTINF:1,\none.ts\n#EXTINF:1,\ntwo.ts\n#EXT-X-ENDLIST\n") + + result = N_m3u8DL().download_segments_url( + manifest, + DownloadRequest( + output_dir=tmp_path / "output", + temporary_dir=tmp_path / "temporary", + file_name="movie", + ffmpeg_path="/not-used-by-download-segments", + ), + ) + + track = result.tracks[0] + assert [path.read_bytes() for path in track.segments] == [b"one", b"two"] + assert track.directory.is_dir() + assert not list((tmp_path / "output").glob("movie.*")) + shutil.rmtree(result.temporary_dir.parent) + assert not result.temporary_dir.exists() + + def test_api_downloads_only_explicit_track_id(tmp_path: Path) -> None: (tmp_path / "video.ts").write_bytes(b"video") (tmp_path / "audio.ts").write_bytes(b"audio") diff --git a/m3u8downloaderd/README.md b/m3u8downloaderd/README.md index 8edbf3f..0166663 100644 --- a/m3u8downloaderd/README.md +++ b/m3u8downloaderd/README.md @@ -27,17 +27,21 @@ python -m m3u8downloaderd \ --host 0.0.0.0 \ --port 8000 \ --download-root /tmp \ - --max-active-downloads 2 + --max-active-downloads 2 \ + --max-active-ffmpeg 2 ``` The service exposes no authentication. Use it only on a trusted network. The folder chooser is restricted to `--download-root` and its descendants. Tasks, progress, history, and settings exist only while the process is running; restarting the service clears them. Downloaded media files remain on disk. -`--max-active-downloads` sets the initial video concurrency. It can be changed -while the service is running from the Settings tab, with a supported range of -1 to 32 videos. Runtime logs are available in the Log tab and remain in memory -only. +`--max-active-downloads` sets the initial video download concurrency. +`--max-active-ffmpeg` sets the initial ffmpeg post-processing concurrency; when +omitted, it defaults to the download concurrency. Both can be changed while the +service is running from the Settings tab, independently, with a supported range +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. Run tests from the repository root: diff --git a/m3u8downloaderd/src/m3u8downloaderd/static/app.js b/m3u8downloaderd/src/m3u8downloaderd/static/app.js index 4872c1d..f95b105 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/static/app.js +++ b/m3u8downloaderd/src/m3u8downloaderd/static/app.js @@ -97,10 +97,13 @@ function renderTask(task) { actions.append(cancel); } if (["partial", "failed", "cancelled"].includes(task.status)) { - const retry = make("button", "secondary-action", "Retry task"); - retry.type = "button"; - retry.addEventListener("click", () => mutate(`/api/tasks/${task.id}/retry`)); - actions.append(retry); + const retryProcessing = make("button", "secondary-action", "Retry processing"); + retryProcessing.type = "button"; + retryProcessing.addEventListener("click", () => mutate(`/api/tasks/${task.id}/retry-processing`)); + const retryVideo = make("button", "secondary-action", "Retry videos"); + retryVideo.type = "button"; + retryVideo.addEventListener("click", () => mutate(`/api/tasks/${task.id}/retry-video`)); + actions.append(retryProcessing, retryVideo); } const progress = taskProgress(task.items); @@ -143,6 +146,7 @@ function renderItem(item) { details.append(make("div", "item-detail", item.stage)); const progress = itemProgress(item); const progressRoot = make("div", "item-progress"); + progressRoot.classList.toggle("is-processing", item.progress_phase === "processing"); const caption = make("div", "progress-caption"); caption.append(make("span", "", progress.caption)); caption.append(make("span", "", progress.total ? `${progress.percent}%` : "Waiting")); @@ -155,11 +159,22 @@ function renderItem(item) { if (item.error) details.append(make("div", "item-detail item-warning", item.error)); row.append(details); const actions = make("div", "item-actions"); - if (["failed", "cancelled", "completed_warning"].includes(item.status)) { - const retry = make("button", "secondary-action", item.status === "completed_warning" ? "Retry cover" : "Retry video"); - retry.type = "button"; - retry.addEventListener("click", () => mutate(`/api/items/${item.id}/retry`)); - actions.append(retry); + if (item.status === "completed_warning") { + const retryCover = make("button", "secondary-action", "Retry cover"); + retryCover.type = "button"; + retryCover.addEventListener("click", () => mutate(`/api/items/${item.id}/retry-cover`)); + actions.append(retryCover); + } else if (["failed", "cancelled"].includes(item.status)) { + if (item.status === "failed" && item.progress_phase === "processing") { + const retryProcessing = make("button", "secondary-action", "Retry processing"); + retryProcessing.type = "button"; + retryProcessing.addEventListener("click", () => mutate(`/api/items/${item.id}/retry-processing`)); + actions.append(retryProcessing); + } + const retryVideo = make("button", "secondary-action", "Retry video"); + retryVideo.type = "button"; + retryVideo.addEventListener("click", () => mutate(`/api/items/${item.id}/retry-video`)); + actions.append(retryVideo); } if (item.output_path) actions.append(make("span", "item-detail", "MP4 ready")); row.append(actions); @@ -176,6 +191,15 @@ function taskProgress(items) { } function itemProgress(item) { + if (item.progress_phase === "processing") { + const total = item.processing_total || 100; + const completed = item.processing_completed || 0; + return { + total, + percent: Math.min(100, Math.round((completed / total) * 100)), + caption: `Processing ${Math.min(100, Math.round((completed / total) * 100))}%`, + }; + } const total = item.total_segments || 0; const completed = item.completed_segments || 0; if (total) { @@ -283,6 +307,7 @@ async function submitSettings(event) { default_directory: $("#setting-directory").value, proxy: $("#setting-proxy").value, max_active_downloads: Number($("#setting-max-downloads").value), + max_active_ffmpeg: Number($("#setting-max-ffmpeg").value), }), }); state.bootstrap.settings = settings; @@ -334,6 +359,7 @@ async function initialize() { $("#setting-directory").value = settings.default_directory; $("#setting-proxy").value = settings.proxy; $("#setting-max-downloads").value = settings.max_active_downloads; + $("#setting-max-ffmpeg").value = settings.max_active_ffmpeg; window.setInterval(() => refreshBootstrap().catch(() => {}), 1000); } diff --git a/m3u8downloaderd/src/m3u8downloaderd/static/index.html b/m3u8downloaderd/src/m3u8downloaderd/static/index.html index 5a9cd07..752dcd7 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/static/index.html +++ b/m3u8downloaderd/src/m3u8downloaderd/static/index.html @@ -60,6 +60,10 @@ Concurrent video downloads +
diff --git a/m3u8downloaderd/src/m3u8downloaderd/static/styles.css b/m3u8downloaderd/src/m3u8downloaderd/static/styles.css index 9906562..e5f1ab2 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/static/styles.css +++ b/m3u8downloaderd/src/m3u8downloaderd/static/styles.css @@ -77,6 +77,7 @@ progress { width: 100%; height: 7px; accent-color: #0d917c; } .item-detail { color: #60797e; font-size: .78rem; overflow-wrap: anywhere; } .item-progress { margin-top: 8px; max-width: 620px; } .item-progress progress { height: 5px; accent-color: #2378a8; } +.item-progress.is-processing progress { accent-color: #b45288; } .item-warning { color: #a3462b; } .item-actions { display: flex; align-items: center; gap: 8px; } .item-actions .secondary-action { min-height: 30px; padding: 0 9px; font-size: .78rem; } diff --git a/m3u8downloaderd/src/m3u8downloaderd/tasks.py b/m3u8downloaderd/src/m3u8downloaderd/tasks.py index 35d17e7..63f41f1 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/tasks.py +++ b/m3u8downloaderd/src/m3u8downloaderd/tasks.py @@ -22,13 +22,14 @@ def _now() -> str: class TaskState: """Keeps settings and task history only for the current process lifetime.""" - def __init__(self, download_root: Path, max_active_downloads: int = 2) -> None: + def __init__(self, download_root: Path, max_active_downloads: int = 2, max_active_ffmpeg: int | None = None) -> None: self._lock = threading.RLock() self._settings = { "default_title_template": DATE_TEMPLATE, "default_directory": str(download_root.resolve()), "proxy": DEFAULT_PROXY, "max_active_downloads": 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._logs: list[dict[str, Any]] = [] @@ -49,6 +50,10 @@ class TaskState: with self._lock: return int(self._settings["max_active_downloads"]) + def max_active_ffmpeg(self) -> int: + with self._lock: + return int(self._settings["max_active_ffmpeg"]) + def logs(self, level: str | None = None) -> list[dict[str, Any]]: normalized = (level or "").lower() if normalized and normalized not in LOG_LEVELS: @@ -104,6 +109,9 @@ class TaskState: "stage": "Queued", "completed_segments": 0, "total_segments": 0, + "progress_phase": "download", + "processing_completed": 0.0, + "processing_total": 0.0, "message": "", "error": None, "warning": None, @@ -159,7 +167,7 @@ class TaskState: now = _now() task["cancel_requested"] = True for item in task["items"]: - if item["status"] == "queued": + if item["status"] in {"queued", "processing_queued"}: item.update(status="cancelled", stage="Cancelled", completed_at=now) self._event_unlocked(task, None, "info", "Cancellation requested") self._finalize_task_unlocked(task) @@ -174,50 +182,105 @@ class TaskState: task = self._tasks[task_id] self._finalize_task_unlocked(task) - def retry_task(self, task_id: int) -> bool: + def retry_task_video(self, task_id: int) -> list[int]: with self._lock: task = self._tasks.get(task_id) if task is None or task["status"] == "running": - return False + return [] retryable = [item for item in task["items"] if item["status"] in {"failed", "cancelled", "completed_warning"}] if not retryable: - return False + return [] for item in retryable: - item.update( - status="queued", - stage="Queued", - completed_segments=0, - total_segments=0, - message="", - error=None, - warning=None, - started_at=None, - completed_at=None, - ) + self._reset_for_video_retry_unlocked(item) task.update(status="queued", cancel_requested=False, started_at=None, completed_at=None, error=None) self._event_unlocked(task, None, "info", "Task retry queued") - return True + return [int(item["id"]) for item in retryable] - def retry_item(self, item_id: int) -> bool: + def retry_item_video(self, item_id: int) -> bool: with self._lock: task, item = self._find_task_and_item_unlocked(item_id) if task["status"] == "running" or item["status"] == "completed": return False - item.update( - status="queued", - stage="Cover retry" if item["status"] == "completed_warning" else "Queued", - completed_segments=0, - total_segments=0, - message="", - error=None, - warning=None, - started_at=None, - completed_at=None, - ) + self._reset_for_video_retry_unlocked(item) task.update(status="queued", cancel_requested=False, started_at=None, completed_at=None, error=None) 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"]: + raise RuntimeError("Task is cancelled") + item.update( + status="processing_queued", + stage=stage, + progress_phase="processing", + processing_completed=0.0, + processing_total=100.0, + message=stage, + error=None, + completed_at=None, + ) + return copy.deepcopy(task), copy.deepcopy(item) + + def retry_item_processing(self, item_id: int) -> tuple[dict[str, Any], dict[str, Any]] | None: + with self._lock: + 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) + 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) + + def retry_task_processing(self, task_id: int) -> list[tuple[dict[str, Any], dict[str, Any]]]: + with self._lock: + task = self._tasks.get(task_id) + if task is None or task["status"] == "running": + return [] + 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) + claimed: list[tuple[dict[str, Any], dict[str, Any]]] = [] + for item in retryable: + item.update(warning=None, started_at=item["started_at"] or _now()) + if item["status"] == "completed_warning": + item.update(status="processing_queued", stage="Waiting for ffmpeg", progress_phase="processing", processing_completed=0.0, processing_total=100.0, message="Waiting for ffmpeg", completed_at=None) + else: + item.update(status="processing_queued", stage="Waiting for ffmpeg", progress_phase="processing", processing_completed=0.0, processing_total=100.0, message="Waiting for ffmpeg", error=None, completed_at=None) + claimed.append((copy.deepcopy(task), copy.deepcopy(item))) + self._event_unlocked(task, None, "info", "Task processing retry queued") + return claimed + + def retry_item_cover(self, item_id: int) -> tuple[dict[str, Any], dict[str, Any]] | None: + with self._lock: + 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) + 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") + + def _reset_for_video_retry_unlocked(self, item: dict[str, Any]) -> None: + item.update( + status="queued", + stage="Queued", + completed_segments=0, + total_segments=0, + progress_phase="download", + processing_completed=0.0, + processing_total=0.0, + message="", + error=None, + warning=None, + output_path=None, + cover_path=None, + started_at=None, + completed_at=None, + ) + def _find_item_unlocked(self, item_id: int) -> dict[str, Any]: return self._find_task_and_item_unlocked(item_id)[1] @@ -274,7 +337,7 @@ class TaskState: if task["status"] not in {"queued", "running"}: return statuses = {item["status"] for item in task["items"]} - if statuses & {"queued", "running"}: + if statuses & {"queued", "running", "processing_queued"}: return if task["cancel_requested"] and statuses <= {"completed", "completed_warning", "cancelled"}: status = "cancelled" diff --git a/m3u8downloaderd/src/m3u8downloaderd/web.py b/m3u8downloaderd/src/m3u8downloaderd/web.py index 7162bc0..63a620a 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/web.py +++ b/m3u8downloaderd/src/m3u8downloaderd/web.py @@ -35,6 +35,7 @@ class ServiceConfig: port: int download_root: Path max_active_downloads: int + max_active_ffmpeg: int | None = None class DownloadService: @@ -42,8 +43,8 @@ class DownloadService: root = config.download_root.expanduser().resolve() if not root.is_dir(): raise RuntimeError(f"Download root does not exist or is not a directory: {root}") - self.config = ServiceConfig(config.host, config.port, root, config.max_active_downloads) - self.store = TaskState(root, self.config.max_active_downloads) + self.config = ServiceConfig(config.host, config.port, root, config.max_active_downloads, config.max_active_ffmpeg) + self.store = TaskState(root, self.config.max_active_downloads, self.config.max_active_ffmpeg) self.runner = TaskRunner(self.store) def start(self) -> None: @@ -67,21 +68,39 @@ class DownloadService: directory = value.get("default_directory", current["default_directory"]) proxy = value.get("proxy", current["proxy"]) max_active_downloads = value.get("max_active_downloads", current["max_active_downloads"]) + max_active_ffmpeg = value.get("max_active_ffmpeg", current["max_active_ffmpeg"]) if not isinstance(template, str): raise ValidationError("Default title must be a string") selected = ensure_within_root(directory, self.config.download_root) normalized_proxy = validate_proxy(proxy) normalized_max_active_downloads = validate_max_active_downloads(max_active_downloads) + normalized_max_active_ffmpeg = validate_max_active_downloads(max_active_ffmpeg) settings = { "default_title_template": template, "default_directory": str(selected), "proxy": normalized_proxy, "max_active_downloads": normalized_max_active_downloads, + "max_active_ffmpeg": normalized_max_active_ffmpeg, } self.store.update_settings(settings) - self.store.log("info", f"Settings updated: concurrent downloads={normalized_max_active_downloads}") + self.store.log("info", f"Settings updated: concurrent downloads={normalized_max_active_downloads}, ffmpeg={normalized_max_active_ffmpeg}") return settings + def retry_task_video(self, task_id: int) -> bool: + return self.runner.retry_task_video(task_id) + + def retry_task_processing(self, task_id: int) -> bool: + return self.runner.retry_task_processing(task_id) + + def retry_item_video(self, item_id: int) -> bool: + return self.runner.retry_item_video(item_id) + + def retry_item_processing(self, item_id: int) -> bool: + return self.runner.retry_item_processing(item_id) + + def retry_item_cover(self, item_id: int) -> bool: + return self.runner.retry_item_cover(item_id) + def logs(self, level: str | None) -> dict[str, Any]: try: return {"logs": self.store.logs(level)} @@ -195,9 +214,48 @@ class ServiceRequestHandler(BaseHTTPRequestHandler): 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): + self._error(HTTPStatus.CONFLICT, "Task cannot be retried") + 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-processing") and self.path.startswith("/api/tasks/"): + task_id = _id_from_path(self.path[: -len("/retry-processing")], "/api/tasks/") + if task_id is None or not self.server.service.retry_task_processing(task_id): + self._error(HTTPStatus.CONFLICT, "Task processing cannot be retried") + 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/items/"): + item_id = _id_from_path(self.path[: -len("/retry-video")], "/api/items/") + if item_id is None or not self.server.service.retry_item_video(item_id): + self._error(HTTPStatus.CONFLICT, "Video cannot be retried") + else: + self._json(HTTPStatus.OK, {"ok": True}) + return + if self.path.endswith("/retry-processing") and self.path.startswith("/api/items/"): + item_id = _id_from_path(self.path[: -len("/retry-processing")], "/api/items/") + if item_id is None or not self.server.service.retry_item_processing(item_id): + self._error(HTTPStatus.CONFLICT, "Video processing cannot be retried") + else: + self._json(HTTPStatus.OK, {"ok": True}) + return + if self.path.endswith("/retry-cover") and self.path.startswith("/api/items/"): + item_id = _id_from_path(self.path[: -len("/retry-cover")], "/api/items/") + if item_id is None or not self.server.service.retry_item_cover(item_id): + self._error(HTTPStatus.CONFLICT, "Cover cannot be retried") + else: + self._json(HTTPStatus.OK, {"ok": True}) + return if self.path.endswith("/retry") and self.path.startswith("/api/tasks/"): task_id = _id_from_path(self.path[: -len("/retry")], "/api/tasks/") - if task_id is None or not self.server.service.store.retry_task(task_id): + if task_id is None or not self.server.service.retry_task_video(task_id): self._error(HTTPStatus.CONFLICT, "Task cannot be retried") else: task = self.server.service.get_task(task_id) @@ -206,10 +264,12 @@ class ServiceRequestHandler(BaseHTTPRequestHandler): return if self.path.endswith("/retry") and self.path.startswith("/api/items/"): item_id = _id_from_path(self.path[: -len("/retry")], "/api/items/") - if item_id is None or not self.server.service.store.retry_item(item_id): + if item_id is None: self._error(HTTPStatus.CONFLICT, "Video cannot be retried") - else: + elif self.server.service.retry_item_video(item_id): self._json(HTTPStatus.OK, {"ok": True}) + else: + self._error(HTTPStatus.CONFLICT, "Video cannot be retried") return self._error(HTTPStatus.NOT_FOUND, "Not found") @@ -311,6 +371,7 @@ def build_argument_parser() -> argparse.ArgumentParser: parser.add_argument("--port", type=int, default=8000) parser.add_argument("--download-root", type=Path, default=Path("/tmp")) parser.add_argument("--max-active-downloads", type=int, default=2) + parser.add_argument("--max-active-ffmpeg", type=int) return parser @@ -319,7 +380,10 @@ def config_from_args(args: argparse.Namespace) -> ServiceConfig: raise ValueError("--port must be between 1 and 65535") if not 1 <= args.max_active_downloads <= MAX_ACTIVE_DOWNLOADS: raise ValueError(f"--max-active-downloads must be between 1 and {MAX_ACTIVE_DOWNLOADS}") - return ServiceConfig(args.host, args.port, args.download_root, args.max_active_downloads) + max_active_ffmpeg = args.max_active_downloads if args.max_active_ffmpeg is None else args.max_active_ffmpeg + if not 1 <= max_active_ffmpeg <= MAX_ACTIVE_DOWNLOADS: + raise ValueError(f"--max-active-ffmpeg must be between 1 and {MAX_ACTIVE_DOWNLOADS}") + return ServiceConfig(args.host, args.port, args.download_root, args.max_active_downloads, max_active_ffmpeg) def _id_from_path(path: str, prefix: str) -> int | None: diff --git a/m3u8downloaderd/src/m3u8downloaderd/worker.py b/m3u8downloaderd/src/m3u8downloaderd/worker.py index a35978d..226b23a 100644 --- a/m3u8downloaderd/src/m3u8downloaderd/worker.py +++ b/m3u8downloaderd/src/m3u8downloaderd/worker.py @@ -1,4 +1,4 @@ -"""Durable task scheduler and media post-processing pipeline.""" +"""Concurrent download scheduling and ffmpeg post-processing.""" from __future__ import annotations @@ -6,16 +6,28 @@ import concurrent.futures import os import re import shutil +import signal import subprocess -import tempfile import threading import time +from collections import deque +from dataclasses import dataclass from pathlib import Path +from typing import Callable from urllib.error import HTTPError, URLError from urllib.parse import urlsplit, urlunsplit from urllib.request import Request -from n_m3u8dl_py import DownloadCancelledError, DownloadEventKind, DownloadRequest, N_m3u8DL, RequestOptions +from n_m3u8dl_py import ( + DownloadCancelledError, + DownloadEventKind, + DownloadRequest, + DownloadedTrackSegments, + MediaType, + N_m3u8DL, + RequestOptions, + SegmentDownloadResult, +) from n_m3u8dl_py.http import build_proxy_opener from n_m3u8dl_py.utils import valid_filename @@ -23,8 +35,18 @@ from .models import MAX_ACTIVE_DOWNLOADS from .tasks import TaskState +@dataclass +class PostprocessJob: + task_id: int + item_id: int + output_dir: Path + output_path: Path + result: SegmentDownloadResult | None = None + cover_only: bool = False + + class TaskRunner: - """Runs queued video downloads while keeping HTTP request handling independent.""" + """Runs download and ffmpeg queues with independent concurrency limits.""" def __init__(self, store: TaskState, worker_capacity: int = MAX_ACTIVE_DOWNLOADS, ffmpeg: str | None = None) -> None: self.store = store @@ -33,244 +55,248 @@ class TaskRunner: if not self.ffmpeg: raise RuntimeError("ffmpeg is required to run m3u8downloaderd") self._stop = threading.Event() - self._dispatcher: threading.Thread | None = None - self._executor = concurrent.futures.ThreadPoolExecutor( - max_workers=self.worker_capacity, - thread_name_prefix="m3u8downloaderd-download", - ) + self._download_dispatcher: threading.Thread | None = None + self._ffmpeg_dispatcher: threading.Thread | None = None + 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._jobs_lock = threading.RLock() + self._processes: set[subprocess.Popen[str]] = set() + self._processes_lock = threading.Lock() def start(self) -> None: - if self._dispatcher is None: - self._dispatcher = threading.Thread(target=self._dispatch, name="m3u8downloaderd-dispatch", daemon=True) - self._dispatcher.start() + if self._download_dispatcher is None: + self._download_dispatcher = threading.Thread(target=self._dispatch_downloads, name="m3u8downloaderd-download-dispatch", daemon=True) + self._download_dispatcher.start() + if self._ffmpeg_dispatcher is None: + self._ffmpeg_dispatcher = threading.Thread(target=self._dispatch_ffmpeg, name="m3u8downloaderd-ffmpeg-dispatch", daemon=True) + self._ffmpeg_dispatcher.start() def stop(self) -> None: self._stop.set() - if self._dispatcher is not None: - self._dispatcher.join(timeout=5) - self._executor.shutdown(wait=False, cancel_futures=True) + self._terminate_ffmpeg_processes() + for dispatcher in (self._download_dispatcher, self._ffmpeg_dispatcher): + if dispatcher is not None: + dispatcher.join(timeout=5) + self._download_executor.shutdown(wait=False, cancel_futures=True) + self._ffmpeg_executor.shutdown(wait=False, cancel_futures=True) + with self._jobs_lock: + jobs = list(self._jobs.values()) + self._jobs.clear() + self._pending.clear() + for job in jobs: + self._cleanup_job(job) - def _dispatch(self) -> None: + def retry_item_video(self, item_id: int) -> bool: + if not self.store.retry_item_video(item_id): + return False + self._discard_job(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 + for item_id in item_ids: + self._discard_job(item_id) + return True + + def retry_item_processing(self, item_id: int) -> bool: + claimed = self.store.retry_item_processing(item_id) + if claimed is None: + return False + task, item = claimed + if not self._has_job(item_id): + self._mark_missing_intermediate(task, item) + return True + self._enqueue_existing(item_id) + return True + + def retry_task_processing(self, task_id: int) -> bool: + claimed = self.store.retry_task_processing(task_id) + if not claimed: + return False + for task, item in claimed: + if self._has_job(int(item["id"])): + self._enqueue_existing(int(item["id"])) + else: + self._mark_missing_intermediate(task, item) + return True + + def retry_item_cover(self, item_id: int) -> bool: + claimed = self.store.retry_item_cover(item_id) + if claimed is None: + return False + task, item = claimed + output_path = Path(str(item["output_path"])) + if not output_path.is_file(): + self._mark_failed(int(task["id"]), item_id, "The existing MP4 for this cover retry no longer exists") + self.store.finalize_task_if_ready(int(task["id"])) + return True + with self._jobs_lock: + job = self._jobs.get(item_id) + if job is None: + job = PostprocessJob(int(task["id"]), item_id, output_path.parent, output_path, cover_only=True) + self._jobs[item_id] = job + else: + job.cover_only = True + self._pending.append(item_id) + return True + + def _dispatch_downloads(self) -> None: active: set[concurrent.futures.Future[None]] = set() while not self._stop.is_set(): - completed = {future for future in active if future.done()} - active.difference_update(completed) - for future in completed: - try: - future.result() - except Exception as error: - self.store.log("error", f"Download worker crashed: {error}") - while len(active) < self.store.max_active_downloads(): + self._collect(active, "Download worker") + while len(active) < self.store.max_active_downloads() and not self._stop.is_set(): claimed = self.store.claim_next_item() if claimed is None: break - task, item = claimed - active.add(self._executor.submit(self._run_item, task, item)) - self._stop.wait(0.2) + active.add(self._download_executor.submit(self._run_download, *claimed)) + self._stop.wait(0.1) - def _run_item(self, task: dict[str, object], item: dict[str, object]) -> None: - item_id = int(item["id"]) - task_id = int(task["id"]) + def _dispatch_ffmpeg(self) -> None: + active: set[concurrent.futures.Future[None]] = set() + while not self._stop.is_set(): + self._collect(active, "ffmpeg worker") + while len(active) < self.store.max_active_ffmpeg() and not self._stop.is_set(): + item_id = self._next_postprocess_item() + if item_id is None: + break + active.add(self._ffmpeg_executor.submit(self._process_item, item_id)) + self._stop.wait(0.1) + + def _collect(self, active: set[concurrent.futures.Future[None]], label: str) -> None: + completed = {future for future in active if future.done()} + active.difference_update(completed) + for future in completed: + try: + future.result() + except Exception as error: + self.store.log("error", f"{label} crashed: {error}") + + def _run_download(self, task: dict[str, object], item: dict[str, object]) -> None: try: - if item["stage"] == "Cover retry" and item["output_path"]: - self._retry_cover(task, item) - return self._download_and_finalize(task, item) - except DownloadCancelledError: - self.store.update_item( - item_id, - status="cancelled", - stage="Cancelled", - completed_at=_timestamp(), - message="Cancellation completed", - ) - self.store.event(task_id, item_id, "info", "Video cancelled") - except Exception as error: - self.store.update_item( - item_id, - status="failed", - stage="Failed", - error=str(error), - completed_at=_timestamp(), - message="Download failed", - ) - self.store.event(task_id, item_id, "error", str(error)) finally: - self.store.finalize_task_if_ready(task_id) + self.store.finalize_task_if_ready(int(task["id"])) def _download_and_finalize(self, task: dict[str, object], item: dict[str, object]) -> None: task_id = int(task["id"]) item_id = int(item["id"]) output_dir = Path(str(task["output_dir"])) - output_dir.mkdir(parents=True, exist_ok=True) - progress: dict[str, tuple[int, int]] = {} - - def on_event(event: object) -> bool | None: - if self.store.cancel_requested(task_id): - return False - kind = getattr(event, "kind", None) - track_id = getattr(event, "track_id", None) or "unknown" - if kind is DownloadEventKind.TRACK_STARTED: - progress[track_id] = (0, int(getattr(event, "total_segments", 0) or 0)) - elif kind is DownloadEventKind.SEGMENT_COMPLETED: - progress[track_id] = ( - int(getattr(event, "completed_segments", 0) or 0), - int(getattr(event, "total_segments", 0) or 0), - ) - elif kind is DownloadEventKind.TRACK_COMPLETED and track_id in progress: - done, total = progress[track_id] - progress[track_id] = (total or done, total or done) - if progress: - completed = sum(value[0] for value in progress.values()) - total = sum(value[1] for value in progress.values()) - self.store.update_item( - item_id, - stage="Downloading", - completed_segments=completed, - total_segments=total, - message=getattr(event, "message", "Downloading"), - ) - message = str(getattr(event, "message", "")) - if message: - level = "debug" if kind is DownloadEventKind.SEGMENT_COMPLETED else "info" - self.store.log(level, message, task_id, item_id) - return None - - self.store.log("info", "Inspecting media playlist", task_id, item_id) - self.store.update_item(item_id, stage="Downloading", message="Inspecting media playlist") - options = RequestOptions( - headers={"Referer": str(item["m3u8_referer"])}, - proxy=str(task["proxy"]) or None, - use_system_proxy=False, - ) - with tempfile.TemporaryDirectory(prefix=f".m3u8downloaderd-{task_id}-{item_id}-", dir=output_dir) as temporary_dir: - request = DownloadRequest( - output_dir=output_dir, - temporary_dir=temporary_dir, - file_name=str(item["title"]), - save_pattern=f".source_{item_id}_.", - thread_count=max(1, os.cpu_count() or 1), - ffmpeg_path=self.ffmpeg, - ) - result = N_m3u8DL(options).download_url(str(item["m3u8_url"]), request, on_event=on_event) - if self.store.cancel_requested(task_id): - raise DownloadCancelledError("Download cancelled by user") - self.store.log("info", "Creating MP4", task_id, item_id) - self.store.update_item(item_id, stage="Muxing", message="Creating MP4") - output_path = self._reserve_output_path(output_dir, str(item["title"])) - mux_temp = output_dir / f".{output_path.stem}.{item_id}.mux.tmp.mp4" + workspace: Path | None = None try: - self._mux_mp4([file.path for file in result.files], mux_temp) - for source in result.files: - source.path.unlink(missing_ok=True) - self._complete_cover_step(task, item, output_path, mux_temp) - except Exception: - output_path.unlink(missing_ok=True) - mux_temp.unlink(missing_ok=True) - raise + output_dir.mkdir(parents=True, exist_ok=True) + progress: dict[str, tuple[int, int]] = {} - def _retry_cover(self, task: dict[str, object], item: dict[str, object]) -> None: - item_id = int(item["id"]) - output_path = Path(str(item["output_path"])) - if not output_path.is_file(): - raise FileNotFoundError("The existing MP4 for this cover retry no longer exists") - self.store.log("info", "Retrying cover", int(task["id"]), item_id) - self.store.update_item(item_id, stage="Cover", message="Retrying cover") - cover_path = Path(str(item["cover_path"])) if item["cover_path"] else None - if cover_path is None or not cover_path.is_file(): - cover_path = self._download_cover( - str(item["image_src"]), - output_path.parent, - item_id, - str(task["proxy"]), - task_id=int(task["id"]), - ) - covered_temp = output_path.parent / f".{output_path.stem}.{item_id}.cover.tmp.mp4" - try: - self._attach_cover(output_path, cover_path, covered_temp) - os.replace(covered_temp, output_path) - cover_path.unlink(missing_ok=True) - self.store.update_item( - item_id, - status="completed", - stage="Completed", - warning=None, - cover_path=None, - completed_at=_timestamp(), - message="Video and cover completed", - ) + def on_event(event: object) -> bool | None: + if self._stop.is_set() or self.store.cancel_requested(task_id): + return False + kind = getattr(event, "kind", None) + track_id = getattr(event, "track_id", None) or "unknown" + if kind is DownloadEventKind.TRACK_STARTED: + progress[track_id] = (0, int(getattr(event, "total_segments", 0) or 0)) + elif kind is DownloadEventKind.SEGMENT_COMPLETED: + progress[track_id] = (int(getattr(event, "completed_segments", 0) or 0), int(getattr(event, "total_segments", 0) or 0)) + elif kind is DownloadEventKind.TRACK_COMPLETED and track_id in progress: + _, total = progress[track_id] + progress[track_id] = (total, total) + if progress: + self.store.update_item(item_id, stage="Downloading", progress_phase="download", completed_segments=sum(value[0] for value in progress.values()), total_segments=sum(value[1] for value in progress.values()), message=getattr(event, "message", "Downloading")) + message = str(getattr(event, "message", "")) + if message: + self.store.log("debug" if kind is DownloadEventKind.SEGMENT_COMPLETED else "info", message, task_id, item_id) + return None + + 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") + 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 workspace is not None: + shutil.rmtree(workspace, ignore_errors=True) + self._cancel_item(task_id, item_id) except Exception as error: - retained_cover = self._retain_cover(cover_path, output_path) - self.store.update_item( - item_id, - status="completed_warning", - stage="Completed with warning", - warning=str(error), - cover_path=str(retained_cover) if retained_cover else None, - completed_at=_timestamp(), - message="Video completed without an embedded cover", - ) - self.store.event(int(task["id"]), item_id, "warning", f"Cover retry failed: {error}") + if workspace is not None: + shutil.rmtree(workspace, ignore_errors=True) + self._mark_failed(task_id, item_id, str(error), "Download failed") finally: - covered_temp.unlink(missing_ok=True) + self.store.finalize_task_if_ready(task_id) - def _complete_cover_step( - self, - task: dict[str, object], - item: dict[str, object], - output_path: Path, - mux_temp: Path, - ) -> None: - item_id = int(item["id"]) - self.store.log("info", "Downloading cover image", int(task["id"]), item_id) - self.store.update_item(item_id, stage="Cover", message="Downloading cover image") - cover_path: Path | None = None - covered_temp = output_path.parent / f".{output_path.stem}.{item_id}.cover.tmp.mp4" + def _process_item(self, item_id: int) -> None: + job = self._job(item_id) + if job is None: + return + task = self.store.get_task(job.task_id) + if task is None or self._stop.is_set() or self.store.cancel_requested(job.task_id): + self._cancel_item(job.task_id, job.item_id) + self._discard_job(item_id) + return try: - cover_path = self._download_cover( - str(item["image_src"]), - output_path.parent, - item_id, - str(task["proxy"]), - task_id=int(task["id"]), - ) - self.store.log("info", "Embedding cover image", int(task["id"]), item_id) - self.store.update_item(item_id, stage="Cover", message="Embedding cover image") - self._attach_cover(mux_temp, cover_path, covered_temp) - os.replace(covered_temp, output_path) - mux_temp.unlink(missing_ok=True) - cover_path.unlink(missing_ok=True) - self.store.update_item( - item_id, - status="completed", - stage="Completed", - output_path=str(output_path), - cover_path=None, - completed_at=_timestamp(), - message="Video and cover completed", - ) - self.store.event(int(task["id"]), item_id, "info", f"Saved {output_path.name}") + if job.cover_only: + self._retry_cover_job(job, task) + else: + self._postprocess_job(job, task) + except DownloadCancelledError: + self._cancel_item(job.task_id, job.item_id) + self._discard_job(item_id) except Exception as error: - os.replace(mux_temp, output_path) - retained_cover = self._retain_cover(cover_path, output_path) - self.store.update_item( - item_id, - status="completed_warning", - stage="Completed with warning", - output_path=str(output_path), - cover_path=str(retained_cover) if retained_cover else None, - warning=str(error), - completed_at=_timestamp(), - message="Video completed without an embedded cover", - ) - self.store.event(int(task["id"]), item_id, "warning", f"Cover step failed: {error}") + self._mark_failed(job.task_id, job.item_id, str(error), "Processing failed") finally: - covered_temp.unlink(missing_ok=True) + self.store.finalize_task_if_ready(job.task_id) - def _mux_mp4(self, source_paths: list[Path], target: Path) -> None: + def _postprocess_job(self, job: PostprocessJob, task: dict[str, object]) -> None: + if job.result is None or not job.result.temporary_dir.is_dir(): + raise RuntimeError("Downloaded segments are unavailable; retry the video") + item = self._item(job.item_id) + tracks = list(job.result.tracks) + weights = self._weights(tracks) + completed = 0.0 + total = sum(weights) + merged: list[Path] = [] + for track, weight in zip(tracks, weights[:-2], strict=True): + self._update_processing(job, f"Merging {track.media_type.value.lower()} track", completed, total) + output = self._track_output(job, track) + self._merge_track(job, track, output, completed, weight, total) + merged.append(output) + completed += weight + mux_temp = job.output_dir / f".{job.output_path.stem}.{job.item_id}.mux.tmp.mp4" + self._update_processing(job, "Creating MP4", completed, total) + self._mux_mp4(job, merged, mux_temp, completed, weights[-2], total) + self._complete_cover_step(job, task, item, mux_temp, completed + weights[-2], weights[-1], total) + + def _merge_track(self, job: PostprocessJob, track: DownloadedTrackSegments, output: Path, base: float, weight: float, total: float) -> None: + files = [path for path in ([track.initialization] if track.initialization else []) + list(track.segments) if path is not None and path.is_file()] + if not files: + raise RuntimeError(f"No segment files remain for {track.track_id}") + if track.media_type is MediaType.SUBTITLES or track.initialization is not None: + self._merge_binary(files, output) + self._update_processing(job, f"Merged {track.media_type.value.lower()} track", base + weight, total) + return + manifest = output.with_suffix(".concat.txt") + manifest.write_text("".join(f"file '{_concat_escape(path)}'\n" for path in files), encoding="utf-8") + command = [self.ffmpeg, "-hide_banner", "-loglevel", "error", "-nostdin", "-y", "-f", "concat", "-safe", "0", "-i", str(manifest), "-c", "copy", str(output)] + try: + self._run_ffmpeg(job, command, "Track merge failed", lambda fraction: self._update_processing(job, f"Merging {track.media_type.value.lower()} track", base + weight * fraction, total), track.duration) + except RuntimeError as error: + self.store.log("warning", f"ffmpeg concat failed for {track.track_id}; using binary merge: {error}", job.task_id, job.item_id) + output.unlink(missing_ok=True) + self._merge_binary(files, output) + self._update_processing(job, f"Merged {track.media_type.value.lower()} track", base + weight, total) + finally: + manifest.unlink(missing_ok=True) + + def _mux_mp4(self, job: PostprocessJob, source_paths: list[Path], target: Path, base: float, weight: float, total: float) -> None: sources = [path for path in source_paths if path.is_file()] if not sources: - raise RuntimeError("Downloader did not produce media files") + raise RuntimeError("No merged media files are available") command = [self.ffmpeg, "-hide_banner", "-loglevel", "error", "-nostdin", "-y"] for source in sources: command.extend(["-i", str(source)]) @@ -280,35 +306,195 @@ class TaskRunner: if any(path.suffix.lower() == ".srt" for path in sources): command.extend(["-c:s", "mov_text"]) command.extend(["-movflags", "+faststart", str(target)]) - self._run_ffmpeg(command, "MP4 muxing failed") + duration = max((track.duration for track in (job.result.tracks if job.result else ())), default=0.0) + self._run_ffmpeg(job, command, "MP4 muxing failed", lambda fraction: self._update_processing(job, "Creating MP4", base + weight * fraction, total), duration) - def _attach_cover(self, video: Path, cover: Path, target: Path) -> None: - command = [ - self.ffmpeg, - "-hide_banner", - "-loglevel", - "error", - "-nostdin", - "-y", - "-i", - str(video), - "-i", - str(cover), - "-map", - "0", - "-map", - "1:v:0", - "-c", - "copy", - "-c:v:1", - "mjpeg", - "-disposition:v:1", - "attached_pic", - "-movflags", - "+faststart", - str(target), - ] - self._run_ffmpeg(command, "Cover embedding failed") + def _complete_cover_step(self, job: PostprocessJob, task: dict[str, object], item: dict[str, object], mux_temp: Path, base: float, weight: float, total: float) -> None: + cover_path: Path | None = None + covered_temp = job.output_dir / f".{job.output_path.stem}.{job.item_id}.cover.tmp.mp4" + try: + self._update_processing(job, "Downloading cover image", base, total) + cover_path = self._download_cover(str(item["image_src"]), job.output_dir, job.item_id, str(task["proxy"]), task_id=job.task_id) + self._update_processing(job, "Embedding cover image", base, total) + self._attach_cover(job, mux_temp, cover_path, covered_temp, base, weight, total) + os.replace(covered_temp, job.output_path) + mux_temp.unlink(missing_ok=True) + cover_path.unlink(missing_ok=True) + self.store.update_item(job.item_id, status="completed", stage="Completed", progress_phase="processing", processing_completed=100.0, processing_total=100.0, output_path=str(job.output_path), cover_path=None, completed_at=_timestamp(), message="Video and cover completed") + self.store.event(job.task_id, job.item_id, "info", f"Saved {job.output_path.name}") + self._discard_job(job.item_id) + except DownloadCancelledError: + raise + except Exception as error: + os.replace(mux_temp, job.output_path) + retained_cover = self._retain_cover(cover_path, job.output_path) + job.cover_only = True + self.store.update_item(job.item_id, status="completed_warning", stage="Completed with warning", progress_phase="processing", processing_completed=100.0, processing_total=100.0, output_path=str(job.output_path), cover_path=str(retained_cover) if retained_cover else None, warning=str(error), completed_at=_timestamp(), message="Video completed without an embedded cover") + self.store.event(job.task_id, job.item_id, "warning", f"Cover step failed: {error}") + finally: + covered_temp.unlink(missing_ok=True) + + def _retry_cover_job(self, job: PostprocessJob, task: dict[str, object]) -> None: + item = self._item(job.item_id) + cover_path = Path(str(item["cover_path"])) if item["cover_path"] else None + covered_temp = job.output_dir / f".{job.output_path.stem}.{job.item_id}.cover.tmp.mp4" + try: + self._update_processing(job, "Retrying cover", 0.0, 100.0) + if cover_path is None or not cover_path.is_file(): + cover_path = self._download_cover(str(item["image_src"]), job.output_dir, job.item_id, str(task["proxy"]), task_id=job.task_id) + self._attach_cover(job, job.output_path, cover_path, covered_temp, 0.0, 100.0, 100.0) + os.replace(covered_temp, job.output_path) + cover_path.unlink(missing_ok=True) + self.store.update_item(job.item_id, status="completed", stage="Completed", progress_phase="processing", processing_completed=100.0, processing_total=100.0, warning=None, cover_path=None, completed_at=_timestamp(), message="Video and cover completed") + self._discard_job(job.item_id) + except DownloadCancelledError: + raise + except Exception as error: + retained_cover = self._retain_cover(cover_path, job.output_path) + self.store.update_item(job.item_id, status="completed_warning", stage="Completed with warning", warning=str(error), cover_path=str(retained_cover) if retained_cover else None, completed_at=_timestamp(), message="Video completed without an embedded cover") + self.store.event(job.task_id, job.item_id, "warning", f"Cover retry failed: {error}") + finally: + covered_temp.unlink(missing_ok=True) + + def _attach_cover(self, job: PostprocessJob, video: Path, cover: Path, target: Path, base: float, weight: float, total: float) -> None: + command = [self.ffmpeg, "-hide_banner", "-loglevel", "error", "-nostdin", "-y", "-i", str(video), "-i", str(cover), "-map", "0", "-map", "1:v:0", "-c", "copy", "-c:v:1", "mjpeg", "-disposition:v:1", "attached_pic", "-movflags", "+faststart", str(target)] + self._run_ffmpeg(job, command, "Cover embedding failed", lambda fraction: self._update_processing(job, "Embedding cover image", base + weight * fraction, total), 0.0) + + def _run_ffmpeg(self, job: PostprocessJob, command: list[str], summary: str, progress: Callable[[float], None], duration: float) -> None: + 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) + output: list[str] = [] + try: + assert process.stdout is not None + for line in process.stdout: + output.append(line) + if self._stop.is_set() or self.store.cancel_requested(job.task_id): + self._terminate_process(process) + raise DownloadCancelledError("Processing cancelled") + 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 process.wait(): + raise RuntimeError(f"{summary}: {''.join(output).strip()[-1500:] or summary}") + progress(1.0) + finally: + with self._processes_lock: + self._processes.discard(process) + + 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) + + @staticmethod + def _weights(tracks: list[DownloadedTrackSegments]) -> list[float]: + durations = [track.duration for track in tracks] + if all(duration > 0 for duration in durations): + mux = max(durations) + return [*durations, mux, max(1.0, (sum(durations) + mux) * 0.02)] + return [*[1.0 for _ in tracks], 1.0, 1.0] + + @staticmethod + def _track_output(job: PostprocessJob, track: DownloadedTrackSegments) -> Path: + suffix = ".srt" if track.media_type is MediaType.SUBTITLES else ".m4a" if track.media_type is MediaType.AUDIO else ".mp4" + return job.result.temporary_dir / f"{valid_filename(track.track_id.replace(':', '_'), 80)}.merged{suffix}" # type: ignore[union-attr] + + @staticmethod + def _merge_binary(files: list[Path], output: Path) -> None: + with output.open("wb") as target: + for source in files: + with source.open("rb") as incoming: + shutil.copyfileobj(incoming, target, 1024 * 1024) + + 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) + + def _enqueue_existing(self, item_id: int) -> None: + with self._jobs_lock: + if item_id in self._jobs: + self._pending.append(item_id) + + def _next_postprocess_item(self) -> int | None: + with self._jobs_lock: + while self._pending: + item_id = self._pending.popleft() + job = self._jobs.get(item_id) + if job is None: + continue + if self.store.cancel_requested(job.task_id): + self._jobs.pop(item_id, None) + self._cleanup_job(job) + continue + return item_id + return None + + def _job(self, item_id: int) -> PostprocessJob | None: + with self._jobs_lock: + return self._jobs.get(item_id) + + def _has_job(self, item_id: int) -> bool: + job = self._job(item_id) + return job is not None and (job.cover_only or job.result is not None and job.result.temporary_dir.is_dir()) + + def _discard_job(self, item_id: int) -> None: + with self._jobs_lock: + job = self._jobs.pop(item_id, None) + if job is not None: + self._cleanup_job(job) + + @staticmethod + def _cleanup_job(job: PostprocessJob) -> None: + if job.result is not None: + shutil.rmtree(job.result.temporary_dir.parent, ignore_errors=True) + if job.output_path.exists() and job.output_path.stat().st_size == 0: + job.output_path.unlink(missing_ok=True) + + def _cancel_item(self, task_id: int, item_id: int) -> None: + self.store.update_item(item_id, status="cancelled", stage="Cancelled", completed_at=_timestamp(), message="Cancellation completed") + self.store.event(task_id, item_id, "info", "Video cancelled") + + def _mark_failed(self, task_id: int, item_id: int, error: str, message: str = "Processing failed") -> None: + self.store.update_item(item_id, status="failed", stage="Failed", error=error, completed_at=_timestamp(), message=message) + self.store.event(task_id, item_id, "error", error) + + def _mark_missing_intermediate(self, task: dict[str, object], item: dict[str, object]) -> None: + self._mark_failed(int(task["id"]), int(item["id"]), "Downloaded segments are unavailable; retry the video", "Processing retry failed") + self.store.finalize_task_if_ready(int(task["id"])) + + def _item(self, item_id: int) -> dict[str, object]: + for task in self.store.list_tasks(): + for item in task["items"]: + if int(item["id"]) == item_id: + return item + raise KeyError(item_id) + + def _terminate_ffmpeg_processes(self) -> None: + with self._processes_lock: + processes = list(self._processes) + for process in processes: + self._terminate_process(process) + + @staticmethod + def _terminate_process(process: subprocess.Popen[str]) -> None: + if process.poll() is not None: + return + try: + os.killpg(process.pid, signal.SIGTERM) + except (AttributeError, ProcessLookupError): + process.terminate() + try: + process.wait(timeout=2) + except subprocess.TimeoutExpired: + try: + os.killpg(process.pid, signal.SIGKILL) + except (AttributeError, ProcessLookupError): + process.kill() def _download_cover(self, url: str, output_dir: Path, item_id: int, proxy: str, *, task_id: int) -> Path: preferred_url = _preferred_cover_url(url) @@ -316,12 +502,7 @@ class TaskRunner: try: return self._download_cover_url(preferred_url, output_dir, item_id, proxy) except Exception as error: - self.store.log( - "warning", - f"Preferred s1080 cover URL download failed; falling back to the original URL: {error}", - task_id, - item_id, - ) + self.store.log("warning", f"Preferred s1080 cover URL download failed; falling back to the original URL: {error}", task_id, item_id) return self._download_cover_url(url, output_dir, item_id, proxy) def _download_cover_url(self, url: str, output_dir: Path, item_id: int, proxy: str) -> Path: @@ -337,8 +518,7 @@ class TaskRunner: if attempt == 3: raise time.sleep(attempt + 1) - extension = _image_extension(data, content_type, url) - path = output_dir / f".m3u8downloaderd-{item_id}.cover{extension}" + path = output_dir / f".m3u8downloaderd-{item_id}.cover{_image_extension(data, content_type, url)}" path.write_bytes(data) return path @@ -366,12 +546,9 @@ class TaskRunner: os.close(descriptor) return candidate - @staticmethod - def _run_ffmpeg(command: list[str], summary: str) -> None: - result = subprocess.run(command, capture_output=True, text=True) - if result.returncode: - detail = result.stderr.strip()[-1500:] or summary - raise RuntimeError(f"{summary}: {detail}") + +def _concat_escape(path: Path) -> str: + return path.as_posix().replace("'", "'\\''") def _read_limited(response: object, limit: int) -> bytes: @@ -406,9 +583,7 @@ _IMG2_SIZE_PATH = re.compile(r"/img2/s\d+/") def _preferred_cover_url(url: str) -> str: parts = urlsplit(url) preferred_path = _IMG2_SIZE_PATH.sub("/img2/s1080/", parts.path) - if preferred_path == parts.path: - return url - return urlunsplit(parts._replace(path=preferred_path)) + return url if preferred_path == parts.path else urlunsplit(parts._replace(path=preferred_path)) def _timestamp() -> str: diff --git a/m3u8downloaderd/tests/test_service.py b/m3u8downloaderd/tests/test_service.py index 65878d3..2a61b2c 100644 --- a/m3u8downloaderd/tests/test_service.py +++ b/m3u8downloaderd/tests/test_service.py @@ -33,6 +33,7 @@ def test_http_api_validates_and_snapshots_settings(tmp_path: Path) -> None: service = DownloadService(config) assert service.bootstrap()["settings"]["proxy"] == DEFAULT_PROXY assert service.bootstrap()["settings"]["max_active_downloads"] == 1 + assert service.bootstrap()["settings"]["max_active_ffmpeg"] == 1 server = DownloadHTTPServer(config, service) thread = threading.Thread(target=server.serve_forever, daemon=True) thread.start() @@ -52,12 +53,15 @@ def test_http_api_validates_and_snapshots_settings(tmp_path: Path) -> None: "default_directory": str(config.download_root), "proxy": "http://127.0.0.1:7890", "max_active_downloads": 2, + "max_active_ffmpeg": 3, }, method="PUT", ) assert settings[0] == 200 assert settings[1]["max_active_downloads"] == 2 + assert settings[1]["max_active_ffmpeg"] == 3 assert service.store.max_active_downloads() == 2 + assert service.store.max_active_ffmpeg() == 3 task = _request(base, "/api/tasks", {"title": "batch", "directory": str(config.download_root), "payload": json.dumps([_item("http://127.0.0.1")])}) assert task[0] == 201 assert task[1]["proxy"] == "http://127.0.0.1:7890" @@ -236,9 +240,11 @@ def test_cover_download_keeps_original_url_without_img2_size_path(tmp_path: Path def test_cli_uses_max_active_downloads() -> None: - args = build_argument_parser().parse_args(["--max-active-downloads", "3"]) + args = build_argument_parser().parse_args(["--max-active-downloads", "3", "--max-active-ffmpeg", "4"]) - assert config_from_args(args).max_active_downloads == 3 + config = config_from_args(args) + assert config.max_active_downloads == 3 + assert config.max_active_ffmpeg == 4 def _get(base: str, path: str) -> tuple[int, dict[str, object]]: