Adds --engine {cpu,cuda} to the parallel driver, forwarding the flag
and NVENC encoding to the insv-stitch fork
209 lines
8.1 KiB
Python
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()
|