diff --git a/fam-edge/config/config.yaml.example b/fam-edge/config/config.yaml.example index 87bfdff..6b58a2a 100644 --- a/fam-edge/config/config.yaml.example +++ b/fam-edge/config/config.yaml.example @@ -69,11 +69,14 @@ models: # cooldown: 900 # - provider: "nvidia" - # enabled: false - # model_name: "nvidia/llama-3.1-nemotron-70b-instruct" + # role: "vision" + # enabled: true + # # Omni 模型原生支持视频输入(video_url),适配器自动按关键帧时间点 + # # 截取片段拼集锦后单次调用;失败自动降级逐帧图片模式 + # model_name: "nvidia/nemotron-3-nano-omni-30b-a3b-reasoning" # api_key: "${NVIDIA_API_KEY}" # base_url: "https://integrate.api.nvidia.com/v1" - # timeout: 30 + # timeout: 120 # reasoning 模型视频推理较慢,勿低于 90 # circuit_breaker: # enabled: true # threshold: 5 diff --git a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py index a7ce75f..e5f83e7 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py @@ -53,9 +53,13 @@ class AIOrchestrator: frame_paths: List[str], frame_timestamps: List[str], known_members_context: str, - rate_limiter=None) -> Dict[str, dict]: + rate_limiter=None, + video_path: str = None) -> Dict[str, dict]: """视觉分析阶段:仅 role=vision 的适配器参与 + 支持 analyze_video 的适配器(如 NVIDIA Omni)优先走原生视频输入, + 失败自动降级回逐帧图片模式。 + orchestrator.mode: - fallback (默认): 按 config 顺序依次尝试,首个成功即采用(单元素 dict) - ensemble: 并行所有健康 vision 模型,全部成功结果都保留(交叉验证) @@ -88,8 +92,20 @@ class AIOrchestrator: continue start = time.time() try: - output = adapter.analyze_frames( - frame_paths, frame_timestamps, known_members_context) + output = None + if video_path and hasattr(adapter, 'analyze_video'): + try: + logger.info(f"[{adapter.provider_name}] 尝试原生视频输入分析") + output = adapter.analyze_video( + video_path, frame_timestamps, known_members_context) + if not output: + logger.warning(f"[{adapter.provider_name}] 视频模式失败,降级逐帧模式") + except Exception as ve: + logger.warning(f"[{adapter.provider_name}] 视频模式异常: {ve},降级逐帧模式") + output = None + if not output: + output = adapter.analyze_frames( + frame_paths, frame_timestamps, known_members_context) duration_ms = int((time.time() - start) * 1000) if output: adapter.get_circuit_breaker().record_success() @@ -354,7 +370,8 @@ class AIOrchestrator: # 3. 并行视觉分析 model_outputs = self.run_visual_analysis( - healthy_adapters, compressed_frames, frame_timestamps, known_members + healthy_adapters, compressed_frames, frame_timestamps, + known_members, video_path=video_path ) if not model_outputs: @@ -450,7 +467,7 @@ class AIOrchestrator: # 3. 并行视觉分析 model_outputs = self.run_visual_analysis( healthy_adapters, compressed_frames, frame_timestamps, - known_members, rate_limiter + known_members, rate_limiter, video_path=video_path ) 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 83cd511..52268c5 100644 --- a/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py +++ b/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py @@ -2,16 +2,23 @@ NvidiaVisionAdapter - NVIDIA NIM 云端 VLM 适配器 provider_name = "nvidia" -模型: meta/llama-3.2-11b-vision-instruct +模型: nvidia/nemotron-3-nano-omni-30b-a3b-reasoning (Omni, 原生视频输入) 角色: vision (视觉分析直出结构化 JSON) + 智能问答 SDK: openai (NIM 兼容 OpenAI API 规范) -限制: NIM 单次请求最多 1 张图 -> 适配器内部逐帧调用,再聚合成 frame_details -熔断器: 启用 + +视频模式 (analyze_video): 按关键帧时间点截取 ±1.5s 片段拼接集锦视频 +(片段左上角叠加原始时间戳),base64 后经 video_url 单次调用 — +模型看到动态画面而非静态帧,动作/轨迹识别显著优于逐帧图片。 + +图片模式 (analyze_frames): 逐帧 image_url 调用(无视频文件时的降级路径)。 +注意: nemotron-omni 是 reasoning 模型,max_tokens 需给足(reasoning 消耗 token)。 """ import os import base64 import json import re +import subprocess +import tempfile from typing import Dict, List, Optional from .base_adapter import BaseModelAdapter @@ -25,16 +32,20 @@ try: except ImportError: OpenAI = None +VIDEO_SEGMENT_PAD = 1.5 # 关键帧前后各截取秒数 +HIGHLIGHT_WIDTH = 640 # 集锦视频宽度(保持宽高比) + class NvidiaVisionAdapter(BaseModelAdapter): - """NVIDIA NIM 云端 VLM 适配器 (逐帧结构化 + 聚合; 文本问答)""" + """NVIDIA NIM 云端 VLM 适配器 (视频集锦单次调用; 逐帧降级; 文本问答)""" def __init__(self, config: dict): super().__init__("nvidia", config) - self.model_name = config.get('model_name', 'meta/llama-3.2-11b-vision-instruct') + self.model_name = config.get( + 'model_name', 'nvidia/nemotron-3-nano-omni-30b-a3b-reasoning') self.api_key = self._resolve_key(config.get('api_key', '')) self.base_url = config.get('base_url', 'https://integrate.api.nvidia.com/v1') - self.timeout = config.get('timeout', 20) + self.timeout = config.get('timeout', 120) cb_cfg = config.get('circuit_breaker', {}) self._cb = CircuitBreaker( threshold=cb_cfg.get('threshold', 3), @@ -67,7 +78,187 @@ class NvidiaVisionAdapter(BaseModelAdapter): return False # ------------------------------------------------------------------ - # 视觉分析:逐帧调用(NIM 限 1 图/请求),聚合为 frame_details + # 视频模式:集锦视频 + video_url 单次调用(主路径) + # ------------------------------------------------------------------ + @staticmethod + def _ts_to_seconds(ts: str) -> float: + """'HH:MM:SS' 或 'HH:MM:SS.mmm' -> 秒""" + parts = str(ts).strip().split(':') + try: + if len(parts) == 3: + return int(parts[0]) * 3600 + int(parts[1]) * 60 + float(parts[2]) + if len(parts) == 2: + return int(parts[0]) * 60 + float(parts[1]) + return float(ts) + except ValueError: + return -1.0 + + def _build_highlight_video(self, video_path: str, + frame_timestamps: List[str]) -> Optional[str]: + """按关键帧时间点截取 ±pad 秒片段,叠加时间戳后拼接集锦视频 + + 时间戳以 00-05-00 形式叠加(连字符避免 ffmpeg drawtext 冒号转义)。 + """ + 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 not clips: + return None + + out_path = os.path.join( + tempfile.mkdtemp(prefix='nim_highlight_'), 'highlight.mp4') + cmd = ['ffmpeg', '-y', '-loglevel', 'error'] + for start, _ in clips: + cmd += ['-ss', f'{start:.2f}', '-t', f'{VIDEO_SEGMENT_PAD * 2}', '-i', video_path] + parts = [] + for i, (_, label) in enumerate(clips): + parts.append( + f"[{i}:v]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]', + '-c:v', 'libx264', '-preset', 'veryfast', '-crf', '28', + '-an', out_path] + try: + subprocess.run(cmd, check=True, capture_output=True, timeout=120) + except subprocess.TimeoutExpired: + logger.warning("NVIDIA 集锦视频生成超时") + return None + except subprocess.CalledProcessError as e: + logger.warning(f"NVIDIA 集锦视频生成失败: {e.stderr.decode()[:200] if e.stderr else e}") + return None + size = os.path.getsize(out_path) + logger.info(f"NVIDIA 集锦视频生成: {len(clips)} 片段, {size // 1024}KB") + if size < 1024: + return None + return out_path + + def analyze_video(self, video_path: str, + frame_timestamps: List[str], + known_members_context: str) -> Optional[Dict]: + """原生视频输入分析: 集锦片段 -> video_url 单次调用""" + if self._cb.is_open(): + logger.warning("NVIDIA 熔断器 OPEN,跳过视频分析") + return None + if self._client is None: + logger.warning("NVIDIA 客户端未初始化,跳过视频分析") + return None + + highlight = self._build_highlight_video(video_path, frame_timestamps) + if not highlight: + logger.warning("NVIDIA 集锦视频不可用,降级逐帧模式") + return None + try: + with open(highlight, 'rb') as f: + b64 = base64.b64encode(f.read()).decode('utf-8') + except Exception as e: + logger.warning(f"NVIDIA 读取集锦视频失败: {e}") + return None + finally: + try: + os.remove(highlight) + os.rmdir(os.path.dirname(highlight)) + except OSError: + pass + + prompt = self._build_video_prompt(frame_timestamps, known_members_context) + try: + resp = self._client.chat.completions.create( + model=self.model_name, + messages=[{"role": "user", "content": [ + {"type": "text", "text": prompt}, + {"type": "video_url", "video_url": { + "url": f"data:video/mp4;base64,{b64}"}} + ]}], + temperature=0.2, + max_tokens=3072, + timeout=self.timeout + ) + content = resp.choices[0].message.content + if not content: + logger.warning("NVIDIA 视频分析返回空 content") + return None + data = self._parse_single_frame_json(content) + if not data or 'frame_details' not in data: + logger.warning(f"NVIDIA 视频 JSON 解析失败: {content[:150]}") + return None + frame_details = self._normalize_frame_details(data, frame_timestamps) + if not frame_details: + return None + self._cb.record_success() + logger.info(f"NVIDIA 视频分析完成,frame_details={len(frame_details)}") + result = {"frame_details": frame_details} + if data.get('global_summary'): + result['global_summary'] = str(data['global_summary']) + if data.get('entities_json'): + result['entities_json'] = data['entities_json'] + return result + except Exception as e: + self._cb.record_failure() + logger.warning(f"NVIDIA 视频分析异常: {e}") + return None + + def _normalize_frame_details(self, data: dict, + frame_timestamps: List[str]) -> List[Dict]: + """归一化模型输出的 frame_details,按已知时间戳对齐""" + details = [] + for i, item in enumerate(data.get('frame_details', []), 1): + if not isinstance(item, dict): + continue + ts = str(item.get('frame_timestamp', + frame_timestamps[i - 1] if i <= len(frame_timestamps) else '')) + details.append({ + "frame_index": i, + "frame_timestamp": ts, + "person": str(item.get('person', '无人')), + "action": str(item.get('action', '')), + "clothing": str(item.get('clothing', '')), + "is_attention_event": bool(item.get('is_attention_event', False)), + "source_providers": ["nvidia"], + }) + return details + + def _build_video_prompt(self, frame_timestamps: List[str], + known_members: str) -> str: + ts_list = '\n'.join(f' 片段{i}: 原始时间 {ts}' for i, ts in enumerate(frame_timestamps, 1)) + return f"""你是家庭监控视频分析助手。下面的视频是由一段长时间监控录像中抽取的片段集锦, +共 {len(frame_timestamps)} 个片段(每个约 3 秒),按顺序拼接。每个片段左上角叠加了 +原始时间戳(ts 后的 00-05-00 表示 00:05:00)。 + +片段时间对照: +{ts_list} + +只输出合法 JSON(不要 markdown、不要解释),结构如下: +{{ + "frame_details": [ + {{ + "frame_timestamp": "<片段原始时间>", + "person": "片段中的人物或'无人'", + "action": "片段中人物的动作(动态观察,如走动/跑动/坐下)", + "clothing": "衣着(颜色+类型)", + "is_attention_event": false + }} + ], + "global_summary": "整段录像的综合摘要", + "entities_json": [{{"person": "人物名或人物X", "action": "行为概括", "clothing": "衣着"}}] +}} + +规则: +1. 只描述客观画面,不猜测。 +2. 已知家庭成员(按特征匹配,匹配到用 real_name,否则用"人物X"): +{known_members or '(暂无已知成员)'} +3. is_attention_event:跌倒、危险、异常哭闹等需关注事件(没有则为 false)。 +4. 无人出现的片段 person 填"无人",action 填""。""" + + # ------------------------------------------------------------------ + # 图片模式:逐帧调用(降级路径,无视频文件时使用) # ------------------------------------------------------------------ def analyze_frames(self, frame_paths: List[str], frame_timestamps: List[str], @@ -184,9 +375,9 @@ class NvidiaVisionAdapter(BaseModelAdapter): 4. 没有人物出现的帧 person 填"无人",action 填""。""" # ------------------------------------------------------------------ - # 智能问答:纯文本 + # 智能问答:纯文本(reasoning 模型,max_tokens 需给足) # ------------------------------------------------------------------ - def chat(self, prompt: str, max_tokens: int = 512) -> Optional[str]: + def chat(self, prompt: str, max_tokens: int = 2048) -> Optional[str]: if self._client is None: logger.warning("NVIDIA 客户端未初始化,跳过问答") return None