Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 13 additions & 2 deletions backend/app/video/ffmpeg.py
Original file line number Diff line number Diff line change
Expand Up @@ -118,8 +118,19 @@ def run(argv: Sequence[str], *, cwd: Optional[str] = None, timeout: int = SHOT_T
f"encoding preset."
) from exc
if result.returncode != 0:
tail = result.stderr.decode(errors="replace")[-800:]
raise FFmpegError(f"ffmpeg failed: {tail}")
stderr = result.stderr.decode(errors="replace")
tail = stderr[-800:].strip()
code = result.returncode
# `-loglevel error` means a process killed from outside says nothing at
# all, so the exit status is the only clue that it was the OOM killer.
if code < 0:
reason = f"killed by signal {-code}"
if code == -9:
reason += " (most likely out of memory)"
else:
reason = f"exit code {code}"
logger.error("ffmpeg %s\nargv: %s\nstderr: %s", reason, " ".join(argv), stderr or "<empty>")
raise FFmpegError(f"ffmpeg failed ({reason}): {tail}".rstrip(": "))


def probe_duration(path: str) -> float:
Expand Down
100 changes: 74 additions & 26 deletions backend/app/video/renderer.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,8 @@
stitch is a `concat -c copy` remux rather than a second full encode, and a shot
is the natural unit of progress. Everything happens in a scratch directory that
is removed on every exit path; only the finished artefacts are moved into the
user's media directory.
user's media directory. The one thing that outlives a failed render is its
encoded segments (see `shot_cache`), so a retry doesn't redo them.
"""

import json
Expand All @@ -20,7 +21,7 @@
from datetime import datetime
from typing import Callable, List, Optional, Sequence, Tuple

from app.video import compose, ffmpeg as F
from app.video import compose, ffmpeg as F, shot_cache
from app.video.narration import (
Cue, NarrationResult, build_narration_chunks, shift_cues, synthesize_shot, write_srt,
)
Expand Down Expand Up @@ -166,6 +167,14 @@ def render(
work_dir = os.path.join(media_dir, WORK_ROOT_NAME, job_id)
os.makedirs(work_dir, exist_ok=True)

# Encoded segments outlive a failed render (see shot_cache), so a retry only
# redoes what actually failed. Housekeeping must never fail a render.
cache = shot_cache.ShotCache(media_dir, user_id)
try:
shot_cache.prune(media_dir)
except Exception:
logger.exception("Could not prune the video shot cache")

try:
progress("Preparing", 2, "Reading the article")
plan: Segmentation = segment(
Expand Down Expand Up @@ -234,6 +243,7 @@ def render(
# ── per-shot render ───────────────────────────────────────────────────
base = 2 + NARRATION_SHARE
shot_files: List[str] = []
reused = 0
# Collected as the shots are rendered and laid out on a timeline
# afterwards, once the transition overlap is known.
durations: List[float] = []
Expand Down Expand Up @@ -322,30 +332,48 @@ def render(

output = f"shot_{index:04d}.mp4"

def _encode(with_options: RenderOptions) -> None:
argv = F.build_shot_command(
def _argv(with_options: RenderOptions) -> List[str]:
return F.build_shot_command(
kind=shot.kind, background=background, audio=narration.path,
duration=duration, output=output, options=with_options,
preview=preview, overlay_png=shot_overlay, subtitle_file=shot_srt,
background_has_audio=has_audio, index=index,
)

def _encode(with_options: RenderOptions) -> None:
# Run inside the work dir so every path in the command is a bare
# filename — which is what keeps the `subtitles=` filter, whose
# argument needs escaping, free of anything that needs escaping.
F.run(argv, cwd=work_dir, timeout=F.shot_timeout(duration))

try:
_encode(options)
except F.FFmpegError as exc:
# Losing a forty-minute render to one expensive segment is a bad
# trade when the two most expensive things in it are also the two
# least important. Retry once without them; the narration is
# already synthesised and cached, so this costs no speech.
logger.warning("Segment %d failed (%s) — retrying it plainer", index + 1, exc)
plain = options.model_copy(deep=True)
plain.ken_burns.effect = "none"
plain.waveform.enabled = False
_encode(plain)
F.run(_argv(with_options), cwd=work_dir, timeout=F.shot_timeout(duration))

# A segment encoded by an earlier attempt that died later on (usually
# at the stitch) is reused rather than encoded again.
key = shot_cache.shot_key(
_argv(options),
[background, narration.path, shot_overlay, shot_srt],
work_dir,
)
fell_back = cache.restore(key, os.path.join(work_dir, output))
if fell_back is not None:
reused += 1
else:
try:
_encode(options)
except F.FFmpegError as exc:
# Losing a forty-minute render to one expensive segment is a bad
# trade when the two most expensive things in it are also the two
# least important. Retry once without them; the narration is
# already synthesised and cached, so this costs no speech.
logger.warning("Segment %d failed (%s) — retrying it plainer", index + 1, exc)
plain = options.model_copy(deep=True)
plain.ken_burns.effect = "none"
plain.waveform.enabled = False
_encode(plain)
fell_back = True
else:
fell_back = False
cache.store(key, os.path.join(work_dir, output), plain=fell_back)
if fell_back:
plan.warnings.append(
f"Segment {index + 1} was rendered without motion or a waveform: "
f"it was too long to render with them."
Expand All @@ -355,6 +383,10 @@ def _encode(with_options: RenderOptions) -> None:
durations.append(duration)
shot_cues.append(narration.cues)

if reused:
logger.info("Video render %s reused %d of %d encoded segments from an earlier attempt",
job_id, reused, len(shots))

# ── timeline ──────────────────────────────────────────────────────────
# A crossfade overlaps each pair of shots, so every shot after the first
# starts before its predecessor ends. Laying the timeline out here, once
Expand All @@ -378,14 +410,27 @@ def _encode(with_options: RenderOptions) -> None:
# A blend needs adjacent shots on screen together, which no per-shot
# filter can express, so this path re-encodes the joined stream once.
# Every other transition still stitches as a remux below.
F.run(
F.build_xfade_command(
shot_files, durations, "stitched.mp4", options=options,
preview=preview, style=options.transition.style, overlap=overlap,
),
cwd=work_dir, timeout=F.shot_timeout(sum(durations)),
)
else:
try:
F.run(
F.build_xfade_command(
shot_files, durations, "stitched.mp4", options=options,
preview=preview, style=options.transition.style, overlap=overlap,
),
cwd=work_dir, timeout=F.shot_timeout(sum(durations)),
)
except F.FFmpegError as exc:
# This is the one step that holds every segment open at once, so
# it is the one a small host runs out of memory in — after all the
# encoding is done. Joining with straight cuts needs almost none,
# so a finished video without the blend beats no video at all.
logger.warning("Crossfade stitch failed (%s) — joining with cuts instead", exc)
plan.warnings.append(
"Crossfade transitions were skipped: blending the segments "
"together failed, so they are joined with straight cuts."
)
overlap = 0.0
chapters, all_cues, timeline = build_timeline(shots, durations, shot_cues, overlap)
if overlap <= 0:
F.write_concat_list(os.path.join(work_dir, "shots.txt"), shot_files)
F.run(F.build_concat_command("shots.txt", "stitched.mp4"), cwd=work_dir, timeout=1800)

Expand Down Expand Up @@ -435,6 +480,9 @@ def _encode(with_options: RenderOptions) -> None:
stem = uuid.uuid4().hex
video_filename = f"{stem}.mp4"
shutil.move(os.path.join(work_dir, final), os.path.join(user_dir, video_filename))
# The segments are inside the finished video now, so their cached copies
# are dead weight; only a render that fails leaves them behind.
cache.discard_used()

subtitle_filename = None
if srt_written and options.subtitles in ("sidecar", "soft", "burn"):
Expand Down
164 changes: 164 additions & 0 deletions backend/app/video/shot_cache.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
"""Finished shot encodes, kept across a failed render so a retry doesn't redo them.

