294 lines
15 KiB
Python
294 lines
15 KiB
Python
"""Isolated, cancellable download worker. Stdout contains JSON lines only."""
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import errno
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
import tempfile
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
|
|
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
|
|
|
from host.common import safe_filename, select_streams, summarize, validate_download, youtube_url
|
|
from host.media import NO_WINDOW, config, encoding_plan, ffmpeg_command, probe, tool
|
|
from host.session import apply_cookies, public_error, is_access_error
|
|
from host.ffmpeg_progress import run_ffmpeg
|
|
|
|
|
|
def emit(kind, **data):
|
|
print(json.dumps({"type": kind, **data}, ensure_ascii=False, allow_nan=False), flush=True)
|
|
|
|
|
|
class QuietLogger:
|
|
def debug(self, message):
|
|
pass
|
|
def info(self, message):
|
|
pass
|
|
def warning(self, message):
|
|
pass
|
|
def error(self, message):
|
|
# The raised DownloadError is reported through the structured channel.
|
|
pass
|
|
|
|
|
|
def ydl_options():
|
|
try:
|
|
node = tool("node")
|
|
except ValueError:
|
|
node = None
|
|
options = {"quiet": True, "no_warnings": True, "logger": QuietLogger(), "noplaylist": True,
|
|
"socket_timeout": 20, "retries": 3, "fragment_retries": 5,
|
|
# Parallel HLS requests can trigger YouTube's fragment throttling;
|
|
# never skip failed fragments and silently publish an incomplete clip.
|
|
"concurrent_fragment_downloads": 1, "skip_unavailable_fragments": False,
|
|
"ffmpeg_location": str(Path(tool("ffmpeg")).parent),
|
|
"windowsfilenames": True, "cachedir": False, "noprogress": True}
|
|
if node:
|
|
options["js_runtimes"] = {"node": {"path": node}}
|
|
return options
|
|
|
|
|
|
def inspect(url, cookies=None):
|
|
from yt_dlp import YoutubeDL
|
|
from yt_dlp.utils import DownloadError
|
|
for attempt in range(2):
|
|
try:
|
|
with YoutubeDL(ydl_options()) as ydl:
|
|
apply_cookies(ydl, cookies)
|
|
return summarize(ydl.extract_info(youtube_url(url), download=False))
|
|
except DownloadError as error:
|
|
if attempt or not is_access_error(error):
|
|
raise
|
|
time.sleep(1.5)
|
|
|
|
|
|
def convert(source, target, options, media, encoders):
|
|
plan = encoding_plan(options, media, encoders)
|
|
duration = float(media.get("format", {}).get("duration") or 0)
|
|
attempts = [(plan, True)]
|
|
if plan["engine"] in {"nvenc", "amf"}:
|
|
attempts.append((plan, False))
|
|
if options["engine"] == "auto":
|
|
attempts.append(({**plan, "engine": "cpu"}, False))
|
|
last_error = ""
|
|
for index, (attempt, cuda_decode) in enumerate(attempts):
|
|
stage = "encoding" if attempt["video"] not in (None, "copy") else ("audio" if attempt["audio"] not in (None, "copy") else "remuxing")
|
|
emit("progress", stage=stage, percent=0, engine=attempt["engine"], codec=attempt["video"] or attempt["audio"],
|
|
message="", speed="", eta=None)
|
|
command = ffmpeg_command(source, target, options, attempt, cuda_decode)
|
|
returncode, errors = run_ffmpeg(command, duration, stage, attempt["engine"], lambda **data: emit("progress", **data))
|
|
if returncode == 0 and target.is_file() and target.stat().st_size:
|
|
# Verify the result really has the requested streams before publishing it.
|
|
output = probe(target)
|
|
types = {s.get("codec_type") for s in output.get("streams", [])}
|
|
required = {"audio"} if options["mode"] == "audio" else ({"video", "audio"} if options["mode"] == "video" else {"video"})
|
|
if not required <= types or (options["mode"] == "video_only" and "audio" in types):
|
|
raise ValueError("Проверка итогового файла не прошла: отсутствует выбранная дорожка.")
|
|
return attempt
|
|
last_error = errors.strip()[-1600:]
|
|
if index + 1 < len(attempts):
|
|
next_plan, next_decode = attempts[index + 1]
|
|
emit("progress", stage="encoding", percent=0, engine=attempt["engine"],
|
|
message="Повтор с обычным декодированием" if next_plan["engine"] == attempt["engine"] else "Аппаратное кодирование недоступно: переход на CPU", speed="", eta=None)
|
|
raise ValueError("FFmpeg не смог обработать видео. " + last_error)
|
|
|
|
|
|
def publish_file(source: Path, folder: Path, name: str, extension: str) -> Path:
|
|
if source.stat().st_dev != folder.stat().st_dev:
|
|
# Copy to an unpublished file on the destination disk, then publish it
|
|
# atomically. Cancellation cannot leave a partial final filename.
|
|
with tempfile.TemporaryDirectory(prefix=".ytdl-", dir=folder) as staging:
|
|
staging_path = Path(staging).resolve()
|
|
emit("workspace", path=str(staging_path), purpose="saving")
|
|
target = staging_path / ("ready." + extension)
|
|
total, copied, last_emit = source.stat().st_size, 0, 0.0
|
|
emit("progress", stage="saving", percent=0, copied=0, total=total, speed=None, eta=None)
|
|
with source.open("rb") as incoming, target.open("xb") as outgoing:
|
|
while chunk := incoming.read(4 * 1024 * 1024):
|
|
outgoing.write(chunk)
|
|
copied += len(chunk)
|
|
if time.monotonic() - last_emit >= 0.2:
|
|
last_emit = time.monotonic()
|
|
emit("progress", stage="saving", percent=min(99, copied / total * 100) if total else 0,
|
|
copied=copied, total=total, speed=None, eta=None)
|
|
outgoing.flush()
|
|
os.fsync(outgoing.fileno())
|
|
destination = publish_file(target, folder, name, extension)
|
|
source.unlink()
|
|
return destination
|
|
# Same-volume hard-link creation is atomic and never overwrites an existing file.
|
|
for index in range(10000):
|
|
suffix = "" if index == 0 else f" ({index})"
|
|
destination = folder / f"{name}{suffix}.{extension}"
|
|
if destination.exists():
|
|
continue
|
|
try:
|
|
os.link(source, destination)
|
|
source.unlink()
|
|
return destination
|
|
except FileExistsError:
|
|
continue
|
|
except OSError as error:
|
|
if destination.exists():
|
|
continue
|
|
if error.errno == errno.EXDEV or getattr(error, "winerror", None) == 17:
|
|
raise ValueError("Не удалось сохранить файл на выбранный диск.") from error
|
|
# Windows rename also refuses to replace existing destinations, including on FAT/exFAT.
|
|
if os.name != "nt":
|
|
raise
|
|
try:
|
|
os.rename(source, destination)
|
|
return destination
|
|
except FileExistsError:
|
|
continue
|
|
except OSError:
|
|
if destination.exists():
|
|
continue
|
|
raise
|
|
raise ValueError("Не удалось выбрать свободное имя файла.")
|
|
|
|
|
|
def download(payload, encoders):
|
|
from yt_dlp import YoutubeDL
|
|
from yt_dlp.postprocessor.ffmpeg import FFmpegMergerPP, FFmpegPostProcessorError
|
|
from yt_dlp.utils import DownloadError
|
|
options = validate_download(payload)
|
|
emit("progress", stage="preparing", percent=None, message="Получаю свежие ссылки на дорожки")
|
|
folder = Path(options["folder"])
|
|
work_folder = Path(options["workFolder"])
|
|
# Only this job's generated directory is removed on exit; final files live outside it.
|
|
# A failed downloader can retain an open .part handle until this worker exits.
|
|
# Preserve its actual error; the native host cleans up again after process exit.
|
|
with tempfile.TemporaryDirectory(prefix=".ytdl-", dir=work_folder, ignore_cleanup_errors=True) as scratch:
|
|
scratch_path = Path(scratch).resolve()
|
|
if scratch_path.parent != work_folder.resolve():
|
|
raise ValueError("Некорректная временная папка.")
|
|
emit("workspace", path=str(scratch_path))
|
|
last_emit = 0.0
|
|
finished_tracks = set()
|
|
track_count = 1
|
|
def format_selector(context):
|
|
nonlocal track_count
|
|
selector = select_streams(context, options)
|
|
track_count = 2 if "+" in selector else 1
|
|
# Keep extraction and downloading in the same YoutubeDL instance so
|
|
# its visitor cookies and request state remain available.
|
|
return ydl.build_format_selector(selector)(context)
|
|
def reject_live(info, *, incomplete=False):
|
|
if info.get("is_live"):
|
|
return "Дождитесь окончания прямой трансляции."
|
|
def progress(event):
|
|
nonlocal last_emit
|
|
if event.get("status") == "finished":
|
|
finished_tracks.add(event.get("filename", ""))
|
|
return
|
|
if event.get("status") != "downloading" or time.monotonic() - last_emit < 0.2:
|
|
return
|
|
last_emit = time.monotonic()
|
|
downloaded = event.get("downloaded_bytes", 0)
|
|
total = event.get("total_bytes") or event.get("total_bytes_estimate") or 0
|
|
track_progress = min(1, downloaded / total) if total else 0
|
|
percent = min(99, (len(finished_tracks) + track_progress) / track_count * 100) if total else None
|
|
stream = event.get("info_dict", {})
|
|
track = "audio" if stream.get("vcodec") == "none" else "video"
|
|
emit("progress", stage="downloading", track=track, percent=percent,
|
|
downloaded=downloaded, total=total, speed=event.get("speed"), eta=event.get("eta"))
|
|
def postprocess(event):
|
|
if event.get("status") == "started" and event.get("postprocessor") == "Merger":
|
|
emit("progress", stage="merging", percent=0, engine="copy", speed="", eta=None,
|
|
processedSeconds=0, durationSeconds=0, written=0, message="Подготовка дорожек к объединению")
|
|
|
|
class ProgressMerger(FFmpegMergerPP):
|
|
def run(self, info):
|
|
self.duration = info.get("duration") or 0
|
|
return super().run(info)
|
|
|
|
def run_ffmpeg_multiple_files(self, inputs, output, opts, **kwargs):
|
|
self.check_version()
|
|
oldest_mtime = min(os.stat(path).st_mtime for path in inputs)
|
|
command = [self.executable, "-hide_banner", "-nostdin", "-loglevel", "info", "-y",
|
|
"-progress", "pipe:1", "-stats_period", "0.4"]
|
|
for path in inputs:
|
|
command += ["-i", str(path)]
|
|
# Preserve yt-dlp's maps and AAC fixup, including HLS audio.
|
|
command += list(opts) + [str(output)]
|
|
code, errors = run_ffmpeg(command, self.duration, "merging", "copy", lambda **data: emit("progress", **data))
|
|
if code:
|
|
raise FFmpegPostProcessorError(errors.strip()[-1600:] or "Не удалось объединить дорожки.")
|
|
self.try_utime(output, oldest_mtime, oldest_mtime)
|
|
|
|
def enable_merge_progress(ydl):
|
|
original_run_pp = ydl.run_pp
|
|
def run_pp(pp, info):
|
|
if isinstance(pp, FFmpegMergerPP) and not isinstance(pp, ProgressMerger):
|
|
pp = ProgressMerger(ydl)
|
|
return original_run_pp(pp, info)
|
|
ydl.run_pp = run_pp
|
|
download_options = {**ydl_options(), "format": format_selector, "match_filter": reject_live,
|
|
"outtmpl": str(scratch_path / "source.%(ext)s"),
|
|
"merge_output_format": "mkv", "progress_hooks": [progress],
|
|
"postprocessor_hooks": [postprocess], "overwrites": False}
|
|
for attempt in range(2):
|
|
try:
|
|
with YoutubeDL(download_options) as ydl:
|
|
enable_merge_progress(ydl)
|
|
apply_cookies(ydl, payload.get("browserCookies"))
|
|
downloaded_info = ydl.extract_info(options["url"], download=True)
|
|
break
|
|
except DownloadError as error:
|
|
if attempt == 0 and (is_access_error(error) or "403" in str(error) or "410" in str(error)):
|
|
finished_tracks.clear()
|
|
emit("progress", stage="preparing", percent=None, message="Обновляю ссылки YouTube и повторяю загрузку")
|
|
time.sleep(1.5)
|
|
else:
|
|
raise
|
|
if not downloaded_info:
|
|
raise ValueError("Ролик недоступен для загрузки.")
|
|
title = str(downloaded_info.get("title", "Видео"))
|
|
name = safe_filename(options["filename"] or title, str(downloaded_info.get("id", "video")))
|
|
candidates = [p for p in scratch_path.glob("source.*") if p.is_file() and p.suffix.lower() in {".mkv", ".mp4", ".webm", ".m4a", ".mp3", ".ogg", ".opus", ".flac", ".wav", ".ts"}]
|
|
if len(candidates) != 1:
|
|
raise ValueError("Не удалось определить загруженный медиафайл.")
|
|
source = candidates[0]
|
|
media = probe(source)
|
|
output = scratch_path / f"final.{options['container']}"
|
|
plan = convert(source, output, options, media, encoders)
|
|
destination = publish_file(output, folder, name, options["container"])
|
|
return {"path": str(destination), "filename": destination.name, "title": title,
|
|
"size": destination.stat().st_size, "engine": plan["engine"],
|
|
"codec": plan["video"] or plan["audio"], "mode": options["mode"]}
|
|
|
|
|
|
def main():
|
|
task = {}
|
|
try:
|
|
task = json.loads(sys.stdin.readline())
|
|
if task["action"] == "inspect":
|
|
result = inspect(task["payload"]["url"], task["payload"].get("browserCookies"))
|
|
elif task["action"] == "download":
|
|
result = download(task["payload"], task.get("encoders", task.get("nvenc", [])))
|
|
elif task["action"] == "dependency_install":
|
|
from host.dependencies import install
|
|
result = install(task["payload"]["dependency"],
|
|
lambda percent, message: emit("progress", stage="installing", percent=percent, message=message),
|
|
lambda path: emit("workspace", path=path))
|
|
elif task["action"] == "helper_update":
|
|
from host.updates import update
|
|
result = update(lambda percent, message: emit("progress", stage="updating", percent=percent, message=message),
|
|
lambda path: emit("workspace", path=path), task["payload"].get("candidates"))
|
|
else:
|
|
raise ValueError("Неизвестная операция.")
|
|
emit("result", result=result)
|
|
except Exception as error:
|
|
emit("error", message=public_error(error, bool(task.get("payload", {}).get("browserCookies"))))
|
|
return 1
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|