split download and ffmpeg procces

This commit is contained in:
root
2026-09-27 12:52:05 +08:00
parent 993fd8f217
commit b6160430a0
13 changed files with 868 additions and 341 deletions
+9
View File
@@ -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,
+4
View File
@@ -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",
]
+120
View File
@@ -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:
+58 -32
View File
@@ -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()
+25
View File
@@ -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")
+9 -5
View File
@@ -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:
@@ -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);
}
@@ -60,6 +60,10 @@
<span>Concurrent video downloads</span>
<input id="setting-max-downloads" name="max_active_downloads" type="number" min="1" max="32" step="1" required>
</label>
<label>
<span>Concurrent ffmpeg processing</span>
<input id="setting-max-ffmpeg" name="max_active_ffmpeg" type="number" min="1" max="32" step="1" required>
</label>
<div class="form-actions">
<button class="primary-action" type="submit">Save settings</button>
<span id="settings-status" class="form-status" role="status"></span>
@@ -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; }
+92 -29
View File
@@ -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,49 +182,104 @@ 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:
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 [int(item["id"]) for item in retryable]
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
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,
)
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
def retry_item(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,
)
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 _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"
+71 -7
View File
@@ -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:
+412 -237
View File
@@ -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():
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
active.add(self._download_executor.submit(self._run_download, *claimed))
self._stop.wait(0.1)
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"Download worker crashed: {error}")
while len(active) < self.store.max_active_downloads():
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)
self.store.log("error", f"{label} crashed: {error}")
def _run_item(self, task: dict[str, object], item: dict[str, object]) -> None:
item_id = int(item["id"])
task_id = int(task["id"])
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"]))
workspace: Path | None = None
try:
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):
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),
)
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)
_, total = progress[track_id]
progress[track_id] = (total, total)
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"),
)
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:
level = "debug" if kind is DownloadEventKind.SEGMENT_COMPLETED else "info"
self.store.log(level, message, task_id, item_id)
self.store.log("debug" if kind is DownloadEventKind.SEGMENT_COMPLETED else "info", 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"<SaveName>.source_{item_id}_<MediaType>.<Ext>",
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):
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")
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"
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
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",
)
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:
+8 -2
View File
@@ -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]]: