Files
sentinel-home-ai/fam-edge/src/fam_edge/video_processor.py
ericwyuan 61db82cb9b refactor(fam-edge): 重构第一阶段 - 人物图片零额外调用 + 运行时稳定性 + 工程质量
人物图片功能重做: 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 <noreply@anthropic.com>
2026-08-22 00:22:57 +08:00

363 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
VideoProcessor - 整视频分析编排
流程(不再切片/抽帧):
1. 从 OracleDB 取当前 known_members_context已命名/合并的人物)
2. 按 vision_order 依次调适配器的 analyze_videoGemini 整视频 -> NVIDIA 整视频)
3. 首个成功结果 -> 归一化 -> 写 OracleDBvideos + events 表,含 person_appearances
4. 把本视频 people_mentioned 更新进 people 表(带 features 特征,供 person_service 合并)
降级: 全部视觉模型失败 -> 标记视频 failed不再本地融合
注: 不依赖 OpenCV/cv2。视频文件校验用 ffprobesubprocess不再生成帧 jpg。
"""
import os
import re
import json
import subprocess
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 _ffprobe_available() -> bool:
"""ffprobe 是否可用ffmpeg 套件自带)。"""
try:
r = subprocess.run(
['ffprobe', '-version'],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=5)
return r.returncode == 0
except (FileNotFoundError, subprocess.TimeoutExpired):
return False
except Exception:
return False
def validate_video(path: str) -> tuple:
"""校验视频文件是否为正常可解码视频。
返回 (ok: bool, error: str, meta: dict|None)
- meta: {fps, frames, duration_sec, width, height}
用 ffprobesubprocess查 stream 信息。无 ffprobe 时仅做大小检查
(与旧 cv2 缺失时行为一致,跳过深度校验)。
不校验会导致空/半成品文件浪费云端配额。
"""
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
if not _ffprobe_available():
return True, "", None # 无 ffprobe 时跳过深度校验(仅大小检查)
# -v error: 只报错;-show_entries: 只取需要的字段;-of json: JSON 输出
r = subprocess.run(
['ffprobe', '-v', 'error', '-show_entries',
'stream=codec_type,avg_frame_rate,nb_frames,duration,width,height',
'-of', 'json', path],
capture_output=True, text=True, timeout=30)
if r.returncode != 0:
return False, f"ffprobe_error: {r.stderr[:200]}", None
try:
data = json.loads(r.stdout or '{}')
except ValueError:
return False, "ffprobe_bad_json", None
streams = data.get('streams') or []
vstream = next((s for s in streams if s.get('codec_type') == 'video'), None)
if not vstream:
return False, "no_video_stream", None
# fps: avg_frame_rate 形如 "25/1" -> 25.0
fps = 0.0
avg_rate = vstream.get('avg_frame_rate', '0/1')
try:
num, den = avg_rate.split('/')
den_f = float(den or '1')
fps = float(num) / den_f if den_f else 0.0
except (ValueError, ZeroDivisionError):
fps = 0.0
frames = 0
try:
frames = int(vstream.get('nb_frames') or 0)
except (ValueError, TypeError):
frames = 0
duration = 0.0
try:
duration = float(vstream.get('duration') or 0)
except (ValueError, TypeError):
duration = 0.0
meta = {
"fps": round(fps, 2),
"frames": frames,
"duration_sec": round(duration, 1) if duration else (
round(frames / max(fps, 0.01), 1) if frames and fps else 0),
"width": int(vstream.get('width') or 0),
"height": int(vstream.get('height') or 0),
}
return True, "", meta
except subprocess.TimeoutExpired:
return False, "ffprobe_timeout", None
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.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.set_event_start_time(video_id, event_start)
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 = []
for ev in events:
abs_ts, _ = _parse_event_ts(ev.get('timestamp'), start_dt)
ev_people = [_clean_person(str(p)) for p in ev.get('people', []) if p]
# 透传 person_appearances含 uid/features/action清洗 uid 字符串
appearances = ev.get('person_appearances') or []
norm_appearances = []
for pa in appearances:
if not isinstance(pa, dict):
continue
uid = _clean_person(str(pa.get('uid', '')))
if not uid:
continue
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,
"description": str(ev.get('description', '')),
"people": ev_people,
"person_appearances": norm_appearances,
"is_attention_event": bool(ev.get('is_attention_event', False)),
})
# 清洗 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 ('无人', '')]
# mark_video_processed 会把 norm_events 里的 person_appearances 落到
# events.person_appearances_json供 person_service 聚合特征
event_ids = self.db.mark_video_processed(video_id, summary, norm_events, people, provider)
# 更新 people 表(标签级 + 特征:从该视频所有 person_appearances 收集每个 uid 的特征)
uid_features = {}
for ev in norm_events:
for pa in ev.get('person_appearances', []):
uid = pa.get('uid')
if not uid or uid in ('无人', ''):
continue
feats = pa.get('features') or {}
if uid not in uid_features:
uid_features[uid] = feats
else:
# 同一 uid 多次出现:合并非 unknown 字段(与 upsert_person 的合并一致)
merged = dict(uid_features[uid])
for k, v in feats.items():
v_str = str(v).strip() if v is not None else ''
if v_str and v_str.lower() != 'unknown':
merged[k] = v_str
elif k not in merged:
merged[k] = v_str or 'unknown'
uid_features[uid] = merged
for p in people:
if p and p not in ('无人', ''):
feats = uid_features.get(p)
if feats:
self.db.upsert_person(p, source='llm', features=feats, display_uid=p)
else:
self.db.upsert_person(p, source='llm')
logger.info(f"[video_id={video_id}] 已落库: summary={len(summary)}字, "
f"events={len(norm_events)}, people={people}, "
f"with_features={len(uid_features)}")