studio / renderer /platform /processor.py
Ava2lon's picture
Upload 52 files
5d782cc verified
Raw
History Blame
33.5 kB
from __future__ import annotations
import json
import math
import mimetypes
import re
import shutil
import zipfile
from pathlib import Path
from typing import Any
from renderer.core.config import Settings
from renderer.core.ingest import AssetIngestor
from renderer.core.models import TaskResult
from renderer.core.utils import safe_filename, temp_workdir, write_json
from renderer.ffmpeg.assets import AssetProbe
from renderer.ffmpeg.command import FFmpegCommand
from renderer.ffmpeg.runner import CommandResult, FFmpegRunner
from renderer.templates import get_platform_profile
TOOLKIT_TASKS = {
"cut",
"trim",
"split",
"concat",
"merge",
"compress",
"normalize",
"resize",
"crop",
"rotate",
"flip",
"scale",
"zoom",
"pan",
"speed",
"speed_ramp",
"time_remap",
"slow_motion",
"fast_motion",
"reverse",
"reverse_playback",
"freeze_frame",
"motion_blur",
"stabilization",
"lens_correction",
"loop",
"extract_audio",
"thumbnail",
"gif",
"frames",
"watermark",
"overlay_text",
"blur_background",
"burn_subtitles",
"convert",
"merge_audio",
"noise_reduction",
"equalizer",
"compressor",
"limiter",
"pitch_shift",
"voice_changer",
"ai_enhancement",
"green_screen",
"chroma_key",
"blue_screen",
"ai_background_removal",
"inspect",
"loudness_analyze",
"silence_detect",
"black_detect",
"scene_detect",
"contact_sheet",
"hls",
}
class PlatformProcessor:
def __init__(self, settings: Settings | None = None, log=None) -> None:
self.settings = settings or Settings()
self.settings.ensure_dirs()
self._commands: list[list[str]] = []
self._logs: list[str] = []
self.runner = FFmpegRunner(self.settings.ffmpeg_timeout_seconds, log=log, on_command=self._record_command)
self.ingest = AssetIngestor(self.settings)
self.assets = AssetProbe(self.settings.metadata_cache)
def ingest_sources(self, sources: list[dict[str, Any]], job_id: str) -> TaskResult:
with temp_workdir(self.settings.temp_dir, f"{job_id}_ingest") as work:
workdir = Path(work)
staged: list[dict[str, Any]] = []
for index, source in enumerate(sources):
url = str(source.get("url") or source.get("source") or "").strip()
if not url:
continue
source_type = str(source.get("type") or _source_type(url))
if source_type == "youtube":
staged.append(
{
"source": url,
"source_type": source_type,
"status": "registered",
"note": "YouTube ingestion is registered for automation; provide a direct downloadable media URL or install yt-dlp for local extraction.",
}
)
continue
resolved = self.ingest.resolve(url, workdir / "inputs", f"source_{index:03d}")
metadata = self.assets.probe(resolved).__dict__
staged.append({"source": url, "source_type": source_type, "path": str(resolved), "metadata": metadata})
output = self._json_artifact(job_id, "ingest_manifest", {"assets": staged, "asset_count": len(staged)})
return self._result(output, {"task": "ingest", "asset_count": len(staged)})
def analyze(self, media: str, job_id: str, *, transcript: str = "", platform: str | None = None) -> TaskResult:
with temp_workdir(self.settings.temp_dir, f"{job_id}_analyze") as work:
workdir = Path(work)
source = self.ingest.resolve(media, workdir / "inputs", "media")
metadata = self.assets.probe(source)
duration = max(0.0, metadata.duration)
words = transcript.split()
words_per_minute = (len(words) / duration * 60) if duration > 0 and words else None
highlights = _highlight_windows(duration)
viral_score = _viral_score(duration, bool(words), metadata.width, metadata.height)
analysis = {
"media": str(source),
"metadata": metadata.__dict__,
"transcript": transcript,
"highlight_moments": highlights,
"viral_score": viral_score,
"hook_quality": _hook_quality(transcript),
"audience_retention_estimate": _retention_estimate(duration, viral_score),
"engagement_prediction": _engagement_prediction(viral_score),
"audience_persona": _persona(transcript),
"platform_recommendations": _platform_recommendations(duration, metadata.width, metadata.height, platform),
"scene_segmentation": highlights,
"speech_pacing": {
"words_per_minute": round(words_per_minute, 1) if words_per_minute else None,
"label": _pacing_label(words_per_minute),
},
"silence_detection": {
"estimated_silence_ratio": 0.0 if transcript else 0.18,
"note": "Heuristic estimate; use transcription with word timestamps for precise silence spans.",
},
}
output = self._json_artifact(job_id, "analysis", analysis)
return self._result(output, {"task": "analyze", "viral_score": viral_score, "duration": duration})
def metadata(self, job_id: str, *, topic: str = "", transcript: str = "", platform: str | None = None) -> TaskResult:
text = transcript or topic or "Untitled video"
title = _title_from_text(text, platform)
tags = _hashtags(text, platform)
payload = {
"title": title,
"description": _description(text, tags),
"hashtags": tags,
"keywords": _keywords(text),
"chapters": _chapters(text),
"seo_tags": _keywords(text) + [platform] if platform else _keywords(text),
"suggested_upload_schedule": _schedule(platform),
"platform": platform or "general",
}
output = self._json_artifact(job_id, "metadata", payload)
return self._result(output, {"task": "metadata", "keyword_count": len(payload["keywords"])})
def publish(self, payload: dict[str, Any], job_id: str) -> TaskResult:
platforms = payload.get("platforms") or [payload.get("platform") or "draft"]
manifest = {
"publish_state": "draft_ready" if payload.get("draft", True) else "credentials_required",
"platforms": platforms,
"scheduled_at": payload.get("scheduled_at"),
"asset": payload.get("asset") or payload.get("media"),
"title": payload.get("title"),
"description": payload.get("description"),
"retry_policy": {"max_attempts": 3, "backoff_seconds": 60},
"note": "Direct publishing requires platform OAuth/API credentials configured outside this CPU render worker.",
}
output = self._json_artifact(job_id, "publish_manifest", manifest)
return self._result(output, {"task": "publish", "platform_count": len(platforms)})
def clips(self, media: str, job_id: str, clips: list[dict[str, Any]] | None = None) -> TaskResult:
with temp_workdir(self.settings.temp_dir, f"{job_id}_clips") as work:
workdir = Path(work)
source = self.ingest.resolve(media, workdir / "inputs", "media")
metadata = self.assets.probe(source)
clip_specs = clips or _highlight_windows(metadata.duration)
outputs: list[Path] = []
for index, clip in enumerate(clip_specs):
start = max(0.0, float(clip.get("start", 0)))
end = float(clip.get("end", start + clip.get("duration", 8)))
duration = max(0.2, end - start)
target = workdir / f"clip_{index + 1:02d}.mp4"
command = (
FFmpegCommand()
.add("-hide_banner", "-ss", start)
.input(source)
.add("-t", duration, "-c", "copy")
.overwrite()
.add(target)
.build()
)
self._run(command)
outputs.append(target)
if len(outputs) == 1:
final = self._export(outputs[0], job_id, outputs[0].name)
else:
final = self.settings.exports_dir / f"{job_id}_clips.zip"
with zipfile.ZipFile(final, "w", zipfile.ZIP_DEFLATED) as archive:
for path in outputs:
archive.write(path, path.name)
return self._result(final, {"task": "clips", "clip_count": len(outputs)})
def thumbnail(self, media: str, job_id: str, *, text: str = "", timestamp: float | None = None, template: str = "bold") -> TaskResult:
with temp_workdir(self.settings.temp_dir, f"{job_id}_thumb") as work:
workdir = Path(work)
source = self.ingest.resolve(media, workdir / "inputs", "media")
metadata = self.assets.probe(source)
target = workdir / "thumbnail.jpg"
seek = timestamp if timestamp is not None else max(0.0, min(metadata.duration * 0.2, 8.0))
vf = "scale=1280:720:force_original_aspect_ratio=increase,crop=1280:720"
if text:
vf += "," + _drawtext_filter(text, template)
command = (
FFmpegCommand()
.add("-hide_banner", "-ss", seek)
.input(source)
.add("-frames:v", 1, "-vf", vf, "-q:v", 2)
.overwrite()
.add(target)
.build()
)
self._run(command)
final = self._export(target, job_id, "thumbnail.jpg")
return self._result(final, {"task": "thumbnail", "timestamp": seek})
def inspect_media(self, media: str, job_id: str) -> TaskResult:
with temp_workdir(self.settings.temp_dir, f"{job_id}_inspect") as work:
source = self.ingest.resolve(media, Path(work) / "inputs", "media")
metadata = self.assets.probe(source).__dict__
metadata["source"] = media
output = self._json_artifact(job_id, "media_inspection", metadata)
return self._result(output, {"task": "inspect", "duration": metadata.get("duration", 0)})
def media_analysis(self, task: str, media: str, job_id: str, params: dict[str, Any]) -> TaskResult:
with temp_workdir(self.settings.temp_dir, f"{job_id}_{task}") as work:
source = self.ingest.resolve(media, Path(work) / "inputs", "media")
command = _analysis_command(task, source, params)
result = self._run(command)
analysis = _parse_analysis(task, result.stderr, params)
analysis.update({"task": task, "source": media})
output = self._json_artifact(job_id, task, analysis)
count = len(analysis.get("events", []))
return self._result(output, {"task": task, "event_count": count})
def toolkit(self, payload: dict[str, Any], job_id: str) -> TaskResult:
task = str(payload.get("task") or payload.get("operation") or "").strip()
if task not in TOOLKIT_TASKS:
raise ValueError(f"Unsupported toolkit task: {task}")
if task == "thumbnail":
return self.thumbnail(str(payload["input"]), job_id, text=str(payload.get("text") or ""), timestamp=payload.get("timestamp"))
if task == "split":
clips = payload.get("clips") if isinstance(payload.get("clips"), list) else payload.get("params", {}).get("clips")
return self.clips(str(payload.get("input") or payload.get("media")), job_id, clips)
if task == "inspect":
return self.inspect_media(str(payload.get("input") or payload.get("media")), job_id)
if task in {"loudness_analyze", "silence_detect", "black_detect", "scene_detect"}:
params = payload.get("params") if isinstance(payload.get("params"), dict) else payload
return self.media_analysis(task, str(payload.get("input") or payload.get("media")), job_id, params)
with temp_workdir(self.settings.temp_dir, f"{job_id}_{task}") as work:
workdir = Path(work)
source_value = payload.get("input") or payload.get("media")
source = self.ingest.resolve(str(source_value), workdir / "inputs", "media") if source_value else workdir / "concat_placeholder.mp4"
params = payload.get("params") if isinstance(payload.get("params"), dict) else payload
output_name = safe_filename(str(payload.get("output_name") or _default_output_name(task)))
output = workdir / output_name
if task == "extract_audio" and output.suffix.lower() != ".mp3":
output = output.with_suffix(".mp3")
if task == "gif" and output.suffix.lower() != ".gif":
output = output.with_suffix(".gif")
if task == "contact_sheet" and output.suffix.lower() not in {".jpg", ".jpeg", ".png"}:
output = output.with_suffix(".jpg")
command = self._toolkit_command(task, source, output, params, workdir)
self._run(command)
if task == "frames":
final = self.settings.exports_dir / f"{job_id}_frames.zip"
with zipfile.ZipFile(final, "w", zipfile.ZIP_DEFLATED) as archive:
for frame in sorted(workdir.glob("frame_*.jpg")):
archive.write(frame, frame.name)
elif task == "hls":
final = self.settings.exports_dir / f"{job_id}_hls.zip"
with zipfile.ZipFile(final, "w", zipfile.ZIP_DEFLATED) as archive:
for path in sorted(workdir.glob("hls_*")):
archive.write(path, path.name)
else:
final = self._export(output, job_id, output.name)
return self._result(final, {"task": task})
def _toolkit_command(self, task: str, source: Path, output: Path, params: dict[str, Any], workdir: Path) -> list[str]:
cmd = FFmpegCommand().add("-hide_banner")
if task == "loop":
cmd.add("-stream_loop", int(params.get("loops", -1)))
if task in {"cut", "trim", "gif", "freeze_frame"} and params.get("start") is not None:
cmd.add("-ss", float(params.get("start", 0)))
cmd.input(source)
if task == "merge_audio":
audio = self.ingest.resolve(str(params.get("audio")), workdir / "inputs", "audio")
cmd.input(audio)
return cmd.add("-map", "0:v", "-map", "1:a", "-c:v", "copy", "-c:a", "aac", "-shortest").overwrite().add(output).build()
if task == "watermark":
image = self.ingest.resolve(str(params.get("watermark") or params.get("image")), workdir / "inputs", "watermark")
cmd.input(image)
return cmd.add("-filter_complex", "[1:v]scale=iw*0.18:-1[wm];[0:v][wm]overlay=W-w-24:H-h-24", "-c:a", "copy").overwrite().add(output).build()
if task in {"concat", "merge"}:
inputs = params.get("inputs")
if not isinstance(inputs, list) or not inputs:
raise ValueError("Concat requires params.inputs")
concat_file = workdir / "concat.txt"
lines: list[str] = []
for index, item in enumerate(inputs):
media = self.ingest.resolve(str(item), workdir / "inputs", f"concat_{index:03d}")
lines.append(f"file '{str(media).replace(chr(39), chr(39) + chr(92) + chr(39) + chr(39))}'")
concat_file.write_text("\n".join(lines), encoding="utf-8")
return FFmpegCommand().add("-hide_banner", "-f", "concat", "-safe", "0").input(concat_file).add("-c", "copy").overwrite().add(output).build()
if task == "contact_sheet":
columns = max(1, min(10, int(params.get("columns", 4))))
rows = max(1, min(10, int(params.get("rows", 4))))
interval = max(0.1, float(params.get("interval_seconds", 5)))
width = max(80, min(1920, int(params.get("thumbnail_width", 320))))
vf = f"fps=1/{interval},scale={width}:-1,tile={columns}x{rows}:padding=4:margin=4"
return cmd.add("-vf", vf, "-frames:v", 1, "-q:v", 2).overwrite().add(output).build()
if task == "hls":
segment_seconds = max(1, min(30, int(params.get("segment_seconds", 6))))
playlist_type = str(params.get("playlist_type", "vod")).lower()
if playlist_type not in {"event", "vod"}:
raise ValueError("HLS playlist_type must be 'event' or 'vod'")
playlist = workdir / "hls_playlist.m3u8"
segments = workdir / "hls_segment_%05d.ts"
return (
cmd.add(
"-c:v", "libx264", "-preset", self.settings.preset, "-crf", int(params.get("crf", self.settings.crf)),
"-c:a", "aac", "-f", "hls", "-hls_time", segment_seconds, "-hls_playlist_type", playlist_type,
"-hls_segment_filename", segments,
)
.overwrite()
.add(playlist)
.build()
)
duration = params.get("duration")
if task in {"cut", "trim", "gif", "loop", "freeze_frame"} and duration is not None:
cmd.add("-t", float(duration))
vf = _video_filter(task, params)
af = _audio_filter(task, params)
if vf:
cmd.add("-vf", vf)
if af:
cmd.add("-af", af)
if task == "extract_audio":
return cmd.add("-vn", "-c:a", "mp3", "-b:a", "192k").overwrite().add(output).build()
if task == "frames":
return cmd.add("-vf", vf or "fps=1", "-q:v", 2).overwrite().add(workdir / "frame_%04d.jpg").build()
if task == "gif":
return cmd.add("-loop", 0).overwrite().add(output).build()
if task in {
"cut",
"trim",
"compress",
"normalize",
"resize",
"crop",
"rotate",
"flip",
"scale",
"zoom",
"pan",
"speed",
"speed_ramp",
"time_remap",
"slow_motion",
"fast_motion",
"reverse",
"reverse_playback",
"freeze_frame",
"motion_blur",
"stabilization",
"lens_correction",
"loop",
"overlay_text",
"blur_background",
"burn_subtitles",
"convert",
"noise_reduction",
"equalizer",
"compressor",
"limiter",
"pitch_shift",
"voice_changer",
"green_screen",
"chroma_key",
"blue_screen",
"ai_background_removal",
}:
cmd.add("-c:v", "libx264", "-preset", self.settings.preset, "-crf", int(params.get("crf", self.settings.crf)), "-c:a", "aac")
return cmd.overwrite().add(output).build()
def _json_artifact(self, job_id: str, name: str, payload: dict[str, Any]) -> Path:
output = self.settings.exports_dir / f"{job_id}_{safe_filename(name)}.json"
write_json(output, payload)
return output
def _export(self, source: Path, job_id: str, output_name: str) -> Path:
target = self.settings.exports_dir / f"{job_id}_{safe_filename(output_name)}"
shutil.copy2(source, target)
return target
def _run(self, command: list[str]) -> CommandResult:
result = self.runner.run(command)
if result.stderr:
self._logs.append(result.stderr[-4000:])
return result
def _record_command(self, command: list[str]) -> None:
self._commands.append(command)
def _result(self, output: Path | None, metrics: dict[str, Any]) -> TaskResult:
return TaskResult(output_path=output, commands=list(self._commands), metrics=metrics, logs=list(self._logs))
def supported_toolkit_tasks() -> list[str]:
return sorted(TOOLKIT_TASKS)
def _source_type(url: str) -> str:
lowered = url.lower()
if "youtube.com" in lowered or "youtu.be" in lowered:
return "youtube"
if "drive.google.com" in lowered:
return "google_drive"
if "s3" in lowered or "amazonaws.com" in lowered:
return "s3"
return "direct_url"
def _highlight_windows(duration: float) -> list[dict[str, Any]]:
if duration <= 0:
return [{"start": 0, "end": 8, "reason": "default opener"}]
windows = [{"start": 0, "end": min(duration, 8), "reason": "opening hook"}]
if duration > 18:
middle = max(0.0, duration * 0.42)
windows.append({"start": round(middle, 2), "end": round(min(duration, middle + 10), 2), "reason": "midpoint payoff"})
if duration > 35:
end = max(0.0, duration - 12)
windows.append({"start": round(end, 2), "end": round(duration, 2), "reason": "closing CTA"})
return windows
def _viral_score(duration: float, has_words: bool, width: int | None, height: int | None) -> int:
score = 48
if 7 <= duration <= 60:
score += 20
elif duration <= 180:
score += 8
if width and height and height >= width:
score += 14
if has_words:
score += 10
return max(1, min(100, score))
def _hook_quality(transcript: str) -> dict[str, Any]:
opener = " ".join(transcript.split()[:18])
signals = sum(1 for token in ("how", "why", "secret", "mistake", "stop", "watch", "you") if token in opener.lower())
return {"score": min(100, 45 + signals * 12), "opening_text": opener}
def _retention_estimate(duration: float, viral_score: int) -> dict[str, Any]:
first_3s = min(96, 55 + viral_score * 0.35)
completion = max(18, min(88, first_3s - math.log(max(duration, 1), 1.8)))
return {"first_3_seconds_percent": round(first_3s, 1), "completion_percent": round(completion, 1)}
def _engagement_prediction(score: int) -> str:
if score >= 80:
return "high"
if score >= 62:
return "medium"
return "needs_work"
def _persona(text: str) -> str:
lowered = text.lower()
if any(word in lowered for word in ("founder", "startup", "product", "launch")):
return "builders and product-led founders"
if any(word in lowered for word in ("money", "sales", "growth", "marketing")):
return "growth-minded operators"
if any(word in lowered for word in ("learn", "tutorial", "how")):
return "learners seeking practical instruction"
return "general short-form viewers"
def _platform_recommendations(duration: float, width: int | None, height: int | None, platform: str | None) -> list[dict[str, Any]]:
vertical = bool(width and height and height >= width)
candidates = ["tiktok", "instagram_reels", "youtube_shorts"] if vertical else ["youtube_1080p", "linkedin_video"]
if platform and platform not in candidates:
candidates.insert(0, platform)
return [{"platform": item, "fit": "strong" if duration <= 90 else "medium"} for item in candidates]
def _pacing_label(wpm: float | None) -> str:
if not wpm:
return "unknown"
if wpm < 125:
return "slow"
if wpm > 185:
return "fast"
return "clear"
def _title_from_text(text: str, platform: str | None) -> str:
words = [word.strip(".,:;!?") for word in text.split() if word.strip(".,:;!?")]
title = " ".join(words[:9]) or "Untitled Video"
suffix = " #Shorts" if platform in {"youtube_shorts", "tiktok", "instagram_reels"} else ""
return f"{title.title()}{suffix}"
def _description(text: str, tags: list[str]) -> str:
summary = " ".join(text.split()[:42])
return f"{summary}\n\n{' '.join(tags)}".strip()
def _hashtags(text: str, platform: str | None) -> list[str]:
base = ["#video", "#content"]
if platform:
base.append(f"#{platform.replace('_', '')}")
for keyword in _keywords(text)[:5]:
tag = "#" + "".join(ch for ch in keyword.title() if ch.isalnum())
if tag not in base:
base.append(tag)
return base[:8]
def _keywords(text: str) -> list[str]:
stop = {"the", "and", "for", "with", "that", "this", "your", "you", "are", "from", "into", "video"}
words = [word.strip(".,:;!?").lower() for word in text.split()]
unique: list[str] = []
for word in words:
if len(word) < 4 or word in stop or word in unique:
continue
unique.append(word)
return unique[:12]
def _chapters(text: str) -> list[dict[str, Any]]:
sentences = [part.strip() for part in text.replace("?", ".").replace("!", ".").split(".") if part.strip()]
return [{"time": f"0:{index * 15:02d}", "title": sentence[:60]} for index, sentence in enumerate(sentences[:6])]
def _schedule(platform: str | None) -> dict[str, str]:
if platform in {"linkedin_video", "youtube_1080p"}:
return {"day": "Tuesday", "time": "09:00 local"}
return {"day": "Thursday", "time": "18:00 local"}
def _video_filter(task: str, params: dict[str, Any]) -> str:
if task in {"compress", "normalize"}:
profile = get_platform_profile(params.get("platform"))
return f"scale={profile.width}:{profile.height}:force_original_aspect_ratio=increase,crop={profile.width}:{profile.height},fps={profile.fps},format=yuv420p"
if task in {"resize", "scale"}:
return f"scale={int(params.get('width', 1080))}:{int(params.get('height', 1920))}"
if task == "crop":
return f"crop={int(params.get('width', 1080))}:{int(params.get('height', 1080))}:{int(params.get('x', 0))}:{int(params.get('y', 0))}"
if task == "rotate":
return {"90": "transpose=1", "180": "hflip,vflip", "270": "transpose=2"}.get(str(params.get("degrees", "90")), "transpose=1")
if task == "flip":
axis = str(params.get("axis", "horizontal"))
return "vflip" if axis in {"vertical", "y"} else "hflip"
if task in {"speed", "speed_ramp", "time_remap", "slow_motion", "fast_motion"}:
default = 0.5 if task == "slow_motion" else 2.0 if task == "fast_motion" else 1.0
factor = max(0.25, min(4.0, float(params.get("factor", default))))
return f"setpts={1 / factor:.4f}*PTS"
if task == "zoom":
factor = max(1.0, min(4.0, float(params.get("factor", 1.2))))
return f"scale=iw*{factor:.3f}:ih*{factor:.3f},crop=iw/{factor:.3f}:ih/{factor:.3f}"
if task == "pan":
width = int(params.get("width", 1080))
height = int(params.get("height", 1920))
x = str(params.get("x", "(iw-ow)/2"))
y = str(params.get("y", "(ih-oh)/2"))
return f"crop={width}:{height}:{x}:{y}"
if task in {"motion_blur", "freeze_frame"}:
return "tmix=frames=3:weights='1 2 1'"
if task == "stabilization":
return "deshake"
if task == "lens_correction":
return f"lenscorrection=k1={float(params.get('k1', -0.15))}:k2={float(params.get('k2', 0.05))}"
if task in {"reverse", "reverse_playback"}:
return "reverse"
if task in {"green_screen", "chroma_key", "ai_background_removal"}:
color = str(params.get("color") or "0x00ff00")
similarity = float(params.get("similarity", 0.18))
blend = float(params.get("blend", 0.08))
return f"chromakey={color}:{similarity}:{blend}"
if task == "blue_screen":
similarity = float(params.get("similarity", 0.18))
blend = float(params.get("blend", 0.08))
return f"chromakey=0x0000ff:{similarity}:{blend}"
if task == "speed":
factor = max(0.25, min(4.0, float(params.get("factor", 1.0))))
return f"setpts={1 / factor:.4f}*PTS"
if task == "gif":
return f"fps={int(params.get('fps', 12))},scale={int(params.get('width', 540))}:-1:flags=lanczos"
if task == "frames":
return f"fps={float(params.get('fps', 1))}"
if task == "overlay_text":
return _drawtext_filter(str(params.get("text") or "Text"), "bold")
if task == "blur_background":
return "gblur=sigma=18"
if task == "burn_subtitles":
subtitles = str(params.get("subtitles") or "").replace("\\", "/").replace(":", r"\:")
return f"subtitles='{subtitles}'"
return ""
def _audio_filter(task: str, params: dict[str, Any]) -> str:
if task in {"speed", "speed_ramp", "time_remap", "slow_motion", "fast_motion"}:
default = 0.5 if task == "slow_motion" else 2.0 if task == "fast_motion" else 1.0
factor = max(0.5, min(2.0, float(params.get("factor", default))))
return f"atempo={factor}"
if task in {"reverse", "reverse_playback"}:
return "areverse"
if task in {"noise_reduction", "ai_enhancement"}:
return "afftdn=nf=-25"
if task == "equalizer":
return f"equalizer=f={float(params.get('frequency', 1000))}:width_type=o:width={float(params.get('width', 1))}:g={float(params.get('gain', 3))}"
if task == "compressor":
return "acompressor=threshold=-18dB:ratio=3:attack=20:release=250"
if task == "limiter":
return "alimiter=limit=0.95"
if task in {"pitch_shift", "voice_changer"}:
factor = max(0.5, min(2.0, float(params.get("factor", 1.0))))
return f"asetrate=48000*{factor:.4f},aresample=48000,atempo={1 / factor:.4f}"
return ""
def _drawtext_filter(text: str, template: str) -> str:
escaped = text.replace("\\", "\\\\").replace(":", r"\:").replace("'", r"\'")
color = "yellow" if template == "bold" else "white"
return (
"drawtext="
f"text='{escaped}':fontcolor={color}:fontsize=58:"
"box=1:boxcolor=black@0.55:boxborderw=24:"
"x=(w-text_w)/2:y=h-(text_h*3)"
)
def _default_output_name(task: str) -> str:
if task == "extract_audio":
return "audio.mp3"
if task == "gif":
return "clip.gif"
if task == "thumbnail":
return "thumbnail.jpg"
if task == "contact_sheet":
return "contact_sheet.jpg"
if task == "hls":
return "stream.zip"
return f"{task}.mp4"
def _analysis_command(task: str, source: Path, params: dict[str, Any]) -> list[str]:
command = FFmpegCommand().add("-hide_banner").input(source)
if task == "loudness_analyze":
return command.add("-vn", "-af", "ebur128=framelog=verbose", "-f", "null", "-").build()
if task == "silence_detect":
noise = float(params.get("noise_db", -35))
duration = float(params.get("duration_seconds", 0.5))
return command.add("-af", f"silencedetect=noise={noise}dB:d={duration}", "-f", "null", "-").build()
if task == "black_detect":
threshold = max(0.0, min(1.0, float(params.get("pixel_threshold", 0.1))))
duration = max(0.1, float(params.get("duration_seconds", 0.1)))
return command.add("-vf", f"blackdetect=d={duration}:pic_th={threshold}", "-an", "-f", "null", "-").build()
if task == "scene_detect":
threshold = max(0.01, min(1.0, float(params.get("threshold", 0.4))))
return command.add("-vf", f"select='gt(scene,{threshold})',showinfo", "-an", "-f", "null", "-").build()
raise ValueError(f"Unsupported analysis task: {task}")
def _parse_analysis(task: str, stderr: str, params: dict[str, Any]) -> dict[str, Any]:
events: list[dict[str, Any]] = []
if task == "loudness_analyze":
integrated = re.findall(r"I:\s*(-?\d+(?:\.\d+)?)\s*LUFS", stderr)
loudness_range = re.findall(r"LRA:\s*(-?\d+(?:\.\d+)?)\s*LU", stderr)
true_peak = re.findall(r"Peak:\s*(-?\d+(?:\.\d+)?)\s*dBFS", stderr)
return {
"integrated_lufs": float(integrated[-1]) if integrated else None,
"loudness_range_lu": float(loudness_range[-1]) if loudness_range else None,
"true_peak_dbfs": float(true_peak[-1]) if true_peak else None,
"target_lufs": float(params.get("target_lufs", -14)),
}
if task == "silence_detect":
starts = re.finditer(r"silence_start:\s*([0-9.]+)", stderr)
ends = list(re.finditer(r"silence_end:\s*([0-9.]+).*?silence_duration:\s*([0-9.]+)", stderr))
end_index = 0
for match in starts:
if end_index < len(ends):
end = ends[end_index]
if float(end.group(1)) >= float(match.group(1)):
events.append({"start": float(match.group(1)), "end": float(end.group(1)), "duration": float(end.group(2))})
end_index += 1
else:
events.append({"start": float(match.group(1))})
return {"events": events, "noise_db": float(params.get("noise_db", -35))}
if task == "black_detect":
for match in re.finditer(r"black_start:([0-9.]+)\s+black_end:([0-9.]+)\s+black_duration:([0-9.]+)", stderr):
events.append({"start": float(match.group(1)), "end": float(match.group(2)), "duration": float(match.group(3))})
return {"events": events, "pixel_threshold": float(params.get("pixel_threshold", 0.1))}
for match in re.finditer(r"pts_time:([0-9.]+)", stderr):
events.append({"time": float(match.group(1))})
return {"events": events, "threshold": float(params.get("threshold", 0.4))}