Lane B / session-5 item 3: default ingest to local queue/ on ultra

The pose engine (HSMR) runs on ultra, so nothing was consuming the m3ultra queue
ingest.py shipped to. ingest.py now delivers to a local queue/ dir by default;
--remote-queue [user@host:/path] restores the m3ultra rsync path for when a real
throughput need appears (ponytail marker + PLAN Ground-rule-1 deviation noted).
CAPTURE.md updated. Verified: inbox drop -> local queue/, ledger row, selfcheck green.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
type-two 2026-07-17 10:02:16 +10:00
parent 9db4fec988
commit 0cf673aa82
3 changed files with 44 additions and 22 deletions

View File

@ -18,6 +18,10 @@ critical path. Lanes EF are follow-ons. Do NOT start E/F before D ships.
training. Python via conda/uv there. MODELBEAST lives here; mocap becomes a sibling training. Python via conda/uv there. MODELBEAST lives here; mocap becomes a sibling
service, same patterns (job dirs, batch keys — see `MESHGOD/scripts/mb_*` for the service, same patterns (job dirs, batch keys — see `MESHGOD/scripts/mb_*` for the
client convention). client convention).
- **v1 deviation (session 5):** HSMR is installed and MPS-verified on **ultra**, not
m3ultra, so v1 runs inference on ultra and `ingest.py` defaults to a **local** `queue/`.
`# ponytail: v1 inference on ultra; migrate to m3ultra when a real throughput need appears`
(`ingest.py --remote-queue` already ships to m3ultra when that day comes).
- **ultra** (this Mac, M1 Ultra 128GB) — GLUE + DCC. Blender headless retarget/cleanup/ - **ultra** (this Mac, M1 Ultra 128GB) — GLUE + DCC. Blender headless retarget/cleanup/
export, orchestration, Unreal verification. Owns the canonical repo. export, orchestration, Unreal verification. Owns the canonical repo.
- **JING5** (M5 MBP) — capture companion + portable inference later. Not in v1 critical path. - **JING5** (M5 MBP) — capture companion + portable inference later. Not in v1 critical path.

View File

@ -25,8 +25,9 @@ uses to auto-trim the clip. No bookends → manual trimming.
- **AirDrop** the clip from the phone to **ultra**, save into - **AirDrop** the clip from the phone to **ultra**, save into
`~/Documents/MOCAPGOD/capture/inbox/`. `~/Documents/MOCAPGOD/capture/inbox/`.
- Then run **`python3 ingest.py`** (or let the Lane D `mocap` CLI do it). Ingest validates, - Then run **`python3 ingest.py`** (or let the Lane D `mocap` CLI do it). Ingest validates,
normalizes to 30fps, and ships it to the m3ultra queue. Rejects (portrait, no feet, garbage) normalizes to 30fps, and puts it in the local pose queue (`queue/` on ultra, where the engine
land in `capture/rejected/` with a reason in `capture/LEDGER.tsv` — check there if a clip vanishes. runs). Rejects (portrait, no feet, garbage) land in `capture/rejected/` with a reason in
`capture/LEDGER.tsv` — check there if a clip vanishes. (`--remote-queue` ships to m3ultra instead.)
## Common rejections (and the fix) ## Common rejections (and the fix)
| Symptom | Fix | | Symptom | Fix |

View File

