Files
L'électron rare a2a7fbd0f9 feat: thread CUDA engine through parallel stitcher
Adds --engine {cpu,cuda} to the parallel driver, forwarding the flag
and NVENC encoding to the insv-stitch fork
2026-08-07 20:21:38 +02:00

209 lines
8.1 KiB
Python

#!/usr/bin/env python3
"""Chunked, resumable, parallel driver around insv-stitch's x5_pipeline.py.
Stitching a full take is CPU-bound and single-threaded: tens of seconds
per frame, i.e. days of wall time for a half-hour capture. This driver
splits the take into fixed-size frame ranges, stitches several of them
concurrently, and concatenates the pieces losslessly.
Resumability is the point as much as the speed: a chunk whose .mp4 is
already present *and* holds exactly the expected number of frames is
skipped, so an interrupted multi-day job restarts where it stopped
rather than from zero.
Usage:
python3 scripts/stitch-parallel.py \
--insv "/Volumes/T7/Balade jeudi/VID_..._00_020.insv" \
--out "/Volumes/T7/Balade jeudi/stitched/balade-020-equirect.mp4" \
--width 3840 --workers 8 --chunk 2000
Only the '_00_' file is passed; x5_pipeline locates its '_10_' sibling
itself for X3-generation dual-file captures.
"""
import argparse
import json
import os
import subprocess
import sys
import time
from pathlib import Path
TOOL = Path.home() / "Documents" / "Projets" / "_tools" / "insv-stitch"
def probe_frame_count(path, count_frames=False):
"""Number of video frames in `path`.
`count_frames` decodes the whole file (slow but exact) and is what
chunk verification needs, since a chunk killed mid-write still
carries a plausible nb_frames tag.
"""
if count_frames:
cmd = ["ffprobe", "-v", "error", "-select_streams", "v:0",
"-count_frames", "-show_entries", "stream=nb_read_frames",
"-of", "default=nokey=1:noprint_wrappers=1", str(path)]
else:
cmd = ["ffprobe", "-v", "error", "-select_streams", "v:0",
"-show_entries", "stream=nb_frames",
"-of", "default=nokey=1:noprint_wrappers=1", str(path)]
r = subprocess.run(cmd, capture_output=True, text=True)
out = r.stdout.strip()
if not out or not out.isdigit():
return 0
return int(out)
def dual_file_sibling(insv):
"""The '_10_' file holding the other lens, for X3 dual-file captures."""
name = insv.name
if "_00_" in name:
return insv.parent / name.replace("_00_", "_10_", 1)
if "_10_" in name:
return insv.parent / name.replace("_10_", "_00_", 1)
return None
def usable_frame_count(insv):
"""Frames the pipeline will actually stitch.
The two lenses of a dual-file capture routinely differ by a frame;
x5_pipeline clamps to the shorter one, so asking for the longer
count would leave the final chunk permanently short and block the
concat.
"""
counts = [probe_frame_count(insv)]
sibling = dual_file_sibling(insv)
if sibling and sibling.exists():
counts.append(probe_frame_count(sibling))
return min(counts)
def chunk_ranges(total, size, start=0, end=None):
end = total if end is None else min(end, total)
return [(a, min(a + size, end)) for a in range(start, end, size)]
def chunk_path(workdir, a, b):
return workdir / f"chunk_{a:07d}_{b:07d}.mp4"
def chunk_done(path, expected):
"""True when `path` already holds a complete chunk."""
return path.exists() and probe_frame_count(path, count_frames=True) == expected
def launch(insv, path, a, b, width, log, engine=None):
cmd = [str(TOOL / ".venv" / "bin" / "python"), str(TOOL / "x5_pipeline.py"),
str(insv), "-o", str(path), "-w", str(width), "--video",
"--start", str(a), "--end", str(b)]
if engine:
cmd.append(f"--engine={engine}")
enum = {"cuda": "nvenc"}.get(engine)
if enum:
cmd.append(f"--encoder={enum}")
env = dict(os.environ)
# One OpenCV/OMP thread per worker: parallelism comes from running
# several chunks, not from threading inside one of them.
env["OPENCV_FOR_THREADS_NUM"] = "1"
env["OMP_NUM_THREADS"] = "1"
fh = open(log, "w")
return subprocess.Popen(cmd, stdout=fh, stderr=subprocess.STDOUT, env=env), fh
def concat(paths, out):
"""Losslessly join chunks — identical encoder settings throughout."""
listing = out.parent / (out.stem + ".concat.txt")
listing.write_text("".join(f"file '{p.as_posix()}'\n" for p in paths))
cmd = ["ffmpeg", "-y", "-v", "error", "-f", "concat", "-safe", "0",
"-i", str(listing), "-c", "copy", str(out)]
subprocess.run(cmd, check=True)
listing.unlink()
def main():
ap = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--insv", required=True, help="the '_00_' .insv file")
ap.add_argument("--out", required=True, help="stitched equirect .mp4")
ap.add_argument("--width", type=int, default=3840, help="equirect output width")
ap.add_argument("--workers", type=int, default=8, help="concurrent chunks")
ap.add_argument("--engine", choices=["cpu", "cuda"], default=None,
help="computing backend to pass to x5_pipeline.py "
"(default: x5_pipeline's own default = cpu)")
ap.add_argument("--chunk", type=int, default=2000, help="frames per chunk")
ap.add_argument("--start", type=int, default=0)
ap.add_argument("--end", type=int, default=None)
ap.add_argument("--workdir", default=None,
help="chunk scratch dir (default: <out>.chunks/)")
ap.add_argument("--concat-only", action="store_true",
help="skip stitching, just join the chunks already present")
args = ap.parse_args()
insv = Path(args.insv)
out = Path(args.out)
out.parent.mkdir(parents=True, exist_ok=True)
workdir = Path(args.workdir) if args.workdir else out.parent / (out.stem + ".chunks")
workdir.mkdir(parents=True, exist_ok=True)
total = usable_frame_count(insv)
ranges = chunk_ranges(total, args.chunk, args.start, args.end)
print(f"{insv.name}: {total} frames -> {len(ranges)} chunks of "
f"{args.chunk}, {args.workers} workers, width {args.width}",
flush=True)
if not args.concat_only:
pending = []
for a, b in ranges:
p = chunk_path(workdir, a, b)
if chunk_done(p, b - a):
print(f" skip {p.name} (already complete)", flush=True)
else:
pending.append((a, b, p))
print(f"{len(pending)} chunks to stitch "
f"({sum(b - a for a, b, _ in pending)} frames)", flush=True)
t0 = time.time()
done_frames = 0
running = [] # (proc, fh, a, b, path)
queue = list(pending)
while queue or running:
while queue and len(running) < args.workers:
a, b, p = queue.pop(0)
log = p.with_suffix(".log")
proc, fh = launch(insv, p, a, b, args.width, log,
engine=args.engine)
running.append((proc, fh, a, b, p))
print(f" start {p.name}", flush=True)
time.sleep(5)
still = []
for proc, fh, a, b, p in running:
if proc.poll() is None:
still.append((proc, fh, a, b, p))
continue
fh.close()
ok = chunk_done(p, b - a)
done_frames += (b - a) if ok else 0
dt = time.time() - t0
rate = done_frames / dt if dt else 0
print(f" {'done' if ok else 'FAILED'} {p.name} rc={proc.returncode} "
f"| {done_frames} frames in {dt/3600:.2f} h "
f"({rate*3600:.0f} frames/h)", flush=True)
if not ok:
print(f" see {p.with_suffix('.log')}", flush=True)
running = still
missing = [chunk_path(workdir, a, b) for a, b in ranges
if not chunk_done(chunk_path(workdir, a, b), b - a)]
if missing:
print(f"{len(missing)} chunk(s) incomplete, not concatenating:", flush=True)
for p in missing[:10]:
print(f" {p.name}", flush=True)
sys.exit(1)
concat([chunk_path(workdir, a, b) for a, b in ranges], out)
print(f"wrote {out} ({probe_frame_count(out)} frames)", flush=True)
if __name__ == "__main__":
main()