fix: atomic compression writes, enforce size limit, job status guard
This commit is contained in:
1 parent
1dbfecaac3
commit
f878110c3c
2 files changed
+259
-63
No files matched your search
@@ -3,9 +3,9 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
import threading
|
||||
|
||||
from dotenv import load_dotenv
|
||||
from flask import (
|
||||
@@ -20,7 +20,7 @@ from pathlib import Path
|
||||
|
||||
from src.job import JobQueue
|
||||
from src.model_registry import scan_models, scan_model_groups, ModelConfig
|
||||
from src.preview import extract_thumbnail, extract_sample_frames
|
||||
from src.preview import extract_thumbnail, extract_sample_frames, probe_video
|
||||
from src.zone_config import load_zone_config, list_zone_presets, build_static_roi, get_active_zone
|
||||
|
||||
load_dotenv()
|
||||
@@ -206,39 +206,75 @@ def jobs_list():
|
||||
return render_template("jobs.html", jobs=jobs)
|
||||
|
||||
|
||||
# Serializes compression encodes: encodes are rare, one global lock is enough.
|
||||
_COMPRESS_LOCK = threading.Lock()
|
||||
_AUDIO_BITRATE = 64_000
|
||||
_MIN_VIDEO_BITRATE = 100_000
|
||||
|
||||
|
||||
def _remove_quiet(path: str) -> None:
|
||||
try:
|
||||
os.remove(path)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
def _compressed_cache_valid(out_path: str, file_path: str) -> bool:
|
||||
"""Cache hit only if the compressed file exists and is newer than source."""
|
||||
try:
|
||||
return (
|
||||
os.path.isfile(out_path)
|
||||
and os.path.getmtime(out_path) >= os.path.getmtime(file_path)
|
||||
)
|
||||
except OSError:
|
||||
return False
|
||||
|
||||
|
||||
def _compress_for_download(file_path: str, max_bytes: int = 250 * 1024 * 1024) -> str | None:
|
||||
"""Compress file_path with ffmpeg if larger than max_bytes.
|
||||
|
||||
Returns path to `<name>_compressed.mp4` next to the input, or None to
|
||||
fall back to the original file (small file, ffmpeg missing/failed).
|
||||
fall back to the original file (small file, ffmpeg missing/failed, or
|
||||
result still over the limit).
|
||||
"""
|
||||
if not os.path.isfile(file_path):
|
||||
return None
|
||||
if os.path.getsize(file_path) <= max_bytes:
|
||||
return None
|
||||
out_path = os.path.splitext(file_path)[0] + "_compressed.mp4"
|
||||
if os.path.isfile(out_path):
|
||||
if _compressed_cache_valid(out_path, file_path):
|
||||
return out_path
|
||||
try:
|
||||
probe = subprocess.run(
|
||||
["ffprobe", "-v", "quiet", "-print_format", "json", "-show_format", file_path],
|
||||
capture_output=True, text=True, timeout=60,
|
||||
with _COMPRESS_LOCK:
|
||||
if _compressed_cache_valid(out_path, file_path):
|
||||
return out_path
|
||||
try:
|
||||
duration = float(probe_video(file_path).get("duration") or 0)
|
||||
except Exception:
|
||||
return None
|
||||
if duration <= 0:
|
||||
return None
|
||||
video_bitrate = max(
|
||||
int((max_bytes * 8 * 0.92) / duration) - _AUDIO_BITRATE,
|
||||
_MIN_VIDEO_BITRATE,
|
||||
)
|
||||
duration = float(json.loads(probe.stdout).get("format", {}).get("duration") or 0)
|
||||
except (OSError, subprocess.SubprocessError, ValueError, TypeError, AttributeError):
|
||||
return None
|
||||
if duration <= 0:
|
||||
return None
|
||||
bitrate = int((max_bytes * 8 * 0.92) / duration)
|
||||
try:
|
||||
subprocess.run(
|
||||
["ffmpeg", "-y", "-i", file_path, "-c:v", "libx264",
|
||||
"-b:v", str(bitrate), "-c:a", "aac", "-b:a", "64k", out_path],
|
||||
capture_output=True, timeout=3600, check=True,
|
||||
)
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
return None
|
||||
return out_path if os.path.isfile(out_path) else None
|
||||
tmp_path = out_path + ".tmp"
|
||||
try:
|
||||
subprocess.run(
|
||||
["ffmpeg", "-y", "-i", file_path, "-c:v", "libx264",
|
||||
"-preset", "veryfast", "-b:v", str(video_bitrate),
|
||||
"-c:a", "aac", "-b:a", str(_AUDIO_BITRATE), "-f", "mp4",
|
||||
tmp_path],
|
||||
capture_output=True, timeout=3600, check=True,
|
||||
stdin=subprocess.DEVNULL,
|
||||
)
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
_remove_quiet(tmp_path)
|
||||
return None
|
||||
if not os.path.isfile(tmp_path) or os.path.getsize(tmp_path) > max_bytes:
|
||||
_remove_quiet(tmp_path)
|
||||
return None
|
||||
os.replace(tmp_path, out_path)
|
||||
return out_path
|
||||
|
||||
|
||||
@app.route("/download/<job_id>/<filename>")
|
||||
@@ -246,6 +282,8 @@ def download(job_id, filename):
|
||||
job = job_queue.get_job(job_id)
|
||||
if job is None:
|
||||
return "Job not found", 404
|
||||
if job.status.name in ("PENDING", "RUNNING"):
|
||||
return "Job still processing", 409
|
||||
safe_filename = secure_filename(filename)
|
||||
output_dir_abs = os.path.abspath(job.output_dir)
|
||||
file_path = os.path.abspath(os.path.join(job.output_dir, safe_filename))
|
||||
|
||||
Reference in new issue
Block a user