Mage-VL / neural_codec /precompute_dcvc_rt.py
Xinjie-Q's picture
Upload Mage-VL: unified codec-native streaming VLM (image+video understanding + proactive gate)
12acbba verified
Raw
History Blame Contribute Delete
12.2 kB
#!/usr/bin/env python3
"""Offline precompute of DCVC-RT patch-selection canvases for the release model.
For each input video this writes, under ``<out_root>/assets/<key>/``:
canvas_000.jpg ... selected-patch canvases (release codec contract)
src_patch_position.npy int32 [total_patches, 3] (t, h, w)
meta.json {fps, canvas_files, cfg, ...}
These directories are consumed at inference time by ``codec_loader.py`` /
``infer_dcvc_rt.py`` (which feed them through the model's existing codec
downstream, i.e. exactly what ``video_backend="codec"`` does for HEVC).
Run in the ``magevl`` conda env (torch + DCVC-RT + cv2). Example::
python neural_codec/precompute_dcvc_rt.py \
--videos /path/videos.jsonl \
--out_root /path/out/dcvc_rt \
--target_canvas 32 --seq_len_frames 64 --max_pixels 150000 --qp 21 \
--cuda_idx 0 1 --num_workers 2
``--videos`` accepts a .jsonl (one obj per line with ``video`` and optional
``key``), a directory of videos, or a single video path.
"""
from __future__ import annotations
import argparse
import json
import os
import sys
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Dict, List, Tuple
import cv2
import numpy as np
_HERE = os.path.dirname(os.path.abspath(__file__))
_MODEL_DEFAULT = os.path.dirname(_HERE) # Mage-VL-Ported/
sys.path.insert(0, _HERE)
from canvas_assembler import AssembleConfig, assemble_canvases # noqa: E402
_VIDEO_EXTS = (".mp4", ".mkv", ".webm", ".avi", ".mov", ".m4v")
@dataclass
class Job:
video: str
key: str
# ------------------------------------------------------------------ video io
def _probe(video_path: str) -> Tuple[int, float]:
cap = cv2.VideoCapture(video_path)
if not cap.isOpened():
return 0, 0.0
total = int(cap.get(cv2.CAP_PROP_FRAME_COUNT) or 0)
fps = float(cap.get(cv2.CAP_PROP_FPS) or 0.0)
cap.release()
if not np.isfinite(fps) or fps <= 0:
fps = 30.0
return max(0, total), fps
def _sample_ids(total: int, seq_len: int) -> List[int]:
total = max(1, int(total))
if seq_len <= 0:
return []
return list(np.linspace(0, total - 1, int(seq_len), dtype=np.int64))
def _maybe_downscale(bgr: np.ndarray, max_side: int) -> np.ndarray:
if max_side <= 0:
return bgr
H, W = bgr.shape[:2]
if max(H, W) <= max_side:
return bgr
s = max_side / float(max(H, W))
return cv2.resize(bgr, (max(1, int(round(W * s))), max(1, int(round(H * s)))),
interpolation=cv2.INTER_AREA)
# ---------------------------------------------------------------- per-video
_ENGINE = None # per-worker DCVCRTEngine
def _process_one(job: Job, args: argparse.Namespace) -> Tuple[str, str]:
from PIL import Image
out_dir = Path(args.out_root) / "assets" / job.key
done = (out_dir / "meta.json").exists() and (out_dir / "src_patch_position.npy").exists()
if done and not args.overwrite:
return "skip", f"{job.key} exists"
if not os.path.exists(job.video):
return "fail", f"{job.key} missing video {job.video}"
total, fps = _probe(job.video)
if total <= 0:
return "fail", f"{job.key} unreadable video"
sampled_ids = _sample_ids(total, args.seq_len_frames)
needed = set(int(x) for x in sampled_ids)
max_fid = max(needed)
# Stream-decode 0..max_fid; DCVC needs contiguous frames. Keep only sampled.
cap = cv2.VideoCapture(job.video)
if not cap.isOpened():
return "fail", f"{job.key} cannot open"
t0 = time.time()
kept_rgb: Dict[int, np.ndarray] = {}
bitmaps: Dict[int, np.ndarray] = {}
seq_started = False
fidx = 0
last_bgr = None
while fidx <= max_fid:
ret, bgr = cap.read()
if not ret or bgr is None:
if last_bgr is None:
break
bgr = last_bgr # pad tail with last good frame
else:
last_bgr = bgr
bgr = _maybe_downscale(bgr, args.dcvc_max_side)
rgb = cv2.cvtColor(bgr, cv2.COLOR_BGR2RGB)
if not seq_started:
_ENGINE.reset_sequence(rgb.shape[0], rgb.shape[1], args.qp_i, args.qp_p)
seq_started = True
bm = _ENGINE.step(fidx, rgb)
if fidx in needed:
kept_rgb[fidx] = rgb
bitmaps[fidx] = bm.float().cpu().numpy()
fidx += 1
cap.release()
if not kept_rgb:
return "fail", f"{job.key} decoded no frames"
# Align sampled order (repeat last available for any missing tail id).
ordered_ids: List[int] = []
ordered_frames: List[np.ndarray] = []
last_ok = None
for fid in sampled_ids:
fid = int(fid)
if fid in kept_rgb:
last_ok = fid
use = fid if fid in kept_rgb else last_ok
if use is None:
continue
ordered_ids.append(fid)
ordered_frames.append(kept_rgb[use])
if fid not in bitmaps and use in bitmaps:
bitmaps[fid] = bitmaps[use]
cfg = AssembleConfig(
patch=args.patch,
spatial_merge_size=args.spatial_merge_size,
target_canvas=args.target_canvas,
seq_len_frames=len(ordered_ids),
max_pixels=args.max_pixels,
canvas_token_side=args.canvas_token_side,
)
images, src_positions = assemble_canvases(ordered_frames, ordered_ids, bitmaps, cfg)
out_dir.mkdir(parents=True, exist_ok=True)
canvas_files = []
for i, im in enumerate(images):
name = f"canvas_{i:03d}.jpg"
im.save(out_dir / name, quality=int(args.jpg_quality))
canvas_files.append(name)
tmp = out_dir / "src_patch_position.npy.tmp"
with open(tmp, "wb") as f:
np.save(f, src_positions, allow_pickle=False)
os.replace(tmp, out_dir / "src_patch_position.npy")
meta = {
"fps": float(fps),
"canvas_files": canvas_files,
"video": job.video,
"key": job.key,
"total_frames": int(total),
"sampled_frame_ids": [int(x) for x in ordered_ids],
"engine": "dcvc-rt",
"cfg": {
"target_canvas": args.target_canvas,
"seq_len_frames": args.seq_len_frames,
"max_pixels": args.max_pixels,
"canvas_token_side": cfg.token_side(),
"patch": args.patch,
"spatial_merge_size": args.spatial_merge_size,
"qp_i": args.qp_i,
"qp_p": args.qp_p,
"intra_period": args.intra_period,
"reset_interval": args.reset_interval,
"dcvc_max_side": args.dcvc_max_side,
},
"elapsed_sec": round(time.time() - t0, 3),
}
tmpj = out_dir / "meta.json.tmp"
with open(tmpj, "w", encoding="utf-8") as f:
json.dump(meta, f, ensure_ascii=False, indent=2)
os.replace(tmpj, out_dir / "meta.json")
return "ok", f"{job.key} n_canvas={len(images)} elapsed={meta['elapsed_sec']}s"
# ------------------------------------------------------------- worker setup
def _init_worker(args: argparse.Namespace, gpu_queue) -> None:
global _ENGINE
os.environ.setdefault("OMP_NUM_THREADS", "1")
try:
cv2.setNumThreads(1)
except Exception:
pass
gpu_id = gpu_queue.get()
from dcvc_rt_engine import DCVCRTEngine
_ENGINE = DCVCRTEngine(
args.intra_tar, args.inter_tar,
device=f"cuda:{gpu_id}" if gpu_id >= 0 else "cpu",
half=not args.fp32,
intra_period=args.intra_period,
reset_interval=args.reset_interval,
)
print(f"[worker] engine ready on cuda:{gpu_id}", flush=True)
def _worker(job_and_args):
job, args = job_and_args
try:
return _process_one(job, args)
except Exception as e: # keep the pool alive
import traceback
traceback.print_exc()
return "fail", f"{job.key} err={e!r}"
# ----------------------------------------------------------------- jobs io
def _load_jobs(spec: str) -> List[Job]:
p = Path(spec)
jobs: List[Job] = []
if p.is_dir():
for fp in sorted(p.rglob("*")):
if fp.suffix.lower() in _VIDEO_EXTS:
jobs.append(Job(str(fp), fp.stem))
elif p.suffix.lower() == ".jsonl":
with open(p, "r", encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
d = json.loads(line)
if not d.get("exists", True):
continue
video = str(d["video"])
key = str(d.get("key") or Path(video).stem)
jobs.append(Job(video, key))
else:
jobs.append(Job(str(p), p.stem))
return jobs
def main():
ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
ap.add_argument("--videos", required=True, help="jsonl / directory / single video")
ap.add_argument("--out_root", required=True)
ap.add_argument("--intra_tar", default=os.path.join(_HERE, "dcvc_rt_intra.tar"))
ap.add_argument("--inter_tar", default=os.path.join(_HERE, "dcvc_rt_inter.tar"))
# selection / canvas
ap.add_argument("--target_canvas", type=int, default=32)
ap.add_argument("--seq_len_frames", type=int, default=64)
ap.add_argument("--max_pixels", type=int, default=150000)
ap.add_argument("--canvas_token_side", type=int, default=None,
help="tokens/side of square canvas (default: derived from max_pixels)")
ap.add_argument("--patch", type=int, default=16,
help="must match the release image processor patch_size (16)")
ap.add_argument("--spatial_merge_size", type=int, default=2)
# dcvc-rt
ap.add_argument("--qp", type=int, default=21, help="shortcut to set qp_i=qp_p")
ap.add_argument("--qp_i", type=int, default=None)
ap.add_argument("--qp_p", type=int, default=None)
ap.add_argument("--intra_period", type=int, default=-1)
ap.add_argument("--reset_interval", type=int, default=32)
ap.add_argument("--dcvc_max_side", type=int, default=512,
help="downscale decoded frames so max(H,W)<=this before DCVC (0=off)")
ap.add_argument("--fp32", action="store_true")
# runtime
ap.add_argument("--cuda_idx", type=int, nargs="+", default=[0])
ap.add_argument("--num_workers", type=int, default=1)
ap.add_argument("--overwrite", action="store_true")
ap.add_argument("--jpg_quality", type=int, default=95)
ap.add_argument("--max_samples", type=int, default=0)
args = ap.parse_args()
args.qp_i = args.qp if args.qp_i is None else args.qp_i
args.qp_p = args.qp if args.qp_p is None else args.qp_p
jobs = _load_jobs(args.videos)
if args.max_samples > 0:
jobs = jobs[: args.max_samples]
Path(args.out_root).mkdir(parents=True, exist_ok=True)
print(f"[info] jobs={len(jobs)} workers={args.num_workers} gpus={args.cuda_idx}")
ok = skip = fail = 0
if args.num_workers <= 1:
import queue
q = queue.Queue()
for g in ([args.cuda_idx[0]] if args.cuda_idx else [0]):
q.put(g)
_init_worker(args, q)
for i, job in enumerate(jobs):
st, msg = _worker((job, args))
ok += st == "ok"; skip += st == "skip"; fail += st == "fail"
print(f"[{i+1}/{len(jobs)}] {st}: {msg}", flush=True)
else:
import multiprocessing as mp
ctx = mp.get_context("spawn")
gpu_q = ctx.Manager().Queue()
# assign one GPU slot per worker (round-robin over cuda_idx)
for w in range(args.num_workers):
gpu_q.put(args.cuda_idx[w % len(args.cuda_idx)])
with ctx.Pool(args.num_workers, initializer=_init_worker, initargs=(args, gpu_q)) as pool:
for i, (st, msg) in enumerate(pool.imap_unordered(_worker, [(j, args) for j in jobs])):
ok += st == "ok"; skip += st == "skip"; fail += st == "fail"
print(f"[{i+1}/{len(jobs)}] {st}: {msg}", flush=True)
print(f"[done] ok={ok} skip={skip} fail={fail} total={len(jobs)}")
if __name__ == "__main__":
main()