#!/usr/bin/env python3 """Run a bounded synthetic FFmpeg workflow from intake through acceptance.""" from __future__ import annotations import argparse import hashlib import json import os import shutil import subprocess from pathlib import Path from typing import Any class WorkflowError(Exception): """A controlled workflow failure with an actionable message.""" def resolve_tool(name: str) -> str: expanded = os.path.expanduser(name) if os.path.sep in expanded: candidate = Path(expanded) if candidate.is_file() and os.access(candidate, os.X_OK): return str(candidate.resolve()) else: resolved = shutil.which(name) if resolved: return resolved raise WorkflowError(f"required executable not found: {name}") def run(argv: list[str], *, timeout: float = 30.0) -> subprocess.CompletedProcess[str]: try: result = subprocess.run( argv, stdin=subprocess.DEVNULL, capture_output=True, text=True, timeout=timeout, check=False, ) except (OSError, subprocess.TimeoutExpired) as exc: raise WorkflowError(f"command could not complete: {argv[0]}: {exc}") from exc if result.returncode: detail = result.stderr.strip()[-2000:] or result.stdout.strip()[-2000:] raise WorkflowError(f"command failed ({result.returncode}): {' '.join(argv)}\n{detail}") return result def write_json(path: Path, value: object) -> None: path.write_text(json.dumps(value, indent=2, sort_keys=True) + "\n", encoding="utf-8") def sha256(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as handle: for block in iter(lambda: handle.read(1024 * 1024), b""): digest.update(block) return f"sha256:{digest.hexdigest()}" def probe(ffprobe: str, media: Path) -> tuple[list[str], dict[str, Any]]: command = [ ffprobe, "-v", "error", "-show_format", "-show_streams", "-show_chapters", "-of", "json", str(media), ] result = run(command) try: document = json.loads(result.stdout) except json.JSONDecodeError as exc: raise WorkflowError(f"ffprobe returned invalid JSON for {media.name}: {exc}") from exc if not isinstance(document, dict): raise WorkflowError(f"ffprobe JSON root for {media.name} must be an object") return command, document def first_line(command: list[str]) -> str: output = run(command).stdout.splitlines() if not output: raise WorkflowError(f"version command returned no output: {command[0]}") return output[0] def duration(document: dict[str, Any]) -> float | None: raw = document.get("format", {}).get("duration") try: return float(raw) except (TypeError, ValueError): return None def stream_types(document: dict[str, Any]) -> list[str]: streams = document.get("streams", []) if not isinstance(streams, list): return [] return [ stream.get("codec_type") for stream in streams if isinstance(stream, dict) and isinstance(stream.get("codec_type"), str) ] def prepare_workspace(path: Path) -> Path: resolved = path.expanduser().resolve() if resolved.exists() and any(resolved.iterdir()): raise WorkflowError(f"workspace must be absent or empty: {resolved}") resolved.mkdir(parents=True, exist_ok=True) return resolved def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("workspace", help="new or empty directory for generated artifacts") parser.add_argument("--ffmpeg", default="ffmpeg", help="ffmpeg executable") parser.add_argument("--ffprobe", default="ffprobe", help="ffprobe executable") parser.add_argument( "--duration", type=float, default=3.0, help="synthetic source duration (1-30 seconds)" ) parser.add_argument( "--timeout", type=float, default=30.0, help="per-command timeout (1-300 seconds)" ) parser.add_argument("--json", action="store_true", help="emit compact summary JSON") return parser def main(argv: list[str] | None = None) -> int: args = build_parser().parse_args(argv) try: if not 1.0 <= args.duration <= 30.0: raise WorkflowError("--duration must be between 1 and 30 seconds") if not 1.0 <= args.timeout <= 300.0: raise WorkflowError("--timeout must be between 1 and 300 seconds") workspace = prepare_workspace(Path(args.workspace)) ffmpeg = resolve_tool(args.ffmpeg) ffprobe = resolve_tool(args.ffprobe) source = workspace / "synthetic-source.mkv" output = workspace / "edited-output.mkv" start = round(args.duration * 0.2, 3) end = round(args.duration * 0.8, 3) expected_duration = round(end - start, 3) commands: list[list[str]] = [] ffmpeg_version = first_line([ffmpeg, "-version"]) ffprobe_version = first_line([ffprobe, "-version"]) generate = [ ffmpeg, "-v", "error", "-n", "-f", "lavfi", "-i", f"testsrc2=size=320x180:rate=24:duration={args.duration:g}", "-f", "lavfi", "-i", f"sine=frequency=880:sample_rate=48000:duration={args.duration:g}", "-map", "0:v:0", "-map", "1:a:0", "-c:v", "mpeg4", "-q:v", "5", "-c:a", "pcm_s16le", "-shortest", str(source), ] run(generate, timeout=args.timeout) commands.append(generate) source_probe_command, source_probe = probe(ffprobe, source) commands.append(source_probe_command) write_json(workspace / "source-probe.json", source_probe) required_types = {"video", "audio"} if not required_types.issubset(stream_types(source_probe)): raise WorkflowError("synthetic source does not contain both video and audio") source_digest = sha256(source) intake = { "schema_version": 1, "workflow_id": "synthetic-editorial-example", "assets": [ { "id": "asset-synthetic-001", "source": source.name, "source_digest": source_digest, "authorization": "generated synthetic media", "probe": { "command": source_probe_command, "result_path": "source-probe.json", }, } ], "tool_context": { "ffmpeg": ffmpeg_version, "ffprobe": ffprobe_version, }, "output_contract": { "container": "matroska", "required_stream_types": ["video", "audio"], "expected_duration_seconds": expected_duration, "duration_tolerance_seconds": 0.15, "downstream_consumer": None, }, "preservation": {"originals_untouched": True, "overwrite_policy": "refuse"}, "privacy_boundary": "synthetic, non-personal fixture only", "assumptions": [], } write_json(workspace / "intake-manifest.json", intake) evidence = { "schema_version": 1, "workflow_id": "synthetic-editorial-example", "asset_id": "asset-synthetic-001", "observations": [ { "class": "observed_artifact", "claim": "source probe reports video and audio streams", "locator": "source-probe.json", "coverage": "container and stream metadata only", } ], "unverified": [ "semantic visual interpretation", "listening quality", "downstream consumer compatibility", ], } write_json(workspace / "evidence-packet.json", evidence) edl = { "schema_version": 1, "workflow_id": "synthetic-editorial-example", "timebase": "seconds", "sources": [ { "asset_id": "asset-synthetic-001", "source": source.name, "duration": duration(source_probe), "digest": source_digest, } ], "events": [ { "id": "event-001", "action": "keep", "asset_id": "asset-synthetic-001", "stream_refs": ["0:v:0", "0:a:0"], "in": start, "out": end, "reason": "exercise a deterministic bounded trim", "evidence": ["evidence-packet.json#observations/0"], "boundary_precision": "decoded", "review_status": "approved_for_synthetic_example", } ], "output": { "mapping": ["video", "audio"], "expected_duration": expected_duration, "tolerance_seconds": 0.15, }, } write_json(workspace / "edit-decision-list.json", edl) render = [ ffmpeg, "-v", "error", "-n", "-i", str(source), "-ss", str(start), "-t", str(expected_duration), "-map", "0:v:0", "-map", "0:a:0", "-c:v", "mpeg4", "-q:v", "5", "-c:a", "pcm_s16le", str(output), ] run(render, timeout=args.timeout) commands.append(render) output_probe_command, output_probe = probe(ffprobe, output) commands.append(output_probe_command) write_json(workspace / "output-probe.json", output_probe) decode = [ ffmpeg, "-v", "error", "-i", str(output), "-map", "0:v:0", "-map", "0:a:0", "-f", "null", "-", ] run(decode, timeout=args.timeout) commands.append(decode) frame_paths: list[str] = [] for index, timestamp in enumerate((0.1, max(0.1, expected_duration - 0.1)), start=1): frame = workspace / f"review-frame-{index}.png" extract = [ ffmpeg, "-v", "error", "-ss", str(timestamp), "-i", str(output), "-frames:v", "1", "-n", str(frame), ] run(extract, timeout=args.timeout) commands.append(extract) frame_paths.append(frame.name) audio_check = [ ffmpeg, "-v", "error", "-i", str(output), "-map", "0:a:0", "-af", "astats=metadata=1:reset=1", "-f", "null", "-", ] run(audio_check, timeout=args.timeout) commands.append(audio_check) actual_duration = duration(output_probe) checks = [ { "criterion": "output_decodes", "boundary": "component", "verdict": "PASS", "evidence": "decode command completed", }, { "criterion": "required_streams", "boundary": "integration", "verdict": "PASS" if required_types.issubset(stream_types(output_probe)) else "FAIL", "evidence": "output-probe.json", }, { "criterion": "duration", "boundary": "integration", "verdict": ( "PASS" if actual_duration is not None and abs(actual_duration - expected_duration) <= 0.15 else "FAIL" ), "expected": expected_duration, "observed": actual_duration, "tolerance": 0.15, "evidence": "output-probe.json", }, { "criterion": "bounded_visual_samples_created", "boundary": "content-sampling", "verdict": "PASS", "evidence": frame_paths, "limitation": "samples establish only the sampled timestamps", }, { "criterion": "audio_processing_path", "boundary": "signal", "verdict": "PASS", "evidence": "bounded astats command completed", "limitation": "no listening review was performed", }, { "criterion": "downstream_consumer", "boundary": "downstream", "verdict": "UNVERIFIED", "evidence": None, }, ] failures = [check for check in checks if check["verdict"] == "FAIL"] acceptance = { "schema_version": 1, "workflow_id": "synthetic-editorial-example", "asset_id": "asset-synthetic-001", "event_ids": ["event-001"], "output": output.name, "output_digest": sha256(output), "render_command": render, "tool_context": intake["tool_context"], "criteria": checks, "overall_verdict": "FAIL" if failures else "PASS_WITH_UNVERIFIED_BOUNDARIES", "unverified_boundaries": [ "semantic visual review", "listening review", "downstream consumer compatibility", ], } write_json(workspace / "acceptance-report.json", acceptance) write_json(workspace / "command-log.json", {"commands": commands}) report = { "ok": not failures, "workflow_id": "synthetic-editorial-example", "workspace": str(workspace), "artifacts": sorted(path.name for path in workspace.iterdir()), "overall_verdict": acceptance["overall_verdict"], "unverified_boundaries": acceptance["unverified_boundaries"], } if args.json: print(json.dumps(report, sort_keys=True, separators=(",", ":"))) else: print(json.dumps(report, indent=2, sort_keys=True)) return 0 if report["ok"] else 1 except WorkflowError as exc: print(json.dumps({"ok": False, "error": str(exc)}, sort_keys=True)) return 2 if __name__ == "__main__": raise SystemExit(main())