from __future__ import annotations import os import shutil import subprocess import threading import time from concurrent.futures import ThreadPoolExecutor, as_completed from dataclasses import dataclass, field from datetime import datetime from pathlib import Path from typing import Callable from m3u_tools import ( OperationLogger, VIDEO_FILE_EXTENSIONS, discover_ffmpeg, discover_ffprobe, ensure_directory, natural_sort_key, slugify_filename, subprocess_no_window_kwargs, ) def _default_parallel_jobs() -> int: cpu_count = os.cpu_count() or 8 # 14600K has 20 threads, RTX 5060 can handle 3 parallel NVENC encodes return min(4, max(1, (cpu_count + 4) // 8)) def auto_rename_series(series_dir: str | Path) -> list[tuple[Path, str]]: """Rename all video files in a directory to sequential numbers (1.mp4, 2.mp4...). Skips empty dirs and non-video files. Returns list of (old_path, new_name).""" dir_path = Path(series_dir) if not dir_path.is_dir(): return [] video_files = sorted( [f for f in dir_path.iterdir() if f.is_file() and f.suffix.lower() in VIDEO_FILE_EXTENSIONS], key=lambda f: natural_sort_key(f.stem), ) if not video_files: return [] renamed: list[tuple[Path, str]] = [] for idx, src in enumerate(video_files, start=1): new_name = f"{idx}{src.suffix}" dst = dir_path / new_name if src != dst: src.rename(dst) renamed.append((src, new_name)) return renamed @dataclass(slots=True) class ConversionConfig: root_dir: str = "." ffmpeg_path: str = field(default_factory=discover_ffmpeg) ffprobe_path: str = field(default_factory=discover_ffprobe) log_dir: str = "logs" parallel_jobs: int = field(default_factory=_default_parallel_jobs) resolution: str = "1280:720" video_codec: str = "hevc_nvenc" video_bitrate: str = "0" video_maxrate: str = "1800k" video_bufsize: str = "3600k" video_cq: int = 30 audio_codec: str = "aac" audio_bitrate: str = "96k" hls_time: int = 10 preset: str = "p3" hwaccel: str = "cuda" max_retries: int = 2 timeout_seconds: int = 3600 input_dir_name: str = "input" output_dir_name: str = "output" episode_folder_mode: bool = False auto_rename: bool = True @dataclass(slots=True) class ConversionJob: name: str input_dir: Path output_dir: Path mode: str @dataclass(slots=True) class ConversionTask: job_name: str input_file: Path output_dir: Path playlist_file: Path segment_prefix: str mode: str @dataclass(slots=True) class ConversionResult: job_name: str input_file: Path output_dir: Path success: bool message: str elapsed_seconds: float @dataclass(slots=True) class ConversionSummary: total: int succeeded: int failed: int log_file: Path elapsed_seconds: float results: list[ConversionResult] def discover_series_jobs(root_dir: str | Path, config: ConversionConfig) -> list[ConversionJob]: """Scan root_dir/input/ subdirs for video series. Each subdir under input/ is treated as one series. Output goes to root_dir/output/// Falls back to legacy mode (numbered input dirs) for backward compat. """ root = Path(root_dir) jobs: list[ConversionJob] = [] input_path = root / config.input_dir_name output_path = root / config.output_dir_name if input_path.is_dir(): for child in sorted(input_path.iterdir(), key=lambda item: natural_sort_key(item.name)): if not child.is_dir(): continue video_count = sum( 1 for f in child.iterdir() if f.is_file() and f.suffix.lower() in VIDEO_FILE_EXTENSIONS ) if video_count == 0: continue jobs.append( ConversionJob( name=child.name, input_dir=child, output_dir=output_path / child.name, mode="series", ) ) if jobs: return jobs # Legacy: numbered input dirs directly under root numbered_inputs = sorted( [path for path in root.iterdir() if path.is_dir() and path.name[:-1] == "input"], key=lambda item: natural_sort_key(item.name), ) for inp in numbered_inputs: suffix = inp.name[5:] # "input" -> suffix after "input" out = root / f"output{suffix}" video_count = sum( 1 for f in inp.iterdir() if f.is_file() and f.suffix.lower() in VIDEO_FILE_EXTENSIONS ) if video_count == 0: continue jobs.append( ConversionJob(name=inp.name, input_dir=inp, output_dir=out, mode="legacy") ) if jobs: return jobs # Single: video files directly in root direct_videos = [ path for path in root.iterdir() if path.is_file() and path.suffix.lower() in VIDEO_FILE_EXTENSIONS ] if direct_videos: jobs.append( ConversionJob( name=root.name, input_dir=root, output_dir=output_path, mode="single", ) ) return jobs def discover_single_series(root_dir: str | Path, series_name: str, config: ConversionConfig) -> ConversionJob | None: """Create a single ConversionJob for a specified series folder. The series folder should be root_dir/input// or an absolute path. """ root = Path(root_dir) series_path = Path(series_name) if not series_path.is_absolute(): series_path = root / config.input_dir_name / series_name if not series_path.is_dir(): return None video_count = sum( 1 for f in series_path.iterdir() if f.is_file() and f.suffix.lower() in VIDEO_FILE_EXTENSIONS ) if video_count == 0: return None output_path = root / config.output_dir_name / series_path.name return ConversionJob( name=series_path.name, input_dir=series_path, output_dir=output_path, mode="series", ) def build_conversion_tasks(jobs: list[ConversionJob], config: ConversionConfig) -> list[ConversionTask]: tasks: list[ConversionTask] = [] for job in jobs: audio_seen: set[str] = set() video_files = sorted( [item for item in job.input_dir.iterdir() if item.is_file() and item.suffix.lower() in VIDEO_FILE_EXTENSIONS], key=lambda item: natural_sort_key(item.name), ) for input_file in video_files: episode_name = slugify_filename(input_file.stem) if config.episode_folder_mode: output_dir = job.output_dir / episode_name playlist_file = output_dir / "index.m3u8" else: output_dir = job.output_dir playlist_file = output_dir / f"{episode_name}.m3u8" tasks.append( ConversionTask( job_name=job.name, input_file=input_file, output_dir=output_dir, playlist_file=playlist_file, segment_prefix=episode_name, mode=job.mode, ) ) return tasks class VideoConverter: def __init__( self, config: ConversionConfig, logger: OperationLogger, event_callback: Callable[[dict], None] | None = None, stop_event: threading.Event | None = None, ) -> None: self.config = config self.logger = logger self.event_callback = event_callback self.stop_event = stop_event or threading.Event() self._active_processes: list[subprocess.Popen] = [] self._process_lock = threading.Lock() def log(self, level: str, message: str) -> None: getattr(self.logger, level, self.logger.info)(message) def register_process(self, process: subprocess.Popen) -> None: with self._process_lock: self._active_processes.append(process) def unregister_process(self, process: subprocess.Popen) -> None: with self._process_lock: if process in self._active_processes: self._active_processes.remove(process) def terminate_active_processes(self) -> None: with self._process_lock: for process in self._active_processes: try: process.terminate() except Exception: pass def build_filter(self) -> str: return f"scale={self.config.resolution}:flags=lanczos,setsar=1:1,format=yuv420p" def probe_duration(self, file_path: Path) -> float: try: cmd = [ self.config.ffprobe_path, "-v", "quiet", "-show_entries", "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", str(file_path), ] result = subprocess.run( cmd, capture_output=True, text=True, timeout=30, **subprocess_no_window_kwargs(), ) if result.returncode == 0 and result.stdout.strip(): return float(result.stdout.strip()) except Exception: pass return 0.0 def cleanup_task_output(self, task: ConversionTask) -> None: if task.playlist_file.exists(): task.playlist_file.unlink() if task.output_dir.exists(): for f in list(task.output_dir.iterdir()): if f.stem.startswith(task.segment_prefix) and f.suffix == '.ts': f.unlink() ensure_directory(task.output_dir) def build_ffmpeg_command(self, task: ConversionTask) -> list[str]: segment_template = task.output_dir / f"{task.segment_prefix}_%03d.ts" command = [ self.config.ffmpeg_path, "-y", "-hide_banner", "-loglevel", "warning", "-threads", "0", "-hwaccel", self.config.hwaccel, "-i", str(task.input_file), "-map", "0:v:0", "-map", "0:a:0?", "-sn", "-dn", "-map_metadata", "-1", "-map_chapters", "-1", "-vf", self.build_filter(), "-c:v", self.config.video_codec, "-preset", self.config.preset, "-rc", "vbr", "-cq", str(self.config.video_cq), "-b:v", self.config.video_bitrate, "-maxrate", self.config.video_maxrate, "-bufsize", self.config.video_bufsize, "-profile:v", "main", "-spatial_aq", "1", "-temporal_aq", "1", "-aq-strength", "8", "-rc-lookahead", "20", "-c:a", self.config.audio_codec, "-ac", "2", "-ar", "48000", "-b:a", self.config.audio_bitrate, "-force_key_frames", f"expr:gte(t,n_forced*{self.config.hls_time})", "-f", "hls", "-hls_time", str(self.config.hls_time), "-hls_playlist_type", "vod", "-hls_flags", "independent_segments", "-hls_list_size", "0", "-hls_segment_filename", str(segment_template), "-progress", "pipe:1", "-nostats", str(task.playlist_file), ] return command def convert_task(self, task: ConversionTask) -> ConversionResult: started = time.perf_counter() duration_seconds = self.probe_duration(task.input_file) last_percent = -1 last_progress_time = time.perf_counter() for attempt in range(self.config.max_retries + 1): if self.stop_event.is_set(): return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=False, message="任务已停止", elapsed_seconds=time.perf_counter() - started, ) self.cleanup_task_output(task) command = self.build_ffmpeg_command(task) self.log( "info", f"开始转码 {task.job_name} / {task.input_file.name} -> {task.playlist_file}", ) process = subprocess.Popen( command, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, text=True, encoding="utf-8", errors="replace", bufsize=1, **subprocess_no_window_kwargs(), ) self.register_process(process) diagnostics: list[str] = [] current_out_time = 0.0 try: while True: if self.stop_event.is_set(): process.terminate() raise RuntimeError("任务已停止") line = process.stdout.readline() if process.stdout else "" if not line and process.poll() is not None: break line = line.strip() if not line: continue if "=" not in line: diagnostics.append(line) diagnostics[:] = diagnostics[-20:] continue key, value = line.split("=", 1) if key == "out_time_ms": try: new_out = float(value) / 1_000_000.0 current_out_time = new_out if duration_seconds > 0: percent = min(int(new_out / duration_seconds * 100), 100) if percent != last_percent: last_percent = percent if self.event_callback: self.event_callback({ "kind": "convert_progress", "job_name": task.job_name, "file_name": task.input_file.name, "percent": percent, "out_time": new_out, "duration": duration_seconds, }) last_progress_time = time.perf_counter() except ValueError: pass elif key == "progress" and value == "end": # ffmpeg signals completion via progress=end pass elif key not in { "frame", "fps", "stream_0_0_q", "bitrate", "total_size", "dup_frames", "drop_frames", "speed", "out_time", "elapsed", "progress", }: diagnostics.append(line) diagnostics[:] = diagnostics[-20:] # Stall detection: kill if no progress for 120s and encoding if duration_seconds > 0 and current_out_time < duration_seconds: stalled = time.perf_counter() - last_progress_time if stalled > 120: self.log("warning", f"FFmpeg 无进度超过 {stalled:.0f}秒,终止重试") process.kill() raise RuntimeError("FFmpeg 进度卡死") process.wait() except subprocess.TimeoutExpired: process.kill() process.wait() if attempt >= self.config.max_retries: return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=False, message="转码超时", elapsed_seconds=time.perf_counter() - started, ) self.log("warning", f"FFmpeg 超时,重试 {attempt+1}/{self.config.max_retries}") time.sleep(2) continue except RuntimeError as error: process.kill() process.wait() if attempt >= self.config.max_retries: return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=False, message=str(error), elapsed_seconds=time.perf_counter() - started, ) self.log("warning", f"FFmpeg 异常: {error},重试 {attempt+1}/{self.config.max_retries}") time.sleep(2) continue except Exception as error: process.kill() process.wait() return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=False, message=str(error), elapsed_seconds=time.perf_counter() - started, ) finally: self.unregister_process(process) if process.returncode == 0: return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=True, message=f"成功(CQ {self.config.video_cq})", elapsed_seconds=time.perf_counter() - started, ) # Diagnose failure playlist_exists = task.playlist_file.exists() segment_files = list(task.output_dir.glob("*.ts")) message = f"FFmpeg 返回码 {process.returncode}" if diagnostics: snippet = " | ".join(diagnostics[-6:]) message += f" | {snippet}" if not playlist_exists: message += " | 未生成 m3u8" elif len(segment_files) >= 2: # Output is partial but non-empty; treat as degraded success self.log("warning", f"FFmpeg 返回码 {process.returncode} 但已生成 {len(segment_files)} 个分片,视为部分成功") return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=True, message=f"部分成功({len(segment_files)} 分片,返回码 {process.returncode})", elapsed_seconds=time.perf_counter() - started, ) if attempt >= self.config.max_retries: return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=False, message=message, elapsed_seconds=time.perf_counter() - started, ) self.log("warning", f"FFmpeg 失败 {message},重试 {attempt+1}/{self.config.max_retries}") time.sleep(2) return ConversionResult( job_name=task.job_name, input_file=task.input_file, output_dir=task.output_dir, success=False, message="未知错误", elapsed_seconds=time.perf_counter() - started, ) def run_conversion_jobs( config: ConversionConfig, event_callback: Callable[[dict], None] | None = None, stop_event: threading.Event | None = None, ) -> ConversionSummary: ensure_directory(config.log_dir) timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") log_file = Path(config.log_dir) / f"convert_{timestamp}.log" logger = OperationLogger( log_file, callback=lambda level, message, line: event_callback( {"kind": "log", "level": level, "message": message, "line": line} ) if event_callback else None, ) converter = VideoConverter(config=config, logger=logger, event_callback=event_callback, stop_event=stop_event) jobs = discover_series_jobs(config.root_dir, config) if not jobs: logger.warning(f"没有在 {Path(config.root_dir).resolve()} 找到可处理的视频任务。") return ConversionSummary( total=0, succeeded=0, failed=0, log_file=log_file, elapsed_seconds=0.0, results=[], ) tasks = build_conversion_tasks(jobs, config) logger.info(f"发现 {len(jobs)} 个剧集目录,待处理 {len(tasks)} 个视频,最大并发 {config.parallel_jobs}") if not tasks: logger.warning("未找到可转码的视频文件。") return ConversionSummary( total=0, succeeded=0, failed=0, log_file=log_file, elapsed_seconds=0.0, results=[], ) started = time.perf_counter() results: list[ConversionResult] = [] finished = 0 with ThreadPoolExecutor(max_workers=config.parallel_jobs) as executor: futures = {executor.submit(converter.convert_task, task): task for task in tasks} for future in as_completed(futures): result = future.result() results.append(result) finished += 1 logger.info( f"[{'PASS' if result.success else 'FAIL'}] {finished}/{len(tasks)} " f"{result.job_name} / {result.input_file.name} | {result.message} | {result.elapsed_seconds:.1f}s" ) if event_callback: event_callback({ "kind": "convert_result", "finished": finished, "total": len(tasks), "result": result, }) if stop_event and stop_event.is_set(): converter.terminate_active_processes() elapsed_seconds = time.perf_counter() - started succeeded = sum(1 for item in results if item.success) failed = len(results) - succeeded logger.info(f"转码完成:成功 {succeeded}/{len(results)},耗时 {elapsed_seconds:.1f} 秒") return ConversionSummary( total=len(results), succeeded=succeeded, failed=failed, log_file=log_file, elapsed_seconds=elapsed_seconds, results=results, ) def run_specific_jobs( jobs: list[ConversionJob], config: ConversionConfig, event_callback: Callable[[dict], None] | None = None, stop_event: threading.Event | None = None, ) -> ConversionSummary: ensure_directory(config.log_dir) timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") log_file = Path(config.log_dir) / f"convert_{timestamp}.log" logger = OperationLogger( log_file, callback=lambda level, message, line: event_callback( {"kind": "log", "level": level, "message": message, "line": line} ) if event_callback else None, ) converter = VideoConverter(config=config, logger=logger, event_callback=event_callback, stop_event=stop_event) tasks = build_conversion_tasks(jobs, config) logger.info(f"处理 {len(jobs)} 个剧集目录,待处理 {len(tasks)} 个视频,最大并发 {config.parallel_jobs}") if not tasks: logger.warning("未找到可转码的视频文件。") return ConversionSummary( total=0, succeeded=0, failed=0, log_file=log_file, elapsed_seconds=0.0, results=[], ) started = time.perf_counter() results: list[ConversionResult] = [] finished = 0 with ThreadPoolExecutor(max_workers=config.parallel_jobs) as executor: futures = {executor.submit(converter.convert_task, task): task for task in tasks} for future in as_completed(futures): result = future.result() results.append(result) finished += 1 logger.info( f"[{'PASS' if result.success else 'FAIL'}] {finished}/{len(tasks)} " f"{result.job_name} / {result.input_file.name} | {result.message} | {result.elapsed_seconds:.1f}s" ) if event_callback: event_callback({ "kind": "convert_result", "finished": finished, "total": len(tasks), "result": result, }) if stop_event and stop_event.is_set(): converter.terminate_active_processes() elapsed_seconds = time.perf_counter() - started succeeded = sum(1 for item in results if item.success) failed = len(results) - succeeded logger.info(f"转码完成:成功 {succeeded}/{len(results)},耗时 {elapsed_seconds:.1f} 秒") return ConversionSummary( total=len(results), succeeded=succeeded, failed=failed, log_file=log_file, elapsed_seconds=elapsed_seconds, results=results, ) def generate_playlist_txt( output_dir: str | Path, url_prefix: str = "http://192.168.10.32:2088", ) -> list[Path]: """Scan output_dir for series folders, generate playlist txt files. Supports both flat mode (m3u8 in series dir) and episode folder mode. One txt per series, placed in output_dir/.""" output = Path(output_dir) if not output.is_dir(): return [] created: list[Path] = [] for series_dir in sorted(output.iterdir(), key=lambda f: natural_sort_key(f.name)): if not series_dir.is_dir(): continue lines: list[str] = [] # Flat mode: m3u8 files directly in series dir flat_m3u8 = sorted(series_dir.glob("*.m3u8"), key=lambda f: natural_sort_key(f.stem)) if flat_m3u8: for m3u8_file in flat_m3u8: ep_name = m3u8_file.stem display_name = series_dir.name + ep_name url = f"{url_prefix}/{series_dir.name}/{m3u8_file.name}" lines.append(f"{display_name},{url}\n") else: # Episode folder mode: m3u8 files in subdirectories episode_dirs = sorted( [d for d in series_dir.iterdir() if d.is_dir()], key=lambda d: natural_sort_key(d.name), ) for ep_dir in episode_dirs: m3u8_files = sorted(ep_dir.glob("*.m3u8")) if not m3u8_files: continue m3u8_name = m3u8_files[0].name display_name = series_dir.name + ep_dir.name url = f"{url_prefix}/{series_dir.name}/{ep_dir.name}/{m3u8_name}" lines.append(f"{display_name},{url}\n") if not lines: continue txt_path = output / f"{series_dir.name}_playlist.txt" txt_path.write_text("".join(lines), encoding="utf-8") created.append(txt_path) return created