132 lines
6.1 KiB
Python
132 lines
6.1 KiB
Python
"""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 ['-readrate', '1'])
|
|
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()
|