A render works in a scratch directory that is removed on every exit path, and a
retry is a brand-new job with a brand-new directory — so a render that died at
the stitch, after every segment had been encoded, used to throw all of that
work away and start over.

Each shot is filed here under a fingerprint of everything that decides its
pixels and samples: the exact ffmpeg command (which carries the filtergraph, the
duration and the codec settings) and the content of every file it reads. It is
deliberately *not* keyed on the job id or on the options as a whole. Retrying
with a different transition, music bed or subtitle setting leaves most shots
byte-for-byte identical, and those are exactly the retries that matter — the
usual reason to retry is to change something that isn't in the shots at all.

Entries are hard links to the scratch file, so keeping one costs no extra disk
while the render runs. A successful render discards the entries it used, since
they are inside the finished video now; whatever a failed or cancelled one left
behind is pruned by age.
"""

import hashlib
import logging
import os
import shutil
import time
from typing import List, Optional, Sequence

logger = logging.getLogger(__name__)

CACHE_ROOT_NAME = "_video_cache"

# Long enough to notice a failure, change a setting and try again; short enough
# that abandoned encodes don't sit on the disk.
MAX_AGE_SECONDS = 48 * 60 * 60

# Bump when a change to the pipeline means an old encode is no longer valid even
# though the command and inputs look the same.
CACHE_VERSION = "1"

_PLAIN_SUFFIX = ".plain"


def shot_key(argv: Sequence[str], inputs: Sequence[Optional[str]], work_dir: str) -> str:
"""Fingerprint of one shot's encode.

