#!/usr/bin/env python3 """ QC Report for all completed files. Checks pipeline standards and lists missing/invalid items. """ import json import os import sys import psycopg2 from pathlib import Path OUTPUT_DIR = os.environ.get("MOMENTRY_OUTPUT_DIR", "/Users/accusys/momentry/output") DATABASE_URL = os.environ.get("DATABASE_URL", "postgresql://accusys@localhost:5432/momentry") # Required processor outputs REQUIRED_PROCESSORS = ["face.json", "asr.json", "asrx.json", "ocr.json", "pose.json", "cut.json", "face_cluster.json", "face_traced.json", "profile.json"] def get_db_connection(): return psycopg2.connect(DATABASE_URL) def check_processor_files(uuid): """Check if all required processor output files exist.""" missing = [] for proc in REQUIRED_PROCESSORS: path = Path(OUTPUT_DIR) / f"{uuid}.{proc}" if not path.exists(): missing.append(proc) return missing def check_trace_profiles(conn, uuid): """Check trace_profiles data quality.""" with conn.cursor() as cur: cur.execute("SELECT COUNT(*), SUM(frame_count), COUNT(CASE WHEN name IS NOT NULL AND name != '' THEN 1 END) FROM public.trace_profiles WHERE file_uuid = %s", (uuid,)) count, total_frames, named_count = cur.fetchone() cur.execute("SELECT COUNT(*) FROM public.trace_profiles WHERE file_uuid = %s AND frame_count <= 0", (uuid,)) zero_frames = cur.fetchone()[0] cur.execute("SELECT COUNT(*) FROM public.trace_profiles WHERE file_uuid = %s AND (vlm_description IS NULL OR vlm_description = '')", (uuid,)) no_vlm_desc = cur.fetchone()[0] issues = [] if count == 0: issues.append("trace_profiles: 0 records") if zero_frames > 0: issues.append(f"trace_profiles: {zero_frames} records with frame_count <= 0") if no_vlm_desc > 0: issues.append(f"trace_profiles: {no_vlm_desc} records missing vlm_description") return issues def check_tkg_nodes(conn, uuid): """Check TKG nodes data.""" with conn.cursor() as cur: cur.execute("SELECT COUNT(*) FROM public.tkg_nodes WHERE file_uuid = %s AND node_type = 'face_track'", (uuid,)) face_tracks = cur.fetchone()[0] cur.execute("SELECT COUNT(*) FROM public.tkg_edges WHERE file_uuid = %s", (uuid,)) edges = cur.fetchone()[0] issues = [] if face_tracks == 0: issues.append("tkg_nodes: 0 face_track nodes") if edges == 0: issues.append("tkg_edges: 0 edges") return issues def check_qdrant_faces(uuid): """Check Qdrant _faces collection using curl.""" try: import subprocess result = subprocess.run( ["curl", "-s", "http://localhost:6333/collections/_faces/points/scroll", "-H", "Content-Type: application/json", "-d", json.dumps({"filter": {"must": [{"key": "file_uuid", "match": {"value": uuid}}]}, "limit": 1})], capture_output=True, text=True, timeout=5 ) if result.returncode == 0: data = json.loads(result.stdout) points = data.get("result", {}).get("points", []) count = len(points) return [] if count > 0 else [f"Qdrant _faces: {count} points"] else: return [f"Qdrant _faces: curl error"] except Exception as e: return [f"Qdrant _faces: error ({e})"] def check_video_metadata(conn, uuid): """Check video metadata completeness.""" with conn.cursor() as cur: cur.execute("SELECT file_name, duration, fps, total_frames, cut_done FROM public.videos WHERE file_uuid = %s", (uuid,)) row = cur.fetchone() if not row: return ["videos: record not found"] name, duration, fps, total_frames, cut_done = row issues = [] if duration <= 0: issues.append(f"videos: duration={duration}") if fps <= 0: issues.append(f"videos: fps={fps}") if total_frames <= 0: issues.append(f"videos: total_frames={total_frames}") if not cut_done: issues.append("videos: cut_done=false") return issues def main(): conn = get_db_connection() with conn.cursor() as cur: cur.execute("SELECT file_uuid, file_name FROM public.videos WHERE status = 'completed' ORDER BY created_at DESC") completed_files = cur.fetchall() if not completed_files: print("No completed files found.") return print("=" * 80) print("QC REPORT FOR COMPLETED FILES") print("=" * 80) print(f"Total completed files: {len(completed_files)}") print() all_issues = [] for uuid, file_name in completed_files: file_issues = [] # Check processor files missing_procs = check_processor_files(uuid) if missing_procs: file_issues.append(f"Missing processor files: {', '.join(missing_procs)}") # Check video metadata file_issues.extend(check_video_metadata(conn, uuid)) # Check trace_profiles file_issues.extend(check_trace_profiles(conn, uuid)) # Check TKG file_issues.extend(check_tkg_nodes(conn, uuid)) # Check Qdrant file_issues.extend(check_qdrant_faces(uuid)) status = "PASS" if not file_issues else "FAIL" print(f"[{status}] {file_name} ({uuid})") for issue in file_issues: print(f" - {issue}") print() if file_issues: all_issues.append((file_name, uuid, file_issues)) print("=" * 80) print(f"SUMMARY: {len(completed_files) - len(all_issues)}/{len(completed_files)} files passed QC") if all_issues: print(f"\nFailed files:") for name, uuid, issues in all_issues: print(f" - {name}: {len(issues)} issue(s)") if __name__ == "__main__": main()