From 61db82cb9b3afe466cf9e7f4bd2c7f2d78aee2ce Mon Sep 17 00:00:00 2001 From: ericwyuan Date: Sat, 22 Aug 2026 00:22:57 +0800 Subject: [PATCH] =?UTF-8?q?refactor(fam-edge):=20=E9=87=8D=E6=9E=84?= =?UTF-8?q?=E7=AC=AC=E4=B8=80=E9=98=B6=E6=AE=B5=20-=20=E4=BA=BA=E7=89=A9?= =?UTF-8?q?=E5=9B=BE=E7=89=87=E9=9B=B6=E9=A2=9D=E5=A4=96=E8=B0=83=E7=94=A8?= =?UTF-8?q?=20+=20=E8=BF=90=E8=A1=8C=E6=97=B6=E7=A8=B3=E5=AE=9A=E6=80=A7?= =?UTF-8?q?=20+=20=E5=B7=A5=E7=A8=8B=E8=B4=A8=E9=87=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 人物图片功能重做: bbox 随核心视频分析那一次 Gemini 调用一并产出(prompts.py 加 person_appearances.bbox 字段, [ymin,xmin,ymax,xmax] 0-1000 归一化), frame_service 直接用存好的 bbox 裁剪头像/事件缩略图, 删除原来"展示时额外调用 Gemini 定位人物"的 整套逻辑(locate_person_bbox/VLM 校验/熔断), 从架构上消除与核心视频分析共抢配额的 问题; 用真实数据验证裁剪结果正确框住人物本体。 NVIDIA 模型修复: 实测原配置的 3 个模型均不可用(asset_id 引用 500/400, 不支持视频), 改用 nemotron-3-nano-omni 的 base64 内嵌视频方式(唯一实测打通), 加 max_base64_mb 防止对大文件做注定失败的编码。 Gemini 多 Key 轮换: 支持 extra_api_keys 配置多个独立项目的 key, 配额用尽时依次 换 key 重试(每换 key 需重新上传, Files API 按项目隔离)。 稳定性加固: CircuitBreaker HALF_OPEN 清空旧失败计数(修复探测一失败就重新 OPEN 的 bug); chat() 统一接入熔断器(原来只有视频分析路径检查); NVIDIA 适配器改用共享 json_parser(原来自己重复实现且不做 schema 校验); Gemini Files API 上传超时也尝试 清理远程孤儿文件; video_processor/video_queue 里直接操作 OracleDB._conn 的裸 SQL 改走新增的 set_event_start_time/mark_video_invalid/reset_video_to_pending 方法; /health 加入队列线程存活状态; 密钥改用 ${ENV_VAR} 引用(.env 已支持自动加载), 不再明文写入 config.yaml。 工程质量: 新增 fam-edge/tests(32 个单元测试, 覆盖熔断器状态机/JSON 解析容错/ 时间戳解析/bbox 坐标换算/多 key 解析), 新增 scripts/smoke_test.py(发版前接口 稳定性检查); 清理死代码(OllamaAdapter.analyze_frames、get_sync_delta 死分支、 未使用的 vision_timeout/max_concurrent_tasks 配置项); 修正 get_events_for_label 排序(改最近优先 + 过滤畸形历史时间戳)。 已部署 Oracle 并跑通 smoke test 全部 6 项检查。 Co-Authored-By: Claude Sonnet 5 --- fam-edge/config/config.yaml | 41 ++-- fam-edge/config/config.yaml.example | 168 ++++++++----- fam-edge/scripts/smoke_test.py | 181 ++++++++++++++ .../fam_edge/ai_orchestrator/json_parser.py | 29 ++- .../src/fam_edge/ai_orchestrator/prompts.py | 6 +- .../src/fam_edge/api_gateway/api_gateway.py | 72 +++++- fam-edge/src/fam_edge/frame_service.py | 232 ++++++++++++++++++ .../model_adapters/adapter_factory.py | 4 +- .../model_adapters/circuit_breaker.py | 1 + .../fam_edge/model_adapters/gemini_adapter.py | 213 +++++++++------- .../fam_edge/model_adapters/nvidia_adapter.py | 172 +++++-------- .../fam_edge/model_adapters/ollama_adapter.py | 84 +------ fam-edge/src/fam_edge/oracle_db.py | 176 ++++++++++++- fam-edge/src/fam_edge/video_processor.py | 12 +- fam-edge/src/fam_edge/video_queue.py | 11 +- fam-edge/tests/conftest.py | 7 + fam-edge/tests/test_circuit_breaker.py | 55 +++++ fam-edge/tests/test_frame_service.py | 22 ++ fam-edge/tests/test_gemini_adapter.py | 40 +++ fam-edge/tests/test_json_parser.py | 87 +++++++ fam-edge/tests/test_video_processor.py | 58 +++++ 21 files changed, 1272 insertions(+), 399 deletions(-) create mode 100644 fam-edge/scripts/smoke_test.py create mode 100644 fam-edge/src/fam_edge/frame_service.py create mode 100644 fam-edge/tests/conftest.py create mode 100644 fam-edge/tests/test_circuit_breaker.py create mode 100644 fam-edge/tests/test_frame_service.py create mode 100644 fam-edge/tests/test_gemini_adapter.py create mode 100644 fam-edge/tests/test_json_parser.py create mode 100644 fam-edge/tests/test_video_processor.py diff --git a/fam-edge/config/config.yaml b/fam-edge/config/config.yaml index b479e80..20fc9b1 100644 --- a/fam-edge/config/config.yaml +++ b/fam-edge/config/config.yaml @@ -11,7 +11,6 @@ server: host: "0.0.0.0" port: 5000 - max_concurrent_tasks: 1 # Google 硬盘同步(rclone 负责同步落地,本段仅描述监听行为) gdrive_sync: @@ -26,9 +25,9 @@ gdrive_sync: oracle_db: path: "/opt/fam-edge/data/oracle.db" -# NAS 拉取同步接口鉴权 token(与 NAS oracle_sync.token 一致;明文直配,不再依赖 .env) +# NAS 拉取同步接口鉴权 token(与 NAS oracle_sync.token 一致,走 .env,不明文入库) sync_api: - token: "MLH92wv5jSDdQtHfcWJgKt-YaStn3IttjlrYxW0DwXA" + token: "${ORACLE_SYNC_TOKEN}" # 人物识别服务 person_service: @@ -56,7 +55,14 @@ models: model_name: "gemini-flash-latest" # 主模型(每日免费配额 20 请求,按模型独立) fallback_models: - "gemini-flash-lite-latest" - api_key: "AQ.Ab8RN6I0l8hC7hLnNHRY6qOXdch5CTWczDNlS4c1XrneGHipUQ" + api_key: "${GEMINI_API_KEY}" + # 多 Key 轮换(各自独立 Google Cloud 项目,配额互不影响):主 key 配额用尽时 + # 依次尝试这些 key,每个 key 都会重新走一遍模型 fallback 链。key 本身放 .env, + # 这里只放环境变量名,不直接写密钥。 + extra_api_keys: + - "${GEMINI_API_KEY_2}" + - "${GEMINI_API_KEY_3}" + - "${GEMINI_API_KEY_4}" timeout: 600 # 模型级独立超时(最终值,不参与编排层 ×2 放大) # gemini-flash-lite 实测 ~22-34s,按用户要求放宽至 8 分钟(480s),避免大视频/排队时过早切断 @@ -70,22 +76,27 @@ models: - provider: "nvidia" role: "vision" enabled: true - # 多模型降级链(实测记录 2026-08-21): - # omni 官方支持视频但 asset_id 引用 500;12b 400;llama-vision 不支持视频; - # cosmos/phi/gemma/kosmos/fuyu/paligemma 均 404 端点不可用。 - # 链机制保留(asset 上传一次,逐个尝试+间隔切换),可用模型出现时自动生效。 + # 模型可用性实测记录(2026-08-21,用真实短视频逐个探测 chat.completions 接口): + # omni(本行 model_name):video_url 只认 base64 data URI,NVCF asset_id 引用 + # 方式对它直接 500("Only base64 data URLs are supported for now")——本适配器 + # 已改为 base64 内嵌整段视频,见 max_base64_mb。确认可用(真实返回结构化 JSON)。 + # nemotron-nano-12b-v2-vl:走 asset_id 引用需要 NVCF-ASSET-DIR/ + # NVCF-FUNCTION-ASSET-IDS 请求头,这两个头的值由 NVCF 服务端按内部路径生成, + # 客户端传什么都 400 "Invalid NVCF-ASSET-DIR",标准 OpenAI 兼容调用打不通,已移除。 + # meta/llama-3.2-11b-vision-instruct:明确不支持视频输入 + # ("At most 0 video(s) may be provided"),只能单图,已移除。 + # base64 方案受请求体大小限制(实测约 25MB 上限),真实监控视频压缩后通常在 + # 20MB 上下,超过 max_base64_mb 直接跳过(不做注定失败的慢速编码),不是本地 + # 故意限制过窄——这是当前唯一能打通的 NVIDIA 视频理解路径。 model_name: "nvidia/nemotron-3-nano-omni-30b-a3b-reasoning" - fallback_models: - - "nvidia/nemotron-nano-12b-v2-vl" - - "meta/llama-3.2-11b-vision-instruct" + fallback_models: [] base_url: "https://integrate.api.nvidia.com/v1" - api_key: "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" + api_key: "${NVIDIA_API_KEY}" timeout: 600 - switch_interval_sec: 5 # 模型切换间隔:一个失败后等待再试下一个 + max_base64_mb: 20 # 超过此大小直接跳过 NVIDIA,不做注定失败的编码+上传 + switch_interval_sec: 5 # 模型切换间隔:一个失败后等待再试下一个(未来加模型时用) model_timeouts: # 模型级独立超时(最终值,不参与 ×2) "nvidia/nemotron-3-nano-omni-30b-a3b-reasoning": 300 - "nvidia/nemotron-nano-12b-v2-vl": 300 - "meta/llama-3.2-11b-vision-instruct": 120 circuit_breaker: enabled: true threshold: 5 diff --git a/fam-edge/config/config.yaml.example b/fam-edge/config/config.yaml.example index 6b58a2a..20fc9b1 100644 --- a/fam-edge/config/config.yaml.example +++ b/fam-edge/config/config.yaml.example @@ -1,83 +1,115 @@ -# FAM-Edge 配置文件 (Oracle 端) -# 复制此文件为 config.yaml 并修改实际值 +# FAM-Edge 配置文件 (Oracle 端) - 新架构 v2 +# +# 新架构(2026-08-21 重构): +# 1. 不再切片/抽帧:整视频直传云端 VLM(Gemini 用 Files API,NVIDIA 用整视频 video_url) +# 2. 视频来源:rclone 从 Google 硬盘实时同步到本地 local_dir,监听目录处理新视频 +# 3. Oracle 自建 SQLite 库存储所有视频摘要/事件/人物,并对外提供同步接口供 NAS 拉取 +# 4. 独立 person_service 汇总全量人物 -> LLM 合并为规范人物表 -> 回灌视频提示 +# 5. NAS 仅作管理后台,每 30 分钟从甲骨文拉增量镜像到本地 MariaDB -# NAS 端回调地址 -nas: - webhook_url: "http://100.x.x.10:8000/api/core/callback/event" - media_base_url: "http://100.x.x.10:8000/media" - media_token: "xxx" - -# Oracle 端服务 +# Oracle 端 HTTP 服务 server: host: "0.0.0.0" port: 5000 - max_concurrent_tasks: 1 -# 关键帧筛选参数(自适应:帧数随视频时长动态计算) -video: - candidate_per_minute: 2 # 每分钟粗抽候选帧数 - candidate_min: 30 # 候选帧下限(短视频保底) - candidate_max: 120 # 候选帧上限(超长视频截断) - key_frame_interval_sec: 150 # 关键帧间隔(秒),每2.5分钟1张 - min_key_frames: 5 # 关键帧下限(帧差不足时补足到此数) - max_key_frames_floor: 8 # 关键帧上限的下限(短视频保底) - max_key_frames_cap: 30 # 关键帧上限(超长视频截断) - mse_threshold: 500 # 帧差阈值 - jpeg_quality: 80 - max_long_edge: 1024 +# Google 硬盘同步(rclone 负责同步落地,本段仅描述监听行为) +gdrive_sync: + enabled: true + local_dir: "/opt/fam-edge/gdrive_videos" # rclone 同步落地目录(video_processing 监听此目录) + watch_interval_sec: 30 # 监听新视频的轮询间隔 + camera_name: "客厅" # 摄像头名称(注入视频提示) + # 文件名解析开始时间:监控文件名含时间戳时使用(如 2026-08-21_081500.mp4) + parse_start_from_filename: true -# 超时(秒) -timeout: - download: 60 - vlm_visual: 240 # 单模型视觉分析超时 - vlm_fusion: 120 - callback: 30 - overall: 600 +# Oracle 本地库(视频摘要/事件/人物) +oracle_db: + path: "/opt/fam-edge/data/oracle.db" -# 模型清单(可扩展,新增模型只需在此数组加一项 + 实现适配器) +# NAS 拉取同步接口鉴权 token(与 NAS oracle_sync.token 一致,走 .env,不明文入库) +sync_api: + token: "${ORACLE_SYNC_TOKEN}" + +# 人物识别服务 +person_service: + enabled: true + schedule_interval_sec: 1800 # 每 30 分钟重新汇总一次人物 + model: "gemini" # 用哪个模型做人物合并(vision 模型也支持纯文本) + +# 视频处理 +video_processing: + max_concurrent: 1 # 消费者线程数(串行处理,避免云端并发超额) + timeout: 900 # 兜底单视频分析超时 + timeout_multiplier: 2 # 模型消费超时倍数:在 models[i].timeout 原值上 ×2(大视频上传+分析耗时) + max_retries: 10 # 单视频失败最大重试次数(配额/过载等瞬时故障给足重试机会) + retry_interval_sec: 3600 # 失败重试最小间隔:距上次失败 ≥1h 才重新入队,等配额恢复 + file_validate: true # 登记入队前用 ffprobe 校验文件可解码;失败标记 invalid 不入队 + stable_window_sec: 60 # 文件 mtime 稳定窗口:写入中(rclone 同步未完成)的文件跳过本轮 + # 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed + vision_order: ["gemini", "nvidia"] + +# 智能问答降级链(与视频分析独立):Gemini -> NVIDIA -> 本地 Ollama models: - - provider: "ollama" - enabled: true - model_name: "llava-phi3" - base_url: "http://localhost:11434" - timeout: 240 - num_predict: 500 # 最大生成 token 数(ARM 上建议限制以控制延迟) - circuit_breaker: - enabled: false # 本地模型不启用熔断 - threshold: 5 - cooldown: 900 - - provider: "gemini" + role: "vision" enabled: true - model_name: "gemini-1.5-flash" - api_key: "${GEMINI_API_KEY}" # 从环境变量读取 - timeout: 8 + model_name: "gemini-flash-latest" # 主模型(每日免费配额 20 请求,按模型独立) + fallback_models: + - "gemini-flash-lite-latest" + api_key: "${GEMINI_API_KEY}" + # 多 Key 轮换(各自独立 Google Cloud 项目,配额互不影响):主 key 配额用尽时 + # 依次尝试这些 key,每个 key 都会重新走一遍模型 fallback 链。key 本身放 .env, + # 这里只放环境变量名,不直接写密钥。 + extra_api_keys: + - "${GEMINI_API_KEY_2}" + - "${GEMINI_API_KEY_3}" + - "${GEMINI_API_KEY_4}" + timeout: 600 + # 模型级独立超时(最终值,不参与编排层 ×2 放大) + # gemini-flash-lite 实测 ~22-34s,按用户要求放宽至 8 分钟(480s),避免大视频/排队时过早切断 + model_timeouts: + "gemini-flash-lite-latest": 480 circuit_breaker: enabled: true threshold: 5 - cooldown: 900 + cooldown: 300 - # v1.1 扩展示例(取消注释并填入 API Key 即启用) - # - provider: "openai" - # enabled: false - # model_name: "gpt-4o" - # api_key: "${OPENAI_API_KEY}" - # timeout: 30 - # circuit_breaker: - # enabled: true - # threshold: 5 - # cooldown: 900 + - provider: "nvidia" + role: "vision" + enabled: true + # 模型可用性实测记录(2026-08-21,用真实短视频逐个探测 chat.completions 接口): + # omni(本行 model_name):video_url 只认 base64 data URI,NVCF asset_id 引用 + # 方式对它直接 500("Only base64 data URLs are supported for now")——本适配器 + # 已改为 base64 内嵌整段视频,见 max_base64_mb。确认可用(真实返回结构化 JSON)。 + # nemotron-nano-12b-v2-vl:走 asset_id 引用需要 NVCF-ASSET-DIR/ + # NVCF-FUNCTION-ASSET-IDS 请求头,这两个头的值由 NVCF 服务端按内部路径生成, + # 客户端传什么都 400 "Invalid NVCF-ASSET-DIR",标准 OpenAI 兼容调用打不通,已移除。 + # meta/llama-3.2-11b-vision-instruct:明确不支持视频输入 + # ("At most 0 video(s) may be provided"),只能单图,已移除。 + # base64 方案受请求体大小限制(实测约 25MB 上限),真实监控视频压缩后通常在 + # 20MB 上下,超过 max_base64_mb 直接跳过(不做注定失败的慢速编码),不是本地 + # 故意限制过窄——这是当前唯一能打通的 NVIDIA 视频理解路径。 + model_name: "nvidia/nemotron-3-nano-omni-30b-a3b-reasoning" + fallback_models: [] + base_url: "https://integrate.api.nvidia.com/v1" + api_key: "${NVIDIA_API_KEY}" + timeout: 600 + max_base64_mb: 20 # 超过此大小直接跳过 NVIDIA,不做注定失败的编码+上传 + switch_interval_sec: 5 # 模型切换间隔:一个失败后等待再试下一个(未来加模型时用) + model_timeouts: # 模型级独立超时(最终值,不参与 ×2) + "nvidia/nemotron-3-nano-omni-30b-a3b-reasoning": 300 + circuit_breaker: + enabled: true + threshold: 5 + cooldown: 300 - # - provider: "nvidia" - # 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: 120 # reasoning 模型视频推理较慢,勿低于 90 - # circuit_breaker: - # enabled: true - # threshold: 5 - # cooldown: 900 + # 本地模型:纯文本 qwen2.5:7b,仅参与智能问答兜底 + - provider: "ollama" + role: "text" + usage: "qa_fallback" + enabled: true + model_name: "qwen2.5:7b" + base_url: "http://localhost:11434" + timeout: 120 + num_predict: 512 + circuit_breaker: + enabled: false diff --git a/fam-edge/scripts/smoke_test.py b/fam-edge/scripts/smoke_test.py new file mode 100644 index 0000000..a1bf3e1 --- /dev/null +++ b/fam-edge/scripts/smoke_test.py @@ -0,0 +1,181 @@ +#!/usr/bin/env python3 +""" +fam-edge 发版前稳定性检查 - 依次调用线上接口做基本断言,输出 pass/fail 清单。 + +不是 CI(项目没有 CI 基础设施),是给人在部署后手动跑的检查脚本。 + +用法: + python3 smoke_test.py --base-url http://127.0.0.1:5000 --token \ + [--video-id 173 --ts "2026-08-21 16:11:25"] [--label 人物A] [--skip-chat] + +video-id/ts/label 是可选的:不传就跳过 /api/oracle/frame 和 /api/oracle/avatar 检查 +(这两个接口需要库里真实存在的数据,不同部署环境的数据不一样,没法硬编码)。 +""" +import argparse +import sys +import time + +import requests + +PASS = "PASS" +FAIL = "FAIL" +SKIP = "SKIP" + + +class Report: + def __init__(self): + self.rows = [] + + def add(self, name, status, detail=""): + self.rows.append((name, status, detail)) + mark = {"PASS": "✓", "FAIL": "✗", "SKIP": "-"}[status] + print(f"[{mark}] {name}: {status}" + (f" ({detail})" if detail else "")) + + def ok(self): + return all(r[1] != FAIL for r in self.rows) + + +def check_health(report, base_url): + try: + t0 = time.time() + resp = requests.get(f"{base_url}/health", timeout=10) + dt = time.time() - t0 + if resp.status_code != 200: + report.add("GET /health", FAIL, f"HTTP {resp.status_code}") + return + data = resp.json() + if data.get("status") != "ok": + report.add("GET /health", FAIL, f"status={data.get('status')}") + return + report.add("GET /health", PASS, f"{dt:.2f}s, processed_videos={data.get('processed_videos')}") + except Exception as e: + report.add("GET /health", FAIL, str(e)) + + +def check_sync(report, base_url, token): + try: + t0 = time.time() + resp = requests.get(f"{base_url}/api/oracle/sync", + params={"since": "", "token": token}, timeout=30) + dt = time.time() - t0 + if resp.status_code != 200: + report.add("GET /api/oracle/sync", FAIL, f"HTTP {resp.status_code}: {resp.text[:150]}") + return + data = resp.json() + missing = [k for k in ("videos", "events", "people", "server_time") if k not in data] + if missing: + report.add("GET /api/oracle/sync", FAIL, f"缺字段: {missing}") + return + report.add("GET /api/oracle/sync", PASS, + f"{dt:.2f}s, videos={len(data['videos'])} events={len(data['events'])}") + except Exception as e: + report.add("GET /api/oracle/sync", FAIL, str(e)) + + +def check_activity(report, base_url, token): + try: + t0 = time.time() + resp = requests.get(f"{base_url}/api/oracle/activity", params={"token": token}, timeout=15) + dt = time.time() - t0 + if resp.status_code != 200: + report.add("GET /api/oracle/activity", FAIL, f"HTTP {resp.status_code}") + return + data = resp.json() + if "queue" not in data or "db" not in data: + report.add("GET /api/oracle/activity", FAIL, "缺 queue/db 字段") + return + report.add("GET /api/oracle/activity", PASS, f"{dt:.2f}s, queue={data.get('queue')}") + except Exception as e: + report.add("GET /api/oracle/activity", FAIL, str(e)) + + +def check_frame(report, base_url, token, video_id, ts): + if not video_id or not ts: + report.add("GET /api/oracle/frame", SKIP, "未提供 --video-id/--ts") + return + try: + t0 = time.time() + resp = requests.get(f"{base_url}/api/oracle/frame", + params={"video_id": video_id, "ts": ts, "token": token, "w": 440}, + timeout=60) + dt = time.time() - t0 + if resp.status_code != 200: + report.add("GET /api/oracle/frame", FAIL, f"HTTP {resp.status_code}") + return + if resp.headers.get("content-type") != "image/jpeg" or len(resp.content) < 100: + report.add("GET /api/oracle/frame", FAIL, "返回内容不是合法 JPEG") + return + report.add("GET /api/oracle/frame", PASS, f"{dt:.2f}s, {len(resp.content)} bytes") + except Exception as e: + report.add("GET /api/oracle/frame", FAIL, str(e)) + + +def check_avatar(report, base_url, token, label): + if not label: + report.add("GET /api/oracle/avatar", SKIP, "未提供 --label") + return + try: + t0 = time.time() + resp = requests.get(f"{base_url}/api/oracle/avatar", + params={"label": label, "token": token, "w": 160}, timeout=60) + dt = time.time() - t0 + if resp.status_code != 200: + report.add("GET /api/oracle/avatar", FAIL, f"HTTP {resp.status_code}") + return + if resp.headers.get("content-type") != "image/jpeg" or len(resp.content) < 100: + report.add("GET /api/oracle/avatar", FAIL, "返回内容不是合法 JPEG") + return + report.add("GET /api/oracle/avatar", PASS, f"{dt:.2f}s, {len(resp.content)} bytes") + except Exception as e: + report.add("GET /api/oracle/avatar", FAIL, str(e)) + + +def check_chat(report, base_url): + try: + t0 = time.time() + resp = requests.post(f"{base_url}/api/edge/chat/ask", + json={"prompt": "只回复两个字:在线", "max_tokens": 30}, timeout=120) + dt = time.time() - t0 + if resp.status_code != 200: + report.add("POST /api/edge/chat/ask", FAIL, f"HTTP {resp.status_code}: {resp.text[:150]}") + return + data = resp.json() + if not data.get("answer"): + report.add("POST /api/edge/chat/ask", FAIL, "answer 为空") + return + report.add("POST /api/edge/chat/ask", PASS, f"{dt:.2f}s, provider={data.get('provider')}") + except Exception as e: + report.add("POST /api/edge/chat/ask", FAIL, str(e)) + + +def main(): + p = argparse.ArgumentParser(description=__doc__) + p.add_argument("--base-url", default="http://127.0.0.1:5000") + p.add_argument("--token", required=True, help="ORACLE_SYNC_TOKEN") + p.add_argument("--video-id", type=int, default=None) + p.add_argument("--ts", default=None, help="例如 '2026-08-21 16:11:25'") + p.add_argument("--label", default=None, help="例如 人物A") + p.add_argument("--skip-chat", action="store_true", help="跳过问答检查(耗时最长且可能扣配额)") + args = p.parse_args() + + report = Report() + check_health(report, args.base_url) + check_sync(report, args.base_url, args.token) + check_activity(report, args.base_url, args.token) + check_frame(report, args.base_url, args.token, args.video_id, args.ts) + check_avatar(report, args.base_url, args.token, args.label) + if not args.skip_chat: + check_chat(report, args.base_url) + else: + report.add("POST /api/edge/chat/ask", SKIP, "--skip-chat") + + print() + n_pass = sum(1 for r in report.rows if r[1] == PASS) + n_fail = sum(1 for r in report.rows if r[1] == FAIL) + n_skip = sum(1 for r in report.rows if r[1] == SKIP) + print(f"结果: {n_pass} PASS / {n_fail} FAIL / {n_skip} SKIP") + sys.exit(0 if report.ok() else 1) + + +if __name__ == "__main__": + main() diff --git a/fam-edge/src/fam_edge/ai_orchestrator/json_parser.py b/fam-edge/src/fam_edge/ai_orchestrator/json_parser.py index 3e439fd..acb9631 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/json_parser.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/json_parser.py @@ -45,6 +45,30 @@ def parse_vlm_json(raw: str) -> dict: raise VLMOutputInvalidError(f"无法从 VLM 输出中解析 JSON: {raw[:200]}") +def _clean_bbox(bbox): + """校验 bbox 是否为合法的 4 元数值数组,非法/缺失一律归一为 None(下游按无 bbox 处理, + 退回整帧兜底,不会因为脏数据崩溃)。""" + if not isinstance(bbox, list) or len(bbox) != 4: + return None + try: + return [float(v) for v in bbox] + except (TypeError, ValueError): + return None + + +def _clean_person_appearances(appearances): + if not isinstance(appearances, list): + return appearances + cleaned = [] + for p in appearances: + if not isinstance(p, dict): + continue + p = dict(p) + p["bbox"] = _clean_bbox(p.get("bbox")) + cleaned.append(p) + return cleaned + + def validate_schema(data: dict) -> dict: """Schema 校验 + 脏数据清洗 @@ -89,8 +113,9 @@ def validate_schema(data: dict) -> dict: "timestamp": str(ev["timestamp"]), "description": str(ev["description"]), "people": [str(p) for p in people if p], - # 人物结构化特征(uid/features/action)——保留透传,供人物合并/特征卡 - "person_appearances": ev.get("person_appearances"), + # 人物结构化特征(uid/features/action/bbox)——保留透传,供人物合并/特征卡/ + # 事件缩略图与头像裁剪(bbox 随本次视频分析一次性产出,避免额外调用模型) + "person_appearances": _clean_person_appearances(ev.get("person_appearances")), "is_attention_event": bool(ev.get("is_attention_event", False)), }) data["events"] = cleaned diff --git a/fam-edge/src/fam_edge/ai_orchestrator/prompts.py b/fam-edge/src/fam_edge/ai_orchestrator/prompts.py index 8a45190..c725504 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/prompts.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/prompts.py @@ -58,7 +58,8 @@ def build_video_prompt(known_members: str, event_start_time: str, "face": "蓄须", "distinguishing": "左手戴手表" }}, - "action": "走向沙发坐下" + "action": "走向沙发坐下", + "bbox": [120, 340, 610, 900] }} ], "is_attention_event": false @@ -102,6 +103,9 @@ def build_video_prompt(known_members: str, event_start_time: str, * face: 面部特征(如 蓄须/戴眼镜/圆脸/unknown) * distinguishing: 辨识点(如 左手戴手表/右脸有痣/跛行/无) - action: 该人物在本时刻的动作(与 description 里该人物动作一致,单独抽出便于检索)。 + - bbox: 该人物在本帧画面中的包围框 [ymin,xmin,ymax,xmax],坐标为 0-1000 的归一化 + 整数(ymin/ymax 相对图片高度,xmin/xmax 相对图片宽度)——用于后续裁剪该人物的 + 缩略图/头像,不需要额外调用模型。看不清/无法定位时填 null,不要瞎猜坐标。 特征硬约束: * 客观描述可见特征,不猜测、不推断、不编造(看不清的字段写 unknown,不要靠常识猜性别/年龄)。 * 同一 uid 在视频多个 event 出现时,features 字段保持一致(衣着变了再如实更新 clothing, diff --git a/fam-edge/src/fam_edge/api_gateway/api_gateway.py b/fam-edge/src/fam_edge/api_gateway/api_gateway.py index 5df3c48..f4dae42 100644 --- a/fam-edge/src/fam_edge/api_gateway/api_gateway.py +++ b/fam-edge/src/fam_edge/api_gateway/api_gateway.py @@ -5,19 +5,18 @@ API-Gateway - Flask 蓝图(新架构 v3) GET /api/oracle/sync NAS 每 30 分钟拉取增量(since + token 校验) POST /api/oracle/people/correct NAS 推送手动命名校正(label -> canonical_name) POST /api/edge/chat/ask 智能问答编排(Gemini -> NVIDIA -> Ollama) - GET /api/oracle/activity 实时服务状态 + 最近活动流 - GET /health 健康检查 + GET /api/oracle/activity 实时服务状态 + 最近活动流 + GET /api/oracle/frame 事件时刻缩略帧(ffmpeg 抽帧 + 磁盘缓存) + GET /api/oracle/avatar 人物头像(按视频分析产出的 bbox 裁剪) + GET /health 健康检查(DB + 队列线程存活) -已移除(v3 去除帧图/avatar 依赖,改用大模型特征值): - /api/oracle/video//thumb, /api/oracle/event//thumb, - /api/oracle/person/avatar —— 不再生成 jpg,UI 读 sync_people.features_json - -已移除(旧推送/分块/队列模式): /video/push, /enqueue, /chunk, /assemble, -/results, /queue/stats, /mark_frames +图片能力(v4): 计算全部在 Oracle 本机 ffmpeg 抽帧/裁剪;人物 bbox 随视频分析那 +一次 Gemini 调用一并产出(见 ai_orchestrator/prompts.py),frame_service 不再 +额外调用任何模型。NAS 经 core 代理读取,不在 NAS 做图像计算。 """ import os -from flask import Blueprint, request, jsonify +from flask import Blueprint, request, jsonify, Response from ..logger import setup_logger from .. import state @@ -165,13 +164,64 @@ def activity(): }), 200 +@api_bp.route('/api/oracle/frame', methods=['GET']) +def oracle_frame(): + """事件时刻缩略帧:video_id + 绝对时间戳 -> jpeg(token 校验,磁盘缓存)""" + if not _check_token(): + return jsonify({"error": "unauthorized"}), 401 + video_id = request.args.get('video_id', type=int) + ts = request.args.get('ts', '') + w = request.args.get('w', default=400, type=int) + if not video_id or not ts: + return jsonify({"error": "缺少 video_id / ts"}), 400 + from ..frame_service import extract_frame + data = extract_frame(state.get_db(), video_id, ts, width=w) + if data is None: + return jsonify({"error": "抽帧失败"}), 404 + return Response(data, mimetype='image/jpeg') + + +@api_bp.route('/api/oracle/avatar', methods=['GET']) +def oracle_avatar(): + """人物头像:label/canonical -> jpeg(VLM 定位人物裁剪,磁盘缓存)""" + if not _check_token(): + return jsonify({"error": "unauthorized"}), 401 + label = request.args.get('label', '') + w = request.args.get('w', default=160, type=int) + if not label: + return jsonify({"error": "缺少 label"}), 400 + from ..frame_service import build_avatar + data = build_avatar(state.get_db(), label, width=w) + if data is None: + return jsonify({"error": "头像生成失败"}), 404 + return Response(data, mimetype='image/jpeg') + + @api_bp.route('/health', methods=['GET']) def health(): - """健康检查""" + """健康检查:DB 连通性 + producer/consumer 线程存活状态。 + + 之前只查一次 DB,队列线程全死了也会报 ok;现在把 VideoQueue.is_alive() 也 + 带上,让 /health 真正能反映处理管线是否还在跑。 + """ try: db = state.get_db() vids = db._conn.execute( "SELECT COUNT(*) c FROM videos WHERE status='done'").fetchone()['c'] - return jsonify({"status": "ok", "processed_videos": vids}), 200 except Exception as e: return jsonify({"status": "error", "error": str(e)}), 500 + + queue = state.get_queue() + queue_alive = None + if queue is not None: + try: + queue_alive = queue.is_alive() + except Exception as e: + logger.warning(f"queue.is_alive 异常: {e}") + + status = "ok" if (queue is None or queue_alive) else "degraded" + return jsonify({ + "status": status, + "processed_videos": vids, + "queue_alive": queue_alive, + }), 200 if status == "ok" else 503 diff --git a/fam-edge/src/fam_edge/frame_service.py b/fam-edge/src/fam_edge/frame_service.py new file mode 100644 index 0000000..8dafedc --- /dev/null +++ b/fam-edge/src/fam_edge/frame_service.py @@ -0,0 +1,232 @@ +""" +Frame-Service - 关键帧抽帧与人物头像裁剪(Oracle 端集中计算) + +新架构 v4(2026-08-21 重构): +- extract_frame: 按 video_id + 绝对时间戳,ffmpeg 精确抽帧,磁盘缓存 +- build_avatar: 用该人物候选事件里随视频分析一次性产出的 bbox(person_appearances + 里的 [ymin,xmin,ymax,xmax])直接裁剪,不再额外调用模型定位人物。 + +v3 曾经的做法是在展示缩略图/头像时另外调用一次 Gemini 做人物定位校验,这会和核心 +视频分析共用同一份 Gemini Key 抢配额(实测两边同时打 429),且当时的坐标解析本身 +也有 bug(错把 Gemini 原生 [ymin,xmin,ymax,xmax]/1000 当成 [x1,y1,x2,y2]/1)。v4 +把 bbox 改成随视频分析那一次 Gemini 调用一并产出(见 ai_orchestrator/prompts.py 的 +person_appearances.bbox 字段),frame_service 只做纯本地的抽帧/裁剪,不再对任何 +模型发起请求,从架构上消除配额争抢。 + +所有产物落到 CACHE_DIR,以 (video_id, offset) 或 label 为 key,避免重复计算。 +NAS 侧只负责代理与展示,不做任何图像计算。 +""" +import os +import shutil +import subprocess +import tempfile +from datetime import datetime + +from .logger import setup_logger + +logger = setup_logger('fam-edge.frame_service') + +try: + from .config_loader import load_config + _CFG = load_config().get('frame_service', {}) +except Exception: + _CFG = {} + +CACHE_DIR = _CFG.get('cache_dir', '/opt/fam-edge/frames_cache') +FFMPEG = shutil.which('ffmpeg') or 'ffmpeg' +FFPROBE = shutil.which('ffprobe') or 'ffprobe' +AVATAR_W = int(_CFG.get('avatar_width', 160)) +FRAME_W = int(_CFG.get('frame_width', 400)) + + +def _ensure_dir(): + os.makedirs(CACHE_DIR, exist_ok=True) + + +def _parse_dt(ts_str: str): + try: + return datetime.strptime(ts_str.strip()[:19], '%Y-%m-%d %H:%M:%S') + except (ValueError, TypeError): + return None + + +def _run_ffmpeg(args, timeout=60) -> bool: + try: + proc = subprocess.run([FFMPEG, '-hide_banner', '-loglevel', 'error', *args], + capture_output=True, timeout=timeout) + return proc.returncode == 0 + except (subprocess.TimeoutExpired, OSError): + return False + + +def _out_size(path): + """ffprobe 读图片宽高 -> (w, h) 或 None""" + try: + out = subprocess.run( + [FFPROBE, '-v', 'error', '-select_streams', 'v:0', + '-show_entries', 'stream=width,height', '-of', 'csv=s=x:p=0', path], + capture_output=True, text=True, timeout=20).stdout.strip() + w, h = out.split('x') + return int(w), int(h) + except Exception: + return None + + +def extract_frame(db, video_id: int, ts: str, width: int = FRAME_W) -> bytes: + """按 video_id + 绝对时间戳抽帧,返回 jpeg bytes(带磁盘缓存)""" + if not FFMPEG: + return None + row = db.get_video_by_id(video_id) + if not row: + return None + local_path = row['local_path'] + if not local_path or not os.path.isfile(local_path): + logger.warning(f"[video_id={video_id}] 视频文件不存在: {local_path}") + return None + start = _parse_dt(row['event_start_time'] or '') + t = _parse_dt(ts) + if t and start: + offset = max(0.0, (t - start).total_seconds()) + else: + offset = 0.0 + + _ensure_dir() + cache = os.path.join(CACHE_DIR, f"frame_{video_id}_{int(offset)}.jpg") + if os.path.isfile(cache) and os.path.getsize(cache) > 0: + with open(cache, 'rb') as f: + return f.read() + + fd, tmp = tempfile.mkstemp(suffix='.jpg', dir=CACHE_DIR) + os.close(fd) + try: + # 粗 seek(-i 前,关键帧快进)+ 精 seek(-i 后,逐帧解码): + # 纯输入端 seek 只能跳到最近关键帧,事件按 3s 密度打点时若 GOP 间隔 + # 大于 3s 会抽到别的关键帧,导致画面与描述对不上。 + coarse = max(0.0, offset - 5.0) + fine = offset - coarse + ok = _run_ffmpeg([ + '-ss', f'{coarse:.3f}', '-i', local_path, + '-ss', f'{fine:.3f}', + '-frames:v', '1', '-vf', f'scale={width}:-2', + '-q:v', '5', '-f', 'image2', '-y', tmp, + ], timeout=120) + if not ok or not os.path.isfile(tmp) or os.path.getsize(tmp) == 0: + return None + with open(tmp, 'rb') as f: + data = f.read() + os.replace(tmp, cache) # 原子落缓存 + return data + finally: + if os.path.exists(tmp): + try: + os.remove(tmp) + except OSError: + pass + + +def _bbox_to_pixels(bbox, width, height): + """bbox 为 [ymin,xmin,ymax,xmax],0-1000 归一化 -> 像素 (x1,y1,x2,y2)。""" + ymin, xmin, ymax, xmax = bbox + x1, y1 = xmin / 1000.0 * width, ymin / 1000.0 * height + x2, y2 = xmax / 1000.0 * width, ymax / 1000.0 * height + return int(x1), int(y1), int(x2), int(y2) + + +def _crop_ffmpeg(img_path: str, bbox_px, target_w) -> bool: + """按像素 bbox 裁剪居中并缩小为正方形,覆盖 img_path。失败返回 False。""" + size = _out_size(img_path) + if not size: + return False + w, h = size + x1, y1, x2, y2 = bbox_px + x1, x2 = max(0, x1), min(w, x2) + y1, y2 = max(0, y1), min(h, y2) + if x2 - x1 <= 0 or y2 - y1 <= 0: + return False + px, py = int((x2 - x1) * 0.3), int((y2 - y1) * 0.3) + x1, y1 = max(0, x1 - px), max(0, y1 - py) + x2, y2 = min(w, x2 + px), min(h, y2 + py) + cw, ch = x2 - x1, y2 - y1 + if cw <= 0 or ch <= 0: + return False + return _apply_filter(img_path, f'crop={cw}:{ch}:{x1}:{y1},scale={target_w}:{target_w}') + + +def _center_square_ffmpeg(img_path: str, target_w) -> bool: + """整帧居中正方形裁剪缩小,当人物没有 bbox 时兜底""" + size = _out_size(img_path) + if not size: + return False + w, h = size + side = min(w, h) + cx, cy = (w - side) // 2, (h - side) // 2 + return _apply_filter(img_path, f'crop={side}:{side}:{cx}:{cy},scale={target_w}:{target_w}') + + +def _apply_filter(img_path: str, vf: str) -> bool: + fd, tmp = tempfile.mkstemp(suffix='.jpg', dir=CACHE_DIR) + os.close(fd) + try: + if not _run_ffmpeg(['-i', img_path, '-vf', vf, + '-q:v', '5', '-frames:v', '1', '-f', 'image2', '-y', tmp]): + return False + os.replace(tmp, img_path) + return True + finally: + if os.path.exists(tmp): + try: + os.remove(tmp) + except OSError: + pass + + +def build_avatar(db, label: str, width: int = AVATAR_W) -> bytes: + """为人物构建头像:用候选事件里已经随视频分析产出的 bbox 直接裁剪; + 没有任何候选事件带 bbox 时,回退整帧居中裁剪。零额外模型调用。 + cache key = avatar_{label}。""" + _ensure_dir() + cache = os.path.join(CACHE_DIR, f"avatar_{label}.jpg") + if os.path.isfile(cache) and os.path.getsize(cache) > 0: + with open(cache, 'rb') as f: + return f.read() + + events = db.get_events_for_label(label, limit=6) + if not events: + return None + + fd, tmp = tempfile.mkstemp(suffix='.jpg', dir=CACHE_DIR) + os.close(fd) + try: + crop_data = None + first_good = None + for ev in events: + src = extract_frame(db, ev['video_id'], ev['ts'], width=600) + if src is None: + continue + with open(tmp, 'wb') as f: + f.write(src) + if first_good is None: + first_good = os.path.getsize(tmp) > 0 + bbox = ev.get('bbox') + if bbox: + size = _out_size(tmp) + if size and _crop_ffmpeg(tmp, _bbox_to_pixels(bbox, *size), width): + crop_data = ('bbox', os.path.getsize(tmp)) + break + # 兜底:首张可用的帧整帧居中(候选事件都没有 bbox 时) + if crop_data is None: + if first_good and _center_square_ffmpeg(tmp, width): + crop_data = ('fallback', os.path.getsize(tmp)) + if crop_data is None: + return None + with open(tmp, 'rb') as f: + data = f.read() + if data: + os.replace(tmp, cache) + return data + finally: + if os.path.exists(tmp): + try: + os.remove(tmp) + except OSError: + pass diff --git a/fam-edge/src/fam_edge/model_adapters/adapter_factory.py b/fam-edge/src/fam_edge/model_adapters/adapter_factory.py index 7a71c93..7134d01 100644 --- a/fam-edge/src/fam_edge/model_adapters/adapter_factory.py +++ b/fam-edge/src/fam_edge/model_adapters/adapter_factory.py @@ -48,6 +48,8 @@ def build_adapters(configs: List[dict]) -> List[BaseModelAdapter]: def register_adapter(provider_name: str, adapter_cls): - """注册新适配器(供扩展使用)""" + """注册新适配器(扩展点,供插件式新增 provider 用,无需改这个文件本身。 + 当前没有调用方——新模型目前都是直接改 _ADAPTER_REGISTRY,保留此函数是为了 + 以后接入第三方/可插拔适配器时不用再改工厂代码)。""" _ADAPTER_REGISTRY[provider_name] = adapter_cls logger.info(f"适配器已注册: {provider_name}") diff --git a/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py b/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py index b494bce..a62d495 100644 --- a/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py +++ b/fam-edge/src/fam_edge/model_adapters/circuit_breaker.py @@ -46,6 +46,7 @@ class CircuitBreaker: with self._lock: if self.state == 'OPEN' and self.last_failure and time.time() - self.last_failure > self.cooldown: self.state = 'HALF_OPEN' + self.failures.clear() # 清空旧失败记录,避免探测调用一失败就被旧记录凑数重新 OPEN return self.state == 'OPEN' def __repr__(self): diff --git a/fam-edge/src/fam_edge/model_adapters/gemini_adapter.py b/fam-edge/src/fam_edge/model_adapters/gemini_adapter.py index 12f9f35..a912850 100644 --- a/fam-edge/src/fam_edge/model_adapters/gemini_adapter.py +++ b/fam-edge/src/fam_edge/model_adapters/gemini_adapter.py @@ -8,6 +8,12 @@ provider_name = "gemini" 熔断器: 启用 整视频分析: 用 Files API 上传完整视频 -> generateContent 直出结构化 JSON (本地不切片、不抽帧;Gemini 原生支持长视频) + +多 Key 轮换(2026-08-21 新增): 不同 Google Cloud 项目的 API Key 各自独立计费/配额, +config 的 `api_key` 为主 Key,`extra_api_keys` 可以再配多个(各自项目的 Key)。 +outer loop 按 Key 顺序尝试,inner loop 才是原有的模型 fallback 链——因为 Gemini +Files API 上传的文件只能被同一个 Key/项目引用,换 Key 必须重新上传,所以每个 Key +都要重走一遍"上传 -> 模型链尝试 -> 删除",不是简单地在同一次上传后换 key 调用。 """ import os import time @@ -34,7 +40,16 @@ class GeminiAdapter(BaseModelAdapter): self.model_name = config.get('model_name', 'gemini-flash-latest') self.model_chain = [self.model_name] + [ m for m in config.get('fallback_models', []) if m and m != self.model_name] - self.api_key = self._resolve_key(config.get('api_key', '')) + # 多 Key 轮换:主 key + extra_api_keys(各自独立项目/配额),去重保序 + raw_keys = [config.get('api_key', '')] + list(config.get('extra_api_keys', []) or []) + seen = set() + self.api_keys = [] + for k in raw_keys: + resolved = self._resolve_key(k) + if resolved and resolved not in seen: + seen.add(resolved) + self.api_keys.append(resolved) + self.api_key = self.api_keys[0] if self.api_keys else '' # 向后兼容单 key 用法 self.timeout = config.get('timeout', 600) # 模型级独立超时(最终值,不参与编排层 ×N 放大): {model_name: seconds} # 例: {"gemini-flash-lite-latest": 90}(按实测耗时 ×4 配置) @@ -84,58 +99,68 @@ class GeminiAdapter(BaseModelAdapter): if self._cb.is_open(): logger.warning("Gemini 熔断器 OPEN,跳过视频分析") return None - if not self.api_key: + if not self.api_keys: logger.warning("Gemini API Key 未配置,跳过视频分析") return None if not os.path.isfile(video_path): logger.warning(f"Gemini 视频文件不存在: {video_path}") return None - file_uri = self._upload_file(video_path) - if not file_uri: - self._cb.record_failure() - return None - prompt = self._build_video_prompt(known_members_context, event_start_time) - try: - # 3 秒密度 + person_appearances 特征使输出 JSON 较大,max_tokens 需足够高防截断 - text = self._generate_video(file_uri, prompt, max_tokens=16384, temperature=0.2) - if text is None: - self._cb.record_failure() - return None + last_err = "no_key_available" + for idx, key in enumerate(self.api_keys): + file_uri, file_name = self._upload_file(video_path, key) + if not file_uri: + self._delete_file(file_name, key) # 即使未等到 ACTIVE,也尽力清理 + last_err = f"key{idx}_upload_failed" + continue try: - result = parse_vlm_json(text) - result = self._normalize(result) - if not result or 'events' not in result: - logger.error(f"Gemini 视频输出缺少 events: {text[:150]}") - self._cb.record_failure() - return None - result['compute_provider'] = 'gemini' - self._cb.record_success() - logger.info(f"Gemini 整视频分析完成,events={len(result.get('events', []))}") - return result - except VLMOutputInvalidError as e: - logger.error(f"Gemini 视频输出无法解析为 JSON: {e}") - self._cb.record_failure() - return None - except requests.Timeout: - logger.warning(f"Gemini 视频分析超时 ({self.timeout}s)") - self._cb.record_failure() - return None - except Exception as e: - logger.error(f"Gemini 视频分析异常: {e}") - self._cb.record_failure() - return None - finally: - self._delete_file(file_uri) + # 3 秒密度 + person_appearances 特征使输出 JSON 较大,max_tokens 需足够高防截断 + text = self._generate_video(file_uri, prompt, key, + max_tokens=16384, temperature=0.2) + if text is None: + last_err = f"key{idx}_all_models_failed" + continue + try: + result = parse_vlm_json(text) + result = self._normalize(result) + if not result or 'events' not in result: + logger.error(f"Gemini 视频输出缺少 events: {text[:150]}") + last_err = f"key{idx}_missing_events" + continue + result['compute_provider'] = 'gemini' + self._cb.record_success() + logger.info(f"Gemini 整视频分析完成(key[{idx}])," + f"events={len(result.get('events', []))}") + return result + except VLMOutputInvalidError as e: + logger.error(f"Gemini 视频输出无法解析为 JSON: {e}") + last_err = f"key{idx}_invalid_json" + continue + except requests.Timeout: + logger.warning(f"Gemini key[{idx}] 视频分析超时 ({self.timeout}s)") + last_err = f"key{idx}_timeout" + except Exception as e: + logger.error(f"Gemini key[{idx}] 视频分析异常: {e}") + last_err = f"key{idx}_exception" + finally: + self._delete_file(file_name, key) + self._cb.record_failure() + logger.error(f"Gemini 全部 {len(self.api_keys)} 个 Key 均失败: {last_err}") + return None - def _upload_file(self, video_path: str) -> Optional[str]: - """用 Files API resumable 可续传协议上传完整视频,返回可引用 URI。""" + def _upload_file(self, video_path: str, api_key: str): + """用 Files API resumable 可续传协议上传完整视频(用指定 api_key 对应的项目)。 + + 返回 (uri, file_name) 元组:uri 在文件 ACTIVE 前不可用时为 None,但只要 Gemini + 端已经创建了文件(file_name 非空),就应该用 file_name 尝试清理,避免超时留下 + 孤儿文件(Google 端存储永远不会被删除)。 + """ name = os.path.basename(video_path) size = os.path.getsize(video_path) upload_timeout = max(self.timeout, 900) # 上传端点必须是 /upload/v1beta/files(/v1beta/files 只是元数据端点,不接受上传协议) - base = f"https://generativelanguage.googleapis.com/upload/v1beta/files?key={self.api_key}" + base = f"https://generativelanguage.googleapis.com/upload/v1beta/files?key={api_key}" # 1) 创建可续传上传会话 try: r0 = requests.post( @@ -153,14 +178,14 @@ class GeminiAdapter(BaseModelAdapter): ) except Exception as e: logger.error(f"Gemini 创建上传会话异常: {e}") - return None + return None, None if r0.status_code not in (200, 201): logger.warning(f"Gemini 创建上传会话失败 HTTP {r0.status_code}: {r0.text[:200]}") - return None + return None, None session_url = r0.headers.get('X-Goog-Upload-URL') if not session_url: logger.warning("Gemini 上传响应缺少 X-Goog-Upload-URL") - return None + return None, None # 2) 上传文件体(流式) try: with open(video_path, 'rb') as f: @@ -177,13 +202,13 @@ class GeminiAdapter(BaseModelAdapter): ) except requests.Timeout: logger.warning(f"Gemini 文件上传超时 ({upload_timeout}s)") - return None + return None, None except Exception as e: logger.error(f"Gemini 文件上传异常: {e}") - return None + return None, None if resp.status_code not in (200, 201): logger.warning(f"Gemini 文件上传失败 HTTP {resp.status_code}: {resp.text[:200]}") - return None + return None, None try: info = resp.json().get('file', {}) uri = info.get('uri') @@ -191,16 +216,16 @@ class GeminiAdapter(BaseModelAdapter): state = info.get('state') except (ValueError, KeyError): logger.warning("Gemini 文件上传响应解析失败") - return None + return None, None if not uri: - return None + return None, file_name # 等待 ACTIVE(大文件可能还在处理) if state != 'ACTIVE' and file_name: - uri = self._wait_active(file_name) - return uri + uri = self._wait_active(file_name, api_key) + return uri, file_name - def _wait_active(self, file_name: str, max_wait: int = 120) -> Optional[str]: - url = f"{self._base_url}/{file_name}?key={self.api_key}" + def _wait_active(self, file_name: str, api_key: str, max_wait: int = 120) -> Optional[str]: + url = f"{self._base_url}/{file_name}?key={api_key}" deadline = time.time() + max_wait while time.time() < deadline: try: @@ -215,16 +240,17 @@ class GeminiAdapter(BaseModelAdapter): logger.warning(f"Gemini 文件 {file_name} 未在 {max_wait}s 内 ACTIVE") return None - def _delete_file(self, file_uri: str): - if not file_uri or 'files/' not in file_uri: + def _delete_file(self, file_name: str, api_key: str): + """按 Files API 的 file_name(如 'files/abc123')删除远程文件,尽力而为。""" + if not file_name: return - name = file_uri.split('files/', 1)[-1] + name = file_name.split('files/', 1)[-1] if 'files/' in file_name else file_name try: - requests.delete(f"{self._base_url}/files/{name}?key={self.api_key}", timeout=15) + requests.delete(f"{self._base_url}/files/{name}?key={api_key}", timeout=15) except Exception: pass - def _generate_video(self, file_uri: str, prompt: str, + def _generate_video(self, file_uri: str, prompt: str, api_key: str, max_tokens: int, temperature: float) -> Optional[str]: """带模型 fallback 链的 generateContent(视频文件引用)调用。""" parts = [ @@ -241,7 +267,7 @@ class GeminiAdapter(BaseModelAdapter): t0 = time.time() try: resp = requests.post( - f"{self._base_url}/models/{model}:generateContent?key={self.api_key}", + f"{self._base_url}/models/{model}:generateContent?key={api_key}", json={"contents": [{"parts": parts}], "generationConfig": { "temperature": temperature, @@ -339,44 +365,53 @@ class GeminiAdapter(BaseModelAdapter): # 智能问答:纯文本 # ------------------------------------------------------------------ def chat(self, prompt: str, max_tokens: int = 512) -> Optional[str]: - if not self.api_key: + if self._cb.is_open(): + logger.warning("Gemini 熔断器 OPEN,跳过问答") + return None + if not self.api_keys: logger.warning("Gemini API Key 未配置,跳过问答") return None try: - return self._generate_text(prompt, max_tokens=max_tokens, temperature=0.3) + result = self._generate_text(prompt, max_tokens=max_tokens, temperature=0.3) except Exception as e: logger.error(f"Gemini 问答异常: {e}") - return None + result = None + if result: + self._cb.record_success() + else: + self._cb.record_failure() + return result def _generate_text(self, text: str, max_tokens: int, temperature: float) -> Optional[str]: - """纯文本 generateContent(复用模型 fallback 链)。""" - for model in self.model_chain: - try: - resp = requests.post( - f"{self._base_url}/models/{model}:generateContent?key={self.api_key}", - json={"contents": [{"parts": [{"text": text}]}], - "generationConfig": { - "temperature": temperature, - "maxOutputTokens": max_tokens}}, - timeout=self.timeout - ) - except requests.Timeout: - logger.warning(f"Gemini [{model}] 问答超时") - continue - except Exception as e: - logger.error(f"Gemini [{model}] 问答异常: {e}") - continue - if resp.status_code == 200: - cands = resp.json().get('candidates', []) - out = ''.join( - p.get('text', '') - for p in (cands[0].get('content', {}) if cands else {}).get('parts', []) - ).strip() if cands else '' - if out: - return out - elif resp.status_code == 429: - logger.warning(f"Gemini [{model}] 429,切换模型") - continue + """纯文本 generateContent,按 key 轮换 × 模型 fallback 链依次尝试。""" + for idx, api_key in enumerate(self.api_keys): + for model in self.model_chain: + try: + resp = requests.post( + f"{self._base_url}/models/{model}:generateContent?key={api_key}", + json={"contents": [{"parts": [{"text": text}]}], + "generationConfig": { + "temperature": temperature, + "maxOutputTokens": max_tokens}}, + timeout=self.timeout + ) + except requests.Timeout: + logger.warning(f"Gemini key[{idx}] [{model}] 问答超时") + continue + except Exception as e: + logger.error(f"Gemini key[{idx}] [{model}] 问答异常: {e}") + continue + if resp.status_code == 200: + cands = resp.json().get('candidates', []) + out = ''.join( + p.get('text', '') + for p in (cands[0].get('content', {}) if cands else {}).get('parts', []) + ).strip() if cands else '' + if out: + return out + elif resp.status_code == 429: + logger.warning(f"Gemini key[{idx}] [{model}] 429,切换下一模型/Key") + continue return None def get_timeout(self) -> int: 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 7b3a79c..6be8c8a 100644 --- a/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py +++ b/fam-edge/src/fam_edge/model_adapters/nvidia_adapter.py @@ -2,17 +2,29 @@ NvidiaVisionAdapter - NVIDIA NIM 云端 VLM 适配器 provider_name = "nvidia" -模型: nvidia/nemotron-nano-12b-v2-vl(NIM 官方支持整视频 video_url 输入,内部自行采样帧) +模型: nvidia/nemotron-3-nano-omni-30b-a3b-reasoning(唯一实测确认可用的视频理解模型) 角色: vision (整视频直出结构化 JSON) + 智能问答 SDK: openai (NIM 兼容 OpenAI API 规范) -整视频分析: 先经 NVIDIA Assets API 上传完整视频拿 asset_id,再以 video_url 引用单次调用 - —— 本地不切片、不抽帧(request payload 有 25MB 上限,base64 直塞不可行,必须走 Assets API) + +整视频分析实测结论(2026-08-21 用真实短视频逐个探测): + - nemotron-3-nano-omni-30b-a3b-reasoning: video_url 只认 base64 data URI + (`data:video/mp4;base64,<...>`),Assets API 的 asset_id 引用方式对它直接 500 + (报错 "Only base64 data URLs are supported for now")——所以本适配器不再走 + Assets API 上传,直接 base64 内嵌整段视频。 + - nemotron-nano-12b-v2-vl: 需要走 NVCF 函数调用协议本身的 NVCF-ASSET-DIR/ + NVCF-FUNCTION-ASSET-IDS 请求头,而这两个头的值是 NVCF 服务端按内部路径生成、 + 不是客户端能自己拼对的(实测传什么都 400 "Invalid NVCF-ASSET-DIR"),标准 + OpenAI 兼容 chat.completions 调用打不通,已从模型链移除。 + - meta/llama-3.2-11b-vision-instruct: 明确不支持视频输入("At most 0 video(s) + may be provided"),只能单图,已移除。 + base64 方案的代价是请求体大小受限(原实现注释称约 25MB 上限),所以本适配器会在 + 上传前检查文件大小,超过 `max_base64_mb`(默认 20MB)直接放弃,不做注定失败的 + 慢速编码+上传。真实监控视频压缩后通常在 20MB 上下,属于"够不到就正常降级到失败 + 重试",不是本地故意限制过窄。 """ +import base64 import os -import json -import re import time -import requests from datetime import datetime, timezone, timedelta from typing import Dict, List, Optional @@ -20,6 +32,7 @@ from .base_adapter import BaseModelAdapter from .circuit_breaker import CircuitBreaker from ..logger import setup_logger from ..ai_orchestrator.prompts import build_video_prompt +from ..ai_orchestrator.json_parser import parse_vlm_json, VLMOutputInvalidError from ..config_loader import load_config logger = setup_logger('fam-edge.nvidia_adapter') @@ -43,12 +56,13 @@ class NvidiaVisionAdapter(BaseModelAdapter): def __init__(self, config: dict): super().__init__("nvidia", config) self.model_name = config.get( - 'model_name', 'nvidia/nemotron-nano-12b-v2-vl') + 'model_name', 'nvidia/nemotron-3-nano-omni-30b-a3b-reasoning') self.model_chain = [self.model_name] + [ m for m in config.get('fallback_models', []) if m and m != self.model_name] 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', 600) + self.max_base64_mb = float(config.get('max_base64_mb', 20)) # 模型级独立超时(最终值,不参与编排层 ×N 放大): {model_name: seconds} self.model_timeouts = { str(k): int(v) for k, v in (config.get('model_timeouts') or {}).items()} @@ -86,63 +100,8 @@ class NvidiaVisionAdapter(BaseModelAdapter): return False # ------------------------------------------------------------------ - # 整视频分析:Assets API 上传 -> video_url(asset_id) 单次调用 + # 整视频分析:base64 内嵌 video_url 单次调用(omni 只认 base64,不认 asset_id 引用) # ------------------------------------------------------------------ - ASSET_API = "https://api.nvcf.nvidia.com/v2/nvcf/assets" - - def _upload_asset(self, video_path: str) -> Optional[str]: - """用 NVIDIA Assets API 上传大视频文件,返回 asset_id 供 video_url 引用。""" - content_type = "video/mp4" - upload_timeout = max(self.timeout, 900) - try: - r = requests.post( - self.ASSET_API, - headers={ - "Authorization": f"Bearer {self.api_key}", - "Content-Type": "application/json", - }, - json={"contentType": content_type, "description": "fam-edge video asset"}, - timeout=60, - ) - except Exception as e: - logger.warning(f"NVIDIA 创建 asset 异常: {e}") - return None - if r.status_code not in (200, 201): - logger.warning(f"NVIDIA 创建 asset 失败 HTTP {r.status_code}: {r.text[:200]}") - return None - try: - j = r.json() - asset_id = j.get("assetId") - upload_url = j.get("uploadUrl") - except ValueError: - logger.warning("NVIDIA asset 响应解析失败") - return None - if not asset_id or not upload_url: - logger.warning("NVIDIA asset 响应缺少 assetId/uploadUrl") - return None - try: - with open(video_path, 'rb') as f: - up = requests.put( - upload_url, - data=f, - # 必须全小写 header 名且 content-type 值与 POST 的 contentType 一致: - # 预签名 S3 URL 签名覆盖这两个值,不一致会 SignatureDoesNotMatch - headers={"content-type": content_type, - "x-amz-meta-nvcf-asset-description": "fam-edge video asset"}, - timeout=upload_timeout, - ) - except requests.Timeout: - logger.warning(f"NVIDIA 上传 asset 超时 ({upload_timeout}s)") - return None - except Exception as e: - logger.warning(f"NVIDIA 上传 asset 异常: {e}") - return None - if up.status_code not in (200, 201): - logger.warning(f"NVIDIA 上传 asset 失败 HTTP {up.status_code}: {up.text[:200]}") - return None - logger.info(f"NVIDIA asset 上传成功: {asset_id}") - return asset_id - def analyze_video(self, video_path: str, known_members_context: str, event_start_time: str = '') -> Optional[Dict]: @@ -156,10 +115,18 @@ class NvidiaVisionAdapter(BaseModelAdapter): logger.warning(f"NVIDIA 视频文件不存在: {video_path}") return None - # asset 只上传一次,模型链内复用同一 assetId - asset_id = self._upload_asset(video_path) - if not asset_id: - self._cb.record_failure() + size_mb = os.path.getsize(video_path) / (1024 * 1024) + if size_mb > self.max_base64_mb: + logger.warning( + f"NVIDIA 视频 {size_mb:.1f}MB 超过 base64 上限 {self.max_base64_mb}MB," + "跳过(不做注定失败的慢速编码)") + return None + + try: + with open(video_path, 'rb') as f: + video_b64 = base64.b64encode(f.read()).decode() + except Exception as e: + logger.warning(f"NVIDIA 读取/编码视频失败: {e}") return None prompt = self._build_video_prompt(known_members_context, event_start_time) @@ -176,7 +143,7 @@ class NvidiaVisionAdapter(BaseModelAdapter): messages=[{"role": "user", "content": [ {"type": "text", "text": prompt}, {"type": "video_url", "video_url": { - "url": f"data:video/mp4;asset_id={asset_id}"}} + "url": f"data:video/mp4;base64,{video_b64}"}} ]}], temperature=0.2, max_tokens=16384, @@ -192,22 +159,19 @@ class NvidiaVisionAdapter(BaseModelAdapter): last_err = f"{model}_empty" self._sleep_switch(idx) continue - data = self._parse_json(content) - if not data or 'events' not in data: + try: + data = parse_vlm_json(content) + except VLMOutputInvalidError as e: self._emit_model_call(model, started, duration, False, "json_parse_failed") - logger.warning(f"NVIDIA [{model}] JSON 解析失败,切换下一模型: {content[:120]}") + logger.warning(f"NVIDIA [{model}] JSON 解析失败,切换下一模型: {e}") last_err = f"{model}_json" self._sleep_switch(idx) continue self._emit_model_call(model, started, duration, True) self._cb.record_success() logger.info(f"NVIDIA [{model}] 整视频分析完成,events={len(data.get('events', []))}") - return { - "global_summary": str(data.get('global_summary', '')), - "events": data.get('events', []), - "people_mentioned": data.get('people_mentioned', []), - "compute_provider": f"nvidia:{model}", - } + data['compute_provider'] = f"nvidia:{model}" + return data except Exception as e: duration = time.time() - t0 self._emit_model_call(model, started, duration, False, str(e)) @@ -224,27 +188,6 @@ class NvidiaVisionAdapter(BaseModelAdapter): logger.info(f"NVIDIA 等待 {self.switch_interval_sec}s 后切换下一模型") time.sleep(self.switch_interval_sec) - @staticmethod - def _parse_json(content: str) -> Optional[dict]: - content = content.strip() - try: - return json.loads(content) - except json.JSONDecodeError: - pass - fence = re.search(r'```(?:json)?\s*(\{.*?\})\s*```', content, re.DOTALL) - if fence: - try: - return json.loads(fence.group(1)) - except json.JSONDecodeError: - pass - brace = re.search(r'\{.*\}', content, re.DOTALL) - if brace: - try: - return json.loads(brace.group(0)) - except json.JSONDecodeError: - pass - return None - def _build_video_prompt(self, known_members: str, event_start_time: str) -> str: camera = load_config().get('gdrive_sync', {}).get('camera_name', '') return build_video_prompt(known_members, event_start_time, camera) @@ -253,22 +196,29 @@ class NvidiaVisionAdapter(BaseModelAdapter): # 智能问答:纯文本 # ------------------------------------------------------------------ def chat(self, prompt: str, max_tokens: int = 2048) -> Optional[str]: + if self._cb.is_open(): + logger.warning("NVIDIA 熔断器 OPEN,跳过问答") + return None if self._client is None: logger.warning("NVIDIA 客户端未初始化,跳过问答") return None - try: - resp = self._client.chat.completions.create( - model=self.model_name, - messages=[{"role": "user", "content": prompt}], - temperature=0.3, - max_tokens=max_tokens, - timeout=self.timeout - ) - content = resp.choices[0].message.content - return content.strip() if content else None - except Exception as e: - logger.warning(f"NVIDIA 问答异常: {e}") - return None + for model in self.model_chain: + try: + resp = self._client.chat.completions.create( + model=model, + messages=[{"role": "user", "content": prompt}], + temperature=0.3, + max_tokens=max_tokens, + timeout=self.timeout + ) + content = resp.choices[0].message.content + if content: + self._cb.record_success() + return content.strip() + except Exception as e: + logger.warning(f"NVIDIA [{model}] 问答异常: {e}") + self._cb.record_failure() + return None def get_timeout(self) -> int: return self.timeout diff --git a/fam-edge/src/fam_edge/model_adapters/ollama_adapter.py b/fam-edge/src/fam_edge/model_adapters/ollama_adapter.py index 5336907..2698c3a 100644 --- a/fam-edge/src/fam_edge/model_adapters/ollama_adapter.py +++ b/fam-edge/src/fam_edge/model_adapters/ollama_adapter.py @@ -7,9 +7,8 @@ provider_name = "ollama" 健康检查: GET /api/tags 不参与视觉分析、不参与视频结构化输出(云端 VLM 直出) """ -import base64 import requests -from typing import Dict, List, Optional +from typing import Dict, Optional from .base_adapter import BaseModelAdapter from .circuit_breaker import CircuitBreaker @@ -54,63 +53,6 @@ class OllamaAdapter(BaseModelAdapter): logger.error(f"Ollama 健康检查异常: {e}") return False - def analyze_frames(self, frame_paths: List[str], - frame_timestamps: List[str], - known_members_context: str) -> Optional[str]: - """调用 Ollama 视觉分析""" - if self._cb.is_open(): - logger.warning("Ollama 熔断器 OPEN,跳过调用") - return None - - # 构建 Prompt - n = len(frame_paths) - prompt = self._build_visual_prompt(n, frame_timestamps, known_members_context) - - # 读取图片并 Base64 编码 - images = [] - for path in frame_paths: - try: - with open(path, 'rb') as f: - images.append(base64.b64encode(f.read()).decode('utf-8')) - except Exception as e: - logger.error(f"读取图片失败 {path}: {e}") - - if not images: - logger.error("没有可用的图片帧") - return None - - try: - resp = requests.post( - f"{self.base_url}/api/generate", - json={ - "model": self.model_name, - "prompt": prompt, - "images": images, - "stream": False, - "options": {"temperature": 0.2, "top_p": 0.8, "num_predict": self.num_predict} - }, - timeout=self.timeout - ) - - if resp.status_code == 200: - output = resp.json().get('response', '') - self._cb.record_success() - logger.info(f"Ollama 视觉分析完成,输出长度={len(output)}") - return output - else: - logger.error(f"Ollama 调用失败: {resp.status_code} {resp.text[:200]}") - self._cb.record_failure() - return None - - except requests.Timeout: - logger.error(f"Ollama 调用超时 ({self.timeout}s)") - self._cb.record_failure() - return None - except Exception as e: - logger.error(f"Ollama 调用异常: {e}") - self._cb.record_failure() - return None - def analyze_video(self, video_path: str, known_members_context: str, event_start_time: str = '') -> Optional[Dict]: @@ -158,27 +100,3 @@ class OllamaAdapter(BaseModelAdapter): logger.error(f"Ollama 问答异常: {e}") self._cb.record_failure() return None - - def _build_visual_prompt(self, n: int, timestamps: List[str], known_members: str) -> str: - """构建视觉分析 Prompt""" - ts_lines = '\n'.join( - f"[图片{i+1}] 时间: {ts}" for i, ts in enumerate(timestamps) - ) - return f"""你是家庭监控视频分析助手。请按时间顺序描述下列 {n} 张图片中可见的内容,只描述客观画面,不要猜测或推测。 - -每张图片对应的时间戳如下: -{ts_lines} - -每张图片需报告: -1. 人物:数量、衣着(颜色+类型)、可见动作 -2. 物品:玩具、奶瓶、家具等显眼物品 -3. 互动:人与人、人与物品之间的互动 - -已知家庭成员清单(按特征匹配,匹配成功用 real_name,未匹配用"人物X"标识): -{known_members or '(暂无已知成员)'} - -输出格式(纯文本,每张图片一段,保留时间戳标记): -[图片1] 时间: {timestamps[0] if timestamps else ''} -内容: ... - -要求简洁、客观。不要输出 JSON,不要输出 markdown。""" diff --git a/fam-edge/src/fam_edge/oracle_db.py b/fam-edge/src/fam_edge/oracle_db.py index 789f05f..e1f09ef 100644 --- a/fam-edge/src/fam_edge/oracle_db.py +++ b/fam-edge/src/fam_edge/oracle_db.py @@ -16,6 +16,7 @@ Oracle 本地库(SQLite) - 视频摘要 / 事件 / 人物 存储 """ import os import json +import re import sqlite3 import threading from datetime import datetime, timezone, timedelta @@ -289,6 +290,170 @@ class OracleDB: (error, now, now, video_id)) self._conn.commit() + def set_event_start_time(self, video_id: int, event_start_time: str): + """回填从文件名解析出的视频开始时间(补录/纠偏用)。""" + self._conn.execute( + "UPDATE videos SET event_start_time=?, updated_at=? WHERE id=?", + (event_start_time, _now_iso(), video_id)) + self._conn.commit() + + def mark_video_invalid(self, video_id: int, error: str = ''): + """文件校验不通过(损坏/非视频等),标记 invalid,producer 不再重试。""" + now = _now_iso() + self._conn.execute( + "UPDATE videos SET status='invalid', file_valid=0, file_error=?, " + "updated_at=? WHERE id=?", + (error or '', now, video_id)) + self._conn.commit() + + def reset_video_to_pending(self, video_id: int): + """文件被重新同步覆盖(mtime 变化)时,清掉旧分析结果重新排队处理。""" + now = _now_iso() + self._conn.execute( + "UPDATE videos SET status='pending', retry_count=0, summary_json=NULL, " + "events_json=NULL, people_json=NULL, compute_provider=NULL, " + "processed_at=NULL, file_valid=1, updated_at=? WHERE id=?", + (now, video_id)) + self._conn.commit() + + def get_first_event_for_label(self, label: str): + """找到某人物(canonical_name 或 UID label)最早一次出现的事件。 + + 返回 dict{video_id, ts, features_text, event_start_time} 或 None。 + features_text 从该事件 person_appearances_json 中对应 uid 的特征拼出, + 供 frame_service 用大模型在画面中定位该人物。 + """ + # canonical_name -> 其下所有 label;否则按 label 本身匹配 + rows = self._conn.execute( + "SELECT label FROM people WHERE canonical_name=?", (label,)).fetchall() + labels = {r['label'] for r in rows} if rows else {label} + + best = None + for lb in labels: + pattern = f'%{lb}%' + row = self._conn.execute( + """SELECT e.video_id, e.ts, e.person_appearances_json, + v.event_start_time + FROM events e JOIN videos v ON v.id=e.video_id + WHERE v.status='done' + AND (e.person_list_json LIKE ? OR e.person_appearances_json LIKE ?) + ORDER BY e.ts ASC LIMIT 1""", + (pattern, pattern)).fetchone() + if row and row['video_id'] and \ + (best is None or (row['ts'] or '') < (best['ts'] or '')): + best = row + if not best: + return None + + features_text = '' + try: + pa = json.loads(best['person_appearances_json'] or '[]') + except (ValueError, TypeError): + pa = [] + if isinstance(pa, list): + for p in pa: + uid = str((p.get('uid') or '')).strip() + if uid and uid in labels and isinstance(p.get('features'), dict): + bits = [str(v) for v in p['features'].values() + if v and str(v).strip().lower() != 'unknown'] + if bits: + features_text = ','.join(bits) + break + return {'video_id': best['video_id'], 'ts': best['ts'], + 'features_text': features_text, + 'event_start_time': best['event_start_time']} + + def get_events_for_label(self, label: str, limit: int = 6): + """该人物(canonical_name 或 UID label)出现的候选事件,按时间倒序(最近优先)。 + + 返回 [{video_id, ts, features_text, bbox}](dict 列表):features_text 是该 + 事件中该人物的结构化特征文本;bbox 是视频分析时随该人物一并产出的包围框 + ([ymin,xmin,ymax,xmax],0-1000 归一化,取不到为 None)——frame_service 直接 + 用它做头像裁剪,不再额外调用模型定位。取最近的事件而不是最早的:bbox 是新 + 加的字段,老事件普遍没有,最近优先能更快用上新数据,也更能反映人物当前样貌。 + """ + rows = self._conn.execute( + "SELECT label FROM people WHERE canonical_name=?", (label,)).fetchall() + labels = {r['label'] for r in rows} | {label} + + seen = set() + out = [] + for lb in labels: + pattern = f'%{lb}%' + rs = self._conn.execute( + """SELECT e.video_id, e.ts, e.person_appearances_json, v.event_start_time + FROM events e JOIN videos v ON v.id=e.video_id + WHERE v.status='done' + AND (e.person_list_json LIKE ? OR e.person_appearances_json LIKE ?) + -- 排除历史遗留的畸形 ts(如缺日期的 "26:21"):这类值既不能 + -- 正确排序(字符串比较会排到最前面),extract_frame 也没法从 + -- 中算出正确偏移,只会抽到视频开头的错误画面 + AND e.ts GLOB '[0-9][0-9][0-9][0-9]-[0-9][0-9]-[0-9][0-9] [0-9][0-9]:[0-9][0-9]:[0-9][0-9]' + ORDER BY e.ts DESC LIMIT ?""", + (pattern, pattern, limit)).fetchall() + for r in rs: + key = (r['video_id'], r['ts']) + if key in seen: + continue + seen.add(key) + d = dict(r) + d['features_text'] = self._features_text_for( + d.get('person_appearances_json'), labels) + d['bbox'] = self._bbox_for_uids(d.get('person_appearances_json'), labels) + out.append(d) + out.sort(key=lambda r: (r['ts'] or ''), reverse=True) + return out[:limit] + + @staticmethod + def _bbox_for_uids(pa_json, uids): + """从 person_appearances_json 里取属于 uids 身份组那个人物的 bbox + ([ymin,xmin,ymax,xmax],0-1000 归一化)。bbox 随视频分析一次性产出, + 取不到/非法一律返回 None(调用方退回整帧兜底,不再额外调用模型定位)。""" + try: + pa = json.loads(pa_json or '[]') + except (ValueError, TypeError): + return None + if not isinstance(pa, list): + return None + for p in pa: + if not isinstance(p, dict): + continue + p_uid = re.sub(r'[((][^()()]*[))]', '', str(p.get('uid', ''))).strip() + if p_uid not in uids: + continue + bbox = p.get('bbox') + if isinstance(bbox, list) and len(bbox) == 4: + try: + return [float(v) for v in bbox] + except (TypeError, ValueError): + return None + return None + + @staticmethod + def _features_text_for(pa_json, uids) -> str: + """从 person_appearances_json 提取属于 uids 身份组的人物特征文本。 + + uid 先剥离括号再匹配(历史事件里存在 '人物A(别名:人物B)' 这类原始输出), + 只取该组人物的特征,避免把同帧其他人的特征混进头像定位 prompt。 + """ + try: + pa = json.loads(pa_json or '[]') + except (ValueError, TypeError): + return '' + if not isinstance(pa, list): + return '' + bits = [] + for p in pa: + if not isinstance(p, dict): + continue + uid = re.sub(r'[((][^()()]*[))]', '', str(p.get('uid', ''))).strip() + if uid not in uids or not isinstance(p.get('features'), dict): + continue + for v in p['features'].values(): + if v and str(v).strip().lower() != 'unknown': + bits.append(str(v)) + return ','.join(bits) + def get_all_videos(self) -> List[sqlite3.Row]: return self._conn.execute( "SELECT * FROM videos WHERE status='done' ORDER BY id ASC").fetchall() @@ -306,6 +471,9 @@ class OracleDB: 覆盖,新非 unknown 字段补齐)。None 时不更新特征列。 display_uid: 大模型给的人物 UID(如 "人物A")。label 本身就是 UID 时可省略。 """ + # 剥离括号后缀(如 '人物A(别名/标识:人物B)' -> '人物A'),防止大模型 + # 带备注的原始输出分裂出垃圾人物行 + label = re.sub(r'[((][^()()]*[))]', '', str(label)).strip() or str(label) now = _now_iso() row = self._conn.execute("SELECT * FROM people WHERE label=?", (label,)).fetchone() # 特征合并(在已有 features_json 基础上) @@ -363,11 +531,12 @@ class OracleDB: return json.dumps(merged, ensure_ascii=False) def set_canonical(self, label: str, canonical_name: str, source: str = 'manual'): - """手动命名:设置规范名(label 可视为别名)。""" - self.upsert_person(label, canonical_name, source='manual') + """设置规范名(label 可视为别名)。source 透传:llm 的可被后续纠正,manual 优先。""" + self.upsert_person(label, canonical_name, source=source) def set_person_appearances(self, label: str, count: int, source: str = 'llm'): """覆盖设置出现次数(reconcile 时用 distinct 视频数校准,避免累加膨胀)。""" + label = re.sub(r'[((][^()()]*[))]', '', str(label)).strip() or str(label) now = _now_iso() row = self._conn.execute("SELECT * FROM people WHERE label=?", (label,)).fetchone() if row: @@ -409,9 +578,8 @@ class OracleDB: videos = self._conn.execute( "SELECT * FROM videos WHERE updated_at > ? ORDER BY id ASC", (since_iso,) ).fetchall() + # events 表本身没有 updated_at 列,变更判断借用所属 video 的 updated_at events = self._conn.execute( - "SELECT * FROM events WHERE updated_at > ? ORDER BY id ASC", (since_iso,) - ).fetchall() if False else self._conn.execute( "SELECT e.* FROM events e JOIN videos v ON e.video_id=v.id " "WHERE v.updated_at > ? ORDER BY e.id ASC", (since_iso,)).fetchall() people = self._conn.execute( diff --git a/fam-edge/src/fam_edge/video_processor.py b/fam-edge/src/fam_edge/video_processor.py index dcb8b74..e1a862d 100644 --- a/fam-edge/src/fam_edge/video_processor.py +++ b/fam-edge/src/fam_edge/video_processor.py @@ -186,7 +186,6 @@ class VideoProcessor: 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( @@ -235,10 +234,7 @@ class VideoProcessor: 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() + self.db.set_event_start_time(video_id, event_start) known = self.db.get_known_members_context() logger.info(f"[video_id={video_id}] 开始整视频分析: {filename} " @@ -307,10 +303,16 @@ class VideoProcessor: feats = pa.get('features') or {} if not isinstance(feats, dict): feats = {} + bbox = pa.get('bbox') + if not (isinstance(bbox, list) and len(bbox) == 4): + bbox = None norm_appearances.append({ "uid": uid, "features": feats, "action": str(pa.get('action', '')), + # 该人物在本帧的包围框([ymin,xmin,ymax,xmax],0-1000 归一化), + # 供 frame_service 裁剪头像用,不再额外调用模型定位 + "bbox": bbox, }) norm_events.append({ "timestamp": abs_ts, diff --git a/fam-edge/src/fam_edge/video_queue.py b/fam-edge/src/fam_edge/video_queue.py index c667858..1ced703 100644 --- a/fam-edge/src/fam_edge/video_queue.py +++ b/fam-edge/src/fam_edge/video_queue.py @@ -91,9 +91,7 @@ class VideoQueue: if not ok: vid = self.db.ensure_video(fn, path, camera_name=self.camera_name) self.db.set_video_file_status(vid, False, verr) - self.db._conn.execute( - "UPDATE videos SET status='invalid' WHERE id=?", (vid,)) - self.db._conn.commit() + self.db.mark_video_invalid(vid, verr) logger.warning(f"文件校验失败,标记 invalid 不入队: {fn} ({verr})") continue if vmeta: @@ -113,12 +111,7 @@ class VideoQueue: # 文件被覆盖(rclone 重新同步/更新):重置 pending 重新分析 logger.info(f"文件内容变更,重置重新分析: {fn} (id={row['id']})") self.db.record_activity('queue', 'reanalyze', f"{fn} (id={row['id']})") - self.db._conn.execute( - "UPDATE videos SET status='pending', retry_count=0, summary_json=NULL, " - "events_json=NULL, people_json=NULL, compute_provider=NULL, " - "processed_at=NULL, file_valid=1 WHERE id=?", - (row['id'],)) - self.db._conn.commit() + self.db.reset_video_to_pending(row['id']) self._enqueue(row['id']) def _file_changed(self, row, path: str) -> bool: diff --git a/fam-edge/tests/conftest.py b/fam-edge/tests/conftest.py new file mode 100644 index 0000000..744b28a --- /dev/null +++ b/fam-edge/tests/conftest.py @@ -0,0 +1,7 @@ +import os +import sys + +# 让测试能直接 `from fam_edge.xxx import yyy`,无需先 pip install -e . +_SRC = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), 'src') +if _SRC not in sys.path: + sys.path.insert(0, _SRC) diff --git a/fam-edge/tests/test_circuit_breaker.py b/fam-edge/tests/test_circuit_breaker.py new file mode 100644 index 0000000..61d3dff --- /dev/null +++ b/fam-edge/tests/test_circuit_breaker.py @@ -0,0 +1,55 @@ +import time + +from fam_edge.model_adapters.circuit_breaker import CircuitBreaker + + +def test_closed_stays_closed_below_threshold(): + cb = CircuitBreaker(threshold=3, cooldown=1) + cb.record_failure() + cb.record_failure() + assert not cb.is_open() + assert cb.state == 'CLOSED' + + +def test_opens_at_threshold(): + cb = CircuitBreaker(threshold=3, cooldown=1) + for _ in range(3): + cb.record_failure() + assert cb.is_open() + assert cb.state == 'OPEN' + + +def test_half_open_after_cooldown_and_clears_stale_failures(): + """进入 HALF_OPEN 时应清空旧失败计数,探测调用失败一次不该立刻把它"凑数"重新 OPEN。""" + cb = CircuitBreaker(threshold=3, cooldown=0.05) + for _ in range(3): + cb.record_failure() + assert cb.is_open() + time.sleep(0.1) + # 冷却期已过,第一次 is_open() 调用把状态转为 HALF_OPEN 并放行一次探测 + assert cb.is_open() is False + assert cb.state == 'HALF_OPEN' + assert len(cb.failures) == 0 + # 探测失败一次:由于 deque 已清空,不应该只凭这一次失败就再次 OPEN + # (state 仍是 HALF_OPEN,不是 OPEN,is_open() 只在 state=='OPEN' 时为 True) + cb.record_failure() + assert cb.state != 'OPEN' + assert not cb.is_open() + + +def test_half_open_probe_success_closes(): + cb = CircuitBreaker(threshold=2, cooldown=0.05) + cb.record_failure() + cb.record_failure() + time.sleep(0.1) + assert cb.is_open() is False # 触发 HALF_OPEN 转换 + cb.record_success() + assert cb.state == 'CLOSED' + assert len(cb.failures) == 0 + + +def test_disabled_never_opens(): + cb = CircuitBreaker(threshold=1, cooldown=1, enabled=False) + cb.record_failure() + cb.record_failure() + assert not cb.is_open() diff --git a/fam-edge/tests/test_frame_service.py b/fam-edge/tests/test_frame_service.py new file mode 100644 index 0000000..ab1ef50 --- /dev/null +++ b/fam-edge/tests/test_frame_service.py @@ -0,0 +1,22 @@ +from fam_edge.frame_service import _bbox_to_pixels + + +def test_bbox_to_pixels_basic(): + # [ymin,xmin,ymax,xmax] 0-1000 归一化 -> 像素 (x1,y1,x2,y2) + # 实测样本:Gemini 对 440x248 帧返回 [0, 690, 203, 725], + # 对应画面右上角门厅处的一个人(今天用真实截图验证过)。 + x1, y1, x2, y2 = _bbox_to_pixels([0, 690, 203, 725], 440, 248) + assert x1 == int(690 / 1000 * 440) + assert y1 == 0 + assert x2 == int(725 / 1000 * 440) + assert y2 == int(203 / 1000 * 248) + + +def test_bbox_to_pixels_full_frame(): + x1, y1, x2, y2 = _bbox_to_pixels([0, 0, 1000, 1000], 400, 300) + assert (x1, y1, x2, y2) == (0, 0, 400, 300) + + +def test_bbox_to_pixels_zero_area(): + x1, y1, x2, y2 = _bbox_to_pixels([500, 500, 500, 500], 400, 300) + assert (x1, y1) == (x2, y2) diff --git a/fam-edge/tests/test_gemini_adapter.py b/fam-edge/tests/test_gemini_adapter.py new file mode 100644 index 0000000..57e8d06 --- /dev/null +++ b/fam-edge/tests/test_gemini_adapter.py @@ -0,0 +1,40 @@ +from fam_edge.model_adapters.gemini_adapter import GeminiAdapter + + +def _cfg(**overrides): + base = { + "provider": "gemini", + "model_name": "gemini-flash-latest", + "api_key": "key-primary", + "circuit_breaker": {"enabled": False}, + } + base.update(overrides) + return base + + +def test_single_key_backward_compat(): + a = GeminiAdapter(_cfg()) + assert a.api_keys == ["key-primary"] + assert a.api_key == "key-primary" + + +def test_extra_keys_appended_in_order(): + a = GeminiAdapter(_cfg(extra_api_keys=["key-2", "key-3", "key-4"])) + assert a.api_keys == ["key-primary", "key-2", "key-3", "key-4"] + + +def test_extra_keys_dedup_against_primary(): + a = GeminiAdapter(_cfg(extra_api_keys=["key-primary", "key-2"])) + assert a.api_keys == ["key-primary", "key-2"] + + +def test_missing_env_var_keys_are_dropped(): + """extra_api_keys 里未设置的 ${ENV_VAR} 解析为空字符串,不应该混进 api_keys 列表。""" + a = GeminiAdapter(_cfg(extra_api_keys=["${SOME_UNSET_GEMINI_KEY_VAR}", "key-2"])) + assert a.api_keys == ["key-primary", "key-2"] + + +def test_no_keys_at_all(): + a = GeminiAdapter(_cfg(api_key="")) + assert a.api_keys == [] + assert a.api_key == "" diff --git a/fam-edge/tests/test_json_parser.py b/fam-edge/tests/test_json_parser.py new file mode 100644 index 0000000..6d993c3 --- /dev/null +++ b/fam-edge/tests/test_json_parser.py @@ -0,0 +1,87 @@ +import pytest + +from fam_edge.ai_orchestrator.json_parser import parse_vlm_json, validate_schema, VLMOutputInvalidError + + +def _base(): + return { + "global_summary": "客厅监控摘要", + "events": [ + {"timestamp": "00:00:03", "description": "人物A走进客厅", + "people": ["人物A"], "is_attention_event": False, + "person_appearances": [ + {"uid": "人物A", "features": {"gender": "男"}, "action": "走动", + "bbox": [10, 20, 500, 400]} + ]} + ], + "people_mentioned": ["人物A"], + } + + +def test_direct_json_parses(): + import json + raw = json.dumps(_base(), ensure_ascii=False) + result = parse_vlm_json(raw) + assert result["global_summary"] == "客厅监控摘要" + assert len(result["events"]) == 1 + assert result["events"][0]["person_appearances"][0]["bbox"] == [10.0, 20.0, 500.0, 400.0] + + +def test_markdown_fence_extraction(): + import json + raw = f"这是模型的解释文字\n```json\n{json.dumps(_base(), ensure_ascii=False)}\n```\n谢谢" + result = parse_vlm_json(raw) + assert result["events"][0]["timestamp"] == "00:00:03" + + +def test_greedy_brace_extraction(): + import json + raw = f"废话前缀 {json.dumps(_base(), ensure_ascii=False)} 废话后缀" + result = parse_vlm_json(raw) + assert result["people_mentioned"] == ["人物A"] + + +def test_invalid_json_raises(): + with pytest.raises(VLMOutputInvalidError): + parse_vlm_json("这不是 JSON,也没有大括号") + + +def test_missing_required_field_raises(): + with pytest.raises(VLMOutputInvalidError): + validate_schema({"events": []}) + + +def test_frame_details_legacy_compat(): + data = { + "global_summary": "旧结构", + "frame_details": [ + {"frame_timestamp": "00:00:05", "action": "走动", "person": "人物A", + "is_attention_event": True} + ], + } + result = validate_schema(data) + assert len(result["events"]) == 1 + assert result["events"][0]["description"] == "走动" + assert result["events"][0]["people"] == ["人物A"] + assert result["events"][0]["is_attention_event"] is True + + +def test_bbox_missing_or_null_becomes_none(): + data = _base() + data["events"][0]["person_appearances"][0]["bbox"] = None + result = validate_schema(data) + assert result["events"][0]["person_appearances"][0]["bbox"] is None + + +def test_bbox_wrong_shape_becomes_none(): + data = _base() + data["events"][0]["person_appearances"][0]["bbox"] = [1, 2, 3] # 长度不对 + result = validate_schema(data) + assert result["events"][0]["person_appearances"][0]["bbox"] is None + + +def test_bbox_non_numeric_becomes_none(): + data = _base() + data["events"][0]["person_appearances"][0]["bbox"] = ["a", "b", "c", "d"] + result = validate_schema(data) + assert result["events"][0]["person_appearances"][0]["bbox"] is None diff --git a/fam-edge/tests/test_video_processor.py b/fam-edge/tests/test_video_processor.py new file mode 100644 index 0000000..7bac70a --- /dev/null +++ b/fam-edge/tests/test_video_processor.py @@ -0,0 +1,58 @@ +from datetime import datetime + +from fam_edge.video_processor import ( + _parse_event_start_from_filename, + _parse_event_ts, + _clean_person, +) + + +def test_parse_filename_pure_digit_format(): + assert _parse_event_start_from_filename( + "Generic_ONVIF-001-20260820-140416-1787205856321-7.mp4" + ) == "2026-08-20 14:04:16" + + +def test_parse_filename_underscore_date_format(): + assert _parse_event_start_from_filename("20260821_081500.mp4") == "2026-08-21 08:15:00" + + +def test_parse_filename_dashed_date_format(): + assert _parse_event_start_from_filename("2026-08-21_081500.mp4") == "2026-08-21 08:15:00" + + +def test_parse_filename_no_match_returns_empty(): + assert _parse_event_start_from_filename("客厅.mp4") == "" + + +def test_parse_event_ts_relative_offset(): + start = datetime(2026, 8, 21, 15, 53, 3) + abs_ts, offset = _parse_event_ts("00:18:22", start) + assert abs_ts == "2026-08-21 16:11:25" + assert offset == 18 * 60 + 22 + + +def test_parse_event_ts_relative_offset_no_start(): + abs_ts, offset = _parse_event_ts("00:01:23", None) + assert abs_ts == "00:01:23" + assert offset == 83.0 + + +def test_parse_event_ts_over_6_hours_falls_back_to_absolute(): + """相对时间 > 6 小时视为模型误输出绝对时间,不强行按偏移定位。""" + start = datetime(2026, 8, 21, 8, 0, 0) + abs_ts, offset = _parse_event_ts("2026-08-21 09:00:00", start) + assert abs_ts == "2026-08-21 09:00:00" + assert offset == 3600.0 + + +def test_clean_person_strips_fullwidth_parens(): + assert _clean_person("人物A(别名/标识:人物B)") == "人物A" + + +def test_clean_person_strips_ascii_parens(): + assert _clean_person("人物A(alias: 人物B)") == "人物A" + + +def test_clean_person_no_parens_unchanged(): + assert _clean_person("汤圆") == "汤圆"