384 lines
19 KiB
Python
384 lines
19 KiB
Python
"""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(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()
|