import functools import io import json import os import subprocess import sys import tempfile import threading import time import unittest from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from unittest.mock import patch from host.common import validate_download from host.main import NativeHost from host.media import NO_WINDOW, ROOT, probe, tool from host.processes import ProcessTree from host.worker import convert, download class ProcessingTests(unittest.TestCase): def test_audio_conversion_and_video_conversion_report_the_actual_engine(self): runtime = ROOT / ".runtime" runtime.mkdir(exist_ok=True) with tempfile.TemporaryDirectory(prefix="processing-stages-", dir=runtime) as directory: source = Path(directory) / "source.mkv" subprocess.run([tool("ffmpeg"), "-nostdin", "-v", "error", "-y", "-f", "lavfi", "-i", "testsrc2=size=160x90:rate=24", "-f", "lavfi", "-i", "sine=frequency=440:sample_rate=48000", "-t", "3", "-c:v", "libx264", "-preset", "ultrafast", "-c:a", "libopus", str(source)], creationflags=NO_WINDOW, stdin=subprocess.DEVNULL, check=True, timeout=30) media = probe(source) for codec, stage in [("auto", "audio"), ("h264", "encoding")]: target = Path(directory) / f"{codec}.mp4" events = [] options = {"mode": "video", "container": "mp4", "codec": codec, "engine": "auto" if codec == "auto" else "cpu", "preset": "fast", "bitrate": 192} with patch("host.worker.emit", side_effect=lambda kind, **data: events.append({"type": kind, **data})): plan = convert(source, target, options, media, {"nvenc": ["h264"]}) self.assertEqual(plan["engine"], "cpu") self.assertEqual(plan["video"], "copy" if codec == "auto" else "h264") self.assertTrue(any(event.get("stage") == stage and event.get("engine") == "cpu" for event in events), events) self.assertFalse(any(event.get("stage") == "merging" for event in events), events) self.assertTrue(any(event.get("stage") == "finalizing" and event.get("percent") == 99 and not event.get("speed") for event in events), events) finalizing = next(index for index, event in enumerate(events) if event.get("stage") == "finalizing") self.assertTrue(all(event.get("stage") == "finalizing" for event in events[finalizing:]), events) self.assertEqual({s.get("codec_name") for s in probe(target)["streams"]}, {"h264", "aac"}) def test_late_completion_cannot_replace_cancellation(self): host = NativeHost(io.BytesIO()) try: host.jobs["late"] = {"id": "late", "status": "running", "stage": "cancelling", "cancelling": True} host.event("late", status="complete", stage="complete") self.assertEqual(host.jobs["late"]["status"], "running") host.event("late", status="cancelled", stage="cancelled", cancelling=False) host.event("late", status="complete", stage="complete") self.assertEqual(host.jobs["late"]["status"], "cancelled") finally: host.close() def test_work_folder_validation_and_default(self): with tempfile.TemporaryDirectory() as folder: payload = {"url": "https://youtu.be/jNQXAC9IVRw", "folder": folder} self.assertEqual(validate_download(payload)["workFolder"], str(Path(folder).resolve())) with self.assertRaisesRegex(ValueError, "Рабочая папка"): validate_download({**payload, "workFolder": str(Path(folder) / "missing")}) @unittest.skipUnless(os.name == "nt", "Native worker process and disk cleanup") def test_native_worker_cancel_removes_both_workspaces_and_stops_ffmpeg(self): runtime = ROOT / ".runtime" runtime.mkdir(exist_ok=True) with tempfile.TemporaryDirectory(dir=runtime) as work, tempfile.TemporaryDirectory() as target: fixture = Path(work) / "host" fixture.mkdir() (fixture / "worker.py").write_text('''import json,subprocess,sys,tempfile,time from pathlib import Path task=json.loads(sys.stdin.readline())['payload'] for folder in [task['workFolder'],task['folder']]: scratch=Path(tempfile.mkdtemp(prefix='.ytdl-',dir=folder)) (scratch/'partial.mp4').write_bytes(b'partial') print(json.dumps({'type':'workspace','path':str(scratch)}),flush=True) child=subprocess.Popen([task['ffmpeg'],'-nostdin','-v','error','-re','-f','lavfi','-i','testsrc2=size=64x64:rate=10','-f','null','-'],stdin=subprocess.DEVNULL,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL) print(json.dumps({'type':'progress','stage':'merging','percent':25,'testPid':child.pid}),flush=True) time.sleep(60) ''', encoding="utf-8") for folder in [work, target]: (Path(folder) / "keep.txt").write_text("user file") with patch("host.main.ROOT", Path(work)): host = NativeHost(io.BytesIO()) host._capabilities = {} host.jobs["native"] = {"id": "native", "status": "running", "stage": "preparing"} thread = threading.Thread(target=host.run_job, args=("native", {"folder": target, "workFolder": work, "ffmpeg": tool("ffmpeg")})) thread.start() try: deadline = time.monotonic() + 10 while "testPid" not in host.jobs["native"] and time.monotonic() < deadline: time.sleep(0.025) pid = host.jobs["native"]["testPid"] host.handle({"id": "stop", "action": "cancel", "payload": {"jobId": "native"}}) thread.join(timeout=10) self.assertFalse(thread.is_alive()) self.assertFalse(self.running(pid)) self.assertFalse(host.processes) for folder in [work, target]: self.assertFalse(list(Path(folder).glob(".ytdl-*"))) self.assertEqual((Path(folder) / "keep.txt").read_text(), "user file") finally: host.close() thread.join(timeout=10) @unittest.skipUnless(os.name == "nt", "Actual Windows process-tree regression") def test_cancel_kills_ffmpeg_after_parent_exit_and_keeps_other_job(self): launcher = """import json,subprocess,sys,time task=json.loads(sys.stdin.readline()) child=subprocess.Popen([task['ffmpeg'],'-hide_banner','-loglevel','error','-nostdin','-re','-f','lavfi','-i','testsrc2=size=64x64:rate=10','-f','null','-'],stdin=subprocess.DEVNULL,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL) print(child.pid,flush=True) if not task['exit']:time.sleep(60) """ processes = [] host = NativeHost(io.BytesIO()) try: children = [] for exits in (True, False): p = ProcessTree.launch([sys.executable, "-c", launcher], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, text=True, creationflags=NO_WINDOW) processes.append(p) p.stdin.write(json.dumps({"ffmpeg": tool("ffmpeg"), "exit": exits}) + "\n") p.stdin.close() children.append(int(p.stdout.readline())) if exits: p.wait(timeout=10) host.jobs["cancel"] = {"id": "cancel", "status": "running", "stage": "merging"} host.processes["cancel"] = processes[0] self.assertTrue(self.running(children[0])) self.assertTrue(self.running(children[1])) host.handle({"id": "request", "action": "cancel", "payload": {"jobId": "cancel"}}) self.assertEqual(host.jobs["cancel"]["status"], "cancelled") self.assertFalse(self.running(children[0]), "Cancelled FFmpeg must actually exit") self.assertTrue(self.running(children[1]), "A different job's FFmpeg must survive") finally: host.close() for p in processes: tree = getattr(p, "ytdl_process_tree", None) if tree: tree.terminate(); tree.close() elif p.poll() is None: p.kill(); p.wait() p.stdout.close() for pid in locals().get("children", []): if self.running(pid): subprocess.run(["taskkill", "/PID", str(pid), "/T", "/F"], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, creationflags=NO_WINDOW, timeout=10) @staticmethod def running(pid): import ctypes from ctypes import wintypes api = ctypes.WinDLL("kernel32", use_last_error=True) api.OpenProcess.argtypes = [wintypes.DWORD, wintypes.BOOL, wintypes.DWORD] api.OpenProcess.restype = wintypes.HANDLE api.WaitForSingleObject.argtypes = [wintypes.HANDLE, wintypes.DWORD] api.WaitForSingleObject.restype = wintypes.DWORD api.CloseHandle.argtypes = [wintypes.HANDLE] handle = api.OpenProcess(0x100000, False, pid) if not handle: return False try: return api.WaitForSingleObject(handle, 0) == 258 finally: api.CloseHandle(handle) def test_actual_download_merger_reports_progress_and_publishes_across_disks(self): from yt_dlp import YoutubeDL runtime = ROOT / ".runtime" runtime.mkdir(exist_ok=True) with tempfile.TemporaryDirectory(prefix="processing-test-", dir=runtime) as working, tempfile.TemporaryDirectory() as output: work, destination = Path(working), Path(output) video, audio = work / "video.mp4", work / "audio.m4a" for inputs, target, codec in [(["-f", "lavfi", "-i", "testsrc2=size=160x90:rate=24"], video, ["-an", "-c:v", "libx264", "-preset", "ultrafast"]), (["-f", "lavfi", "-i", "sine=frequency=440:sample_rate=48000"], audio, ["-vn", "-c:a", "aac"])]: subprocess.run([tool("ffmpeg"), "-nostdin", "-v", "error", "-y", *inputs, "-t", "12", *codec, str(target)], creationflags=NO_WINDOW, stdin=subprocess.DEVNULL, check=True, timeout=30) class Handler(SimpleHTTPRequestHandler): def log_message(self, *args): pass server = ThreadingHTTPServer(("127.0.0.1", 0), functools.partial(Handler, directory=str(work))) threading.Thread(target=server.serve_forever, daemon=True).start() events = [] origin = f"http://127.0.0.1:{server.server_port}/" def extract(ydl, url, download=True): return ydl.process_ie_result({"id": "jNQXAC9IVRw", "title": "Processing regression", "duration": 12, "webpage_url": "https://youtu.be/jNQXAC9IVRw", "extractor": "fixture", "extractor_key": "Fixture", "formats": [ {"format_id": "137", "url": origin + "video.mp4", "ext": "mp4", "protocol": "http", "vcodec": "avc1", "acodec": "none", "height": 90, "width": 160, "fps": 24}, {"format_id": "140", "url": origin + "audio.m4a", "ext": "m4a", "protocol": "http", "vcodec": "none", "acodec": "mp4a.40.2", "abr": 128}]}, download=download) try: with patch.object(YoutubeDL, "extract_info", extract), patch("host.worker.emit", side_effect=lambda kind, **data: events.append({"type": kind, **data})): result = globals()["download"]({"url": "https://youtu.be/jNQXAC9IVRw", "folder": str(destination), "workFolder": str(work), "mode": "video", "container": "mp4", "codec": "copy"}, []) finally: server.shutdown(); server.server_close() saved = Path(result["path"]) self.assertEqual(saved.parent, destination.resolve()) self.assertEqual({s["codec_type"] for s in probe(saved)["streams"]}, {"video", "audio"}) self.assertTrue(any(e.get("stage") == "merging" and e.get("processedSeconds", 0) > 0 and e.get("written", 0) > 0 and e.get("percent") is not None for e in events), events) self.assertTrue(any(e.get("stage") == "remuxing" and e.get("engine") == "copy" for e in events), events) stages = [e.get("stage") for e in events if e.get("type") == "progress"] self.assertLess(stages.index("merging"), stages.index("remuxing")) self.assertNotIn("merging", stages[stages.index("remuxing"):]) if work.stat().st_dev != destination.stat().st_dev: self.assertTrue(any(e.get("stage") == "saving" and e.get("copied", 0) > 0 for e in events), events) self.assertFalse(list(work.glob(".ytdl-*"))) self.assertFalse(list(destination.glob(".ytdl-*"))) self.assertEqual(result["size"], saved.stat().st_size) if __name__ == "__main__": unittest.main()