""" 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 /api/oracle/frame 事件时刻缩略帧(ffmpeg 抽帧 + 磁盘缓存) GET /api/oracle/avatar 人物头像(按视频分析产出的 bbox 裁剪) GET /health 健康检查(DB + 队列线程存活) 图片能力(v4): 计算全部在 Oracle 本机 ffmpeg 抽帧/裁剪;人物 bbox 随视频分析那 一次 Gemini 调用一并产出(见 ai_orchestrator/prompts.py),frame_service 不再 额外调用任何模型。NAS 经 core 代理读取,不在 NAS 做图像计算。 """ import json import os from flask import Blueprint, request, jsonify, Response from ..logger import setup_logger from .. import state from ..qa import QAOrchestrator logger = setup_logger('fam-edge.api_gateway') api_bp = Blueprint('api_gateway', __name__) _qa = None def get_qa(): global _qa if _qa is None: _qa = QAOrchestrator() return _qa def _check_token() -> bool: expected = _sync_token() token = request.args.get('token') or request.form.get('token') or \ (request.get_json(silent=True) or {}).get('token', '') return bool(expected) and token == expected _SYNC_TOK = None def _sync_token(): global _SYNC_TOK if _SYNC_TOK is None: from ..config_loader import load_config _SYNC_TOK = load_config().get('sync_api', {}).get('token', '${ORACLE_SYNC_TOKEN}') if _SYNC_TOK.startswith('${') and _SYNC_TOK.endswith('}'): _SYNC_TOK = os.environ.get(_SYNC_TOK[2:-1], '') return _SYNC_TOK or '' def _now_iso_str() -> str: from datetime import datetime, timezone, timedelta return datetime.now(timezone(timedelta(hours=8))).strftime('%Y-%m-%d %H:%M:%S') @api_bp.route('/api/oracle/sync', methods=['GET']) def sync_pull(): """NAS 拉取增量数据。since=ISO 时间字符串(默认 '' 拉全量)。 返回: {videos:[...], events:[...], people:[...], server_time} """ if not _check_token(): return jsonify({"error": "unauthorized"}), 401 since = request.args.get('since', '') try: delta = state.get_db().get_sync_delta(since) except Exception as e: logger.error(f"sync_pull 异常: {e}") return jsonify({"error": str(e)}), 500 return jsonify(delta), 200 @api_bp.route('/api/oracle/people/correct', methods=['POST']) def people_correct(): """NAS 手动命名校正推送。 请求: {"label": "人物A", "canonical_name": "张三", "token": "..."} 更新 people 表(manual 优先,不被 LLM 覆盖),立即重算 known_members_context。 """ if not _check_token(): return jsonify({"error": "unauthorized"}), 401 data = request.get_json(silent=True) if not data: return jsonify({"error": "Invalid JSON"}), 400 label = (data.get('label') or '').strip() canonical = (data.get('canonical_name') or '').strip() if not label or not canonical: return jsonify({"error": "缺少 label / canonical_name"}), 400 try: state.get_db().set_canonical(label, canonical, source='manual') except Exception as e: logger.error(f"people_correct 异常: {e}") return jsonify({"error": str(e)}), 500 return jsonify({"status": "ok", "label": label, "canonical_name": canonical}), 200 @api_bp.route('/api/oracle/identity/correct', methods=['POST']) def identity_correct(): """事件时间轴/人物管理"这个人识别错了"纠错入口(比 people/correct 粒度更细)。 请求: {"video_id": 123, "current_name": "爷爷", "new_name": "爸爸", "token": "..."} 只改这一段视频里被错误识别的那个人,不影响同名字符串在其他视频里的映射—— 人物 uid 只在单次视频分析内稳定,同一个"人物A"字符串在不同视频里可能是不同 真人,纠错必须落到 (video_id, 当前展示名) 这一粒度,不能按全局 label 改。 写 manual 来源,受保护不会被后续自动识别覆盖回去;立即重写这段视频的展示数据。 """ if not _check_token(): return jsonify({"error": "unauthorized"}), 401 data = request.get_json(silent=True) if not data: return jsonify({"error": "Invalid JSON"}), 400 video_id = data.get('video_id') current_name = (data.get('current_name') or '').strip() new_name = (data.get('new_name') or '').strip() if not video_id or not current_name or not new_name: return jsonify({"error": "缺少 video_id / current_name / new_name"}), 400 try: video_id = int(video_id) except (TypeError, ValueError): return jsonify({"error": "video_id 必须是数字"}), 400 try: state.get_db().correct_video_identity(video_id, current_name, new_name) except Exception as e: logger.error(f"identity_correct 异常: {e}") return jsonify({"error": str(e)}), 500 return jsonify({"status": "ok", "video_id": video_id, "current_name": current_name, "new_name": new_name}), 200 @api_bp.route('/api/oracle/video/delete', methods=['POST']) def video_delete(): """删除视频会话(事件时间轴"删除"入口,NAS 转发)。 请求: {"video_id": 123, "token": "..."} 删 Oracle 端 events + videos 行 + 磁盘上的运动片段文件;不清理对应的 ss_motion_events 源事件(独立生命周期,见 oracle_db.delete_video 注释)。 """ if not _check_token(): return jsonify({"error": "unauthorized"}), 401 data = request.get_json(silent=True) if not data: return jsonify({"error": "Invalid JSON"}), 400 video_id = data.get('video_id') if not video_id: return jsonify({"error": "缺少 video_id"}), 400 try: video_id = int(video_id) except (TypeError, ValueError): return jsonify({"error": "video_id 必须是数字"}), 400 try: local_path = state.get_db().delete_video(video_id) except Exception as e: logger.error(f"video_delete 异常: {e}") return jsonify({"error": str(e)}), 500 if local_path is None: return jsonify({"error": "视频不存在"}), 404 state.get_db().record_activity('video', 'delete', f"video_id={video_id}") return jsonify({"status": "ok", "video_id": video_id}), 200 @api_bp.route('/api/edge/chat/ask', methods=['POST']) def chat_ask(): """智能问答编排:Gemini → NVIDIA → 本地 Ollama(两云端都失败才用本地兜底) 请求: {"prompt": "...", "max_tokens": 512} 响应: {"answer": "...", "provider": "gemini"|"nvidia"|"ollama"} """ data = request.get_json(silent=True) if not data or 'prompt' not in data: return jsonify({"error": "缺少必填字段: prompt"}), 400 prompt = data['prompt'] # 默认值从 1024 提到 3072(2026-08-23):问答链路现在优先用推理类模型 # (NVIDIA Nemotron-3 系列),回答前会先输出一段思考过程再给最终答案, # 1024 经常在思考阶段就被截断,永远看不到真正的回答。 max_tokens = int(data.get('max_tokens', 3072)) answer, provider = get_qa().run_qa(prompt, max_tokens=max_tokens) if answer is None: return jsonify({ "error": "所有模型均不可用(Gemini / NVIDIA / Ollama 全部失败)" }), 503 return jsonify({"answer": answer, "provider": provider}), 200 @api_bp.route('/api/edge/chat/ask/stream', methods=['POST']) def chat_ask_stream(): """智能问答编排(流式版):SSE 逐块推送,边生成边显示,不用等全量回答。 请求同 /api/edge/chat/ask。响应 Content-Type: text/event-stream, 每行 `data: \\n\\n`,json 结构见 qa.QAOrchestrator.run_qa_stream 注释。 """ data = request.get_json(silent=True) if not data or 'prompt' not in data: return jsonify({"error": "缺少必填字段: prompt"}), 400 prompt = data['prompt'] # 默认值从 1024 提到 3072(2026-08-23):问答链路现在优先用推理类模型 # (NVIDIA Nemotron-3 系列),回答前会先输出一段思考过程再给最终答案, # 1024 经常在思考阶段就被截断,永远看不到真正的回答。 max_tokens = int(data.get('max_tokens', 3072)) def generate(): for event in get_qa().run_qa_stream(prompt, max_tokens=max_tokens): yield f"data: {json.dumps(event, ensure_ascii=False)}\n\n" # mimetype 显式带 charset=utf-8:响应体本身一直是 UTF-8,但不声明的话下游 # (fam-core 转发这一跳、或任何用 requests 消费这个流的客户端)会自己猜 # 编码,猜错就是中文乱码——跟 gemini_adapter.py chat_stream() 那个坑同源。 return Response(generate(), mimetype='text/event-stream; charset=utf-8', headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}) @api_bp.route('/api/oracle/activity', methods=['GET']) def activity(): """实时服务状态 + 最近活动流(token 校验)。 返回各服务当前状态(队列/当前处理视频/rclone 最近同步/人物合并/模型调用) 与最近 50 条活动(service_activity,只保留 7 天)。 """ if not _check_token(): return jsonify({"error": "unauthorized"}), 401 db = state.get_db() queue = state.get_queue() # 队列实时状态 q_status = None if queue is not None: try: q_status = queue.status() except Exception as e: logger.warning(f"queue.status 异常: {e}") # 各服务最近活动(rclone / person / model 切换) def _last_activity(service): row = db._conn.execute( "SELECT service, action, detail, ts FROM service_activity " "WHERE service=? ORDER BY id DESC LIMIT 1", (service,)).fetchone() return dict(row) if row else None # 最近模型调用(实时模型卡) model_calls = db._conn.execute( "SELECT id, provider, model, filename, started_at, duration_sec, " "success, error FROM model_calls ORDER BY id DESC LIMIT 5").fetchall() # 运动片段分割状态(服务状态页「视频分割」卡) segment = db.get_segment_status() try: import os from ..config_loader import load_config seg_dir = load_config().get('motion_segment', {}).get( 'clips_dir', '/opt/fam-edge/motion_clips') segment['clips_files'] = len(os.listdir(seg_dir)) if os.path.isdir(seg_dir) else 0 except Exception: segment['clips_files'] = 0 # 事件↔片段一致性对账(数量可验证) segment['consistency'] = db.get_segment_consistency() return jsonify({ "queue": q_status, "db": db.get_queue_status(), "segment": segment, "rclone": _last_activity('rclone'), "person": _last_activity('person'), "model_calls": [dict(m) for m in model_calls], "activities": db.get_recent_activities(50), "ts": _now_iso_str(), }), 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('/api/ss/motion', methods=['POST']) def ss_motion(): """接收 NAS 推送的运动侦测事件(单向:NAS -> Oracle)。 请求体: {"token": "...", "events": [ {"event_id": 25349, "camera_id": 2, "event_type": 10, "start_time": 1787318023, "duration": 149, "thumbnail_url": "107471,12075"} ]} start_time/duration 为 Unix epoch(与 SS 同源,时区无关)。events 可以是空数组 ——NAS 侧也会定期用空 events 单纯调一次这个接口当心跳,证明推送链路还活着 (鉴权通过就刷新心跳,不要求 stored>0),供 has_motion_in_range_local() 判断 "无运动"结论是否可信。 落库 ss_motion_events(按 event_id 幂等),供 video_processor 本地预过滤使用。 """ if not _check_token(): return jsonify({"error": "unauthorized"}), 401 data = request.get_json(silent=True) if not data or 'events' not in data: return jsonify({"error": "缺少 events"}), 400 events = data.get('events') or [] if not isinstance(events, list): return jsonify({"error": "events 必须是数组"}), 400 try: stored = state.get_db().record_motion_events(events) state.get_db().record_motion_heartbeat() except Exception as e: logger.error(f"ss_motion 落库异常: {e}") return jsonify({"error": str(e)}), 500 if stored: state.get_db().record_activity( 'motion', 'push', f"接收 NAS 运动事件 {stored} 条") return jsonify({"status": "ok", "received": len(events), "stored": stored}), 200 @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'] 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