Files
videoconvert/video_pipeline.py
2026-06-12 10:59:32 +08:00

776 lines
27 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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/<series_name>/<episode>/
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/<series_name>/ 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