refactor: 全量迁云——fam-core 直读甲骨文 SQLite,NAS 只剩推送进程

前一天刚把登录迁到甲骨文,隔天 NAS 上的 fam-core 又挂了导致数据接口 502。
盘点后确认:甲骨文的 SQLite 才是权威数据源(videos 3113 / events 16159 /
people 60 / model_calls 9876),NAS 的 MariaDB 全是它的镜像——前端读的数据
本来就产自甲骨文,绕了一圈回家又绕回来。

改动:
- db_layer.py 从 725 行重写成 377 行:MySQL 镜像查询改为直读 fam-edge 的
  SQLite。5 个 upsert_sync_*(约 300 行去重逻辑,8/29 和 9/3 两次 1062 事故的
  发源地)连同 oracle_sync.py 整个删除。SQL 方言:JSON_CONTAINS -> json_each
  (前置 json_valid,历史脏数据不会把查询搞崩)、LEFT() -> substr()、%s -> ?。
  函数名 get_sync_* 一并改掉——已经没有 sync 这回事了
- 新增 edge_client.py:写操作(改名/删除)、帧图头像、服务状态都打给同机
  fam-edge,全走 127.0.0.1
- 新增 fam-notifier/:motion_notifier 从 fam-core 拆出独立成服务,游标从
  MariaDB 换成本地 JSON 文件。NAS 上从此没有 Flask、没有数据库、没有监听端口
