#!/usr/bin/env python3 """ Backfill trace profiles from Qdrant _faces collection. For each (file_uuid, trace_id) group in Qdrant: 1. Compute frame_count, start_frame, end_frame, avg_confidence 2. Pick representative frame (highest confidence) 3. Extract key_frame.jpg from video via ffmpeg 4. Crop key_face.jpg from key_frame using representative bbox 5. Write output/{file_uuid}/trace_{N}/trace_profile.json Usage: python3 backfill_trace_profiles.py [--file-uuid UUID] [--dry-run] """ import argparse import json import os import subprocess import sys import urllib.request import urllib.error from collections import defaultdict OUTPUT_DIR = os.environ.get("MOMENTRY_OUTPUT_DIR", "/Users/accusys/momentry/output") QDRANT_URL = os.environ.get("QDRANT_URL", "http://localhost:6333") QDRANT_API_KEY = os.environ.get("QDRANT_API_KEY", "Test3200Test3200Test3200") FACES_COLLECTION = "_faces" BATCH_SIZE = 1000 def qdrant_scroll(filter_dict, limit=BATCH_SIZE, offset=None, with_payload=None): """Scroll Qdrant collection with filter.""" body = {"limit": limit, "filter": filter_dict, "with_vector": False} if offset: body["offset"] = offset if with_payload: body["with_payload"] = with_payload url = f"{QDRANT_URL}/collections/{FACES_COLLECTION}/points/scroll" data = json.dumps(body).encode() req = urllib.request.Request(url, data=data, method="POST") req.add_header("Content-Type", "application/json") req.add_header("Api-Key", QDRANT_API_KEY) with urllib.request.urlopen(req) as resp: return json.loads(resp.read()) def scroll_all(filter_dict, with_payload=None): """Scroll all matching points.""" all_points = [] offset = None while True: result = qdrant_scroll(filter_dict, offset=offset, with_payload=with_payload) points = result.get("result", {}).get("points", []) if not points: break all_points.extend(points) offset = result.get("result", {}).get("next_page_offset") if not offset or len(points) < BATCH_SIZE: break return all_points def get_video_path(file_uuid): """Get video file path from database.""" psql = "/opt/homebrew/Cellar/libpq/18.4/bin/psql" result = subprocess.run( [psql, "-U", "accusys", "-d", "momentry", "-t", "-A", "-c", f"SELECT file_path FROM videos WHERE file_uuid = '{file_uuid}'"], capture_output=True, text=True ) if result.returncode == 0 and result.stdout.strip(): return result.stdout.strip() return None def extract_key_frame(video_path, frame_num, fps, output_path): """Extract a specific frame from video using ffmpeg.""" if fps <= 0: return False timestamp = frame_num / fps try: result = subprocess.run( ["ffmpeg", "-y", "-ss", f"{timestamp:.3f}", "-i", video_path, "-vframes", "1", "-vf", "scale=640:-1", "-q:v", "5", output_path], capture_output=True, timeout=30 ) return result.returncode == 0 and os.path.exists(output_path) except Exception as e: print(f" key_frame extraction failed: {e}", file=sys.stderr) return False def crop_key_face(key_frame_path, bbox, output_path): """Crop key_face from key_frame using bbox.""" x, y, w, h = bbox["x"], bbox["y"], bbox["width"], bbox["height"] if w <= 0 or h <= 0: return False try: result = subprocess.run( ["ffmpeg", "-y", "-i", key_frame_path, "-vf", f"crop={w}:{h}:{x}:{y}", "-q:v", "2", output_path], capture_output=True, timeout=10 ) return result.returncode == 0 and os.path.exists(output_path) except Exception as e: print(f" key_face crop failed: {e}", file=sys.stderr) return False def build_trace_profiles(file_uuid=None, dry_run=False): """Build trace profiles from Qdrant _faces data.""" # Get all unique file_uuids with trace_id >= 0 if file_uuid: file_uuids = [file_uuid] else: print("Scanning Qdrant for all file_uuids with trace_id >= 0...") points = scroll_all( {"must": [{"key": "trace_id", "range": {"gte": 0}}]}, with_payload={"include": ["file_uuid"]} ) file_uuids = sorted(set(p["payload"]["file_uuid"] for p in points)) print(f"Found {len(file_uuids)} files with trace data") total_profiles = 0 for fid in file_uuids: print(f"\n--- {fid} ---") # Scroll all points for this file with trace_id >= 0 points = scroll_all( { "must": [ {"key": "file_uuid", "match": {"value": fid}}, {"key": "trace_id", "range": {"gte": 0}}, ] }, with_payload={"include": ["frame", "trace_id", "bbox", "confidence"]} ) if not points: print(" No points with trace_id >= 0") continue # Group by trace_id traces = defaultdict(list) for p in points: pl = p["payload"] tid = pl.get("trace_id", 0) traces[tid].append({ "frame": pl["frame"], "bbox": pl.get("bbox", {}), "confidence": pl.get("confidence", 0.0), }) print(f" {len(points)} points, {len(traces)} traces") # Get video path video_path = get_video_path(fid) if not video_path or not os.path.exists(video_path): print(f" Video not found, skipping key_frame extraction") video_path = None # Get FPS from DB fps = 30.0 if video_path: psql = "/opt/homebrew/Cellar/libpq/18.4/bin/psql" result = subprocess.run( [psql, "-U", "accusys", "-d", "momentry", "-t", "-A", "-c", f"SELECT COALESCE(fps, 30.0) FROM videos WHERE file_uuid = '{fid}'"], capture_output=True, text=True ) if result.returncode == 0 and result.stdout.strip(): try: fps = float(result.stdout.strip()) except ValueError: pass for tid, faces in sorted(traces.items()): if tid < 0: continue frames = [f["frame"] for f in faces] confidences = [f["confidence"] for f in faces] frame_count = len(faces) start_frame = min(frames) end_frame = max(frames) avg_confidence = sum(confidences) / frame_count if frame_count > 0 else 0.0 # Representative frame: highest confidence best = max(faces, key=lambda f: f["confidence"]) best_frame = best["frame"] best_bbox = best["bbox"] trace_dir = os.path.join(OUTPUT_DIR, fid, f"trace_{tid}") profile_path = os.path.join(trace_dir, "trace_profile.json") kf_path = os.path.join(trace_dir, "key_frame.jpg") face_path = os.path.join(trace_dir, "key_face.jpg") profile = { "version": "1.0", "file_uuid": fid, "trace_id": tid, "label": "", "frame_count": frame_count, "start_frame": start_frame, "end_frame": end_frame, "avg_confidence": round(avg_confidence, 6), "key_frame": "key_frame.jpg" if os.path.exists(kf_path) else None, "key_face": "key_face.jpg" if os.path.exists(face_path) else None, "status": "pending", } if dry_run: print(f" trace_{tid}: {frame_count} frames [{start_frame}-{end_frame}] " f"conf={avg_confidence:.3f} best_frame={best_frame}") continue os.makedirs(trace_dir, exist_ok=True) # Extract key_frame.jpg if not exists if not os.path.exists(kf_path) and video_path: extract_key_frame(video_path, best_frame, fps, kf_path) if os.path.exists(kf_path): profile["key_frame"] = "key_frame.jpg" # Crop key_face.jpg from key_frame if not exists if not os.path.exists(face_path) and os.path.exists(kf_path) and best_bbox: crop_key_face(kf_path, best_bbox, face_path) if os.path.exists(face_path): profile["key_face"] = "key_face.jpg" # Write trace_profile.json with open(profile_path, "w") as f: json.dump(profile, f, indent=2, ensure_ascii=False) total_profiles += 1 if not dry_run: print(f" Created {len([t for t in traces if t >= 0])} trace profiles") print(f"\nDone: {total_profiles} trace profiles created") def main(): parser = argparse.ArgumentParser(description="Backfill trace profiles from Qdrant") parser.add_argument("--file-uuid", help="Process only this file UUID") parser.add_argument("--dry-run", action="store_true", help="Show what would be created") args = parser.parse_args() build_trace_profiles(file_uuid=args.file_uuid, dry_run=args.dry_run) if __name__ == "__main__": main()