""" VideoProcessor - 整视频分析编排 流程(不再切片/抽帧): 1. 从 OracleDB 取当前 known_members_context(已命名/合并的人物) 2. 按 vision_order 依次调适配器的 analyze_video(Gemini 整视频 -> NVIDIA 整视频) 3. 首个成功结果 -> 归一化 -> 写 OracleDB(videos + events 表) 4. 把本视频 people_mentioned 更新进 people 表(供 person_service 后续合并) 降级: 全部视觉模型失败 -> 标记视频 failed(不再本地融合) """ import os import re from datetime import datetime, timedelta, timezone from typing import Dict, List, Optional from .logger import setup_logger from .config_loader import load_config from .model_adapters.adapter_factory import build_adapters from .model_adapters.base_adapter import BaseModelAdapter from . import oracle_db logger = setup_logger('fam-edge.video_processor') def _parse_event_start_from_filename(filename: str) -> str: """从监控文件名解析开始时间(北京时间)。 支持格式: - 2026-08-21_081500.mp4 / 20260821_081500.mp4(带/不带分隔符日期) - Generic_ONVIF-001-20260820-140416-xxx.mp4(纯数字 YYYYMMDD-HHMMSS) """ # 纯数字: YYYYMMDD-HHMMSS 或 YYYYMMDDHHMMSS(监控录像文件名格式) m0 = re.search(r'(\d{4})(\d{2})(\d{2})[-_]?(\d{2})(\d{2})(\d{2})', filename) if m0: y, mo, d, hh, mm, ss = m0.groups() try: dt = datetime(int(y), int(mo), int(d), int(hh), int(mm), int(ss)) return dt.strftime('%Y-%m-%d %H:%M:%S') except ValueError: pass m = re.search(r'(\d{4})[-_](\d{2})[-_](\d{2})[_-]?(\d{2})(\d{2})(\d{2})', filename) if m: y, mo, d, hh, mm, ss = m.groups() try: dt = datetime(int(y), int(mo), int(d), int(hh), int(mm), int(ss)) return dt.strftime('%Y-%m-%d %H:%M:%S') except ValueError: pass # 退而求其次: 2026-08-21 08-15-00 等 m2 = re.search(r'(\d{4}-\d{2}-\d{2})[ _T-]+(\d{2})[-:](\d{2})[-:](\d{2})', filename) if m2: return f"{m2.group(1)} {m2.group(2)}:{m2.group(3)}:{m2.group(4)}" return '' def _parse_event_ts(ts: str, start_dt): """解析事件时间戳 -> (绝对时间显示串, 视频内偏移秒)。 优先识别"视频内相对时间" HH:MM:SS(新 prompt 要求,定位最准); 兼容旧数据的绝对时间 YYYY-MM-DD HH:MM:SS(偏移=绝对-视频开始)。 """ ts = (ts or '').strip() m = re.match(r'^(\d{1,2}):(\d{2}):(\d{2})$', ts) if m: off = int(m.group(1)) * 3600 + int(m.group(2)) * 60 + int(m.group(3)) # 启发式:监控单段通常 ≤1h,相对时间超过 6h 视为模型误输出绝对时间(无日期),不强行定位 if off <= 6 * 3600: if start_dt is not None: abs_ts = (start_dt + timedelta(seconds=off)).strftime('%Y-%m-%d %H:%M:%S') return abs_ts, float(off) return ts, float(off) if start_dt is not None: try: ev_dt = datetime.strptime(ts[:19], '%Y-%m-%d %H:%M:%S') return ts, (ev_dt - start_dt).total_seconds() except ValueError: pass return ts, 0.0 if start_dt is not None: try: ev_dt = datetime.strptime(ts[:19], '%Y-%m-%d %H:%M:%S') return ts, (ev_dt - start_dt).total_seconds() except ValueError: pass return ts, 0.0 def validate_video(path: str) -> tuple: """校验视频文件是否为正常可解码视频。 返回 (ok: bool, error: str, meta: dict|None) - meta: {fps, frames, duration_sec, width, height} 用 OpenCV 打开并读取至少 1 帧(不校验会导致空/半成品文件浪费云端配额)。 """ meta = None try: if not path or not os.path.isfile(path): return False, "file_missing", None if os.path.getsize(path) == 0: return False, "file_empty", None try: import cv2 except ImportError: return True, "", None # 无 cv2 时跳过深度校验(仅大小检查) cap = cv2.VideoCapture(path) try: if not cap.isOpened(): return False, "cannot_open", None ok, frame = cap.read() if not ok or frame is None: return False, "no_decodable_frame", None fps = float(cap.get(cv2.CAP_PROP_FPS) or 0) frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT) or 0) meta = { "fps": round(fps, 2), "frames": frames, "duration_sec": round(frames / max(fps, 0.01), 1), "width": int(cap.get(cv2.CAP_PROP_FRAME_WIDTH) or 0), "height": int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT) or 0), } finally: cap.release() return True, "", meta except Exception as e: return False, f"validate_exc: {e}", None def _clean_person(s: str) -> str: """清洗人物标识:去掉括号注释(如 "人物A(别名/标识:人物B)" -> "人物A")。""" s = (s or '').strip() for sep in ('(', '('): if sep in s: s = s.split(sep, 1)[0].strip() break return s class VideoProcessor: def __init__(self, db: oracle_db.OracleDB): self.config = load_config() self.db = db self.vision_order = self.config.get('video_processing', {}).get( 'vision_order', ['gemini', 'nvidia']) self.vision_timeout = self.config.get('video_processing', {}).get('timeout', 900) self.file_validate = bool(self.config.get('video_processing', {}).get( 'file_validate', True)) self.parse_start = self.config.get('gdrive_sync', {}).get( 'parse_start_from_filename', True) adapters = build_adapters(self.config.get('models', [])) self.vision_adapters: Dict[str, BaseModelAdapter] = { a.provider_name: a for a in adapters if a.get_role() == 'vision'} def _ordered_vision_adapters(self) -> List[BaseModelAdapter]: ordered = [] for name in self.vision_order: if name in self.vision_adapters: ordered.append(self.vision_adapters[name]) # 追加未在顺序里但启用的视觉适配器 for name, a in self.vision_adapters.items(): if name not in self.vision_order: ordered.append(a) return ordered def process_video(self, video_id: int, filename: str, local_path: str, timeout_multiplier: float = 1.0) -> bool: """处理一个视频记录,返回是否成功。 timeout_multiplier: 云端模型消费的超时放大倍数(如 2 = 在配置 timeout 上 ×2)。 每次调用前临时放大对应 adapter.timeout,调用后恢复,避免影响其他调用方。 """ if not os.path.isfile(local_path): logger.error(f"[video_id={video_id}] 文件不存在,跳过: {local_path}") self.db.mark_video_failed(video_id, "file_missing") return False # 处理前二次确认文件有效性(防止登记后文件被破坏/截断;校验结果落库) if self.file_validate: ok, verr, vmeta = validate_video(local_path) if not ok: logger.error(f"[video_id={video_id}] 文件校验失败({verr}),标记 failed: {local_path}") self.db.set_video_file_status(video_id, False, verr) self.db.mark_video_failed(video_id, f"invalid_file:{verr}") return False if vmeta: self.db.set_video_file_status(video_id, True, '', vmeta) camera_name = self.db.get_video_by_filename(filename)['camera_name'] or '' event_start = '' if self.parse_start: event_start = _parse_event_start_from_filename(filename) # 回写解析到的开始时间 if event_start: self.db._conn.execute( "UPDATE videos SET event_start_time=? WHERE id=?", (event_start, video_id)) self.db._conn.commit() known = self.db.get_known_members_context() logger.info(f"[video_id={video_id}] 开始整视频分析: {filename} " f"(event_start={event_start}, known_members={'有' if known else '无'})") last_err = "no_vision_adapter" for adapter in self._ordered_vision_adapters(): # 模型调用统计 hook(带当前 video_id/filename,前端展示用) adapter.model_call_hook = ( lambda p, m, s, d, ok, e, _vid=video_id, _fn=filename: self.db.record_model_call(p, m, _vid, _fn, s, d, ok, e)) # 按模型原配置超时 × multiplier(默认 1x;队列消费默认 2x) orig_timeout = adapter.get_timeout() if timeout_multiplier != 1.0: adapter.timeout = int(orig_timeout * timeout_multiplier) logger.info(f"[video_id={video_id}] {adapter.provider_name} 超时 " f"{orig_timeout}s -> {adapter.timeout}s (×{timeout_multiplier})") try: logger.info(f"[video_id={video_id}] 尝试 {adapter.provider_name} 整视频分析") result = adapter.analyze_video(local_path, known, event_start) except Exception as e: logger.error(f"[video_id={video_id}] {adapter.provider_name} 异常: {e}") last_err = str(e) continue finally: adapter.timeout = orig_timeout if result: self._store_result(video_id, result) return True else: last_err = f"{adapter.provider_name}_failed" logger.warning(f"[video_id={video_id}] {adapter.provider_name} 未返回结果,降级下一模型") logger.error(f"[video_id={video_id}] 所有视觉模型失败,标记 failed: {last_err}") self.db.mark_video_failed(video_id, last_err) return False def _store_result(self, video_id: int, result: Dict): events = result.get('events', []) people = result.get('people_mentioned', []) summary = result.get('global_summary', '') provider = result.get('compute_provider', 'unknown') # 视频开始时间(绝对时间由后端精确计算:开始时间 + 相对偏移) start_dt = None vrow = self.db.get_video_by_id(video_id) if vrow and vrow['event_start_time']: try: start_dt = datetime.strptime(vrow['event_start_time'], '%Y-%m-%d %H:%M:%S') except ValueError: pass norm_events = [] offsets = [] for ev in events: abs_ts, off = _parse_event_ts(ev.get('timestamp'), start_dt) norm_events.append({ "timestamp": abs_ts, "description": str(ev.get('description', '')), "people": [_clean_person(str(p)) for p in ev.get('people', []) if p], "is_attention_event": bool(ev.get('is_attention_event', False)), }) offsets.append(off) # 清洗 people_mentioned(去掉括号注释串,防污染人物表/合并) people = [_clean_person(str(p)) for p in people if p] people = [p for p in people if p and p not in ('无人', '无')] event_ids = self.db.mark_video_processed(video_id, summary, norm_events, people, provider) # 缩略图 + 每个事件对应时间点的画面截图(用相对偏移直接定位,避免模型绝对时间误差) if vrow and vrow['local_path']: self._generate_thumb(video_id, vrow['local_path']) self._generate_event_thumbs(video_id, vrow['local_path'], offsets, event_ids) # 更新 people 表(标签级,待 person_service 合并) for p in people: if p and p not in ('无人', '无'): self.db.upsert_person(p, source='llm') logger.info(f"[video_id={video_id}] 已落库: summary={len(summary)}字, " f"events={len(norm_events)}, people={people}") def _thumbs_dir(self) -> str: db_path = self.config.get('oracle_db', {}).get( 'path', '/opt/fam-edge/data/oracle.db') d = os.path.abspath(os.path.join(os.path.dirname(db_path), '..', 'thumbs')) os.makedirs(d, exist_ok=True) return d def _generate_thumb(self, video_id: int, video_path: str) -> bool: """抽视频首帧生成 JPEG 缩略图(/opt/fam-edge/thumbs/{video_id}.jpg)""" try: import cv2 out = os.path.join(self._thumbs_dir(), f"{video_id}.jpg") cap = cv2.VideoCapture(video_path) try: ok, frame = cap.read() finally: cap.release() if not ok or frame is None: logger.warning(f"抽帧失败 video_id={video_id}: 无法读取首帧") return False h, w = frame.shape[:2] if w > 640: frame = cv2.resize(frame, (640, int(h * 640 / w))) cv2.imwrite(out, frame, [cv2.IMWRITE_JPEG_QUALITY, 65]) logger.info(f"缩略图已生成: {out}") return True except Exception as e: logger.warning(f"抽帧异常 video_id={video_id}: {e}") return False def _generate_event_thumbs(self, video_id: int, video_path: str, offsets: List[float], event_ids: List[int]): """按事件在视频内的偏移秒定位帧,生成事件画面截图 ev_{event_id}.jpg""" try: import cv2 except Exception as e: logger.warning(f"事件截图依赖缺失 video_id={video_id}: {e}") return try: thumbs = self._thumbs_dir() cap = cv2.VideoCapture(video_path) try: for off, eid in zip(offsets, event_ids): if off < 0: off = 0.0 cap.set(cv2.CAP_PROP_POS_MSEC, int(off * 1000)) ok, frame = cap.read() if not ok or frame is None: cap.set(cv2.CAP_PROP_POS_FRAMES, 0) ok, frame = cap.read() if not ok or frame is None: logger.warning(f"事件截图失败 ev_{eid}: 无法读取 offset={off:.0f}s") continue h, w = frame.shape[:2] if w > 640: frame = cv2.resize(frame, (640, int(h * 640 / w))) out = os.path.join(thumbs, f"ev_{eid}.jpg") cv2.imwrite(out, frame, [cv2.IMWRITE_JPEG_QUALITY, 65]) logger.info(f"事件截图已生成 ev_{eid}.jpg (offset={off:.0f}s)") finally: cap.release() except Exception as e: logger.warning(f"事件截图异常 video_id={video_id}: {e}")