@ -1,15 +1,17 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""Lane A ingest — capture/inbox/*.{mp4,mov} → validate → 30fps h264 → rename → rsync to """Lane A ingest — capture/inbox/*.{mp4,mov} → validate → 30fps h264 → rename → pose queue
the m3ultra pose queue capture/LEDGER.tsv. capture/LEDGER.tsv.
python3 ingest.py # scan capture/inbox/ and ship everything python3 ingest.py # scan capture/inbox/ → local queue/ (the v1 default)
python3 ingest.py clip.mov --name dj_scratch_01 # one file, explicit clip name python3 ingest.py clip.mov --name dj_scratch_01 # one file, explicit clip name
python3 ingest.py --remote-queue # ship to the m3ultra queue instead (rsync)
python3 ingest.py --selfcheck # run built-in probe/validate self-test, no network python3 ingest.py --selfcheck # run built-in probe/validate self-test, no network
The pose engine (HSMR) currently runs on ultra, so ingest defaults to the LOCAL queue/ dir.
--remote-queue [user@host:/path] rsyncs to m3ultra instead (default target: MOCAP_REMOTE).
Rejects (portrait / malformed) move to capture/rejected/ with a reason in the ledger. Rejects (portrait / malformed) move to capture/rejected/ with a reason in the ledger.
Successful originals move to capture/done/; on rsync failure the normalized file is kept Successful originals move to capture/done/; on a failed remote ship the normalized file is
in capture/staged/ so a re-run retries the ship without re-encoding. kept in capture/staged/ so a re-run retries without re-encoding.
Remote defaults to m3ultra; override with MOCAP_REMOTE=user@host:/path.
""" """
import argparse, datetime, json, os, re, shutil, subprocess, sys, tempfile import argparse, datetime, json, os, re, shutil, subprocess, sys, tempfile
from pathlib import Path from pathlib import Path
@ -20,7 +22,9 @@ STAGED = ROOT / "capture" / "staged"
DONE = ROOT / "capture" / "done" DONE = ROOT / "capture" / "done"
REJECTED = ROOT / "capture" / "rejected" REJECTED = ROOT / "capture" / "rejected"
LEDGER = ROOT / "capture" / "LEDGER.tsv" LEDGER = ROOT / "capture" / "LEDGER.tsv"
REMOTE = os.environ.get("MOCAP_REMOTE", "m3ultra@100.89.131.57:~/Documents/MOCAPGOD/queue") QUEUE = ROOT / "queue" # v1 default: pose engine (HSMR) runs on ultra, next to the repo
# ponytail: v1 inference on ultra; migrate to m3ultra (--remote-queue) when a real throughput need appears
REMOTE_DEFAULT = os.environ.get("MOCAP_REMOTE", "m3ultra@100.89.131.57:~/Documents/MOCAPGOD/queue")
VIDEO_EXT = {".mp4", ".mov", ".m4v"} VIDEO_EXT = {".mp4", ".mov", ".m4v"}
HEADERS = ["utc", "original", "output", "fps_in", "dur_s", "resolution", "audio", "status", "note"] HEADERS = ["utc", "original", "output", "fps_in", "dur_s", "resolution", "audio", "status", "note"]
@ -90,13 +94,21 @@ def normalize(src, dst):
raise ValueError("ffmpeg normalize failed: " + tail) raise ValueError("ffmpeg normalize failed: " + tail)
def ship(staged): def deliver(staged, remote):
host, _, rpath = REMOTE.partition(":") """Put the normalized clip in the pose queue. Local queue/ dir by default; rsync to a
run(["ssh", host, f"mkdir -p {rpath}"]) # queue dir may not exist yet remote (user@host:/path) when --remote-queue is used. Returns the destination string."""
r = run(["rsync", "-az", str(staged), REMOTE + "/"]) if remote is None:
QUEUE.mkdir(parents=True, exist_ok=True)
dst = QUEUE / staged.name
shutil.move(str(staged), dst) # same-fs move; local delivery can't half-fail
return str(dst)
host, _, rpath = remote.partition(":")
run(["ssh", host, f"mkdir -p {rpath}"]) # remote queue dir may not exist yet
r = run(["rsync", "-az", str(staged), remote + "/"])
if r.returncode != 0: if r.returncode != 0:
tail = (r.stderr.strip().splitlines() or ["?"])[-1] tail = (r.stderr.strip().splitlines() or ["?"])[-1]
raise ConnectionError(f"rsync to {REMOTE} failed: {tail}") raise ConnectionError(f"rsync to {remote} failed: {tail}")
return f"{remote}/{staged.name}"
def ledger(row): def ledger(row):
@ -112,7 +124,7 @@ def slug(s):
return re.sub(r"[^A-Za-z0-9_-]+", "_", s).strip("_") or "clip" return re.sub(r"[^A-Za-z0-9_-]+", "_", s).strip("_") or "clip"
def process_one(src, name=None): def process_one(src, name=None, remote=None):
src = Path(src) src = Path(src)
out_name = f"{datetime.date.today():%Y%m%d}_{slug(name or src.stem)}.mp4" out_name = f"{datetime.date.today():%Y%m%d}_{slug(name or src.stem)}.mp4"
ts = datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") ts = datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
@ -139,9 +151,9 @@ def process_one(src, name=None):
res, audio = f"{info['w']}x{info['h']}", "yes" if info["audio"] else "no" res, audio = f"{info['w']}x{info['h']}", "yes" if info["audio"] else "no"
staged = STAGED / out_name staged = STAGED / out_name
try: try:
if not staged.exists(): # skip re-encode when retrying a failed ship if not staged.exists(): # skip re-encode when retrying a failed remote ship
normalize(src, staged) normalize(src, staged)
ship(staged) dest = deliver(staged, remote)
except (ValueError, ConnectionError) as e: except (ValueError, ConnectionError) as e:
ledger([ts, src.name, out_name, info["fps"], f"{info['dur']:.1f}", res, audio, "FAIL", str(e)]) ledger([ts, src.name, out_name, info["fps"], f"{info['dur']:.1f}", res, audio, "FAIL", str(e)])
print(f" ! {src.name}: {e} (left in inbox; staged kept for retry)") print(f" ! {src.name}: {e} (left in inbox; staged kept for retry)")
@ -152,7 +164,7 @@ def process_one(src, name=None):
staged.unlink(missing_ok=True) staged.unlink(missing_ok=True)
ledger([ts, src.name, out_name, info["fps"], f"{info['dur']:.1f}", res, audio, ledger([ts, src.name, out_name, info["fps"], f"{info['dur']:.1f}", res, audio,
"OK", "" if info["audio"] else "audio stripped (none present)"]) "OK", "" if info["audio"] else "audio stripped (none present)"])
print(f"{src.name}{REMOTE}/{out_name} ({res} {info['fps']}fps {info['dur']:.1f}s)") print(f"{src.name}{dest} ({res} {info['fps']}fps {info['dur']:.1f}s)")
return True return True
@ -174,24 +186,29 @@ def selfcheck():
def main(): def main():
ap = argparse.ArgumentParser(description="Lane A: ingest phone videos into the m3ultra pose queue.") ap = argparse.ArgumentParser(description="Lane A: ingest phone videos into the pose queue.")
ap.add_argument("video", nargs="?", help="single file to ingest (default: scan capture/inbox/)") ap.add_argument("video", nargs="?", help="single file to ingest (default: scan capture/inbox/)")
ap.add_argument("--name", help="clip name for the output (default: source filename)") ap.add_argument("--name", help="clip name for the output (default: source filename)")
ap.add_argument("--remote-queue", nargs="?", const=REMOTE_DEFAULT, default=None,
metavar="USER@HOST:PATH",
help=f"rsync to a remote pose queue instead of the local queue/ dir "
f"(default target: {REMOTE_DEFAULT})")
ap.add_argument("--selfcheck", action="store_true", help="run the built-in self-test and exit") ap.add_argument("--selfcheck", action="store_true", help="run the built-in self-test and exit")
args = ap.parse_args() args = ap.parse_args()
if args.selfcheck: if args.selfcheck:
return selfcheck() return selfcheck()
dest_label = args.remote_queue if args.remote_queue else str(QUEUE)
if args.video: if args.video:
return 0 if process_one(args.video, args.name) else 1 return 0 if process_one(args.video, args.name, args.remote_queue) else 1
INBOX.mkdir(parents=True, exist_ok=True) INBOX.mkdir(parents=True, exist_ok=True)
vids = sorted(p for p in INBOX.iterdir() if p.suffix.lower() in VIDEO_EXT) vids = sorted(p for p in INBOX.iterdir() if p.suffix.lower() in VIDEO_EXT)
if not vids: if not vids:
print(f"inbox empty: {INBOX}") print(f"inbox empty: {INBOX}")
return 0 return 0
print(f"ingesting {len(vids)} file(s) from {INBOX}{REMOTE}") print(f"ingesting {len(vids)} file(s) from {INBOX}{dest_label}")
n_ok = sum(process_one(p) for p in vids) n_ok = sum(process_one(p, remote=args.remote_queue) for p in vids)
print(f"done: {n_ok}/{len(vids)} ok") print(f"done: {n_ok}/{len(vids)} ok")
return 0 if n_ok == len(vids) else 1 return 0 if n_ok == len(vids) else 1