feat: video analysis switched to push mode (upload whole video, sync response)
Rationale: Oracle cannot reach NAS (Tailscale userspace mode on NAS, no TUN), the old pull+webhook design requires Edge to download video from NAS and callback to NAS - both blocked. New design is one-way NAS -> Oracle: - FAM-Edge: new POST /api/edge/video/push endpoint accepts multipart video upload, reuses existing OpenCV scene-change keyframe selection, analyzes synchronously and returns the result payload directly in the HTTP response (no webhook callback). Old /api/edge/video/analyze kept for compatibility. - FAM-Edge: VideoPreprocessor.save_upload() saves the uploaded file - FAM-Edge: AIOrchestrator.process_push_task() runs the full pipeline (health check -> extract -> select -> compress -> VLM -> fusion) and returns callback-style payload dict - FAM-Core: Dispatcher rewritten to push mode - reads local video file, uploads with task metadata (camera_name, event_start_time from file mtime, known_members_context), applies the result to DB via shared event_receiver.apply_success_event() - FAM-Core: event_receiver success logic extracted into reusable apply_success_event() (used by both webhook route and dispatcher) - config: edge_url -> /api/edge/video/push, push_timeout 1800s, gunicorn Edge timeout raised to 1800s for long synchronous analysis
This commit is contained in:
@@ -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):
|
||||
"""调度重试"""
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user