- fam-core 移到甲骨文 /opt/fam-core(systemd,gunicorn -w 2,只绑
  127.0.0.1:5401——5400 被 chat-relay 占了)。Caddy 的 /api/* 从"frp 隧道
  回源 NAS"改成同机反代,forward_auth 闸门不变
- 前端删掉侧边栏同步面板、统计页同步状态、服务状态页的"NAS 同步"卡片与
  "立即同步"按钮(背后的镜像层已不存在);换成"NAS 运动推送"卡片,读
  fam-edge activity 新增的 motion 段(心跳年龄 + 最近事件)
- 顺带修掉一个隐蔽 bug:镜像表为保外键稳定用的是 NAS 本地自增 id,而帧图接口
  要的是甲骨文的 id,两边在 9/3 那次 id 重排后就对不上了。现在只有一套 id

测试:fam-core 21(新增 12 个 db_layer 用例:脏 JSON 不崩、人物精确匹配不误伤
"人物B"、日期过滤、统计口径、chat_history 懒建表)、fam-notifier 6、
fam-edge 157,全绿。

生产验证:甲骨文 /api/ui/stats 返回 videos 2965 / events 16159 / people 59;
NAS 侧 fam-notifier 已推送成功(事件 33070-33072 落库,心跳新鲜);
chat_history 19 条经 scripts/import_chat_history.py 迁移完成。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
ericwyuan
2026-09-13 08:03:32 +08:00
parent dc284f4bf5
commit 1964e976f4
33 changed files with 1178 additions and 1187 deletions

View File

@@ -1,10 +1,20 @@
"""
FAM-Core 主应用 - Flask 单进程(新架构 v2
FAM-Core 主应用 - Flask2026-09-13 从 NAS 迁到甲骨文
承载: Oracle-Sync每 30 分钟拉取增量镜像)+ Member-Manager + Chat-Handler
承载: UI-API + Member-Manager + Chat-Handler + Img-Proxy —— 全是查询与转发,
**没有任何后台线程**,进程随时可重启,不持有任何状态。
NAS 不再处理视频:无 Scheduler / Dispatcher / Poller / Event-Receiver / Video-Server。
所有视频分析在 Oracle 完成NAS 仅作管理后台拉取展示CPU 占用大幅降低
迁云前它跑在 NAS 上,还扛着两个后台线程,现在都不在这里了:
- Oracle-Sync每 30 分钟把甲骨文数据拉一份镜像进 MariaDB—— 整个删除
本服务现在与 fam-edge 同机,直接读它的 SQLite见 db_layer.py 开头)。
- MotionNotifier轮询 Surveillance Station 推运动事件)—— 留在 NAS
拆成独立的 fam-notifier摄像头插在 NAS 上,这部分搬不走)。
两件跟安全有关的事:
- 本服务只监听 127.0.0.1config.yaml 的 server.host唯一的客户端是同机
Caddy。不像迁云前那样经 frp 把 :8000 暴露到公网。
- 登录校验也不在这里Caddy 用 forward_auth 打 fam-edge 的 /api/auth/verify
通过了才反代进来。
"""
import os
import sys
@@ -16,87 +26,52 @@ sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from .config_loader import load_config
from .logger import setup_logger
from .oracle_sync import get_sync
from .motion_notifier.motion_notifier import get_motion_notifier
from .chat_handler.chat_handler import chat_bp
from .member_manager.member_manager import member_bp
from .img_proxy import img_bp
from .motion_bp import motion_bp
from .ui_api import ui_bp
from .auth import auth_bp, init_auth
from . import db_layer, edge_client
logger = setup_logger('fam-core.app')
app = Flask(__name__)
# 注册蓝图:全部 /api/* 路由(前端已迁云服务器 Caddy:80NAS 不再托管 SPA 静态文件)
app.register_blueprint(chat_bp)
app.register_blueprint(member_bp)
app.register_blueprint(img_bp)
app.register_blueprint(motion_bp)
app.register_blueprint(ui_bp)
app.register_blueprint(auth_bp)
# 登录校验全局拦截(页面未登录 302 /login/api/* 未登录 401
init_auth(app)
# 健康检查
@app.route('/health', methods=['GET'])
def health():
return jsonify({"status": "ok", "service": "fam-core"}), 200
# 初始化后台同步线程NAS 唯一常驻线程)
_sync = None
try:
_sync = get_sync()
_sync.start()
logger.info("Oracle-Sync 已启动")
# 启动后立刻拉一次,前端无需等待首个周期即有数据
try:
_sync.trigger_now()
logger.info("启动首次同步完成")
except Exception as e:
logger.warning(f"启动首次同步失败(后续周期会重试): {e}")
except Exception as e:
logger.error(f"Oracle-Sync 启动失败: {e}")
# 初始化运动监测通知服务NAS 轮询 SS EventCenter 主路径)
_notifier = None
try:
_notifier = get_motion_notifier()
_notifier.start()
logger.info("MotionNotifier 已初始化")
except Exception as e:
logger.error(f"MotionNotifier 初始化失败: {e}")
@app.route('/api/status', methods=['GET'])
def status():
"""系统状态"""
"""服务自检库能不能读、fam-edge 在不在。
迁云前这里报告的是镜像同步状态(游标/周期/上次增量条数),镜像没了之后
那些字段不再存在,前端侧边栏的同步面板也一并去掉了。
"""
db_ok, db_err = True, None
try:
conn = db_layer.get_conn()
try:
conn.execute("SELECT 1 FROM videos LIMIT 1").fetchone()
finally:
conn.close()
except Exception as e:
db_ok, db_err = False, str(e)
return jsonify({
"service": "fam-core",
"sync": _sync.status() if _sync else {"running": False, "error": "未初始化"},
"motion": _notifier.status() if _notifier else {"running": False, "error": "未初始化"},
"db": {"ok": db_ok, "error": db_err},
"edge_base_url": edge_client.base_url(),
}), 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)
app.run(host='0.0.0.0', port=port, debug=False)
server = cfg.get('server', {})
app.run(host=server.get('host', '127.0.0.1'),
port=server.get('port', 5401), debug=False)

View File

@@ -47,7 +47,7 @@ def _call_edge_qa(prompt: str) -> str:
"""调用 FAM-Edge 问答编排端点Gemini → NVIDIA → 本地 Ollama 兜底)"""
cfg = load_config()
qa_url = cfg.get('chat_handler', {}).get(
'qa_url', 'http://129.146.26.249:5000/api/edge/chat/ask'
'qa_url', 'http://127.0.0.1:5000/api/edge/chat/ask'
)
timeout = cfg.get('chat_handler', {}).get('timeout', 120)
@@ -79,7 +79,7 @@ def chat_ask():
logger.info(f"Chat: person={queried_person}, date={queried_date}, question={question}")
rows = db_layer.query_sync_events_for_person_date(queried_person, queried_date)
rows = db_layer.query_events_for_person_date(queried_person, queried_date)
if len(rows) == 0:
answer = f"今天没有观察到{queried_person}"
@@ -87,7 +87,7 @@ def chat_ask():
else:
context = _format_events(rows)
context_summary = f"查询 sync_events {len(rows)}"
members = db_layer.get_sync_known_members_context()
members = db_layer.get_known_members_context()
prompt = build_chat_prompt(
context=context,
members=members or queried_person,
@@ -134,7 +134,7 @@ def chat_ask_stream():
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)
rows = db_layer.query_events_for_person_date(queried_person, queried_date)
def sse(obj):
return f"data: {json.dumps(obj, ensure_ascii=False)}\n\n"
@@ -157,14 +157,14 @@ def chat_ask_stream():
yield sse({"type": "context", "count": len(rows), "summary": context_summary,
"preview": context[:800]})
members = db_layer.get_sync_known_members_context()
members = db_layer.get_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.26.249:5000/api/edge/chat/ask/stream')
'qa_stream_url', 'http://127.0.0.1:5000/api/edge/chat/ask/stream')
timeout = cfg.get('chat_handler', {}).get('timeout', 120)
full_answer = []

View File

@@ -1,11 +1,42 @@
"""
配置加载器 - 从 config.yaml 读取配置
配置加载器 - 从 config.yaml 读取配置,支持 ${ENV_VAR} 解析
"""
import os
import re
import yaml
def _load_env_file():
"""加载部署目录下的 .env支持 export KEY=VALUE 格式)。
迁云前是 start_core.sh 负责 source .env 再起 gunicorn现在由 systemd 拉起,
没有那一步,所以在这里兜底加载(跟 fam-edge 的做法一致)。
已存在的环境变量不覆盖。
"""
path = os.environ.get('FAM_ENV_FILE') or os.path.join(
os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))),
'.env')
if not os.path.isfile(path):
return
with open(path, 'r', encoding='utf-8') as f:
for line in f:
line = line.strip()
if not line or line.startswith('#'):
continue
if line.startswith('export '):
line = line[7:].strip()
if '=' not in line:
continue
key, _, value = line.partition('=')
key = key.strip()
value = value.strip().strip('"').strip("'")
if key and key not in os.environ:
os.environ[key] = value
_load_env_file()
def _resolve_env_vars(value):
"""递归解析字符串中的 ${ENV_VAR} 引用"""
if isinstance(value, str):

View File

@@ -1,300 +1,176 @@
"""
数据库访问层 - MariaDB 连接管理与同步镜像 CRUD
DB-Layer - 直读 fam-edge 的 SQLite 库2026-09-13 迁云重写)
新架构 (2026-08-21 重构):
NAS 不再处理视频,仅作为管理后台。
Oracle (FAM-Edge) 处理整视频分析后存 SQLiteNAS 每 30 分钟拉增量,
镜像到本地三张表:
sync_videos : 视频会话(全局摘要 + 事件 JSON + 人物 JSON
sync_events : 视频拆出的时间点事件(描述 + 涉及人物 + 是否关注)
sync_people : 规范人物表label + canonical_nameOracle 维护)
sync_cursor : 同步游标(上次成功拉取到的 server_time
**这个文件以前是什么样**fam-core 跑在 NAS 上,每 30 分钟把甲骨文的数据拉一份
镜像进 MariaDB`sync_videos`/`sync_events`/…前端读镜像。725 行里有一半是
镜像 upsert 的去重逻辑——8/29 和 9/3 两次线上事故1062 主键冲突、游标卡死)
都出在那一半。
本层只服务同步镜像 + 问答历史,旧 process_tasks/event_details/monitor_events/
family_members 相关逻辑已全部移除(视频处理职责已迁移至 Oracle
**现在**fam-core 跟 fam-edge 同机,直接打开它的 SQLite 库读,镜像层整个不存在了。
连带消失的还有一个隐蔽 bug镜像表为了保外键稳定用的是 NAS 本地自增 id而帧图
接口要的是甲骨文的 id两边在 9/3 那次 id 重排后就对不上了。现在只有一套 id。
要点:
- 库文件是 fam-edge 的(`database.path`**本模块只读它写的表**,唯一写的表是
`chat_history`(问答历史,迁云时从 NAS MariaDB 搬过来的fam-edge 不碰)。
- fam-edge 那边开了 WAL读不会阻塞它的写这边同样设 busy_timeout 兜底。
- 视频/人物的写操作(改名、删除)不在这里做,走 `edge_client` 打给 fam-edge——
它除了改库还要合并人物、删磁盘素材。
"""
import json
import pymysql
from datetime import datetime
from typing import Optional, List, Dict, Any
import sqlite3
from datetime import datetime, timedelta, timezone
from typing import Dict, List, Optional
from .config_loader import load_config
from .logger import setup_logger
logger = setup_logger('fam-core.db')
_config = None
_db_path = None
_chat_schema_ready = False
def get_config():
global _config
if _config is None:
_config = load_config()
return _config
def _now() -> str:
return datetime.now(timezone(timedelta(hours=8))).strftime('%Y-%m-%d %H:%M:%S')
def get_conn():
"""获取数据库连接(单 worker gunicorn无需连接池"""
cfg = get_config().get('database', {})
kwargs = dict(
host=cfg.get('host', '127.0.0.1'),
port=cfg.get('port', 3306),
user=cfg.get('user', 'root'),
password=cfg.get('password', ''),
database=cfg.get('database', 'sentinel_home_ai'),
charset='utf8mb4',
autocommit=False
)
unix_socket = cfg.get('unix_socket')
if unix_socket:
kwargs['unix_socket'] = unix_socket
return pymysql.connect(**kwargs)
def get_config() -> Dict:
return load_config()
def _path() -> str:
global _db_path
if _db_path is None:
_db_path = load_config().get('database', {}).get(
'path', '/opt/fam-edge/data/oracle.db')
return _db_path
def get_conn() -> sqlite3.Connection:
"""每次调用开一个连接(请求级,跟改写前的 MySQL 用法一致)。
busy_timeoutfam-edge 的写事务提交时会短暂持锁,这里等而不是立刻报
`database is locked`9/3 那次 Oracle 过载时刷过一片这个错)。
"""
conn = sqlite3.connect(_path(), timeout=10)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA busy_timeout=10000")
return conn
def _rows(cur) -> List[Dict]:
return [dict(r) for r in cur.fetchall()]
def _row(cur) -> Optional[Dict]:
r = cur.fetchone()
return dict(r) if r else None
# 人物命中判断person_list_json 是 JSON 数组文本,用 json_each 展开精确匹配。
# 必须先 json_valid——历史数据里有非 JSON 的脏值,直接 json_each 会整条查询报错。
_PERSON_HIT = ("{col} IS NOT NULL AND json_valid({col}) "
"AND EXISTS(SELECT 1 FROM json_each({col}) WHERE json_each.value = ?)")
# 日期维度:录制时间优先(文件名解析出来的),回退分析时间
_DATE_EXPR = "COALESCE(NULLIF({a}.event_start_time,''), {a}.processed_at, {a}.updated_at, {a}.created_at)"
# 只展示"有内容"的会话:运动片段,或含事件的视频。
# 整段素材分割 0 段的空会话(历史素材无运动事件)不展示,避免淹没时间轴。
_CONTENT_FILTER = ("substr(v.filename, 1, 7) = 'motion_' "
"OR EXISTS(SELECT 1 FROM events se WHERE se.video_id = v.id)")
# ============================================================
# 同步镜像sync_videos
# videos / events
# ============================================================
def upsert_sync_videos(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 videos 增量。rows 为 Oracle 端 dict 列表。
v2 修复 (2026-09-03):不再以 Oracle 的 videos.id 作为 NAS 主键。Oracle 库重建/
重排后 id 会复用/跳变(实测 2303 段跳到 3408/4600+ 段),旧写法按 id 插入时与
已存在的同 filename 行在 UNIQUE filename 上二次冲突报 1062整批中止、游标不推进。
现改为:以 filename 为业务唯一键去重,命中则就地 UPDATE保留 NAS 原 id防止
sync_events/model_calls/identity_map 的 video_id 外键失效Oracle id 仅落
oracle_id 列溯源;未命中则插入(优先用 Oracle id 作主键以对齐子表引用,主键冲突时
回退本地自增,避免 1062
"""
if not rows:
return 0
conn = get_conn()
n = 0
try:
cur = conn.cursor()
for r in rows:
fn = r.get('filename')
oracle_id = r.get('id')
cur.execute("SELECT id FROM sync_videos WHERE filename=%s", (fn,))
row = cur.fetchone()
if row:
# 命中业务键:就地更新,保留原 NAS id外键不失效
cur.execute(
"""UPDATE sync_videos SET
oracle_id=%s, drive_file_id=%s, camera_name=%s,
duration_sec=%s, event_start_time=%s, status=%s,
summary_json=%s, events_json=%s, people_json=%s,
compute_provider=%s, created_at=%s, updated_at=%s,
processed_at=%s, synced_at=NOW()
WHERE id=%s""",
(oracle_id, r.get('drive_file_id'), r.get('camera_name'),
r.get('duration_sec') or 0, r.get('event_start_time'),
r.get('status'), r.get('summary_json'), r.get('events_json'),
r.get('people_json'), r.get('compute_provider'),
r.get('created_at'), r.get('updated_at'),
r.get('processed_at'), row[0]))
else:
# 未命中:新文件。优先用 Oracle id 作主键(对齐子表 video_id 引用)
try:
cur.execute(
"""INSERT INTO sync_videos
(id, oracle_id, filename, camera_name, duration_sec,
event_start_time, status, summary_json, events_json,
people_json, compute_provider, created_at, updated_at,
processed_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, oracle_id, fn, r.get('camera_name'),
r.get('duration_sec') or 0, r.get('event_start_time'),
r.get('status'), r.get('summary_json'), r.get('events_json'),
r.get('people_json'), r.get('compute_provider'),
r.get('created_at'), r.get('updated_at'),
r.get('processed_at')))
except Exception:
# 主键 id 被复用(极端情况):回退本地自增,避免 1062
cur.execute(
"""INSERT INTO sync_videos
(oracle_id, filename, camera_name, duration_sec,
event_start_time, status, summary_json, events_json,
people_json, compute_provider, created_at, updated_at,
processed_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, fn, r.get('camera_name'),
r.get('duration_sec') or 0, r.get('event_start_time'),
r.get('status'), r.get('summary_json'), r.get('events_json'),
r.get('people_json'), r.get('compute_provider'),
r.get('created_at'), r.get('updated_at'),
r.get('processed_at')))
n += 1
conn.commit()
return n
finally:
conn.close()
def get_sync_videos(limit=15, offset=0, date_filter=None) -> List[Dict]:
"""获取视频会话列表(已完成优先),支持日期筛选与分页。
排序/日期维度按视频实际录制时间event_start_time文件名解析
为空回退 processed_at/updated_at/created_at。
date_filter 形如 '2026-08-21'
只返回"有内容"的会话运动片段filename 前缀 motion_或含事件的视频
整段素材分割 0 段的空会话(历史素材无运动事件)不展示,避免淹没时间轴。
"""
def get_videos(limit=15, offset=0, date_filter=None) -> List[Dict]:
"""视频会话列表(事件时间轴左侧),按录制时间倒序,支持日期筛选与分页。"""
d = _DATE_EXPR.format(a='v')
sql = f"""SELECT v.id, v.filename, v.camera_name, v.event_start_time, v.status,
v.summary_json, v.events_json, v.people_json, v.compute_provider,
v.processed_at, v.updated_at,
(SELECT COUNT(*) FROM events se WHERE se.video_id = v.id) AS event_count
FROM videos v
WHERE v.status='done' AND ({_CONTENT_FILTER}) {{extra}}
ORDER BY {d} DESC
LIMIT ? OFFSET ?"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
# 录制时间优先,回退分析时间
date_expr = "COALESCE(NULLIF(event_start_time,''), processed_at, updated_at, created_at)"
# 注意pymysql execute 用 % 做参数占位符SQL 字面量不能含 %,故用 LEFT 判断 motion_ 前缀
content_filter = ("LEFT(v.filename, 7) = 'motion_' "
"OR EXISTS(SELECT 1 FROM sync_events se WHERE se.video_id = v.id)")
if date_filter:
cur.execute(
f"""SELECT v.id, v.filename, v.camera_name, v.event_start_time, v.status,
v.summary_json, v.events_json, v.people_json, v.compute_provider,
v.processed_at, v.updated_at,
(SELECT COUNT(*) FROM sync_events se WHERE se.video_id = v.id) AS event_count
FROM sync_videos v
WHERE v.status='done' AND ({content_filter}) AND {date_expr} LIKE %s
ORDER BY {date_expr} DESC
LIMIT %s OFFSET %s""",
(f'{date_filter}%', limit, offset))
cur = conn.execute(sql.format(extra=f"AND {d} LIKE ?"),
(f'{date_filter}%', limit, offset))
else:
cur.execute(
f"""SELECT v.id, v.filename, v.camera_name, v.event_start_time, v.status,
v.summary_json, v.events_json, v.people_json, v.compute_provider,
v.processed_at, v.updated_at,
(SELECT COUNT(*) FROM sync_events se WHERE se.video_id = v.id) AS event_count
FROM sync_videos v
WHERE v.status='done' AND ({content_filter})
ORDER BY {date_expr} DESC
LIMIT %s OFFSET %s""",
(limit, offset))
return cur.fetchall()
cur = conn.execute(sql.format(extra=''), (limit, offset))
return _rows(cur)
finally:
conn.close()
def get_sync_video(video_id: int) -> Optional[Dict]:
def get_video(video_id: int) -> Optional[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute("SELECT * FROM sync_videos WHERE id=%s", (video_id,))
return cur.fetchone()
return _row(conn.execute("SELECT * FROM videos WHERE id=?", (video_id,)))
finally:
conn.close()
def delete_sync_video(video_id: int):
"""删除本地镜像里的这条视频events 先删,再删 videos
增量同步get_sync_delta只做 upsert感知不到 Oracle 那边的物理删除,
所以删除操作不能走"转发 Oracle + trigger_now() 拉增量"这条老路,必须在
Oracle 确认删除成功后由调用方显式清理本地镜像。
"""
def get_events_for_video(video_id: int) -> List[Dict]:
conn = get_conn()
try:
cur = conn.cursor()
cur.execute("DELETE FROM sync_events WHERE video_id=%s", (video_id,))
cur.execute("DELETE FROM sync_videos WHERE id=%s", (video_id,))
conn.commit()
finally:
conn.close()
# ============================================================
# 同步镜像sync_events
# ============================================================
def upsert_sync_events(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 events 增量。"""
if not rows:
return 0
conn = get_conn()
n = 0
try:
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_events
(id, video_id, ts, description, person_list_json,
person_appearances_json, is_attention_event,
updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
video_id=VALUES(video_id),
ts=VALUES(ts),
description=VALUES(description),
person_list_json=VALUES(person_list_json),
person_appearances_json=VALUES(person_appearances_json),
is_attention_event=VALUES(is_attention_event),
updated_at=VALUES(updated_at),
synced_at=NOW()""",
(r.get('id'), r.get('video_id'), r.get('ts'), r.get('description'),
r.get('person_list_json'), r.get('person_appearances_json'),
1 if r.get('is_attention_event') else 0,
r.get('updated_at'))
)
n += 1
conn.commit()
return n
finally:
conn.close()
def get_sync_events_for_video(video_id: int) -> List[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
return _rows(conn.execute(
"""SELECT id, video_id, ts, description, person_list_json,
is_attention_event, updated_at
FROM sync_events WHERE video_id=%s ORDER BY ts ASC""",
(video_id,))
return cur.fetchall()
is_attention_event
FROM events WHERE video_id=? ORDER BY ts ASC""", (video_id,)))
finally:
conn.close()
def get_sync_people_clips(label: str, limit: int = 10) -> List[Dict]:
"""某人物label 或 canonical_name出现过的运动片段列表。
匹配label 本身 + 其 canonical_name 下全部 label命中 sync_events 的
person_list_jsonJSON_CONTAINS→ 关联 video优先运动片段 motion_
返回按 event_start_time 倒序,含 first_ts该人物在片段内最早事件 ts
前端缩略图定位用)与 clip_events片段内该人物事件数
"""
def get_attention_events(limit: int = 200) -> List[Dict]:
"""需关注事件(统计图表页用),按录制日期倒序。"""
conn = get_conn()
try:
return _rows(conn.execute(
"""SELECT COALESCE(NULLIF(v.event_start_time,''), v.processed_at) AS ev_date,
e.person_list_json
FROM events e JOIN videos v ON e.video_id=v.id
WHERE e.is_attention_event = 1
ORDER BY ev_date DESC LIMIT ?""", (limit,)))
finally:
conn.close()
def get_people_clips(label: str, limit: int = 10) -> List[Dict]:
"""某人物出现过的运动片段列表。
匹配 label 本身 + 同一 canonical_name 下的全部 label返回含 first_ts
(该人物在片段内最早事件时间点,前端缩略图定位用)与 clip_events。
"""
hit_e = _PERSON_HIT.format(col='e.person_list_json')
hit_e2 = _PERSON_HIT.format(col='e2.person_list_json')
hit_e3 = _PERSON_HIT.format(col='e3.person_list_json')
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
labels = {label}
cur.execute(
"SELECT label, canonical_name FROM sync_people WHERE label=%s OR canonical_name=%s",
(label, label))
for r in cur.fetchall():
cn = r.get('canonical_name')
if cn:
cur.execute("SELECT label FROM sync_people WHERE canonical_name=%s", (cn,))
labels.update(x['label'] for x in cur.fetchall())
clips = []
seen = set()
for r in _rows(conn.execute(
"SELECT label, canonical_name FROM people WHERE label=? OR canonical_name=?",
(label, label))):
if r.get('canonical_name'):
labels.update(x['label'] for x in _rows(conn.execute(
"SELECT label FROM people WHERE canonical_name=?", (r['canonical_name'],))))
clips, seen = [], set()
for lb in sorted(labels):
cur.execute(
"""SELECT DISTINCT v.id AS video_id, v.filename, v.event_start_time,
v.duration_sec, v.summary_json, v.camera_name,
(SELECT MIN(e2.ts) FROM sync_events e2
WHERE e2.video_id=v.id AND e2.person_list_json IS NOT NULL
AND JSON_CONTAINS(e2.person_list_json, JSON_QUOTE(%s), '$')) AS first_ts,
(SELECT COUNT(*) FROM sync_events e3
WHERE e3.video_id=v.id AND e3.person_list_json IS NOT NULL
AND JSON_CONTAINS(e3.person_list_json, JSON_QUOTE(%s), '$')) AS clip_events
FROM sync_events e
JOIN sync_videos v ON e.video_id=v.id
WHERE e.person_list_json IS NOT NULL
AND JSON_CONTAINS(e.person_list_json, JSON_QUOTE(%s), '$')
AND v.status='done'""",
cur = conn.execute(
f"""SELECT DISTINCT v.id AS video_id, v.filename, v.event_start_time,
v.duration_sec, v.summary_json, v.camera_name,
(SELECT MIN(e2.ts) FROM events e2
WHERE e2.video_id=v.id AND {hit_e2}) AS first_ts,
(SELECT COUNT(*) FROM events e3
WHERE e3.video_id=v.id AND {hit_e3}) AS clip_events
FROM events e JOIN videos v ON e.video_id=v.id
WHERE {hit_e} AND v.status='done'""",
(lb, lb, lb))
for row in cur.fetchall():
for row in _rows(cur):
if row['video_id'] not in seen:
seen.add(row['video_id'])
clips.append(row)
@@ -304,397 +180,177 @@ def get_sync_people_clips(label: str, limit: int = 10) -> List[Dict]:
conn.close()
def query_sync_events_for_person_date(person: str, date_str: str) -> List[Dict]:
def query_events_for_person_date(person: str, date_str: str) -> List[Dict]:
"""问答上下文:某人在某天的事件。
说明: Oracle 事件 ts 视频内相对时间点(如 00:01:23不是绝对日期
因此按所属视频的录制日期event_start_time回退 processed_at过滤
再按 person_list_json 命中人名。
person 可为真名或抽象标签Oracle 回灌上下文用真名,但历史标签也保留)。
事件 ts 视频内相对时间点(如 00:01:23不是绝对日期所以按所属视频的
录制日期过滤,再按 person_list_json 命中人名。person 可为真名或抽象标签。
"""
hit = _PERSON_HIT.format(col='e.person_list_json')
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"""SELECT e.ts, e.description, e.person_list_json, e.is_attention_event,
v.camera_name, v.filename, v.event_start_time, v.processed_at
FROM sync_events e
JOIN sync_videos v ON e.video_id = v.id
WHERE COALESCE(NULLIF(v.event_start_time,''), v.processed_at) LIKE %s
AND e.person_list_json IS NOT NULL
AND JSON_CONTAINS(e.person_list_json, JSON_QUOTE(%s), '$')
ORDER BY COALESCE(NULLIF(v.event_start_time,''), v.processed_at) ASC, e.ts ASC""",
(f'{date_str}%', person))
return cur.fetchall()
return _rows(conn.execute(
f"""SELECT e.ts, e.description, e.person_list_json, e.is_attention_event,
v.camera_name, v.filename, v.event_start_time, v.processed_at
FROM events e JOIN videos v ON e.video_id = v.id
WHERE COALESCE(NULLIF(v.event_start_time,''), v.processed_at) LIKE ?
AND {hit}
ORDER BY COALESCE(NULLIF(v.event_start_time,''), v.processed_at) ASC, e.ts ASC""",
(f'{date_str}%', person)))
finally:
conn.close()
# ============================================================
# 同步镜像sync_people
# people
# ============================================================
def upsert_sync_people(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 people 增量。
v2 修复 (2026-09-03):同 videos以 label 为业务唯一键去重,命中就地 UPDATE
保留 NAS 原 idpeople.id 无外键引用仅作展示排序Oracle id 落 oracle_id 列
溯源;未命中插入(优先 Oracle id主键冲突回退自增避免 1062
"""
if not rows:
return 0
conn = get_conn()
n = 0
try:
cur = conn.cursor()
for r in rows:
label = r.get('label')
oracle_id = r.get('id')
cur.execute("SELECT id FROM sync_people WHERE label=%s", (label,))
row = cur.fetchone()
if row:
cur.execute(
"""UPDATE sync_people SET
oracle_id=%s, canonical_name=%s, first_seen=%s,
appearances=%s, source=%s, features_json=%s,
display_uid=%s, updated_at=%s, synced_at=NOW()
WHERE id=%s""",
(oracle_id, r.get('canonical_name'), r.get('first_seen'),
r.get('appearances') or 0, r.get('source'),
r.get('features_json'), r.get('display_uid'),
r.get('updated_at'), row[0]))
else:
try:
cur.execute(
"""INSERT INTO sync_people
(id, oracle_id, label, canonical_name, first_seen,
appearances, source, features_json, display_uid,
updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, oracle_id, label, r.get('canonical_name'),
r.get('first_seen'), r.get('appearances') or 0,
r.get('source'), r.get('features_json'),
r.get('display_uid'), r.get('updated_at')))
except Exception:
cur.execute(
"""INSERT INTO sync_people
(oracle_id, label, canonical_name, first_seen,
appearances, source, features_json, display_uid,
updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())""",
(oracle_id, label, r.get('canonical_name'),
r.get('first_seen'), r.get('appearances') or 0,
r.get('source'), r.get('features_json'),
r.get('display_uid'), r.get('updated_at')))
n += 1
conn.commit()
return n
finally:
conn.close()
def get_sync_people() -> List[Dict]:
def get_people() -> List[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
return _rows(conn.execute(
"SELECT id, label, canonical_name, first_seen, appearances, source, "
"features_json, display_uid, updated_at "
"FROM sync_people ORDER BY id ASC")
return cur.fetchall()
"features_json, display_uid, updated_at FROM people ORDER BY id ASC"))
finally:
conn.close()
def upsert_sync_model_calls(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 model_calls 增量(幂等,重复覆盖)。"""
if not rows:
return 0
def get_named_members() -> List[str]:
"""已命名成员的真名列表UI 下拉用)。"""
conn = get_conn()
n = 0
try:
cur = conn.cursor()
for r in rows:
cur.execute(
"""INSERT INTO sync_model_calls
(id, provider, model, video_id, filename, started_at,
duration_sec, success, error, created_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
provider=VALUES(provider),
model=VALUES(model),
video_id=VALUES(video_id),
filename=VALUES(filename),
started_at=VALUES(started_at),
duration_sec=VALUES(duration_sec),
success=VALUES(success),
error=VALUES(error),
created_at=VALUES(created_at),
synced_at=NOW()""",
(r.get('id'), r.get('provider'), r.get('model'),
r.get('video_id'), r.get('filename'), r.get('started_at'),
r.get('duration_sec') or 0, 1 if r.get('success') else 0,
(r.get('error') or '')[:500], r.get('created_at')))
n += 1
conn.commit()
return n
return [r['canonical_name'] for r in _rows(conn.execute(
"SELECT DISTINCT canonical_name FROM people "
"WHERE canonical_name IS NOT NULL AND canonical_name != '' "
"ORDER BY canonical_name ASC"))]
finally:
conn.close()
def upsert_sync_identity_map(rows: List[Dict]) -> int:
"""批量 upsert Oracle 传来的 person_identity_map 增量(幂等,重复覆盖)。
v2 修复 (2026-08-29):不再以 Oracle 的 person_identity_map.id 作为 NAS 主键。
该 id 只是 Oracle 端代理主键,一旦 Oracle 库重建/恢复id 会被复用(实测出现
id=744 同时对应 NAS 上的 (1302,'人物E') 与 Oracle 现在的 (1073,'汤圆')
旧写法会先按 id 撞主键、再 UPDATE 成新组合,触发 uq_video_raw_uid 二次冲突
报 1062。现改为NAS 本地自增 id 做主键,业务键 (video_id, raw_uid) 唯一,
Oracle 的 id 只落 oracle_id 列作溯源参考。
"""
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
(oracle_id, video_id, raw_uid, canonical_name, source, updated_at, synced_at)
VALUES (%s,%s,%s,%s,%s,%s, NOW())
ON DUPLICATE KEY UPDATE
oracle_id=VALUES(oracle_id),
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_known_members_context() -> str:
"""人物清单文本,注入问答 Prompt让模型用真名指代。"""
rows = get_people()
parts = []
for r in rows:
name = r['canonical_name'] or r['label']
if r['canonical_name'] and r['canonical_name'] != r['label']:
parts.append(f"- {name}(标识:{r['label']}")
else:
parts.append(f"- {name}")
return "\n".join(parts)
def get_sync_model_calls(limit: int = 200) -> List[Dict]:
"""最近模型调用记录(前端统计展示)。"""
# ============================================================
# model_calls
# ============================================================
def get_model_calls(limit: int = 200) -> List[Dict]:
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
return _rows(conn.execute(
"SELECT id, provider, model, video_id, filename, started_at, "
"duration_sec, success, error, created_at "
"FROM sync_model_calls ORDER BY id DESC LIMIT %s", (limit,))
return cur.fetchall()
"FROM model_calls ORDER BY id DESC LIMIT ?", (limit,)))
finally:
conn.close()
def get_sync_model_calls_stats() -> Dict:
"""模型调用统计:按 provider+model 聚合成功/失败/平均耗时。"""
def get_model_calls_stats() -> List[Dict]:
"""按 provider+model 聚合成功/失败/平均耗时。"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
return _rows(conn.execute(
"""SELECT provider, model,
SUM(CASE WHEN success=1 THEN 1 ELSE 0 END) AS ok_cnt,
SUM(CASE WHEN success=0 THEN 1 ELSE 0 END) AS fail_cnt,
ROUND(AVG(duration_sec), 1) AS avg_duration,
MAX(created_at) AS last_call
FROM sync_model_calls
GROUP BY provider, model ORDER BY provider, model""")
return cur.fetchall()
finally:
conn.close()
def get_sync_named_members() -> List[str]:
"""已命名成员的真名列表(供 UI 下拉 / 快捷选择)。"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"SELECT DISTINCT canonical_name FROM sync_people "
"WHERE canonical_name IS NOT NULL AND canonical_name != '' "
"ORDER BY canonical_name ASC")
return [r['canonical_name'] for r in cur.fetchall()]
finally:
conn.close()
def get_sync_known_members_context() -> str:
"""获取人物清单文本,注入问答 Prompt让模型用真名指代。"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"SELECT label, canonical_name FROM sync_people ORDER BY id ASC")
rows = cur.fetchall()
if not rows:
return ""
parts = []
for r in rows:
name = r['canonical_name'] or r['label']
if r['canonical_name'] and r['canonical_name'] != r['label']:
parts.append(f"- {name}(标识:{r['label']}")
else:
parts.append(f"- {name}")
return "\n".join(parts)
FROM model_calls GROUP BY provider, model ORDER BY provider, model"""))
finally:
conn.close()
# ============================================================
# 同步游标
# 概览统计
# ============================================================
def get_sync_cursor() -> str:
def get_stats(date_str: str = None) -> Dict:
"""视频数 / 事件数 / 关注事件数 / 出现人物数(可按日期过滤)。"""
d = _DATE_EXPR.format(a='sv')
conn = get_conn()
try:
cur = conn.cursor()
cur.execute("SELECT `value` FROM sync_cursor WHERE `key`='last_since'")
row = cur.fetchone()
return row[0] if row else ''
finally:
conn.close()
def set_sync_cursor(value: str):
conn = get_conn()
try:
cur = conn.cursor()
cur.execute(
"""INSERT INTO sync_cursor (`key`, `value`) VALUES ('last_since', %s)
ON DUPLICATE KEY UPDATE `value`=VALUES(`value`)""",
(value,))
conn.commit()
finally:
conn.close()
# ============================================================
# 运动通知游标MotionNotifier 增量去重用,复用 sync_cursor 表)
# ============================================================
def get_motion_cursor() -> str:
"""返回上次推送到甲骨文的最大 SS 事件 id字符串无则返回空。"""
conn = get_conn()
try:
cur = conn.cursor()
cur.execute("SELECT `value` FROM sync_cursor WHERE `key`='motion_last_event_id'")
row = cur.fetchone()
return row[0] if row else ''
finally:
conn.close()
def set_motion_cursor(value: str):
conn = get_conn()
try:
cur = conn.cursor()
cur.execute(
"""INSERT INTO sync_cursor (`key`, `value`) VALUES ('motion_last_event_id', %s)
ON DUPLICATE KEY UPDATE `value`=VALUES(`value`)""",
(value,))
conn.commit()
finally:
conn.close()
# ============================================================
# 统计
# ============================================================
def get_attention_events(limit: int = 200) -> List[Dict]:
"""需关注事件列表(供统计图表页展示日期 + 涉及人物),按录制日期倒序。"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
cur.execute(
"""SELECT COALESCE(NULLIF(v.event_start_time,''), v.processed_at) AS ev_date,
e.person_list_json
FROM sync_events e JOIN sync_videos v ON e.video_id=v.id
WHERE e.is_attention_event = 1
ORDER BY ev_date DESC LIMIT %s""",
(limit,))
return cur.fetchall()
finally:
conn.close()
def get_sync_stats(date_str: str = None) -> Dict:
"""概览统计:视频数 / 事件数 / 关注事件数 / 出现人物数(按 date 可选过滤)。
人物数通过对 sync_events.person_list_json LIKE 统计MariaDB 不支持 JSON 数组展开)。
"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
# 视频/事件/关注数(日期维度=录制时间 event_start_time回退 processed_at
_D = "COALESCE(NULLIF(sv.event_start_time,''), sv.processed_at)"
if date_str:
cur.execute(
"""SELECT
(SELECT COUNT(*) FROM sync_videos sv
WHERE sv.status='done'
AND (LEFT(sv.filename, 7) = 'motion_'
OR EXISTS(SELECT 1 FROM sync_events se WHERE se.video_id=sv.id))
AND {d} LIKE %s) AS videos,
(SELECT COUNT(*) FROM sync_events se
JOIN sync_videos sv ON se.video_id=sv.id
WHERE {d} LIKE %s) AS events,
(SELECT COALESCE(SUM(se.is_attention_event),0) FROM sync_events se
JOIN sync_videos sv ON se.video_id=sv.id
WHERE {d} LIKE %s) AS attention""".format(d=_D),
(f'{date_str}%', f'{date_str}%', f'{date_str}%'))
like = f'{date_str}%'
stat = _row(conn.execute(
f"""SELECT
(SELECT COUNT(*) FROM videos sv
WHERE sv.status='done'
AND (substr(sv.filename,1,7)='motion_'
OR EXISTS(SELECT 1 FROM events se WHERE se.video_id=sv.id))
AND {d} LIKE ?) AS videos,
(SELECT COUNT(*) FROM events se JOIN videos sv ON se.video_id=sv.id
WHERE {d} LIKE ?) AS events,
(SELECT COALESCE(SUM(se.is_attention_event),0) FROM events se
JOIN videos sv ON se.video_id=sv.id
WHERE {d} LIKE ?) AS attention""",
(like, like, like))) or {}
else:
cur.execute(
stat = _row(conn.execute(
"""SELECT
(SELECT COUNT(*) FROM sync_videos
WHERE status='done'
AND (LEFT(filename, 7) = 'motion_'
OR EXISTS(SELECT 1 FROM sync_events se WHERE se.video_id=sync_videos.id))) AS videos,
(SELECT COUNT(*) FROM sync_events) AS events,
(SELECT COALESCE(SUM(is_attention_event),0) FROM sync_events) AS attention""")
stat = cur.fetchone() or {}
(SELECT COUNT(*) FROM videos sv
WHERE sv.status='done'
AND (substr(sv.filename,1,7)='motion_'
OR EXISTS(SELECT 1 FROM events se WHERE se.video_id=sv.id))) AS videos,
(SELECT COUNT(*) FROM events) AS events,
(SELECT COALESCE(SUM(is_attention_event),0) FROM events) AS attention""")) or {}
# 人物数distinct label 命中 sync_events
cur.execute("SELECT id, label, canonical_name FROM sync_people")
people = cur.fetchall()
# 基于 person_list_json 命中计数:逐 label 统计命中事件数
person_hits = 0
for p in people:
# 出现人物数:逐 label/真名在事件里查有没有命中(沿用改写前的口径
hits = 0
for p in _rows(conn.execute("SELECT label, canonical_name FROM people")):
name = p['canonical_name'] or p['label']
cur.execute(
"SELECT COUNT(*) c FROM sync_events WHERE person_list_json LIKE %s",
(f'%{name}%',))
if cur.fetchone()['c'] > 0:
person_hits += 1
stat['people'] = person_hits
c = conn.execute("SELECT COUNT(*) AS c FROM events "
"WHERE person_list_json LIKE ?", (f'%{name}%',)).fetchone()
if c and c['c'] > 0:
hits += 1
stat['people'] = hits
return stat
finally:
conn.close()
# ============================================================
# chat_history保留:问答历史
# chat_history本模块唯一写的表2026-09-13 从 NAS MariaDB 迁入
# ============================================================
def _ensure_chat_schema(conn):
global _chat_schema_ready
if _chat_schema_ready:
return
conn.executescript("""
CREATE TABLE IF NOT EXISTS chat_history (
chat_id INTEGER PRIMARY KEY AUTOINCREMENT,
user_question TEXT NOT NULL,
ai_answer TEXT NOT NULL,
context_summary TEXT,
queried_date TEXT,
queried_person TEXT,
created_at TEXT
);
CREATE INDEX IF NOT EXISTS idx_chat_created ON chat_history(created_at);
""")
conn.commit()
_chat_schema_ready = True
def insert_chat_history(user_question: str, ai_answer: str,
context_summary: str, queried_date: str,
queried_person: str) -> int:
"""插入对话记录"""
conn = get_conn()
try:
cur = conn.cursor()
cur.execute(
_ensure_chat_schema(conn)
cur = conn.execute(
"""INSERT INTO chat_history
(user_question, ai_answer, context_summary, queried_date, queried_person)
VALUES (%s, %s, %s, %s, %s)""",
(user_question, ai_answer, context_summary, queried_date, queried_person)
)
(user_question, ai_answer, context_summary, queried_date,
queried_person, created_at)
VALUES (?, ?, ?, ?, ?, ?)""",
(user_question, ai_answer, context_summary, queried_date,
queried_person, _now()))
conn.commit()
return cur.lastrowid
finally:
@@ -702,24 +358,20 @@ def insert_chat_history(user_question: str, ai_answer: str,
def get_chat_history(limit=20, offset=0, date_filter=None, person_filter=None) -> List[Dict]:
"""获取对话历史(分页)"""
conn = get_conn()
try:
cur = conn.cursor(pymysql.cursors.DictCursor)
conditions = []
params = []
_ensure_chat_schema(conn)
conds, params = [], []
if date_filter:
conditions.append("queried_date = %s")
conds.append("queried_date = ?")
params.append(date_filter)
if person_filter:
conditions.append("queried_person = %s")
conds.append("queried_person = ?")
params.append(person_filter)
where = f"WHERE {' AND '.join(conditions)}" if conditions else ""
where = f"WHERE {' AND '.join(conds)}" if conds else ""
params.extend([limit, offset])
cur.execute(
f"SELECT * FROM chat_history {where} ORDER BY created_at DESC LIMIT %s OFFSET %s",
params
)
return cur.fetchall()
return _rows(conn.execute(
f"SELECT * FROM chat_history {where} ORDER BY created_at DESC LIMIT ? OFFSET ?",
params))
finally:
conn.close()

View File

@@ -0,0 +1,102 @@
"""
Edge-Client - 访问同机 fam-edge 的薄客户端2026-09-13 迁云后新增)
迁云之前这些调用散在 `oracle_sync.py` 里NAS 跨公网访问甲骨文),随镜像层一起
删掉了。现在 fam-core 与 fam-edge 同机,全部走 127.0.0.1:不出公网、无 TLS、
无隧道token 仍然带着fam-edge 侧的鉴权没变,且它同时对公网监听 :5000
只保留三类调用:
- push_name_correct / push_video_delete —— 写操作仍由 fam-edge 处理,因为它
除了改库还要做人物合并、删磁盘素材,不是单纯一条 SQL
- get_activity —— 服务状态页
- fetch_image —— 事件帧图 / 人物头像(图源和裁剪都在 fam-edge
"""
import requests
from .config_loader import load_config
from .logger import setup_logger
logger = setup_logger('fam-core.edge_client')
_cfg = None
def _conf():
global _cfg
if _cfg is None:
c = load_config().get('edge', {})
_cfg = {
'base_url': (c.get('base_url') or 'http://127.0.0.1:5000').rstrip('/'),
'token': c.get('token') or '',
'timeout': int(c.get('timeout', 60)),
}
return _cfg
def base_url():
return _conf()['base_url']
def token():
return _conf()['token']
def _post(rel: str, payload: dict):
"""返回 (ok, err)。写操作统一用这个,失败时把原因带回给前端。"""
c = _conf()
body = dict(payload)
if c['token']:
body['token'] = c['token']
try:
resp = requests.post(f"{c['base_url']}{rel}", json=body, timeout=(5, c['timeout']))
except requests.RequestException as e:
logger.error(f"fam-edge {rel} 请求失败: {e}")
return False, str(e)
if resp.status_code != 200:
logger.warning(f"fam-edge {rel} 返回 {resp.status_code}: {resp.text[:160]}")
return False, f"HTTP {resp.status_code} {resp.text[:120]}"
return True, ''
def push_name_correct(label: str, canonical_name: str):
return _post('/api/oracle/people/correct',
{'label': label, 'canonical_name': canonical_name})
def push_identity_correct(video_id: int, current_name: str, new_name: str):
return _post('/api/oracle/identity/correct',
{'video_id': video_id, 'current_name': current_name, 'new_name': new_name})
def push_video_delete(video_id: int):
return _post('/api/oracle/video/delete', {'video_id': video_id})
def get_activity():
"""服务状态页用fam-edge 的队列 / 分割 / 模型活动快照。返回 (data, err)。"""
c = _conf()
try:
r = requests.get(f"{c['base_url']}/api/oracle/activity",
params={'token': c['token']}, timeout=(5, 15))
except requests.RequestException as e:
return None, f"连接 fam-edge 失败: {e}"
if r.status_code != 200:
return None, f"fam-edge activity HTTP {r.status_code}"
return r.json(), None
def fetch_image(rel: str, params: dict):
"""帧图 / 头像透传,失败返回 None调用方给 404"""
c = _conf()
p = dict(params)
if c['token']:
p['token'] = c['token']
try:
resp = requests.get(f"{c['base_url']}{rel}", params=p, timeout=(5, c['timeout']))
except requests.RequestException as e:
logger.error(f"fam-edge {rel} 请求失败: {e}")
return None
if resp.status_code == 200 and resp.content:
return resp.content
logger.warning(f"fam-edge {rel} 返回 {resp.status_code}")
return None

View File

@@ -1,14 +1,12 @@
"""
Img-Proxy - NAS 端图片代理(图源在 Oracle计算全部在 Oracle
Img-Proxy - 图片代理(图源与裁剪都在 fam-edge
浏览器不直连 Oracle避免公网暴露 5000 端口与 token 外泄),
而是访问 NAS fam-core 的 /api/proxy/*,由 NAS 出网到 Oracle 拉取 jpeg 回传
Oracle 侧已有磁盘缓存 / VLM 人物定位裁剪NAS 端仅透传,不做图像计算。
浏览器不直连 fam-edge避免 token 外泄),而是访问 /api/proxy/*,由本服务转一手。
迁云后这一跳是同机 127.0.0.1,纯透传,不做图像计算
"""
from flask import Blueprint, Response, request
from .oracle_sync import get_sync
from . import edge_client
from .logger import setup_logger
logger = setup_logger('fam-core.img_proxy')
@@ -17,21 +15,7 @@ img_bp = Blueprint('img_proxy', __name__)
def _fetch(rel: str, params: dict):
import requests
sync = get_sync() # 复用 oracle_sync 已解析好的 base_url/token/timeout避免两处配置各读一份
params = dict(params)
if sync.token:
params['token'] = sync.token
try:
resp = requests.get(f"{sync.base_url}{rel}", params=params,
timeout=(10, sync.timeout))
if resp.status_code == 200 and resp.content:
return resp.content
logger.warning(f"Oracle {rel} 返回 {resp.status_code}: {resp.text[:120]}")
return None
except Exception as e:
logger.error(f"Oracle {rel} 请求失败: {e}")
return None
return edge_client.fetch_image(rel, params)
@img_bp.route('/api/proxy/frame', methods=['GET'])

View File

@@ -14,7 +14,7 @@ from flask import Blueprint, request, jsonify
from ..logger import setup_logger
from .. import db_layer
from ..oracle_sync import get_sync
from .. import edge_client
logger = setup_logger('fam-core.member_manager')
@@ -24,7 +24,7 @@ member_bp = Blueprint('member_manager', __name__)
@member_bp.route('/api/member/unnamed', methods=['GET'])
def list_unnamed():
"""列出未命名人物canonical_name 为空)"""
members = db_layer.get_sync_people()
members = db_layer.get_people()
result = []
for m in members:
canonical = m.get('canonical_name')
@@ -40,7 +40,7 @@ def list_unnamed():
@member_bp.route('/api/member/list', methods=['GET'])
def list_members():
"""列出所有人物(按 canonical_name 或 label 展示)"""
members = db_layer.get_sync_people()
members = db_layer.get_people()
result = []
for m in members:
canonical = m.get('canonical_name')
@@ -58,7 +58,7 @@ def list_members():
@member_bp.route('/api/member/name', methods=['POST'])
def name_member():
"""命名人物(回推 Oracle + 立即拉回本地镜像
"""命名人物(交给 fam-edge 落库 + 合并人物
请求: {"label": "人物A", "canonical_name": "张三"}
"""
@@ -71,18 +71,13 @@ def name_member():
if not label or not canonical_name:
return jsonify({"error": "缺少必填字段: label, canonical_name"}), 400
logger.info(f"命名: {label} -> {canonical_name}(回推 Oracle")
ok, err = get_sync().push_name_correct(label, canonical_name)
logger.info(f"命名: {label} -> {canonical_name}(回推 fam-edge")
ok, err = edge_client.push_name_correct(label, canonical_name)
if not ok:
return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502
return jsonify({"error": f"回推 fam-edge 失败: {err}"}), 502
# 立即拉回最新 people 镜像,前端无需等待下一个 30 分钟周期
try:
get_sync().trigger_now()
except Exception as e:
logger.warning(f"命名后即时拉回失败(下一个周期会自动同步): {e}")
members = db_layer.get_sync_people()
members = db_layer.get_people()
return jsonify({
"status": "ok",
"label": label,
@@ -114,24 +109,20 @@ def merge_member():
return jsonify({"error": "source 与 target 不能相同"}), 400
# 解析 target 的规范名
members = {m['label']: m for m in db_layer.get_sync_people()}
members = {m['label']: m for m in db_layer.get_people()}
target_row = members.get(target)
if target_row and target_row.get('canonical_name'):
canonical = target_row['canonical_name']
else:
canonical = target # target 未命名 -> 以 label 作为规范名
logger.info(f"合并: {source} -> {canonical}(回推 Oracle")
ok, err = get_sync().push_name_correct(source, canonical)
logger.info(f"合并: {source} -> {canonical}(回推 fam-edge")
ok, err = edge_client.push_name_correct(source, canonical)
if not ok:
return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502
return jsonify({"error": f"回推 fam-edge 失败: {err}"}), 502
try:
get_sync().trigger_now()
except Exception as e:
logger.warning(f"合并后即时拉回失败(下一个周期会自动同步): {e}")
members = db_layer.get_sync_people()
members = db_layer.get_people()
return jsonify({
"status": "ok",
"source": source,
@@ -163,15 +154,11 @@ def identity_correct():
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)
logger.info(f"人物纠错: video_id={video_id} {current_name} -> {new_name}(回推 fam-edge")
ok, err = edge_client.push_identity_correct(video_id, current_name, new_name)
if not ok:
return jsonify({"error": f"回推 Oracle 失败: {err}"}), 502
return jsonify({"error": f"回推 fam-edge 失败: {err}"}), 502
try:
get_sync().trigger_now()
except Exception as e:
logger.warning(f"纠错后即时拉回失败(下一个周期会自动同步): {e}")
return jsonify({
"status": "ok", "video_id": video_id,

View File

@@ -1,20 +0,0 @@
"""
运动监测状态接口fam-core
运动事件采集走 MotionNotifier 轮询 SS EventCenter.Event.List真实
event_id/start_time/duration不再提供 /api/ss/webhook 接收端点
SS Webhook 行动规则已弃用2026-08-25 移除)。
"""
from flask import Blueprint, jsonify
from .logger import setup_logger
from .motion_notifier.motion_notifier import get_motion_notifier
logger = setup_logger('fam-core.motion_bp')
motion_bp = Blueprint('motion_bp', __name__)
@motion_bp.route('/api/ss/status', methods=['GET'])
def ss_status():
return jsonify(get_motion_notifier().status()), 200

View File

@@ -1,4 +0,0 @@
"""MotionNotifier 包NAS 端运动监测通知服务。"""
from .motion_notifier import MotionNotifier, get_motion_notifier
__all__ = ["MotionNotifier", "get_motion_notifier"]

View File

@@ -1,370 +0,0 @@
"""
MotionNotifier - NAS 端运动监测通知服务轮询主路径2026-08-22 定型)
职责:
在 NAS 本机**轮询**群晖 Surveillance Station 的运动侦测事件
SYNO.SurveillanceStation.EventCenter.Event method=List参数名下划线风格
camera_ids/start_time/end_timeevent_type=10 即运动),增量推送到甲骨文
FAM-Edge 的 /api/ss/motion 接口。轮询能拿到真实 event_id/start_time/duration。
数据方向(关键约束): NAS -> Oracle单向。甲骨文不再反向访问 NAS。
- 原 fam-edge 的 dsm_motion_client甲骨文主动查 SS API已停用
- 推送后由甲骨文本地落库 ss_motion_events供 video_processor 本地预过滤使用。
可靠性设计:
- 游标MariaDB sync_cursor.motion_last_event_id持久化重启优先续用 DB 游标,
停机期间的事件由窗口回看补推;仅首次部署时才初始化为 SS 当前最大 id。
- 推送失败的批次不推进游标,下一轮窗口回看重试(不丢事件)。
- 心跳线程enabled 即跑,与轮询无关):定期空 POST /api/ss/motion证明
NAS->Oracle 推送链路存活,供甲骨文侧判断"无运动"结论是否可信。
"""
import os
import re
import time
import threading
from datetime import datetime, timezone
import requests
from ..logger import setup_logger
from ..config_loader import load_config
from .. import db_layer
logger = setup_logger('fam-core.motion_notifier')
_MOTION_NOTIFIER = None
def get_motion_notifier():
"""模块级单例app.py 启动时创建并 start状态接口经此获取"""
global _MOTION_NOTIFIER
if _MOTION_NOTIFIER is None:
_MOTION_NOTIFIER = MotionNotifier()
return _MOTION_NOTIFIER
class MotionNotifier:
def __init__(self):
cfg = load_config().get('motion_notifier', {})
self.enabled = bool(cfg.get('enabled', False))
self.dsm_host = cfg.get('dsm_host', '192.168.50.64')
self.dsm_port = int(cfg.get('dsm_port', 5000))
self.dsm_account = self._resolve(cfg.get('dsm_account', ''))
self.dsm_password = self._resolve(cfg.get('dsm_password', ''))
self.camera_ids = cfg.get('camera_ids', [2])
self.oracle_base_url = cfg.get('oracle_base_url',
'http://129.146.26.249:5000').rstrip('/')
self.oracle_token = self._resolve(cfg.get('oracle_token', '${ORACLE_SYNC_TOKEN}'))
self.poll_interval_sec = int(cfg.get('poll_interval_sec', 60))
self.poll_window_hours = int(cfg.get('poll_window_hours', 2))
self.batch_size = int(cfg.get('batch_size', 100))
self.timeout = int(cfg.get('timeout_sec', 10))
# 是否启用轮询(默认开启):轮询 EventCenter.Event.List 为主数据源,拿到
# 真实 event_id/start_time/durationWebhook 仅作为可选的低延迟补充
# SS 行动规则未配置时不会触发,不影响主路径)
self.poll_enabled = bool(cfg.get('poll_enabled', True))
# 心跳:跟"轮询 SS"是两回事——不查 SS只是定期空 POST 一下甲骨文,证明
# NAS->Oracle 这条推送链路本身还活着。enabled=true 时始终跑(不受
# poll_enabled 影响),供甲骨文侧 has_motion_in_range_local() 判断
# "这段时间没收到运动事件"是真的没运动,还是推送链路已经挂了。
self.heartbeat_interval_sec = int(cfg.get('heartbeat_interval_sec', 300))
# 轮询路径poll_enabled=true 时)关注的摄像头
self.camera_ids = cfg.get('camera_ids', [2])
self._base = f"http://{self.dsm_host}:{self.dsm_port}/webapi"
self._sid = None
self._running = False
self._thread = None
self._last_event_id = None
# 推送过但 duration<=0动作进行中SS 尚未回填真实时长)的事件 id。
# 下轮窗口回查时若已结束duration>0补推覆盖 Oracle避免分割永远跳过。
self._zero_dur_ids = set()
self._last_poll_at = None
self._last_error = None
self._pushed_total = 0
self._hb_running = False
self._hb_thread = None
self._last_heartbeat_at = None
self._last_heartbeat_error = None
@staticmethod
def _resolve(v):
"""解析 ${ENV} 引用;非字符串或不含 ${...} 原样返回。"""
if isinstance(v, str) and v.startswith('${') and v.endswith('}'):
return os.environ.get(v[2:-1], '')
return v
# ------------------------------------------------------------------
# Surveillance Station 登录 / 事件查询
# ------------------------------------------------------------------
def _login(self):
try:
resp = requests.get(
f"{self._base}/auth.cgi",
params={"api": "SYNO.API.Auth", "version": 6, "method": "login",
"account": self.dsm_account, "passwd": self.dsm_password,
"session": "SurveillanceStation", "format": "sid"},
timeout=self.timeout)
data = resp.json()
if data.get('success'):
return data['data']['sid']
logger.warning(f"SS 登录失败: {data.get('error')}")
except Exception as e:
logger.warning(f"SS 登录异常: {e}")
return None
def _fetch_events(self, start_ts: int, end_ts: int):
"""查询 [start_ts, end_ts] 窗口内的 SS 事件。返回事件列表或 None查询失败"""
if self._sid is None:
self._sid = self._login()
if self._sid is None:
return None
for attempt in range(2):
try:
resp = requests.get(
f"{self._base}/entry.cgi",
params={"api": "SYNO.SurveillanceStation.EventCenter.Event",
"version": 1, "method": "List",
"camera_ids": ",".join(str(c) for c in self.camera_ids),
"start_time": start_ts, "end_time": end_ts,
"limit": 1000, "_sid": self._sid},
timeout=self.timeout)
data = resp.json()
except Exception as e:
logger.warning(f"SS 事件查询异常: {e}")
return None
if not data.get('success'):
code = (data.get('error') or {}).get('code')
if code in (106, 107, 119) and attempt == 0:
# session 过期/被顶掉,重新登录重试一次
self._sid = self._login()
if self._sid is None:
return None
continue
logger.warning(f"SS 事件查询失败: {data.get('error')}")
return None
# 响应按 ds_idCMS 多机场景的服务器 id单机固定 "0")分组,拉平
events = [e for grp in (data.get('data') or {}).values()
for e in (grp or [])]
return events
return None
# ------------------------------------------------------------------
# 推送到甲骨文
# ------------------------------------------------------------------
def push_events_to_oracle(self, events) -> int:
"""把标准化后的事件列表推送到 Oracle /api/ss/motion。返回成功推送条数。"""
if not events:
return 0
norm = []
for e in events:
eid = e.get('id') or e.get('event_id')
if eid is None:
continue
norm.append({
"event_id": int(eid),
"camera_id": e.get('camera_id'),
"event_type": e.get('event_type'),
"start_time": e.get('start_time'),
"duration": e.get('duration'),
"thumbnail_url": e.get('thumbnail_url') or e.get('thumbnail_dir'),
})
if not norm:
return 0
try:
resp = requests.post(
f"{self.oracle_base_url}/api/ss/motion",
json={"token": self.oracle_token, "events": norm},
timeout=(10, 30))
except requests.RequestException as e:
logger.error(f"推送运动事件到 Oracle 失败: {e}")
self._last_error = str(e)
return 0
if resp.status_code != 200:
logger.error(f"推送运动事件到 Oracle 返回 {resp.status_code}: {resp.text[:200]}")
self._last_error = f"HTTP {resp.status_code}"
return 0
try:
stored = resp.json().get('stored', 0)
except ValueError:
stored = len(norm)
self._pushed_total += stored
self._last_error = None
logger.info(f"运动事件推送成功: {len(norm)} 条 -> Oracle 存储 {stored}")
return stored
def send_heartbeat(self) -> bool:
"""空 events 调一次 /api/ss/motion只为证明 NAS->Oracle 推送链路还活着。
跟 push_events_to_oracle 分开一个方法,是因为那个方法 events 为空时直接
return 0不发请求——心跳恰恰就是要在没有真实事件时也发一次请求。
"""
try:
resp = requests.post(
f"{self.oracle_base_url}/api/ss/motion",
json={"token": self.oracle_token, "events": []},
timeout=(10, 30))
except requests.RequestException as e:
logger.warning(f"运动心跳推送失败: {e}")
self._last_heartbeat_error = str(e)
return False
if resp.status_code != 200:
logger.warning(f"运动心跳推送返回 {resp.status_code}: {resp.text[:200]}")
self._last_heartbeat_error = f"HTTP {resp.status_code}"
return False
self._last_heartbeat_at = datetime.now()
self._last_heartbeat_error = None
return True
def _heartbeat_run(self):
logger.info(f"运动心跳线程启动,间隔 {self.heartbeat_interval_sec}s")
while self._hb_running:
try:
self.send_heartbeat()
except Exception as e:
self._last_heartbeat_error = str(e)
logger.error(f"运动心跳异常: {e}", exc_info=True)
for _ in range(self.heartbeat_interval_sec):
if not self._hb_running:
break
time.sleep(1)
# ------------------------------------------------------------------
# 增量轮询主循环
# ------------------------------------------------------------------
def _init_cursor(self):
"""启动时初始化 last_event_id。
优先使用 DB 已保存的游标:服务停机/重启期间产生的新事件,会在重启后
被下一轮窗口回看捞到并补推(不丢事件)。
仅当 DB 无有效游标(首次部署)时,取 SS 当前最大事件 id 作为起点,
避免把历史事件全部回灌一遍。
"""
saved = db_layer.get_motion_cursor()
if saved and int(saved) > 0:
self._last_event_id = int(saved)
logger.info(f"运动通知游标初始化DB: last_event_id={self._last_event_id}")
return
now = int(datetime.now(timezone.utc).timestamp())
events = self._fetch_events(now - 3600, now) # 最近 1h 用于定位最大 id
if events:
self._last_event_id = max(
int(e.get('id', 0)) for e in events if e.get('id'))
else:
self._last_event_id = 0
logger.info(f"运动通知游标初始化SS 当前最大): last_event_id={self._last_event_id}")
def _poll_once(self):
now = int(datetime.now(timezone.utc).timestamp())
start_ts = now - int(self.poll_window_hours * 3600)
events = self._fetch_events(start_ts, now)
if events is None:
return # 查询失败,下一轮重试
by_id = {int(e.get('id')): e for e in events if e.get('id')}
# 1) 补推:先前推送时 duration<=0动作进行中的事件现在若已结束
# SS 回填真实 duration>0重推覆盖 OracleON CONFLICT UPDATE
# 重启兜底:游标 DB 续用 → 窗口回看会把历史事件整体重推,届时同样覆盖。
if self._zero_dur_ids:
recheck = []
for eid in list(self._zero_dur_ids):
e = by_id.get(eid)
if e and int(e.get('duration') or 0) > 0:
recheck.append(e)
self._zero_dur_ids.discard(eid)
if recheck:
self.push_events_to_oracle(recheck)
logger.info(f"补推 {len(recheck)} 条已结束事件的真实 duration")
# 2) 增量推送新事件id 大于游标)
new = [e for e in events
if e.get('id') and int(e.get('id')) > (self._last_event_id or 0)]
if not new:
return
new.sort(key=lambda e: int(e.get('id', 0)))
cursor = self._last_event_id
for i in range(0, len(new), self.batch_size):
batch = new[i:i + self.batch_size]
stored = self.push_events_to_oracle(batch)
if stored > 0:
# 只有推送成功的批次才推进游标;失败批次保持原地,
# 下一轮窗口回看会重新捞到并重试(不丢事件)
cursor = max(int(e.get('id', 0)) for e in batch)
# 记录推送时仍在进行中的事件duration<=0待下轮补推真实时长
for e in batch:
if int(e.get('duration') or 0) <= 0:
self._zero_dur_ids.add(int(e.get('id')))
self._last_event_id = cursor
db_layer.set_motion_cursor(str(self._last_event_id))
def _run(self):
logger.info(f"MotionNotifier 启动,轮询间隔 {self.poll_interval_sec}s"
f"目标 SS {self.dsm_host}:{self.dsm_port}")
try:
self._init_cursor()
except Exception as e:
logger.warning(f"运动通知游标初始化失败(从 0 开始): {e}")
self._last_event_id = 0
while self._running:
try:
self._poll_once()
self._last_poll_at = datetime.now()
except Exception as e:
self._last_error = str(e)
logger.error(f"运动通知轮询异常: {e}", exc_info=True)
# 分段休眠,便于 stop 快速唤醒
for _ in range(self.poll_interval_sec):
if not self._running:
break
time.sleep(1)
# ------------------------------------------------------------------
def start(self):
if not self.enabled:
logger.info("MotionNotifier 未启用motion_notifier.enabled=false")
return
# 心跳线程跟轮询是否开启无关:只要 MotionNotifier 整体 enabled就该
# 持续证明推送链路活着,哪怕当前正好没有真实运动事件可推送。
if not self._hb_running:
self._hb_running = True
self._hb_thread = threading.Thread(
target=self._heartbeat_run, daemon=True, name='motion-heartbeat')
self._hb_thread.start()
if not self.poll_enabled:
logger.info("MotionNotifier 轮询已禁用poll_enabled=false仅作为 "
"Webhook 推送客户端 + 心跳 + 摄像头名映射使用")
return
if self._running:
return
self._running = True
self._thread = threading.Thread(target=self._run, daemon=True,
name='motion-notifier')
self._thread.start()
def is_alive(self):
return self._thread is not None and self._thread.is_alive()
def stop(self):
self._running = False
if self._thread:
self._thread.join(timeout=5)
self._hb_running = False
if self._hb_thread:
self._hb_thread.join(timeout=5)
def status(self) -> dict:
return {
"heartbeat_running": self._hb_thread is not None and self._hb_thread.is_alive(),
"last_heartbeat_at": self._last_heartbeat_at.isoformat() if self._last_heartbeat_at else None,
"last_heartbeat_error": self._last_heartbeat_error,
"running": self.is_alive(),
"enabled": self.enabled,
"poll_enabled": self.poll_enabled,
"last_event_id": self._last_event_id,
"zero_dur_pending": len(self._zero_dur_ids),
"last_poll_at": self._last_poll_at.isoformat() if self._last_poll_at else None,
"last_error": self._last_error,
"pushed_total": self._pushed_total,
"oracle": self.oracle_base_url,
"dsm": f"{self.dsm_host}:{self.dsm_port}",
}

View File

@@ -1,234 +0,0 @@
"""
Oracle-Sync - NAS 端唯一后台线程
职责:
1. 每 interval_sec默认 1800s = 30 分钟)从甲骨文 FAM-Edge 拉取增量:
GET {base_url}/api/oracle/sync?since=<cursor>&token=<token>
返回 {videos, events, people, server_time},写入本地 MariaDB 镜像表
(sync_videos / sync_events / sync_people),并推进 sync_cursor。
2. 接收命名校正回推: POST {base_url}/api/oracle/people/correct
{label, canonical_name, token} —— 手动命名manual 优先,不被 LLM 覆盖)。
数据流向(新架构 v2:
Google 硬盘 --rclone--> 甲骨文本地 --> 整视频分析 --> Oracle SQLite
--> [本线程每 30 分钟拉增量] --> NAS MariaDB 镜像 --> fam-ui 读取展示
NAS 不再处理任何视频CPU 占用显著降低。
"""
import time
import threading
import requests
from datetime import datetime
from .logger import setup_logger
from .config_loader import load_config
from . import db_layer
logger = setup_logger('fam-core.oracle_sync')
_SYNC_INSTANCE = None
def get_sync():
"""模块级单例app.py 启动时创建并 start其余模块经此获取"""
global _SYNC_INSTANCE
if _SYNC_INSTANCE is None:
_SYNC_INSTANCE = OracleSync()
return _SYNC_INSTANCE
class OracleSync:
def __init__(self):
cfg = load_config().get('oracle_sync', {})
self.base_url = cfg.get('base_url', 'http://129.146.26.249:5000').rstrip('/')
self.token = cfg.get('token', '')
self.interval_sec = int(cfg.get('interval_sec', 1800))
self.timeout = int(cfg.get('timeout', 120))
self._running = False
self._thread = None
self._last_sync_at = None
self._last_error = None
self._last_count = None
# 拉取互斥锁trigger_now命名后即时拉回与后台 _run 并发时只允许一个执行
self._pull_lock = threading.Lock()
# ------------------------------------------------------------------
def _pull_once(self) -> bool:
"""执行一次增量拉取(加锁防并发双拉)。返回是否成功。"""
with self._pull_lock:
return self._pull_once_locked()
def _pull_once_locked(self) -> bool:
since = db_layer.get_sync_cursor() or ''
params = {'since': since, 'token': self.token}
try:
resp = requests.get(
f"{self.base_url}/api/oracle/sync",
params=params, timeout=(10, self.timeout))
except requests.RequestException as e:
self._last_error = f"请求失败: {e}"
logger.error(f"拉取同步失败: {e}")
return False
if resp.status_code == 401:
self._last_error = "token 校验失败"
logger.error("同步 token 校验失败 (401),请检查 oracle_sync.token 配置")
return False
if resp.status_code != 200:
self._last_error = f"HTTP {resp.status_code}"
logger.error(f"同步返回异常: {resp.status_code} {resp.text[:200]}")
return False
try:
data = resp.json()
except ValueError:
self._last_error = "非 JSON 响应"
logger.error("同步返回非 JSON 响应")
return False
videos = data.get('videos', []) or []
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, n_identity)
logger.info(
f"同步完成: videos+{n_videos} events+{n_events} people+{n_people} "
f"model_calls+{n_calls} identity_map+{n_identity} since={since!r} "
f"-> server_time={server_time}")
return True
# ------------------------------------------------------------------
def push_name_correct(self, label: str, canonical_name: str):
"""回推命名校正到 Oracle手动命名优先级最高不被 LLM 覆盖)。
返回 (success: bool, error: str)
"""
label = (label or '').strip()
canonical_name = (canonical_name or '').strip()
if not label or not canonical_name:
return False, "缺少 label / canonical_name"
try:
resp = requests.post(
f"{self.base_url}/api/oracle/people/correct",
json={"label": label, "canonical_name": canonical_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 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 push_video_delete(self, video_id: int):
"""回推事件时间轴"删除该视频"到 Oracle。返回 (success: bool, error: str)。"""
try:
video_id = int(video_id)
except (TypeError, ValueError):
return False, "video_id 必须是数字"
try:
resp = requests.post(
f"{self.base_url}/api/oracle/video/delete",
json={"video_id": video_id, "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}")
while self._running:
try:
self._pull_once()
except Exception as e:
self._last_error = str(e)
logger.error(f"同步异常: {e}", exc_info=True)
# 分段休眠,便于 stop 快速唤醒
for _ in range(self.interval_sec):
if not self._running:
break
time.sleep(1)
def start(self):
if self._running:
return
self._running = True
self._thread = threading.Thread(target=self._run, daemon=True, name='oracle-sync')
self._thread.start()
def is_alive(self):
return self._thread is not None and self._thread.is_alive()
def stop(self):
self._running = False
if self._thread:
self._thread.join(timeout=5)
def trigger_now(self) -> bool:
"""立即触发一次同步(命名后即时拉回 / 手动)。"""
return self._pull_once()
def status(self) -> dict:
return {
"running": self.is_alive(),
"last_sync_at": self._last_sync_at.isoformat() if self._last_sync_at else None,
"last_error": self._last_error,
"last_count": self._last_count,
"cursor": db_layer.get_sync_cursor(),
"interval_sec": self.interval_sec,
}

View File

@@ -7,9 +7,8 @@ UI-API - Vue 前端数据接口(新架构 v3Vue SPA 取代 Streamlit
<id>`(删除视频会话),因为语义上属于 videos 这个资源,比塞进 member_manager.py
更清晰;写法沿用项目里"先回推 Oracle成功后处理本地状态"的既有模式。
/api/ui/service-status 需要代理 Oracle 的 /api/oracle/activity浏览器不直连
Oracle避免 token 暴露),写法照抄 img_proxy.py 的模式:复用 oracle_sync.get_sync()
已解析好的 base_url/token不再单独存一份配置。
/api/ui/service-status 代理同机 fam-edge 的 /api/oracle/activity浏览器不直连
fam-edge避免 token 暴露),走 edge_client 统一读 base_url/token。
"""
import re
from datetime import datetime
@@ -19,7 +18,7 @@ from flask import Blueprint, request, jsonify
from .logger import setup_logger
from . import db_layer
from .oracle_sync import get_sync
from . import edge_client
logger = setup_logger('fam-core.ui_api')
@@ -50,7 +49,7 @@ def videos():
date_filter = request.args.get('date') or None
page = max(0, request.args.get('page', 0, type=int))
page_size = 15
rows = db_layer.get_sync_videos(limit=page_size, offset=page * page_size,
rows = db_layer.get_videos(limit=page_size, offset=page * page_size,
date_filter=date_filter)
return jsonify({"videos": _ser(rows), "page": page, "page_size": page_size}), 200
@@ -58,10 +57,10 @@ def videos():
@ui_bp.route('/api/ui/videos/<int:video_id>', methods=['GET'])
def video_detail(video_id):
"""单个视频会话详情 + 时间线事件列表(事件时间轴右侧)。"""
video = db_layer.get_sync_video(video_id)
video = db_layer.get_video(video_id)
if not video:
return jsonify({"error": "视频不存在"}), 404
events = db_layer.get_sync_events_for_video(video_id)
events = db_layer.get_events_for_video(video_id)
return jsonify({"video": _ser(video), "events": _ser(events)}), 200
@@ -69,16 +68,13 @@ def video_detail(video_id):
def video_delete(video_id):
"""删除视频会话(事件时间轴"删除"入口)。
先回推 Oracle 物理删除events + videos 行 + 磁盘文件),成功后再清理本地
MariaDB 镜像——增量同步get_sync_delta只做 upsert 感知不到删除,不能像
命名纠错那样靠 trigger_now() 拉增量顺带清理,必须显式调 delete_sync_video
Oracle 回推失败时不清理本地镜像,避免"Oracle 还留着、NAS 却以为删了"
数据不一致,让用户看错误提示后重试。
交给 fam-edge 处理:它删 events + videos 行,还要删磁盘上的运动片段文件。
迁云前这里还要再清一次 NAS 本地镜像(增量同步只 upsert感知不到删除
现在没有镜像了fam-edge 删完就是最终状态
"""
ok, err = get_sync().push_video_delete(video_id)
ok, err = edge_client.push_video_delete(video_id)
if not ok:
return jsonify({"error": f"删除失败: {err}"}), 502
db_layer.delete_sync_video(video_id)
return jsonify({"status": "ok", "video_id": video_id}), 200
@@ -86,7 +82,7 @@ def video_delete(video_id):
def stats():
"""统计卡:视频/事件/人物/关注数(可选按日期过滤)。"""
date_filter = request.args.get('date') or None
return jsonify(_ser(db_layer.get_sync_stats(date_filter))), 200
return jsonify(_ser(db_layer.get_stats(date_filter))), 200
def _clean_person(s: str) -> str:
@@ -100,7 +96,7 @@ def people():
对应旧 Streamlit 版本 app.py 里的分组逻辑,原样搬到服务端。
"""
rows = db_layer.get_sync_people()
rows = db_layer.get_people()
groups = {}
for p in rows:
key = p.get('canonical_name') or p['label']
@@ -128,7 +124,7 @@ def people():
def people_clips():
"""某人物出现过的运动片段列表(人物卡「运动片段」区块)。
按 label/canonical_name 匹配 sync_events.person_list_json → 关联视频
按 label/canonical_name 匹配 events.person_list_json → 关联视频
(运动片段优先)。返回片段 video_id/filename/event_start_time/duration/
summary/camera_name + first_ts缩略图定位+ clip_events片段内事件数
"""
@@ -136,7 +132,7 @@ def people_clips():
if not label:
return jsonify({"error": "缺少 label 参数"}), 400
limit = min(20, request.args.get('limit', 10, type=int) or 10)
clips = db_layer.get_sync_people_clips(label, limit)
clips = db_layer.get_people_clips(label, limit)
return jsonify({"label": label, "clips": _ser(clips)}), 200
@@ -163,41 +159,25 @@ def attention_events():
@ui_bp.route('/api/ui/named-members', methods=['GET'])
def named_members():
"""已命名成员真名列表AI 对话页快捷选择)。"""
return jsonify({"members": db_layer.get_sync_named_members()}), 200
return jsonify({"members": db_layer.get_named_members()}), 200
@ui_bp.route('/api/ui/model-stats', methods=['GET'])
def model_stats():
"""云端模型调用统计:按模型聚合 + 最近调用明细。"""
agg = db_layer.get_sync_model_calls_stats()
calls = db_layer.get_sync_model_calls(limit=100)
agg = db_layer.get_model_calls_stats()
calls = db_layer.get_model_calls(limit=100)
return jsonify({"aggregate": _ser(agg), "recent_calls": _ser(calls)}), 200
@ui_bp.route('/api/ui/service-status', methods=['GET'])
def service_status():
"""服务状态页:NAS 同步状态 + Oracle 实时活动代理token 不下发浏览器)。"""
sync = get_sync()
nas_status = sync.status()
"""服务状态页:fam-edge 的队列/分割/模型活动代理token 不下发浏览器)。
oracle_data = None
oracle_error = None
if sync.base_url and sync.token:
import requests
try:
r = requests.get(f"{sync.base_url}/api/oracle/activity",
params={'token': sync.token}, timeout=15)
if r.status_code == 200:
oracle_data = r.json()
else:
oracle_error = f"Oracle activity HTTP {r.status_code}"
except Exception as e:
oracle_error = f"连接 Oracle 失败: {e}"
else:
oracle_error = "未配置 oracle_sync.base_url/token"
return jsonify({
"nas_sync": nas_status,
"oracle": oracle_data,
"oracle_error": oracle_error,
}), 200
迁云前这里还有一块 "NAS 同步状态"(镜像拉取的游标/周期/上次条数),随镜像层
一起删了——现在前端读的就是 fam-edge 写的那份库,没有"同步"这个中间状态。
NAS 侧只剩运动事件推送,它的心跳在 fam-edge 的 service_activity 里,
已经包含在 activity 快照中。
"""
data, err = edge_client.get_activity()
return jsonify({"oracle": data, "oracle_error": err}), 200