[阶段2] FAM-Edge 重构为整视频分析+同步接口+人物服务 - 移除切片/抽帧/队列,新增 oracle_db/person_service/qa/watch_processor/video_processor,api_gateway 提供 /api/oracle/sync 与 /api/oracle/people/correct

This commit is contained in:
ericwyuan
2026-08-21 10:38:15 +08:00
parent ae1c13589f
commit 9b1cc8f93b
23 changed files with 1097 additions and 2335 deletions

View File

@@ -1,534 +0,0 @@
"""
AI-Orchestrator - 多模型编排
视频分析链路(新框架):
1. 加载所有启用的模型适配器
2. 健康检查
3. 抽帧 + 关键帧筛选 + 压缩
4. 云端 VLM 视觉分析Gemini 主 / NVIDIA 兜底),直出结构化 JSON
5. format_cloud_result对云端结果做**格式化/校验**(无本地模型调用,不汇总摘要)
6. 同步返回 NAS → 落库
智能问答链路(新框架):
- run_qaGemini → NVIDIA → 本地 Ollama仅当两云端都失败才用本地兜底
"""
import time
import json
import base64
import requests
from datetime import datetime, timedelta
from concurrent.futures import ThreadPoolExecutor, as_completed, TimeoutError as FuturesTimeout
from typing import Dict, List, Optional, Tuple
from ..logger import setup_logger, log_task
from ..config_loader import load_config
from ..model_adapters.adapter_factory import build_adapters
from ..model_adapters.base_adapter import BaseModelAdapter
from ..video_preprocessor.preprocessor import VideoPreprocessor
from .json_parser import VLMOutputInvalidError, validate_schema
logger = setup_logger('fam-edge.orchestrator')
class AIOrchestrator:
"""AI 编排器"""
def __init__(self):
self.config = load_config()
self.adapters: List[BaseModelAdapter] = build_adapters(self.config.get('models', []))
self.timeout_cfg = self.config.get('timeout', {})
def health_check_all(self) -> List[BaseModelAdapter]:
"""健康检查,返回健康的适配器列表"""
healthy = []
for adapter in self.adapters:
try:
if adapter.health_check():
healthy.append(adapter)
except Exception as e:
logger.error(f"适配器 {adapter.provider_name} 健康检查异常: {e}")
return healthy
def run_visual_analysis(self, adapters: List[BaseModelAdapter],
frame_paths: List[str],
frame_timestamps: List[str],
known_members_context: str,
rate_limiter=None,
video_path: str = None,
event_start_time: str = '') -> Dict[str, dict]:
"""视觉分析阶段:仅 role=vision 的适配器参与
支持 analyze_video 的适配器(如 NVIDIA Omni优先走原生视频输入
失败自动降级回逐帧图片模式。
orchestrator.mode:
- fallback (默认): 按 config 顺序依次尝试,首个成功即采用(单元素 dict
- ensemble: 并行所有健康 vision 模型,全部成功结果都保留(交叉验证)
rate_limiter: 可选 RateLimiter 实例,按 provider 限速2x burst
"""
vision_adapters = [a for a in adapters if getattr(a, 'role', 'vision') == 'vision']
if not vision_adapters:
logger.error("没有 vision 角色的可用适配器")
return {}
mode = self.config.get('orchestrator', {}).get('mode', 'fallback')
if mode == 'ensemble':
return self._run_visual_ensemble(
vision_adapters, frame_paths, frame_timestamps,
known_members_context, rate_limiter)
# fallback: 顺序降级,首个成功即采用
model_outputs = {}
for adapter in vision_adapters:
if adapter.get_circuit_breaker().is_open():
logger.warning(f"[{adapter.provider_name}] 熔断器 OPEN跳过")
continue
# 速率限制:按 provider 获取 token2x burst
if rate_limiter:
acquired = rate_limiter.acquire(adapter.provider_name, timeout=300)
if not acquired:
logger.warning(f"[{adapter.provider_name}] 速率限制超时,跳过")
continue
start = time.time()
try:
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,
event_start_time=event_start_time)
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()
log_task(logger, 0, f'model_{adapter.provider_name}',
f'视觉分析成功', duration_ms=duration_ms)
model_outputs[adapter.provider_name] = output
logger.info(f"fallback 采用 [{adapter.provider_name}],停止降级")
break
else:
adapter.get_circuit_breaker().record_failure()
logger.warning(f"[{adapter.provider_name}] 视觉分析返回空,降级下一模型")
except Exception as e:
logger.error(f"[{adapter.provider_name}] 视觉分析异常: {e}")
adapter.get_circuit_breaker().record_failure()
return model_outputs
def _run_visual_ensemble(self, vision_adapters, frame_paths,
frame_timestamps, known_members_context,
rate_limiter=None) -> Dict[str, dict]:
"""并行调用所有健康 vision 模型,保留全部成功结果(交叉验证)"""
model_outputs = {}
max_timeout = max((a.get_timeout() for a in vision_adapters), default=240)
with ThreadPoolExecutor(max_workers=len(vision_adapters)) as pool:
futures = {}
for adapter in vision_adapters:
if adapter.get_circuit_breaker().is_open():
logger.warning(f"[{adapter.provider_name}] 熔断器 OPEN跳过")
continue
# 速率限制:按 provider 获取 token2x burst
if rate_limiter:
acquired = rate_limiter.acquire(adapter.provider_name, timeout=300)
if not acquired:
logger.warning(f"[{adapter.provider_name}] 速率限制超时,跳过")
continue
future = pool.submit(
adapter.analyze_frames,
frame_paths, frame_timestamps, known_members_context)
futures[future] = adapter.provider_name
for future in as_completed(futures, timeout=max_timeout + 10):
provider = futures[future]
start = time.time()
try:
adapter = next(a for a in vision_adapters if a.provider_name == provider)
output = future.result(timeout=adapter.get_timeout())
duration_ms = int((time.time() - start) * 1000)
if output:
model_outputs[provider] = output
adapter.get_circuit_breaker().record_success()
log_task(logger, 0, f'model_{provider}',
f'视觉分析成功,输出长度={len(output)}', duration_ms=duration_ms)
else:
adapter.get_circuit_breaker().record_failure()
logger.warning(f"[{provider}] 视觉分析返回空")
except FuturesTimeout:
logger.warning(f"[{provider}] 视觉分析超时")
adapter = next(a for a in vision_adapters if a.provider_name == provider)
adapter.get_circuit_breaker().record_failure()
except Exception as e:
logger.error(f"[{provider}] 视觉分析异常: {e}")
adapter = next(a for a in vision_adapters if a.provider_name == provider)
adapter.get_circuit_breaker().record_failure()
return model_outputs
@staticmethod
def _attach_frame_images(frame_details: List[dict], frame_paths: List[str]) -> None:
"""把关键帧图片 base64 附加到 frame_details按位置对齐视觉分析输入帧
附带人脸红框标记与 face_countNAS 落盘 meta.jsonUI 据此挑有人像的头像)
"""
from ..frame_marker import mark_jpeg
for i, fd in enumerate(frame_details):
if i >= len(frame_paths):
break
try:
with open(frame_paths[i], 'rb') as f:
raw = f.read()
marked, faces = mark_jpeg(raw)
fd['frame_image'] = base64.b64encode(marked).decode('ascii')
fd['face_count'] = faces
except OSError as e:
logger.warning(f"关键帧图片读取失败: {frame_paths[i]}: {e}")
def format_cloud_result(self, provider: str, raw_result: dict,
known_members_context: str = '',
task_id: int = 0) -> dict:
"""格式化云端 VLM 直出的结构化结果(**无本地模型调用**)。
- 云端模型已产出结构化数据frame_details / 可选 global_summary / entities_json
- 本方法仅做字段归一化、source_providers 与 compute_provider 填充、
entities 推导、global_summary 缺失时格式化生成
- 解析/校验失败抛 VLMOutputInvalidError
"""
if not isinstance(raw_result, dict):
raise VLMOutputInvalidError("云端视觉模型未返回结构化数据(dict)")
data = dict(raw_result)
frame_details = data.get('frame_details')
if not isinstance(frame_details, list) or not frame_details:
raise VLMOutputInvalidError("云端结果缺少非空的 frame_details")
# 归一化每条 frame_detail
normalized = []
for f in frame_details:
if not isinstance(f, dict):
continue
sp = f.get('source_providers')
if not isinstance(sp, list) or not sp:
sp = [provider]
normalized.append({
"frame_index": int(f.get("frame_index", len(normalized) + 1)),
"frame_timestamp": str(f.get("frame_timestamp", "")),
"person": str(f.get("person", "无人")),
"action": str(f.get("action", "")),
"clothing": str(f.get("clothing", "")),
"is_attention_event": bool(f.get("is_attention_event", False)),
"source_providers": [str(p) for p in sp],
})
if not normalized:
raise VLMOutputInvalidError("frame_details 解析后为空")
data['frame_details'] = normalized
# compute_provider本次实际成功的云端模型
data['compute_provider'] = [provider]
# entities_json缺失时由 frame_details 推导(按人物去重)
if not data.get('entities_json'):
seen = set()
ents = []
for f in normalized:
p = f['person']
if p and p != '无人' and p not in seen:
seen.add(p)
ents.append({
"person": p,
"action": f['action'],
"clothing": f['clothing'],
})
data['entities_json'] = ents
# global_summary云端未给则格式化生成非 LLM 汇总,仅拼接事实)
if not data.get('global_summary'):
data['global_summary'] = self._build_summary_from_frames(normalized)
return validate_schema(data)
def _build_summary_from_frames(self, frame_details: List[dict]) -> str:
"""当云端模型未提供 global_summary 时,由 frame_details 格式化生成摘要。
注意:这是确定性事实拼接,非 LLM 二次汇总。"""
persons = {}
has_attention = False
for f in frame_details:
p = f['person']
if p and p != '无人':
persons.setdefault(p, set()).add(f['action'])
if f.get('is_attention_event'):
has_attention = True
if not persons:
summary = "整个时段内画面中未检测到人物出现,主要为环境静态画面。"
else:
parts = []
for p, acts in persons.items():
acts_desc = "".join(sorted(a for a in acts if a)) or "无明显动作"
parts.append(f"{p}{acts_desc}")
summary = f"时段内检测到:{''.join(parts)}"
if has_attention:
summary += " ⚠️ 存在需关注的异常事件。"
return summary
def run_qa(self, prompt: str, max_tokens: int = 512) -> Tuple[Optional[str], Optional[str]]:
"""智能问答编排Gemini → NVIDIA → 本地 Ollama仅当两云端都失败才用本地兜底
返回 (answer, provider);全部失败返回 (None, None)。
"""
qa_order = ['gemini', 'nvidia', 'ollama']
for name in qa_order:
adapter = next((a for a in self.adapters if a.provider_name == name), None)
if adapter is None:
logger.warning(f"[qa] 未配置模型 {name},跳过")
continue
try:
if adapter.get_circuit_breaker().is_open():
logger.warning(f"[qa] {name} 熔断器 OPEN跳过")
continue
answer = adapter.chat(prompt, max_tokens=max_tokens)
if answer:
logger.info(f"[qa] 由 {name} 回答(长度={len(answer)}")
return answer, name
logger.warning(f"[qa] {name} 返回空")
except Exception as e:
logger.error(f"[qa] {name} 调用异常: {e}")
return None, None
def send_callback(self, webhook_url: str, task_id: int,
result: dict, camera_name: str = '',
event_start_time: str = '', event_end_time: str = ''):
"""回调 NAS"""
payload = {
"task_id": task_id,
"status": "success",
"event_start_time": event_start_time,
"event_end_time": event_end_time,
"camera_name": camera_name,
"global_summary": result.get('global_summary', ''),
"entities_json": result.get('entities_json', []),
"frame_details": result.get('frame_details', []),
"compute_provider": result.get('compute_provider', []),
"error_message": None
}
callback_timeout = self.timeout_cfg.get('callback', 30)
max_retries = 3
for attempt in range(max_retries):
try:
resp = requests.post(webhook_url, json=payload, timeout=callback_timeout)
if resp.status_code == 200:
log_task(logger, task_id, 'callback', '回调成功')
return
else:
logger.warning(f"[task_id={task_id}] 回调返回 {resp.status_code},重试 {attempt+1}/{max_retries}")
except Exception as e:
logger.warning(f"[task_id={task_id}] 回调异常: {e},重试 {attempt+1}/{max_retries}")
raise Exception(f"回调失败,已重试 {max_retries}")
def send_failure_callback(self, webhook_url: str, task_id: int,
failure_stage: str, error_message: str):
"""发送失败回调"""
payload = {
"task_id": task_id,
"status": "failed",
"failure_stage": failure_stage,
"error_message": error_message
}
try:
requests.post(webhook_url, json=payload, timeout=30)
except Exception as e:
logger.error(f"[task_id={task_id}] 失败回调也失败: {e}")
def process_task(self, task_data: dict):
"""端到端处理任务拉取模式webhook 回调)"""
task_id = task_data.get('task_id')
video_url = task_data.get('video_url')
webhook_url = task_data.get('webhook_url')
known_members = task_data.get('known_members_context', '')
logger.info(f"[task_id={task_id}] ====== 开始处理任务 ======")
start_time = time.time()
# 1. 健康检查
healthy_adapters = self.health_check_all()
if not healthy_adapters:
logger.error(f"[task_id={task_id}] 所有模型不健康,返回 503")
self.send_failure_callback(webhook_url, task_id, 'vlm_visual', 'All models unhealthy')
return 503
# 2. 下载 + 抽帧
preprocessor = VideoPreprocessor(task_id)
try:
# 下载
video_path = preprocessor.download_video(video_url)
# 抽帧
candidate_frames = preprocessor.extract_candidate_frames(video_path)
if not candidate_frames:
raise Exception("抽帧失败,无候选帧")
# 关键帧筛选
key_frames = preprocessor.select_key_frames(candidate_frames)
# 压缩
compressed_frames = preprocessor.compress_frames(key_frames)
if not compressed_frames:
raise Exception("压缩后无可用帧")
# 计算时间戳
event_start_time = task_data.get('event_start_time', '')
frame_timestamps = preprocessor.compute_timestamps(
video_path, len(compressed_frames), event_start_time
)
# 3. 并行视觉分析
model_outputs = self.run_visual_analysis(
healthy_adapters, compressed_frames, frame_timestamps,
known_members, video_path=video_path,
event_start_time=event_start_time
)
if not model_outputs:
raise Exception('All models failed in visual analysis')
# 4. 云端直出结果格式化(无本地融合)
provider = next(iter(model_outputs))
fusion_result = self.format_cloud_result(
provider, model_outputs[provider], known_members, task_id)
# 5. 回调
# 从视频文件名推断 camera_name
camera_name = task_data.get('camera_name', '')
event_end_time = task_data.get('event_end_time', '')
self.send_callback(
webhook_url, task_id, fusion_result,
camera_name=camera_name,
event_start_time=event_start_time,
event_end_time=event_end_time
)
total_ms = int((time.time() - start_time) * 1000)
log_task(logger, task_id, 'overall', f'任务完成', duration_ms=total_ms)
except VLMOutputInvalidError as e:
logger.error(f"[task_id={task_id}] 云端结果格式化失败: {e}")
self.send_failure_callback(webhook_url, task_id, 'vlm_fusion', str(e))
except Exception as e:
logger.error(f"[task_id={task_id}] 任务处理失败: {e}", exc_info=True)
self.send_failure_callback(webhook_url, task_id, 'download', str(e))
finally:
# 6. 清理
if 'preprocessor' in locals():
preprocessor.cleanup()
return 200
def process_push_task(self, task_data: dict, video_path: str,
preprocessor: 'VideoPreprocessor',
rate_limiter=None) -> dict:
"""推送模式:同步处理上传的视频,结果直接返回(无 webhook 回调)
rate_limiter: 可选 RateLimiter 实例,按 provider 限速2x burst
返回 payload 结构与原 webhook 回调一致:
- 成功: {task_id, status: "success", event_start_time, ..., frame_details, ...}
- 失败: {task_id, status: "failed", failure_stage, error_message}
"""
task_id = task_data.get('task_id')
known_members = task_data.get('known_members_context', '')
event_start_time = task_data.get('event_start_time', '')
logger.info(f"[task_id={task_id}] ====== 开始处理推送任务 ======")
start_time = time.time()
try:
# 1. 健康检查
healthy_adapters = self.health_check_all()
if not healthy_adapters:
logger.error(f"[task_id={task_id}] 所有模型不健康")
return {
"task_id": task_id, "status": "failed",
"failure_stage": "vlm_visual",
"error_message": "All models unhealthy"
}
# 2. 抽帧(视频已由调用方保存到本地,无需下载)
candidate_frames = preprocessor.extract_candidate_frames(video_path)
if not candidate_frames:
raise Exception("抽帧失败,无候选帧")
key_frames = preprocessor.select_key_frames(candidate_frames)
compressed_frames = preprocessor.compress_frames(key_frames)
if not compressed_frames:
raise Exception("压缩后无可用帧")
frame_timestamps = preprocessor.compute_timestamps(
video_path, len(compressed_frames), event_start_time
)
# event_end_time 未提供时,用 start + 视频时长推算DB 列 NOT NULL
event_end_time = task_data.get('event_end_time', '')
if not event_end_time and event_start_time and preprocessor.video_duration > 0:
try:
start_dt = datetime.strptime(event_start_time, '%Y-%m-%d %H:%M:%S')
event_end_time = (
start_dt + timedelta(seconds=int(preprocessor.video_duration))
).strftime('%Y-%m-%d %H:%M:%S')
except ValueError:
pass
# 3. 并行视觉分析
model_outputs = self.run_visual_analysis(
healthy_adapters, compressed_frames, frame_timestamps,
known_members, rate_limiter, video_path=video_path,
event_start_time=event_start_time
)
if not model_outputs:
raise Exception('All models failed in visual analysis')
# 4. 云端直出结果格式化(无本地融合)
provider = next(iter(model_outputs))
fusion_result = self.format_cloud_result(
provider, model_outputs[provider], known_members, task_id)
# 5. 附加关键帧图片NAS 落盘后供 UI 时间轴展示)
frame_details = fusion_result.get('frame_details', [])
self._attach_frame_images(frame_details, compressed_frames)
total_ms = int((time.time() - start_time) * 1000)
log_task(logger, task_id, 'overall', '推送任务完成', duration_ms=total_ms)
return {
"task_id": task_id,
"status": "success",
"event_start_time": event_start_time,
"event_end_time": event_end_time,
"camera_name": task_data.get('camera_name', ''),
"global_summary": fusion_result.get('global_summary', ''),
"entities_json": fusion_result.get('entities_json', []),
"frame_details": frame_details,
"compute_provider": fusion_result.get('compute_provider', []),
"error_message": None
}
except VLMOutputInvalidError as e:
logger.error(f"[task_id={task_id}] VLM 输出解析失败: {e}")
return {
"task_id": task_id, "status": "failed",
"failure_stage": "vlm_fusion", "error_message": str(e)
}
except Exception as e:
logger.error(f"[task_id={task_id}] 推送任务处理失败: {e}", exc_info=True)
return {
"task_id": task_id, "status": "failed",
"failure_stage": "process", "error_message": str(e)
}