"""Bounded, ephemeral H.264 playback; original recordings remain untouched.""" from contextlib import contextmanager import http.client import os import shutil import subprocess import threading import tempfile import time import urllib.request from pathlib import Path SLOTS = threading.BoundedSemaphore(2) MAX_SOURCE_BYTES = 256 * 1024 * 1024 MAX_SESSION_SECONDS = 300 class PlaybackError(Exception): pass @contextmanager def cached_source(source, authorization, reserve_bytes): # A private, bounded scratch file decouples MediaMTX's HTTP write timeout # from browser buffering and pause. Use the data disk, not a small /tmp RAM disk. base = Path(os.environ.get('VISION_DATA', tempfile.gettempdir())) / 'playback-tmp' base.mkdir(mode=0o700, parents=True, exist_ok=True) for old in base.glob('buffer-*'): try: if old.is_symlink() or not old.is_dir() or time.time() - old.stat().st_mtime < 3600: continue (old / 'source.mp4').unlink(missing_ok=True) old.rmdir() except OSError: pass with tempfile.TemporaryDirectory(prefix='buffer-', dir=base) as folder: if shutil.disk_usage(folder).free < reserve_bytes + MAX_SOURCE_BYTES: raise PlaybackError('磁盘剩余空间不足以缓冲回放,请先释放空间') request = urllib.request.Request(source, headers={'Authorization': authorization}) target = Path(folder) / 'source.mp4' deadline = time.monotonic() + 30 size = 0 with urllib.request.urlopen(request, timeout=10) as response, target.open('wb') as output: while chunk := response.read(256 * 1024): size += len(chunk) if size > MAX_SOURCE_BYTES: raise PlaybackError('该录像片段过大,请缩短回放时段') if time.monotonic() > deadline: raise PlaybackError('回放缓冲超时,请稍后重试') if shutil.disk_usage(folder).free < reserve_bytes + len(chunk): raise PlaybackError('磁盘预留空间不足,已停止缓冲回放') output.write(chunk) if not size: raise PlaybackError('该时段没有可读取的录像') yield str(target) def command(source, authorization, duration): executable = shutil.which('ffmpeg') if not executable: raise PlaybackError('回放转换需要服务器安装 FFmpeg') device = os.environ.get('VISION_VAAPI_DEVICE', '/dev/dri/renderD128') acceleration = [] if os.access(device, os.R_OK | os.W_OK): acceleration = ['-hwaccel', 'vaapi', '-hwaccel_device', device, '-hwaccel_output_format', 'vaapi', '-extra_hw_frames', '8'] video = ['-vf', 'fps=15,scale_vaapi=w=1920:h=1080:format=nv12:force_original_aspect_ratio=decrease:force_divisible_by=2', '-c:v', 'h264_vaapi', '-profile:v', 'main', '-level:v', '4.1', '-rc_mode', 'CQP', '-qp', '25', '-g', '30', '-bf', '0'] else: video = ['-vf', 'fps=10,scale=768:432:force_original_aspect_ratio=decrease:force_divisible_by=2', '-c:v', 'libx264', '-threads', '1', '-preset', 'veryfast', '-tune', 'zerolatency', '-profile:v', 'baseline', '-pix_fmt', 'yuv420p', '-crf', '26', '-g', '20', '-bf', '0'] input_options = (['-rw_timeout', '15000000', '-headers', 'Authorization: ' + authorization + '\r\n'] if authorization else []) return [executable, '-hide_banner', '-loglevel', 'error', '-nostdin', '-threads', '1', '-filter_threads', '1'] + acceleration + input_options + [ '-analyzeduration', '1000000', '-probesize', '1000000', # Recordings are finite HTTP downloads. Pacing their input like a # live source stalls upstream writes and truncates long responses. '-i', source, '-t', str(duration), '-map', '0:v:0', '-map', '0:a:0?', '-sn', '-dn'] + video + [ '-c:a', 'aac', '-ar', '48000', '-ac', '1', '-b:a', '64k', '-movflags', '+frag_keyframe+empty_moov+default_base_moof', '-frag_duration', '1000000', '-f', 'mp4', 'pipe:1'] @contextmanager def stream(source, authorization, duration, reserve_bytes=8 * 1024**3): if not 0 < duration <= MAX_SESSION_SECONDS: raise PlaybackError('请刷新页面,使用每段最多 5 分钟的自动连续回放') if not SLOTS.acquire(blocking=False): raise PlaybackError('已有两个回放正在进行,请关闭其他回放后重试') process = timer = None try: if not shutil.which('ffmpeg'): raise PlaybackError('回放转换需要服务器安装 FFmpeg') with cached_source(source, authorization, reserve_bytes) as local_source: # Credentials never appear in the converter's process arguments. process = subprocess.Popen(command(local_source, '', duration), stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL) timer = threading.Timer(duration + 45, process.kill) timer.daemon = True timer.start() try: yield process.stdout finally: if process.poll() is None: process.terminate() try: process.wait(timeout=3) except subprocess.TimeoutExpired: process.kill() process.wait() except (OSError, http.client.HTTPException) as exc: raise PlaybackError('录像转换服务暂不可用') from exc finally: if timer: timer.cancel() if process: if process.poll() is None: process.terminate() try: process.wait(timeout=3) except subprocess.TimeoutExpired: process.kill() process.wait() process.stdout.close() SLOTS.release()