arthur/clips/extraction.py
2026-10-03 00:39:07 -04:00

396 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Upload a video once, then turn it into the two things the app actually reads.
WHAT CHANGED AND WHY. This used to decode one PNG per source frame and store every
one of them. A 7.6-second 1440x1920 take is 112MB that way, and the 900-frame limit
is 1.1GB — for pixels whose only consumer was a canvas that MediaPipe then read
once. The page now detects from the video itself (see `frontend/src/arthur/flow/
ingest.cljs`), so this produces:
THE PROXY. One browser-safe H.264/yuv420p MP4, CFR, `+faststart`. The same take
is 6MB. This is the analysis source, and it is re-encoded RATHER THAN KEPT AS
UPLOADED even when the upload is already H.264, for two reasons that are both
about not guessing: an iPhone's HEVC is not decodable in every browser, and the
footage's identity is the digest of this file — one produced by one ffmpeg
invocation, not one that depends on which branch the source happened to take.
THE TRACING STILLS. One JPEG per frame, long edge capped, for the tracing editor
to draw over. Reference images; nothing measures them. They are not in the
footage digest — see `models.Footage`.
The proxy is probed after it is written rather than before. `width`, `height` and
`frames` are properties of the file the browser will decode, and taking them from
the source instead is how a scaler or a dropped frame becomes a silent one-frame
offset between the landmarks and the audio.
"""
import hashlib
import json
import subprocess
import tempfile
import threading
import time
from fractions import Fraction
from pathlib import Path
from django.db import close_old_connections, transaction
from . import blobs
from .models import Blob, Extraction, Footage, FootageFrame
_active = set()
_lock = threading.Lock()
TIMEOUT = 3600
# Visually lossless enough that landmarks do not move: measured against the same
# frames as PNGs, IMAGE-mode landmarks shifted at most 0.0033 of frame width.
PROXY_CRF = "18"
# The long edge of a tracing still. The proxy keeps full resolution because the
# detector reads it; a still only has to be good enough to draw a cel over.
TRACING_EDGE = 1280
TRACING_QUALITY = "4"
def _command(args):
result = subprocess.run(args, capture_output=True, text=True, timeout=TIMEOUT)
if result.returncode:
raise ValueError((result.stderr or result.stdout or "media tool failed")[-1200:])
return result.stdout
def _run_with_progress(job, args, root, name, total, span):
"""Run one ffmpeg and publish its live frame count as `span` of the job.
ffmpeg's `-progress` file is the only honest source for this: parsing its
stderr means parsing a format that is explicitly not an interface, and a
spinner that is not attached to frames is a spinner that lies on a long take.
"""
progress_path = root / f"{name}.progress"
log_path = root / f"{name}.log"
first, last = span
args = ["ffmpeg", "-hide_banner", "-loglevel", "error", "-y",
"-stats_period", "0.25", "-progress", str(progress_path)] + args
with open(log_path, "wb") as log:
proc = subprocess.Popen(args, stdout=log, stderr=subprocess.STDOUT)
deadline = time.monotonic() + TIMEOUT
try:
while proc.poll() is None:
if time.monotonic() >= deadline:
raise TimeoutError(f"{name} timed out")
if progress_path.exists():
lines = progress_path.read_text(errors="replace").splitlines()
count = next((int(line[6:].strip()) for line in reversed(lines)
if line.startswith("frame=") and
line[6:].strip().isdigit()), 0)
if count and total:
reached = first + int((last - first) * min(1.0, count / total))
if reached > job.progress:
job.progress = reached
job.save(update_fields=["progress", "updated"])
time.sleep(0.2)
finally:
if proc.poll() is None:
proc.kill()
proc.wait()
if proc.returncode:
raise ValueError(log_path.read_text(errors="replace")[-1200:] or f"{name} failed")
def _encode_proxy(job, source_path, proxy_path, facts, root):
"""The uploaded video -> one H.264 file every browser can decode and seek."""
total = facts.get("reported_frames") or round(facts["duration"] * facts["fps"])
_run_with_progress(
job,
["-i", str(source_path), "-an",
# Constant frame rate at the rate `probe` chose. This RESAMPLES rather
# than asserts: the upload is allowed to be variable, and this is the
# step that makes the thing the page measures not be.
"-fps_mode", "cfr", "-r", facts.get("rate") or str(facts["fps"]),
"-c:v", "libx264", "-preset", "veryfast", "-crf", PROXY_CRF,
# NO B-FRAMES, AND THIS IS THE LOAD-BEARING FLAG. It is what makes
# decode order presentation order, so the page can treat access unit k
# of the elementary stream as frame k without demuxing a container or
# consulting a timestamp. With them x264 has a
# two-frame reordering delay, ffmpeg compensates by writing an edit list
# (`elst` media_time 1024 at timebase 1/15360 — exactly two frames), and
# the browser then lives on two timelines at once: `currentTime` obeys the
# edit list and the `mediaTime` reported by requestVideoFrameCallback does
# not. Seek to frame 0 and the browser correctly hands back a frame whose
# mediaTime says 2. Software decoding hides it; hardware decoding does
# not, which is the worst possible way for it to be wrong. Without
# B-frames DTS equals PTS, no edit list is written, and the two timelines
# are the same one. It also makes decode order presentation order, should
# this ever be fed to a WebCodecs VideoDecoder.
"-bf", "0",
# yuv420p and an even frame size are what makes this playable everywhere
# rather than only in the browser that happened to be tested.
"-pix_fmt", "yuv420p", "-vf", "scale=trunc(iw/2)*2:trunc(ih/2)*2",
"-movflags", "+faststart", str(proxy_path)],
root, "proxy", total, (0, 55))
def _elementary_stream(proxy_path, out_path):
"""The proxy's video, unwrapped into a raw Annex-B H.264 stream.
A STREAM COPY, not a second encode: the same coded frames as the MP4, with
the container's length-prefixed NAL units rewritten as start-code-delimited
ones. It costs a file read and nothing else.
This exists because the page decodes with WebCodecs, and `VideoDecoder` takes
demuxed chunks rather than a container. Handing it Annex-B means the client
needs no demuxer: NAL start codes are findable in a loop, and because the
proxy is encoded with no B-frames, decode order is presentation order — so
access unit k IS frame k, with no container timing to consult and no clock to
reconcile. That is the whole reason this file is worth the bytes it costs.
"""
_command(["ffmpeg", "-hide_banner", "-loglevel", "error", "-y",
"-i", str(proxy_path), "-an", "-c:v", "copy",
"-bsf:v", "h264_mp4toannexb", "-f", "h264", str(out_path)])
def _extract_stills(job, proxy_path, frames_dir, frames, root):
"""The proxy -> one tracing JPEG per frame, long edge capped."""
_run_with_progress(
job,
["-i", str(proxy_path), "-fps_mode", "passthrough",
"-vf", f"scale='if(gt(iw,ih),min({TRACING_EDGE},iw),-2)':"
f"'if(gt(iw,ih),-2,min({TRACING_EDGE},ih))'",
"-q:v", TRACING_QUALITY, str(frames_dir / "%04d.jpg")],
root, "stills", frames, (55, 85))
MAX_RATE = 120 # a capture rate; past this the container is describing something else
def probe_audio(path):
"""The length of an uploaded sound in seconds, refusing a file with no audio."""
data = json.loads(_command(["ffprobe", "-v", "error", "-show_streams",
"-show_format", "-of", "json", str(path)]))
if not any(s.get("codec_type") == "audio" for s in data.get("streams", [])):
raise ValueError("the uploaded file has no audio stream")
duration = float(data.get("format", {}).get("duration") or 0)
if duration <= 0:
raise ValueError("the sound's length is unknown")
return duration
def probe(path):
"""What the upload is, as far as choosing a proxy rate goes.
IT NO LONGER REFUSES VARIABLE-FRAME-RATE INPUT, and the reason is the proxy.
That refusal was written when the page measured the source's own frames, where
a wandering frame duration really does break `frame = floor(t * fps)`. Nothing
measures the source now: ffmpeg resamples it onto a constant rate, and the
proxy — constant by construction, and re-probed after it is written — is the
only timeline anything downstream sees.
Keeping the check would have been worse than useless, because the thing it
tested is not reliable. Ordinary iPhone footage, shot straight from the camera
app, reports `avg_frame_rate` 8670/299 and `nb_frames` 289 on a stream whose
decoded timestamps are 280 frames exactly 1/30s apart. The container's summary
of itself disagreed with the container's own contents, so the guard rejected
CFR video for being variable.
THE RATE IS THE NOMINAL ONE. `r_frame_rate` is the rate every timestamp in the
stream can be expressed at, which is the rate that keeps every distinct source
frame; resampling to the average would drop some. Duration is preserved either
way — ffmpeg's CFR conversion is driven by timestamps, so the audio stays in
sync at any rate — so this trades a possible duplicated frame against a
certainly lost one.
"""
data = json.loads(_command(["ffprobe", "-v", "error", "-show_streams",
"-show_format", "-of", "json", str(path)]))
video = next((s for s in data.get("streams", []) if s.get("codec_type") == "video"), None)
if not video:
raise ValueError("the uploaded file has no video stream")
nominal = Fraction(video.get("r_frame_rate") or "0")
average = Fraction(video.get("avg_frame_rate") or "0")
if nominal <= 0 and average <= 0:
raise ValueError("the video's frame rate is unknown")
rate = nominal if 0 < nominal <= MAX_RATE else average
if not 0 < rate <= MAX_RATE:
raise ValueError(f"the video reports a frame rate of {float(rate):g}, which is "
"not a rate footage can be measured at")
duration = float(data.get("format", {}).get("duration") or 0)
# if duration > 0 and duration * float(rate) > 901:
# raise ValueError("video is longer than the 900-frame footage limit")
frames = video.get("nb_frames")
return {"fps": float(rate),
# The exact rate, for ffmpeg. 30000/1001 is not a float, and handing
# `-r` a rounded one is how a long take drifts out of sync.
"rate": f"{rate.numerator}/{rate.denominator}",
"nominal_fps": float(nominal), "average_fps": float(average),
"width": int(video["width"]), "height": int(video["height"]),
"duration": duration,
# KEPT, AND NO LONGER TRUSTED AS A COUNT. See the docstring: this is
# the container's claim about itself, it is wrong on ordinary phone
# footage, and `run` checks the proxy's DURATION instead.
"reported_frames": int(frames) if frames and frames.isdigit() else None,
"has_audio": any(s.get("codec_type") == "audio" for s in data.get("streams", [])),
"vfr": nominal != average}
def _refuse_a_shifted_timeline(path):
"""The proxy must put frame `i` at `i / fps` on BOTH of the browser's clocks.
Asserted rather than assumed, because the failure is silent and the symptom is
unrecognisable. An encoder delay makes ffmpeg write an edit list, `currentTime`
then obeys it while `requestVideoFrameCallback`'s `mediaTime` does not, and the
page's frame walk is uniformly off by the delay — on hardware decoding only. It
cost two wrong diagnoses to find, so it does not get to come back silently if
somebody changes an encoder flag.
"""
data = json.loads(_command(["ffprobe", "-v", "error", "-select_streams", "v:0",
"-show_streams", "-of", "json", str(path)]))
stream = data["streams"][0]
if int(stream.get("has_b_frames") or 0):
raise ValueError(
"the proxy was encoded with B-frames, whose reordering delay makes the "
"browser's seek clock and its frame-timestamp clock disagree")
if float(stream.get("start_time") or 0) != 0:
raise ValueError(f"the proxy starts at {stream['start_time']}s rather than 0")
def count_frames(path):
"""How many frames a file really holds, counted rather than reported.
`nb_frames` is a container's claim. This is the decoder's answer, and it is
what the page will get when it walks the proxy — so a disagreement between the
two has to be settled before the count reaches a manifest, not after it has
become a one-frame audio offset nobody can find.
"""
text = _command(["ffprobe", "-v", "error", "-select_streams", "v:0",
"-count_frames", "-show_entries", "stream=nb_read_frames",
"-of", "default=nokey=1:noprint_wrappers=1", str(path)])
counted = text.strip()
if not counted.isdigit():
raise ValueError("could not count the proxy's frames")
return int(counted)
def extraction_key(source, settings):
# Scheme 3: the extraction now also produces the elementary stream the page
# decodes, so a job run under scheme 2 did not make everything this one does.
text = json.dumps({"scheme": 3, "source": source.blob_id, "settings": settings},
sort_keys=True, separators=(",", ":"))
return "sha256:" + hashlib.sha256(text.encode()).hexdigest()
def _register(job, proxy_path, stream_path, stills, audio_path, facts):
proxy_digest, proxy_size = blobs.adopt(proxy_path)
stream_digest, stream_size = blobs.adopt(stream_path)
audio_digest, audio_size = blobs.adopt(audio_path)
still_blobs = [(index, *blobs.adopt(path)) for index, path in enumerate(stills)]
width, height, fps, frames = facts["width"], facts["height"], facts["fps"], facts["frames"]
# The footage's own identity: the bytes the page will measure, the audio it
# will clock against, and the rate that ties them together. Scheme 2 — scheme
# 1 hashed a PNG per frame, and those footages name pixels this no longer has.
h = hashlib.sha256()
h.update(f"arthur-footage-2/{fps}/{frames}/{width}x{height}\n".encode())
h.update(proxy_digest.encode())
h.update(audio_digest.encode())
with transaction.atomic():
proxy_blob, _ = Blob.objects.get_or_create(
digest=proxy_digest, defaults={"size": proxy_size, "media_type": "video/mp4"})
stream_blob, _ = Blob.objects.get_or_create(
digest=stream_digest, defaults={"size": stream_size, "media_type": "video/h264"})
audio_blob, _ = Blob.objects.get_or_create(
digest=audio_digest, defaults={"size": audio_size, "media_type": "audio/wav"})
footage, created = Footage.objects.get_or_create(
digest=h.hexdigest(),
defaults={"label": job.source.filename[:200], "source": job.source.filename[:200],
"fps": fps, "frames": frames, "width": width, "height": height,
"audio": audio_blob, "video": proxy_blob, "stream": stream_blob})
if not created and not footage.stream_id:
# The same footage by identity, extracted before the elementary
# stream existed. Its digest is over the proxy and the audio, which
# have not changed — so this is the same footage gaining a file it
# was always entitled to, not a different one.
footage.stream = stream_blob
if not footage.video_id:
footage.video = proxy_blob
footage.save(update_fields=["stream", "video"])
if created:
rows = []
for index, digest, size in still_blobs:
blob, _ = Blob.objects.get_or_create(
digest=digest, defaults={"size": size, "media_type": "image/jpeg"})
rows.append(FootageFrame(footage=footage, index=index, blob=blob))
FootageFrame.objects.bulk_create(rows)
return footage
def run(key):
close_old_connections()
try:
job = Extraction.objects.select_related("source", "source__blob").get(key=key)
job.state, job.progress, job.error = "running", 0, ""
job.save(update_fields=["state", "progress", "error", "updated"])
facts = job.source.probe
source_path = blobs.path_for(job.source.blob_id)
with tempfile.TemporaryDirectory(prefix="arthur-extract-") as directory:
root = Path(directory)
proxy_path = root / "proxy.mp4"
_encode_proxy(job, source_path, proxy_path, facts, root)
# Everything downstream describes the PROXY, not the upload.
proxy_facts = probe(proxy_path)
_refuse_a_shifted_timeline(proxy_path)
frames = count_frames(proxy_path)
#if not 1 <= frames <= 900:
# raise ValueError(f"the proxy holds {frames} frames; the limit is 1–900")
# CHECKED AS A DURATION, not as a frame count. The page's clock is
# `frame = floor(audio.currentTime * fps)`, so what must not drift is
# how long the picture lasts against how long the audio lasts — and
# the source's own frame count is a number this has already caught
# lying. A resample to a constant rate legitimately changes the count
# and must not change the duration.
drift = abs(frames / proxy_facts["fps"] - facts["duration"])
if facts["duration"] > 0 and drift > 0.5:
raise ValueError(
f"the proxy runs {frames / proxy_facts['fps']:.2f}s and the upload "
f"runs {facts['duration']:.2f}s; refusing footage whose picture and "
"audio would drift")
proxy_facts["frames"] = frames
stream_path = root / "proxy.h264"
_elementary_stream(proxy_path, stream_path)
frames_dir = root / "stills"
frames_dir.mkdir()
_extract_stills(job, proxy_path, frames_dir, frames, root)
stills = sorted(frames_dir.glob("*.jpg"))
if len(stills) != frames:
raise ValueError(f"wrote {len(stills)} tracing stills for {frames} frames")
job.progress = 85
job.save(update_fields=["progress", "updated"])
audio_path = root / "audio.wav"
if facts["has_audio"]:
_command(["ffmpeg", "-hide_banner", "-loglevel", "error", "-y",
"-i", str(source_path), "-vn", "-ac", "1", "-ar", "44100",
str(audio_path)])
else:
_command(["ffmpeg", "-hide_banner", "-loglevel", "error", "-y",
"-f", "lavfi", "-i", "anullsrc=r=44100:cl=mono",
"-t", str(frames / proxy_facts["fps"]), "-c:a", "pcm_s16le",
str(audio_path)])
footage = _register(job, proxy_path, stream_path, stills, audio_path, proxy_facts)
job.footage, job.state, job.progress = footage, "done", 100
job.save(update_fields=["footage", "state", "progress", "updated"])
except Exception as exc:
Extraction.objects.filter(key=key).update(state="failed", error=str(exc)[:2000])
finally:
with _lock:
_active.discard(key)
close_old_connections()
def enqueue(key):
with _lock:
if key in _active:
return
_active.add(key)
threading.Thread(target=run, args=(key,), daemon=True,
name=f"arthur-extract-{key[7:15]}").start()