diff --git a/fam-core/config/config.yaml b/fam-core/config/config.yaml index bdd6bdc..319cec8 100644 --- a/fam-core/config/config.yaml +++ b/fam-core/config/config.yaml @@ -1,6 +1,7 @@ # FAM-Core 配置文件 (NAS 端) - 实际部署配置 -# Tailscale: NAS=100.70.234.39, Oracle=100.74.137.126 # 注: Tailscale 防火墙待修复,当前 edge_url 使用 Oracle 公网 IP +# 推送模式: NAS 上传视频到 Edge /api/edge/video/push,结果同步随响应返回, +# 无需 Oracle 反向访问 NAS(webhook/Video-Server 拉取均不再使用) # Ollama 未对外暴露,chat_handler 通过 FAM-Edge 代理 server: @@ -24,12 +25,12 @@ scheduler: dispatcher: poll_interval: 30 - edge_url: "http://129.146.203.203:5000/api/edge/video/analyze" - webhook_url: "http://100.70.234.39:8000/api/core/callback/event" + edge_url: "http://129.146.203.203:5000/api/edge/video/push" max_retries: 3 + push_timeout: 1800 video_server: - base_url: "http://100.70.234.39:8000/media" + base_url: "http://127.0.0.1:8000/media" token: "sentinel-media-2026" video_dir: "/volume1/surveillance" diff --git a/fam-core/config/config.yaml.example b/fam-core/config/config.yaml.example index 3e834fd..9ea3182 100644 --- a/fam-core/config/config.yaml.example +++ b/fam-core/config/config.yaml.example @@ -21,9 +21,9 @@ scheduler: dispatcher: poll_interval: 30 # 轮询间隔(秒) - edge_url: "http://100.x.x.20:5000/api/edge/video/analyze" - webhook_url: "http://100.x.x.10:8000/api/core/callback/event" + edge_url: "http://100.x.x.20:5000/api/edge/video/push" # 推送模式端点(视频上传,同步返回结果) max_retries: 3 + push_timeout: 1800 # 推送+分析同步超时(秒) video_server: base_url: "http://100.x.x.10:8000/media" diff --git a/fam-core/src/fam_core/dispatcher/dispatcher.py b/fam-core/src/fam_core/dispatcher/dispatcher.py index 71f7a29..cf8a3a5 100644 --- a/fam-core/src/fam_core/dispatcher/dispatcher.py +++ b/fam-core/src/fam_core/dispatcher/dispatcher.py @@ -1,9 +1,15 @@ """ -Dispatcher - 30s 轮询 PENDING 任务,下发至 Edge +Dispatcher - 30s 轮询 PENDING 任务,推送视频至 Edge(推送模式) + +流程(纯单向通信,NAS → Oracle,无需 Oracle 反向访问 NAS): +1. 读取任务对应的本地视频文件 +2. multipart 上传至 Edge /api/edge/video/push(附已知成员清单等元数据) +3. Edge 同步抽帧+分析,结果直接随 HTTP 响应返回 +4. Dispatcher 收到响应后直接写 monitor_events/event_details,任务标记 SUCCESS 退避重试: min(60 * (retry_count + 1) * 2, 600) 秒 -payload 注入 family_members 表的已命名+未命名成员清单 """ +import os import time import threading import requests @@ -12,18 +18,22 @@ from datetime import datetime, timedelta from ..logger import setup_logger, log_task from ..config_loader import load_config from .. import db_layer +from ..event_receiver.event_receiver import apply_success_event logger = setup_logger('fam-core.dispatcher') class Dispatcher: - """任务下发器,30s 轮询""" + """任务下发器,30s 轮询(推送模式)""" def __init__(self): cfg = load_config() self.poll_interval = cfg.get('dispatcher', {}).get('poll_interval', 30) - self.edge_url = cfg.get('dispatcher', {}).get('edge_url', 'http://localhost:5000/api/edge/video/analyze') + self.edge_url = cfg.get('dispatcher', {}).get('edge_url', 'http://localhost:5000/api/edge/video/push') self.max_retries = cfg.get('dispatcher', {}).get('max_retries', 3) + # 推送+分析全程同步,耗时较长(大视频上传 + ARM 多帧分析) + self.push_timeout = cfg.get('dispatcher', {}).get('push_timeout', 1800) + self.camera_name = cfg.get('scheduler', {}).get('camera_name', '默认摄像头') self._running = False self._thread = None @@ -42,41 +52,83 @@ class Dispatcher: return True def _build_payload(self, task): - """构建下发 payload,注入已知成员清单""" - known_members = db_layer.get_known_members_context() - return { - "task_id": task['task_id'], - "video_url": task['video_url'], - "webhook_url": load_config().get('dispatcher', {}).get( - 'webhook_url', 'http://localhost:8000/api/core/callback/event' - ), - "known_members_context": known_members + """构建推送元数据,注入已知成员清单与事件时间""" + video_path = task['video_path'] + payload = { + "task_id": str(task['task_id']), + "camera_name": self.camera_name, + "known_members_context": db_layer.get_known_members_context(), + "event_start_time": "", + "event_end_time": "", } + # 用文件 mtime 近似事件开始时间 + try: + mtime = os.path.getmtime(video_path) + start_dt = datetime.fromtimestamp(mtime) + payload["event_start_time"] = start_dt.strftime('%Y-%m-%d %H:%M:%S') + except OSError: + pass + return payload def _dispatch_one(self, task): - """下发单个任务""" + """推送单个任务:上传视频 → 同步等结果 → 直接写库""" task_id = task['task_id'] + video_path = task['video_path'] + + # 本地视频必须存在,否则直接失败(重试无意义) + if not video_path or not os.path.isfile(video_path): + db_layer.update_task_status( + task_id, 'FAILED', + error_message=f"视频文件不存在: {video_path}", + failure_stage='upload' + ) + logger.error(f"[task_id={task_id}] 视频文件不存在,标记 FAILED: {video_path}") + return + payload = self._build_payload(task) + size_mb = os.path.getsize(video_path) / (1024 * 1024) + db_layer.update_task_status(task_id, 'PROCESSING') + log_task(logger, task_id, 'dispatcher', + f'推送视频至 Edge: {self.edge_url} ({size_mb:.1f}MB)') try: - log_task(logger, task_id, 'dispatcher', f'下发至 Edge: {self.edge_url}') - resp = requests.post(self.edge_url, json=payload, timeout=30) - - if resp.status_code == 202: - db_layer.update_task_status(task_id, 'PROCESSING') - log_task(logger, task_id, 'dispatcher', 'Edge 接受任务,状态切换为 PROCESSING') - elif resp.status_code == 429: - logger.warning(f"[task_id={task_id}] Edge 队列已满 (429),稍后重试") - elif resp.status_code == 503: - logger.warning(f"[task_id={task_id}] Edge Ollama 不可用 (503),退避重试") - self._schedule_retry(task) - else: - logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}") - self._schedule_retry(task) - + with open(video_path, 'rb') as fh: + resp = requests.post( + self.edge_url, + data=payload, + files={'video': (os.path.basename(video_path), fh, 'video/mp4')}, + timeout=self.push_timeout + ) except requests.RequestException as e: - logger.error(f"[task_id={task_id}] 下发失败: {e}") + logger.error(f"[task_id={task_id}] 推送失败: {e}") self._schedule_retry(task) + return + + if resp.status_code == 429: + logger.warning(f"[task_id={task_id}] Edge 忙 (429),回到 PENDING 稍后重试") + db_layer.update_task_status(task_id, 'PENDING') + return + + if resp.status_code != 200: + logger.error(f"[task_id={task_id}] Edge 返回异常状态码: {resp.status_code}") + self._schedule_retry(task) + return + + result = resp.json(silent=True) or {} + if result.get('status') == 'success': + try: + event_id = apply_success_event(task_id, result) + log_task(logger, task_id, 'dispatcher', f'推送任务完成: event_id={event_id}') + except Exception as e: + logger.error(f"[task_id={task_id}] 结果落库失败: {e}", exc_info=True) + db_layer.update_task_status( + task_id, 'FAILED', error_message=str(e), failure_stage='callback') + else: + error_message = result.get('error_message', 'unknown') + failure_stage = result.get('failure_stage', 'edge') + logger.error(f"[task_id={task_id}] Edge 分析失败: stage={failure_stage}, error={error_message}") + db_layer.update_task_status( + task_id, 'FAILED', error_message=error_message, failure_stage=failure_stage) def _schedule_retry(self, task): """调度重试""" diff --git a/fam-core/src/fam_core/event_receiver/event_receiver.py b/fam-core/src/fam_core/event_receiver/event_receiver.py index 74b8da5..7e6a093 100644 --- a/fam-core/src/fam_core/event_receiver/event_receiver.py +++ b/fam-core/src/fam_core/event_receiver/event_receiver.py @@ -45,6 +45,49 @@ def _upsert_abstract_members(frame_details: list): logger.info(f"upsert family_member: {label} (feature={feature})") +def apply_success_event(task_id, data: dict) -> int: + """将成功结果写库,返回 event_id + + 供两条路径复用: + - webhook 回调路由 (拉取模式) + - Dispatcher 收到推送模式同步响应后直接落库 + """ + # 1. 插入 monitor_events + event_id = db_layer.insert_event( + task_id=task_id, + event_start_time=data['event_start_time'], + event_end_time=data['event_end_time'], + camera_name=data.get('camera_name', ''), + global_summary=data.get('global_summary', ''), + entities_json=data.get('entities_json', []), + compute_provider=data.get('compute_provider', []) + ) + + # 2. 遍历 frame_details 逐条插入 + frame_details = data.get('frame_details', []) + for frame in frame_details: + db_layer.insert_event_detail( + event_id=event_id, + task_id=task_id, + frame_index=frame.get('frame_index', 0), + frame_timestamp=frame.get('frame_timestamp', ''), + camera_name=frame.get('camera_name', data.get('camera_name', '')), + person=frame.get('person', '未知'), + action=frame.get('action', ''), + clothing=frame.get('clothing', ''), + is_attention_event=frame.get('is_attention_event', False), + source_providers=frame.get('source_providers', []) + ) + + # 3. 对未命名的 abstract_label 自动 upsert + _upsert_abstract_members(frame_details) + + # 4. 更新任务状态 + db_layer.update_task_status(task_id, 'SUCCESS') + logger.info(f"[task_id={task_id}] 事件处理完成: event_id={event_id}, frame_details={len(frame_details)}条") + return event_id + + @event_bp.route('/api/core/callback/event', methods=['POST']) def receive_event(): """接收 Edge 回调""" @@ -58,40 +101,7 @@ def receive_event(): if status == 'success': try: - # 1. 插入 monitor_events - event_id = db_layer.insert_event( - task_id=task_id, - event_start_time=data['event_start_time'], - event_end_time=data['event_end_time'], - camera_name=data.get('camera_name', ''), - global_summary=data.get('global_summary', ''), - entities_json=data.get('entities_json', []), - compute_provider=data.get('compute_provider', []) - ) - - # 2. 遍历 frame_details 逐条插入 - frame_details = data.get('frame_details', []) - for frame in frame_details: - db_layer.insert_event_detail( - event_id=event_id, - task_id=task_id, - frame_index=frame.get('frame_index', 0), - frame_timestamp=frame.get('frame_timestamp', ''), - camera_name=frame.get('camera_name', data.get('camera_name', '')), - person=frame.get('person', '未知'), - action=frame.get('action', ''), - clothing=frame.get('clothing', ''), - is_attention_event=frame.get('is_attention_event', False), - source_providers=frame.get('source_providers', []) - ) - - # 3. 对未命名的 abstract_label 自动 upsert - _upsert_abstract_members(frame_details) - - # 4. 更新任务状态 - db_layer.update_task_status(task_id, 'SUCCESS') - logger.info(f"[task_id={task_id}] 事件处理完成: event_id={event_id}, frame_details={len(frame_details)}条") - + event_id = apply_success_event(task_id, data) return jsonify({"status": "ok", "event_id": event_id}), 200 except Exception as e: diff --git a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py index 18d5a0e..58c8a80 100644 --- a/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py +++ b/fam-edge/src/fam_edge/ai_orchestrator/orchestrator.py @@ -246,7 +246,7 @@ class AIOrchestrator: 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') @@ -325,3 +325,82 @@ class AIOrchestrator: preprocessor.cleanup() return 200 + + def process_push_task(self, task_data: dict, video_path: str, + preprocessor: 'VideoPreprocessor') -> dict: + """推送模式:同步处理上传的视频,结果直接返回(无 webhook 回调) + + 返回 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 + ) + + # 3. 并行视觉分析 + model_outputs = self.run_visual_analysis( + healthy_adapters, compressed_frames, frame_timestamps, known_members + ) + if not model_outputs: + raise Exception('All models failed in visual analysis') + + # 4. 文本融合 + fusion_result = self.run_text_fusion(model_outputs, known_members, task_id) + + 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": task_data.get('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": fusion_result.get('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) + } 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 ef613f6..f58211d 100644 --- a/fam-edge/src/fam_edge/api_gateway/api_gateway.py +++ b/fam-edge/src/fam_edge/api_gateway/api_gateway.py @@ -9,6 +9,7 @@ from flask import Blueprint, request, jsonify from ..logger import setup_logger from ..ai_orchestrator.orchestrator import AIOrchestrator +from ..video_preprocessor.preprocessor import VideoPreprocessor logger = setup_logger('fam-edge.api_gateway') @@ -71,6 +72,63 @@ def receive_task(): return jsonify({"status": "accepted", "task_id": task_id}), 202 +@api_bp.route('/api/edge/video/push', methods=['POST']) +def receive_push_task(): + """推送模式:接收 multipart 视频上传,同步分析,结果随 HTTP 响应返回 + + NAS 无法被 Oracle 反向访问(Tailscale 不通),因此改为 NAS 主动上传视频, + Edge 用 OpenCV 场景变化检测抽帧后分析,摘要直接放在响应里带回。 + """ + global _currently_processing + + task_id_raw = request.form.get('task_id') + file = request.files.get('video') + if not task_id_raw or not file: + return jsonify({"error": "缺少必填字段: task_id, video"}), 400 + + try: + task_id = int(task_id_raw) + except ValueError: + return jsonify({"error": "task_id 必须是整数"}), 400 + + logger.info(f"[task_id={task_id}] 收到推送任务: {file.filename}") + + # 并发控制(同步处理,占用整个请求周期) + with _current_task_lock: + if _currently_processing: + logger.warning(f"[task_id={task_id}] 已有任务处理中,返回 429") + return jsonify({"error": "Queue full", "retry_after": 60}), 429 + _currently_processing = True + + preprocessor = None + try: + preprocessor = VideoPreprocessor(task_id) + video_path = preprocessor.save_upload(file) + + task_data = { + "task_id": task_id, + "camera_name": request.form.get('camera_name', ''), + "event_start_time": request.form.get('event_start_time', ''), + "event_end_time": request.form.get('event_end_time', ''), + "known_members_context": request.form.get('known_members_context', ''), + } + + result = get_orchestrator().process_push_task(task_data, video_path, preprocessor) + return jsonify(result), 200 + + except Exception as e: + logger.error(f"[task_id={task_id}] 推送任务异常: {e}", exc_info=True) + return jsonify({ + "task_id": task_id, "status": "failed", + "failure_stage": "upload", "error_message": str(e) + }), 200 + finally: + if preprocessor is not None: + preprocessor.cleanup() + with _current_task_lock: + _currently_processing = False + + @api_bp.route('/health', methods=['GET']) def health(): """健康检查""" diff --git a/fam-edge/src/fam_edge/video_preprocessor/preprocessor.py b/fam-edge/src/fam_edge/video_preprocessor/preprocessor.py index cd6d243..8e3e261 100644 --- a/fam-edge/src/fam_edge/video_preprocessor/preprocessor.py +++ b/fam-edge/src/fam_edge/video_preprocessor/preprocessor.py @@ -68,6 +68,17 @@ class VideoPreprocessor: log_task(logger, self.task_id, 'download', f'下载完成: {size_mb:.1f}MB', duration_ms=duration_ms) return self.video_path + def save_upload(self, file_storage) -> str: + """保存推送模式上传的视频文件(multipart),替代 download_video""" + os.makedirs(self.work_dir, exist_ok=True) + start = time.time() + file_storage.save(self.video_path) + duration_ms = int((time.time() - start) * 1000) + size_mb = os.path.getsize(self.video_path) / (1024 * 1024) + log_task(logger, self.task_id, 'upload', + f'保存上传视频: {size_mb:.1f}MB', duration_ms=duration_ms) + return self.video_path + def _get_video_duration(self, video_path: str) -> float: """用 ffprobe 获取视频时长(秒)""" try: