diff --git a/fam-core/src/fam_core/dispatcher/dispatcher.py b/fam-core/src/fam_core/dispatcher/dispatcher.py index 650fe23..67849de 100644 --- a/fam-core/src/fam_core/dispatcher/dispatcher.py +++ b/fam-core/src/fam_core/dispatcher/dispatcher.py @@ -17,6 +17,7 @@ Dispatcher - 30s 轮询 PENDING 任务,上传视频至 Edge 异步队列 """ import os import io +import re import time import math import shutil @@ -90,6 +91,22 @@ class Dispatcher: return False return True + def _probe_duration(self, video_path): + """用 ffmpeg 解析视频时长(NAS 无独立 ffprobe),失败返回 0""" + if not self.ffmpeg: + return 0 + try: + r = subprocess.run( + [self.ffmpeg, '-i', video_path], + capture_output=True, timeout=30) + m = re.search(r'Duration:\s*(\d+):(\d+):(\d+(?:\.\d+)?)', + r.stderr.decode('utf-8', 'ignore')) + if m: + return int(m.group(1)) * 3600 + int(m.group(2)) * 60 + float(m.group(3)) + except (subprocess.TimeoutExpired, OSError): + pass + return 0 + def _build_payload(self, task): """构建推送元数据,注入已知成员清单与事件时间""" video_path = task['video_path'] @@ -102,7 +119,9 @@ class Dispatcher: } try: mtime = os.path.getmtime(video_path) - start_dt = datetime.fromtimestamp(mtime) + # mtime 是录制结束时刻,开始时间 = 结束时间 - 视频时长 + duration = self._probe_duration(video_path) + start_dt = datetime.fromtimestamp(mtime - duration) payload["event_start_time"] = start_dt.strftime('%Y-%m-%d %H:%M:%S') except OSError: pass @@ -238,7 +257,7 @@ class Dispatcher: self.edge_url, data=payload, files={'video': (os.path.basename(video_path), fh, 'video/mp4')}, - timeout=(10, 120) + timeout=(10, 300) ) except requests.RequestException as e: logger.error(f"[task_id={task_id}] 上传失败: {e}") diff --git a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py index e5f83e7..eb99ccd 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py @@ -54,7 +54,8 @@ class AIOrchestrator: frame_timestamps: List[str], known_members_context: str, rate_limiter=None, - video_path: str = None) -> Dict[str, dict]: + video_path: str = None, + event_start_time: str = '') -> Dict[str, dict]: """视觉分析阶段:仅 role=vision 的适配器参与 支持 analyze_video 的适配器(如 NVIDIA Omni)优先走原生视频输入, @@ -97,7 +98,8 @@ class AIOrchestrator: try: logger.info(f"[{adapter.provider_name}] 尝试原生视频输入分析") output = adapter.analyze_video( - video_path, frame_timestamps, known_members_context) + video_path, frame_timestamps, known_members_context, + event_start_time=event_start_time) if not output: logger.warning(f"[{adapter.provider_name}] 视频模式失败,降级逐帧模式") except Exception as ve: @@ -371,7 +373,8 @@ class AIOrchestrator: # 3. 并行视觉分析 model_outputs = self.run_visual_analysis( healthy_adapters, compressed_frames, frame_timestamps, - known_members, video_path=video_path + known_members, video_path=video_path, + event_start_time=event_start_time ) if not model_outputs: @@ -467,7 +470,8 @@ class AIOrchestrator: # 3. 并行视觉分析 model_outputs = self.run_visual_analysis( healthy_adapters, compressed_frames, frame_timestamps, - known_members, rate_limiter, video_path=video_path + known_members, rate_limiter, video_path=video_path, + event_start_time=event_start_time ) if not model_outputs: raise Exception('All models failed in visual analysis') diff --git a/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py b/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py index 6f68779..7898b29 100644 --- a/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py +++ b/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py @@ -98,20 +98,42 @@ class NvidiaVisionAdapter(BaseModelAdapter): except ValueError: return -1.0 - def _build_highlight_video(self, video_path: str, - frame_timestamps: List[str]) -> Optional[str]: + @staticmethod + def _probe_duration(video_path: str) -> float: + """ffprobe 解析视频时长,失败返回 0""" + try: + r = subprocess.run( + ['ffprobe', '-v', 'quiet', '-show_entries', 'format=duration', + '-of', 'csv=p=0', video_path], + capture_output=True, timeout=30) + return float(r.stdout.decode().strip() or 0) + except (subprocess.TimeoutExpired, OSError, ValueError): + return 0.0 + + def _build_highlight_video(self, video_path: str, frame_timestamps: List[str], + event_start_time: str = '') -> Optional[str]: """按关键帧时间点截取 ±pad 秒片段,叠加时间戳后拼接集锦视频 - 时间戳以 00-05-00 形式叠加(连字符避免 ffmpeg drawtext 冒号转义)。 + 时间戳以 2026-08-20 04-34-10 形式叠加(连字符避免 ffmpeg drawtext 冒号转义)。 + 偏移换算: 绝对时间戳 - 视频开始时间(event_start_time 缺失时, + 时间戳值本身须已是视频内偏移,如 HH:MM:SS 相对时间)。 """ + start_sec = self._ts_to_seconds(event_start_time) if event_start_time else 0.0 + duration = self._probe_duration(video_path) clips = [] for ts in frame_timestamps: sec = self._ts_to_seconds(ts) if sec < 0: continue - start = max(0.0, sec - VIDEO_SEGMENT_PAD) - label = str(ts).replace(':', '-') - clips.append((start, label)) + if start_sec > 0: + sec -= start_sec + if sec < 0: + sec += 86400 # 跨午夜 + if duration > 0 and (sec < -VIDEO_SEGMENT_PAD + or sec > duration - 0.5): + logger.info(f"NVIDIA 跳过超界片段: {ts} -> {sec:.1f}s (视频 {duration:.0f}s)") + continue + clips.append((sec, str(ts).replace(':', '-'))) if not clips: return None @@ -123,12 +145,13 @@ class NvidiaVisionAdapter(BaseModelAdapter): parts = [] for i, (_, label) in enumerate(clips): parts.append( - f"[{i}:v]scale={HIGHLIGHT_WIDTH}:-2," + f"[{i}:v]fps=15,scale={HIGHLIGHT_WIDTH}:-2," f"drawtext=text='ts {label}':x=8:y=8:fontsize=22:" f"fontcolor=white:box=1:boxcolor=black@0.6[v{i}]") concat_in = ''.join(f'[v{i}]' for i in range(len(clips))) parts.append(f'{concat_in}concat=n={len(clips)}:v=1:a=0[out]') cmd += ['-filter_complex', ';'.join(parts), '-map', '[out]', + '-r', '15', '-c:v', 'libx264', '-preset', 'veryfast', '-crf', '28', '-an', out_path] try: @@ -147,7 +170,8 @@ class NvidiaVisionAdapter(BaseModelAdapter): def analyze_video(self, video_path: str, frame_timestamps: List[str], - known_members_context: str) -> Optional[Dict]: + known_members_context: str, + event_start_time: str = '') -> Optional[Dict]: """原生视频输入分析: 集锦片段 -> video_url 单次调用""" if self._cb.is_open(): logger.warning("NVIDIA 熔断器 OPEN,跳过视频分析") @@ -156,7 +180,8 @@ class NvidiaVisionAdapter(BaseModelAdapter): logger.warning("NVIDIA 客户端未初始化,跳过视频分析") return None - highlight = self._build_highlight_video(video_path, frame_timestamps) + highlight = self._build_highlight_video( + video_path, frame_timestamps, event_start_time) if not highlight: logger.warning("NVIDIA 集锦视频不可用,降级逐帧模式") return None