feat: 事件时间轴缩略帧 + 人物管理头像 + 人物合并硬规则校验
## 新架构:Oracle 集中计算 + NAS 代理展示 ### Oracle 端 (fam-edge) - 新增 frame_service: ffmpeg 视频抽帧 + VLM 人物定位裁剪头像(磁盘缓存) - 新增 /api/oracle/frame: 按 video_id+ts 抽帧返回 jpeg(带 token) - 新增 /api/oracle/avatar: 按 label 生成人物头像(VLM 定位人物 + 兜底整帧居中) - 新增 person_identifier: 人物身份识别模块 - Gemini 适配器支持 flash/flash-lite 双模型切换,429 自动降级 - frame_service VLM 全模型 429 时进入 10 分钟熔断,避免每次请求白打配额 - 兜底头像不落缓存,配额恢复后自动重试 VLM 精确定位 ### 人物合并硬规则校验(框架级修复) - person_service: LLM 合并结果落库前加硬冲突检测 - 性别冲突 → 绝不合并 - 年龄档跨未成年/成年 → 绝不合并(防止把爷爷/宝宝并进同一人) - oracle_db: upsert_person 入口剥离括号后缀(人物A(别名:人物B) → 人物A),消灭垃圾人物行 - 修复 set_canonical 丢弃 source 参数的 bug(旧代码硬编码 'manual' 导致错误合并被永久固化) - get_events_for_label: 只提取该身份组的特征文本,头像定位更精准 ### NAS 端 (fam-core) - 新增 img_proxy: /api/proxy/frame 和 /api/proxy/avatar 代理 Oracle 图片 - app.py 注册 img_bp 蓝图 - oracle_sync / db_layer / member_manager 同步人物表 ### UI 端 (fam-ui) - 事件时间轴: 每条事件卡片加时间点缩略帧 - 人物管理: 每人卡片加头像(150x150 圆角) - parse_persons: 剥离括号备注,与 Oracle 归一化一致 - 新增 EventItem 组件、Timeline 页改造 - Chat / ServiceStatus 页相应调整 ### 数据库 - scripts/ddl.sql: 同步表结构更新 - Oracle people 表: features_json / display_uid / source 字段完善
This commit is contained in:
@@ -90,6 +90,21 @@ def status():
|
||||
}), 200
|
||||
|
||||
|
||||
@app.route('/api/sync/trigger', methods=['POST'])
|
||||
def sync_trigger():
|
||||
"""手动立即触发一次甲骨文增量同步(服务状态页"立即同步"按钮)。
|
||||
|
||||
正常情况下后台线程每 30 分钟自动拉一次;这个接口给用户想立刻看到最新数据
|
||||
时用,跟后台线程共用同一把拉取锁(oracle_sync._pull_lock),不会并发重复拉。
|
||||
"""
|
||||
if not _sync:
|
||||
return jsonify({"error": "同步服务未初始化"}), 503
|
||||
ok = _sync.trigger_now()
|
||||
if not ok:
|
||||
return jsonify({"status": "failed", "error": _sync.status().get("last_error")}), 502
|
||||
return jsonify({"status": "ok", **_sync.status()}), 200
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
cfg = load_config()
|
||||
port = cfg.get('server', {}).get('port', 8000)
|
||||
|
||||
@@ -13,7 +13,7 @@ Chat-Handler - Flask 蓝图,接收用户问答(新架构 v2)
|
||||
"""
|
||||
import json
|
||||
import requests
|
||||
from flask import Blueprint, request, jsonify
|
||||
from flask import Blueprint, request, jsonify, Response, stream_with_context
|
||||
|
||||
from ..logger import setup_logger
|
||||
from ..config_loader import load_config
|
||||
@@ -115,6 +115,95 @@ def chat_ask():
|
||||
}), 200
|
||||
|
||||
|
||||
@chat_bp.route('/api/chat/ask/stream', methods=['POST'])
|
||||
def chat_ask_stream():
|
||||
"""流式问答:SSE 逐块推送,边生成边显示。
|
||||
|
||||
先立即推一条 context 事件(用了哪些 sync_events,本地查询很快,不用等
|
||||
大模型);再把甲骨文 /api/edge/chat/ask/stream 的分块原样转发给前端;
|
||||
最后一次性把拼好的完整回答写进 chat_history(跟非流式版一致)。
|
||||
"""
|
||||
data = request.get_json(silent=True)
|
||||
if not data:
|
||||
return jsonify({"error": "Invalid JSON"}), 400
|
||||
|
||||
question = data.get('question', '')
|
||||
queried_person = data.get('queried_person', '')
|
||||
queried_date = data.get('queried_date', '')
|
||||
if not question or not queried_person or not queried_date:
|
||||
return jsonify({"error": "缺少必填字段: question, queried_person, queried_date"}), 400
|
||||
|
||||
logger.info(f"Chat(stream): person={queried_person}, date={queried_date}, question={question}")
|
||||
rows = db_layer.query_sync_events_for_person_date(queried_person, queried_date)
|
||||
|
||||
def sse(obj):
|
||||
return f"data: {json.dumps(obj, ensure_ascii=False)}\n\n"
|
||||
|
||||
def generate():
|
||||
if len(rows) == 0:
|
||||
answer = f"今天没有观察到{queried_person}。"
|
||||
context_summary = "查询 sync_events 0 条"
|
||||
yield sse({"type": "context", "count": 0, "summary": context_summary})
|
||||
yield sse({"type": "chunk", "provider": None, "text": answer})
|
||||
yield sse({"type": "done", "provider": None})
|
||||
db_layer.insert_chat_history(
|
||||
user_question=question, ai_answer=answer,
|
||||
context_summary=context_summary,
|
||||
queried_date=queried_date, queried_person=queried_person)
|
||||
return
|
||||
|
||||
context = _format_events(rows)
|
||||
context_summary = f"查询 sync_events {len(rows)} 条"
|
||||
yield sse({"type": "context", "count": len(rows), "summary": context_summary,
|
||||
"preview": context[:800]})
|
||||
|
||||
members = db_layer.get_sync_known_members_context()
|
||||
prompt = build_chat_prompt(
|
||||
context=context, members=members or queried_person,
|
||||
question=question, queried_person=queried_person)
|
||||
|
||||
cfg = load_config()
|
||||
stream_url = cfg.get('chat_handler', {}).get(
|
||||
'qa_stream_url', 'http://129.146.203.203:5000/api/edge/chat/ask/stream')
|
||||
timeout = cfg.get('chat_handler', {}).get('timeout', 120)
|
||||
|
||||
full_answer = []
|
||||
provider_used = None
|
||||
try:
|
||||
resp = requests.post(stream_url, json={"prompt": prompt},
|
||||
timeout=(10, timeout), stream=True)
|
||||
if resp.status_code != 200:
|
||||
raise Exception(f"HTTP {resp.status_code}")
|
||||
for line in resp.iter_lines(decode_unicode=True):
|
||||
if not line or not line.startswith('data: '):
|
||||
continue
|
||||
yield line + '\n\n' # 原样转发给前端(已经是同样的 SSE 格式)
|
||||
try:
|
||||
obj = json.loads(line[len('data: '):])
|
||||
except ValueError:
|
||||
continue
|
||||
if obj.get('type') == 'chunk' and obj.get('text'):
|
||||
full_answer.append(obj['text'])
|
||||
provider_used = obj.get('provider') or provider_used
|
||||
elif obj.get('type') == 'done':
|
||||
provider_used = obj.get('provider') or provider_used
|
||||
except Exception as e:
|
||||
logger.error(f"流式问答代理失败: {e}")
|
||||
if not full_answer:
|
||||
yield sse({"type": "error", "message": "AI 服务暂时不可用,请稍后重试"})
|
||||
return
|
||||
|
||||
answer = ''.join(full_answer).strip() or "抱歉,暂时无法生成回答。"
|
||||
logger.info(f"流式问答由 {provider_used} 提供回答(长度={len(answer)})")
|
||||
db_layer.insert_chat_history(
|
||||
user_question=question, ai_answer=answer,
|
||||
context_summary=context_summary,
|
||||
queried_date=queried_date, queried_person=queried_person)
|
||||
|
||||
return Response(stream_with_context(generate()), mimetype='text/event-stream',
|
||||
headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'})
|
||||
|
||||
|
||||
@chat_bp.route('/api/chat/history', methods=['GET'])
|
||||
def chat_history():
|
||||
"""获取对话历史"""
|
||||
|
||||
@@ -370,6 +370,39 @@ def upsert_sync_model_calls(rows: List[Dict]) -> int:
|
||||
conn.close()
|
||||
|
||||
|
||||
def upsert_sync_identity_map(rows: List[Dict]) -> int:
|
||||
"""批量 upsert Oracle 传来的 person_identity_map 增量(幂等,重复覆盖)。
|
||||
|
||||
甲骨文不稳定,提取出的有效数据(含闭集人物识别结果/纠错记录)都要同步到
|
||||
NAS 防丢失——这张表跟 sync_people/sync_model_calls 走一样的镜像模式。
|
||||
"""
|
||||
if not rows:
|
||||
return 0
|
||||
conn = get_conn()
|
||||
n = 0
|
||||
try:
|
||||
cur = conn.cursor()
|
||||
for r in rows:
|
||||
cur.execute(
|
||||
"""INSERT INTO sync_identity_map
|
||||
(id, video_id, raw_uid, canonical_name, source, updated_at, synced_at)
|
||||
VALUES (%s,%s,%s,%s,%s,%s, NOW())
|
||||
ON DUPLICATE KEY UPDATE
|
||||
video_id=VALUES(video_id),
|
||||
raw_uid=VALUES(raw_uid),
|
||||
canonical_name=VALUES(canonical_name),
|
||||
source=VALUES(source),
|
||||
updated_at=VALUES(updated_at),
|
||||
synced_at=NOW()""",
|
||||
(r.get('id'), r.get('video_id'), r.get('raw_uid'),
|
||||
r.get('canonical_name'), r.get('source'), r.get('updated_at')))
|
||||
n += 1
|
||||
conn.commit()
|
||||
return n
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def get_sync_model_calls(limit: int = 200) -> List[Dict]:
|
||||
"""最近模型调用记录(前端统计展示)。"""
|
||||
conn = get_conn()
|
||||
|
||||
@@ -142,3 +142,38 @@ def merge_member():
|
||||
"display_name": m.get('canonical_name') or m['label'],
|
||||
} for m in members]
|
||||
}), 200
|
||||
|
||||
|
||||
@member_bp.route('/api/member/identity-correct', methods=['POST'])
|
||||
def identity_correct():
|
||||
"""事件时间轴"这个人识别错了"纠错入口(比 /api/member/name 粒度更细)。
|
||||
|
||||
请求: {"video_id": 123, "current_name": "爷爷", "new_name": "爸爸"}
|
||||
只改这一段视频里被错误识别的那个人,不影响同名字符串在其他视频里的映射
|
||||
(人物 uid 只在单次视频分析内稳定,同一字符串在不同视频里可能是不同真人,
|
||||
不能像 /api/member/name 那样按全局 label 改)。
|
||||
"""
|
||||
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
|
||||
|
||||
logger.info(f"人物纠错: video_id={video_id} {current_name} -> {new_name}(回推 Oracle)")
|
||||
ok, err = get_sync().push_identity_correct(video_id, current_name, new_name)
|
||||
if not ok:
|
||||
return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502
|
||||
|
||||
try:
|
||||
get_sync().trigger_now()
|
||||
except Exception as e:
|
||||
logger.warning(f"纠错后即时拉回失败(下一个周期会自动同步): {e}")
|
||||
|
||||
return jsonify({
|
||||
"status": "ok", "video_id": video_id,
|
||||
"current_name": current_name, "new_name": new_name,
|
||||
}), 200
|
||||
|
||||
@@ -90,22 +90,25 @@ class OracleSync:
|
||||
events = data.get('events', []) or []
|
||||
people = data.get('people', []) or []
|
||||
model_calls = data.get('model_calls', []) or []
|
||||
identity_map = data.get('identity_map', []) or []
|
||||
server_time = data.get('server_time', '') or ''
|
||||
|
||||
n_videos = db_layer.upsert_sync_videos(videos)
|
||||
n_events = db_layer.upsert_sync_events(events)
|
||||
n_people = db_layer.upsert_sync_people(people)
|
||||
n_calls = db_layer.upsert_sync_model_calls(model_calls)
|
||||
n_identity = db_layer.upsert_sync_identity_map(identity_map)
|
||||
|
||||
if server_time:
|
||||
db_layer.set_sync_cursor(server_time)
|
||||
|
||||
self._last_sync_at = datetime.now()
|
||||
self._last_error = None
|
||||
self._last_count = (n_videos, n_events, n_people, n_calls)
|
||||
self._last_count = (n_videos, n_events, n_people, n_calls, n_identity)
|
||||
logger.info(
|
||||
f"同步完成: videos+{n_videos} events+{n_events} people+{n_people} "
|
||||
f"model_calls+{n_calls} since={since!r} -> server_time={server_time}")
|
||||
f"model_calls+{n_calls} identity_map+{n_identity} since={since!r} "
|
||||
f"-> server_time={server_time}")
|
||||
return True
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
@@ -133,6 +136,39 @@ class OracleSync:
|
||||
logger.error(f"命名校正回推失败: {msg}")
|
||||
return False, msg
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
def push_identity_correct(self, video_id: int, current_name: str, new_name: str):
|
||||
"""回推事件时间轴/人物管理"这个人识别错了"纠错到 Oracle。
|
||||
|
||||
跟 push_name_correct 的区别:这个按 (video_id, 当前展示名) 定位,只改这
|
||||
一段视频里错认的那个人,不影响同名字符串在其他视频里的映射(人物 uid
|
||||
只在单次视频分析内稳定,跨视频复用同一字符串完全可能是不同真人)。
|
||||
|
||||
返回 (success: bool, error: str)
|
||||
"""
|
||||
try:
|
||||
video_id = int(video_id)
|
||||
except (TypeError, ValueError):
|
||||
return False, "video_id 必须是数字"
|
||||
current_name = (current_name or '').strip()
|
||||
new_name = (new_name or '').strip()
|
||||
if not current_name or not new_name:
|
||||
return False, "缺少 current_name / new_name"
|
||||
try:
|
||||
resp = requests.post(
|
||||
f"{self.base_url}/api/oracle/identity/correct",
|
||||
json={"video_id": video_id, "current_name": current_name,
|
||||
"new_name": new_name, "token": self.token},
|
||||
timeout=(10, 30))
|
||||
except requests.RequestException as e:
|
||||
logger.error(f"人物纠错回推失败: {e}")
|
||||
return False, str(e)
|
||||
if resp.status_code == 200:
|
||||
return True, ""
|
||||
msg = f"HTTP {resp.status_code}: {resp.text[:200]}"
|
||||
logger.error(f"人物纠错回推失败: {msg}")
|
||||
return False, msg
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
def _run(self):
|
||||
logger.info(f"OracleSync 线程启动,间隔 {self.interval_sec}s,目标 {self.base_url}")
|
||||
|
||||
Reference in New Issue
Block a user