`argv` is the command that would produce it, minus its trailing output name
(which only carries the shot's position, so a shot that merely moved would
otherwise miss). `inputs` are the files it reads: generated ones inside
`work_dir` are hashed by content, since they are rebuilt every run; the
user's own media, which can be hundreds of megabytes, by size and mtime.
"""
digest = hashlib.sha256(CACHE_VERSION.encode())
for arg in argv[:-1]:
digest.update(arg.encode("utf-8", "replace"))
digest.update(b"\0")
for name in inputs:
if not name:
continue
digest.update(_signature(name, work_dir).encode())
digest.update(b"\0")
return digest.hexdigest()


def _signature(name: str, work_dir: str) -> str:
full = name if os.path.isabs(name) else os.path.join(work_dir, name)
try:
stat = os.stat(full)
except OSError:
return "missing"
root = os.path.abspath(work_dir) + os.sep
if not os.path.abspath(full).startswith(root):
return f"{stat.st_size}:{stat.st_mtime_ns}"
digest = hashlib.sha256()
with open(full, "rb") as f:
for block in iter(lambda: f.read(1 << 20), b""):
digest.update(block)
return digest.hexdigest()


def _link_or_copy(src: str, dest: str) -> None:
"""Hard link where the filesystem allows it, else copy. Lands atomically."""
tmp = f"{dest}.{os.getpid()}.tmp"
try:
try:
os.link(src, tmp)
except OSError:
shutil.copy2(src, tmp)
os.replace(tmp, dest)
finally:
try:
os.remove(tmp)
except OSError:
pass


class ShotCache:
"""One user's cache directory, plus a note of which entries this render touched."""

def __init__(self, media_dir: str, user_id: str) -> None:
self.dir = os.path.join(media_dir, CACHE_ROOT_NAME, user_id)
self._used: List[str] = []

def _path(self, key: str) -> str:
return os.path.join(self.dir, f"{key}.mp4")

def restore(self, key: str, dest: str) -> Optional[bool]:
"""Put a cached encode at `dest`.

None on a miss. On a hit, whether that encode had to fall back to the
plainer settings, so the render can repeat the warning it earned then.
"""
path = self._path(key)
if not os.path.isfile(path):
return None
try:
_link_or_copy(path, dest)
os.utime(path) # a retry in progress keeps its entries alive
except OSError as exc:
logger.warning("Could not reuse cached shot %s: %s", key[:12], exc)
return None
self._used.append(key)
return os.path.isfile(path + _PLAIN_SUFFIX)

def store(self, key: str, src: str, *, plain: bool) -> None:
"""File a finished encode. Best effort: a full disk must not fail a render."""
try:
os.makedirs(self.dir, exist_ok=True)
path = self._path(key)
_link_or_copy(src, path)
if plain:
open(path + _PLAIN_SUFFIX, "wb").close()
self._used.append(key)
except OSError as exc:
logger.warning("Could not cache shot %s: %s", key[:12], exc)

def discard_used(self) -> None:
"""The render succeeded, so what it used is inside the finished video."""
for key in self._used:
for path in (self._path(key), self._path(key) + _PLAIN_SUFFIX):
try:
os.remove(path)
except OSError:
pass
self._used = []


def prune(media_dir: str, max_age_seconds: int = MAX_AGE_SECONDS) -> None:
"""Delete cache entries nobody has touched for `max_age_seconds`."""
root = os.path.join(media_dir, CACHE_ROOT_NAME)
cutoff = time.time() - max_age_seconds
for folder, _dirs, files in os.walk(root, topdown=False):
for name in files:
path = os.path.join(folder, name)
try:
if os.path.getmtime(path) < cutoff:
os.remove(path)
except OSError:
pass
if folder != root:
try:
os.rmdir(folder) # only succeeds once empty
except OSError:
pass
13 changes: 13 additions & 0 deletions backend/tests/test_video_ffmpeg.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import math
import os
import tempfile
import types

import pytest

Expand Down Expand Up @@ -264,6 +265,18 @@ def test_concat_list_quotes_each_entry():
assert "'\\''" in body.splitlines()[1]


@pytest.mark.parametrize("code,stderr,expected", [
(-9, b"", "killed by signal 9 (most likely out of memory)"),
(1, b"Invalid argument", "exit code 1): Invalid argument"),
])
def test_run_failure_reports_how_the_process_died(monkeypatch, code, stderr, expected):
result = types.SimpleNamespace(returncode=code, stderr=stderr)
monkeypatch.setattr(F.subprocess, "run", lambda *a, **k: result)
with pytest.raises(F.FFmpegError) as err:
F.run(["ffmpeg", "-i", "x"])
assert expected in str(err.value)


# ── chunking ─────────────────────────────────────────────────────────────────

def test_short_text_is_one_chunk():
Expand Down
Loading
Loading