refactor(fam-core): 重构第二阶段 - img_proxy 收尾 + 稳定性小修 + 工程质量
img_proxy 收尾: 提交此前未提交的图片代理蓝图(app.py 注册),改为复用 oracle_sync
已解析好的 base_url/token/timeout,不再独立读一份配置(避免配置改动时两处不同步)。
稳定性/正确性: chat_handler 问答失败时不再把 Oracle 内部 HTTP 状态码等细节透传给
客户端,改为通用错误信息(详细原因仍记服务端日志); 删除 logger.py 里旧任务架构
遗留的死函数 log_task; 更新 config.yaml.example 到当前 v2 架构(原文件还是重构前
的 scheduler/dispatcher/video_server 旧结构,当前代码完全不读这些字段)。
工程质量: 新增 fam-core/tests(9 个单元测试,覆盖 _format_events 上下文格式化和
config_loader 的 ${ENV_VAR} 解析)。
部署时发现并修复一个和这次改动无关的运维问题: NAS .env 文件缺 export 关键字,
plain source 只在当前 shell 生效不会被子进程(gunicorn)继承,导致 Oracle-Sync
token 校验失败;用 set -a/set +a 强制导出重启,非代码改动。
已部署 NAS 并验证:/health、/api/status、/api/proxy/frame、/api/proxy/avatar
全部通过;浏览器实测事件时间轴、人物管理页正常渲染。
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
@@ -1,5 +1,9 @@
|
|||||||
# FAM-Core 配置文件 (NAS 端)
|
# FAM-Core 配置文件 (NAS 端) - 新架构 v2
|
||||||
# 复制此文件为 config.yaml 并修改实际值
|
# 复制此文件为 config.yaml 并修改实际值
|
||||||
|
#
|
||||||
|
# NAS 仅作管理后台,不再处理视频。唯一后台线程 Oracle-Sync 每 30 分钟
|
||||||
|
# 从甲骨文 FAM-Edge 拉取增量镜像到本地 MariaDB(sync_videos/events/people)。
|
||||||
|
# 所有视频分析在 Oracle 完成。
|
||||||
|
|
||||||
server:
|
server:
|
||||||
host: "0.0.0.0"
|
host: "0.0.0.0"
|
||||||
@@ -11,32 +15,18 @@ database:
|
|||||||
user: "root"
|
user: "root"
|
||||||
password: ""
|
password: ""
|
||||||
database: "sentinel_home_ai"
|
database: "sentinel_home_ai"
|
||||||
|
unix_socket: "/run/mysqld/mysqld10.sock"
|
||||||
|
|
||||||
scheduler:
|
# 甲骨文同步(每 30 分钟拉增量镜像)
|
||||||
scan_interval: 60 # 扫描间隔(秒)
|
oracle_sync:
|
||||||
video_dir: "/volume1/surveillance" # 视频录制目录
|
# FAM-Edge 对外同步接口地址(端口同其 server.port=5000)
|
||||||
video_extensions: [".mp4", ".mkv", ".avi"]
|
base_url: "http://<oracle-public-ip>:5000"
|
||||||
file_stable_seconds: 60 # 文件稳定判定秒数
|
# 与 Oracle 端 sync_api.token 一致
|
||||||
camera_name: "默认摄像头"
|
token: ""
|
||||||
|
interval_sec: 1800 # 拉取间隔(秒),默认 30 分钟
|
||||||
dispatcher:
|
timeout: 120 # 单次拉取超时(秒)
|
||||||
poll_interval: 30 # 轮询间隔(秒)
|
|
||||||
edge_url: "http://100.x.x.20:5000/api/edge/video/enqueue" # 异步队列端点(上传入队,Poller 拉取结果)
|
|
||||||
max_retries: 5 # 文件级重试次数(分块级重试另计,每块3次)
|
|
||||||
stale_timeout: 1800 # PROCESSING 僵尸回收(秒),需大于压缩+上传+Edge队列积压总时长
|
|
||||||
# ffmpeg_path: "/var/packages/CodecPack/target/bin/ffmpeg41" # 可选,默认自动探测(Synology 需 CodecPack 版,系统版无 h264 编码)
|
|
||||||
compress_timeout: 3600 # 单个视频预压缩超时(秒)
|
|
||||||
|
|
||||||
video_server:
|
|
||||||
base_url: "http://100.x.x.10:8000/media"
|
|
||||||
token: "your-secret-token-here" # 视频访问 token
|
|
||||||
video_dir: "/volume1/surveillance"
|
|
||||||
|
|
||||||
chat_handler:
|
chat_handler:
|
||||||
ollama_url: "http://100.x.x.20:11434/api/generate"
|
# 智能问答统一走 FAM-Edge 编排端点(Gemini → NVIDIA → 本地 Ollama 兜底)
|
||||||
model_name: "llava-phi3"
|
qa_url: "http://<oracle-public-ip>:5000/api/edge/chat/ask"
|
||||||
timeout: 120
|
timeout: 120
|
||||||
|
|
||||||
storage:
|
|
||||||
# 关键帧落盘目录(event_receiver 写入,fam-ui 读取展示时间轴)
|
|
||||||
frame_image_dir: "/volume1/web/sentinel-home-ai/fam-ui/static/frames"
|
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ from .logger import setup_logger
|
|||||||
from .oracle_sync import get_sync
|
from .oracle_sync import get_sync
|
||||||
from .chat_handler.chat_handler import chat_bp
|
from .chat_handler.chat_handler import chat_bp
|
||||||
from .member_manager.member_manager import member_bp
|
from .member_manager.member_manager import member_bp
|
||||||
|
from .img_proxy import img_bp
|
||||||
|
|
||||||
logger = setup_logger('fam-core.app')
|
logger = setup_logger('fam-core.app')
|
||||||
|
|
||||||
@@ -27,6 +28,7 @@ app = Flask(__name__)
|
|||||||
# 注册蓝图
|
# 注册蓝图
|
||||||
app.register_blueprint(chat_bp)
|
app.register_blueprint(chat_bp)
|
||||||
app.register_blueprint(member_bp)
|
app.register_blueprint(member_bp)
|
||||||
|
app.register_blueprint(img_bp)
|
||||||
|
|
||||||
# 健康检查
|
# 健康检查
|
||||||
@app.route('/health', methods=['GET'])
|
@app.route('/health', methods=['GET'])
|
||||||
|
|||||||
@@ -98,7 +98,7 @@ def chat_ask():
|
|||||||
answer = _call_edge_qa(prompt)
|
answer = _call_edge_qa(prompt)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"问答编排调用失败: {e}")
|
logger.error(f"问答编排调用失败: {e}")
|
||||||
return jsonify({"error": f"AI 调用失败: {e}"}), 503
|
return jsonify({"error": "AI 服务暂时不可用,请稍后重试"}), 503
|
||||||
|
|
||||||
chat_id = db_layer.insert_chat_history(
|
chat_id = db_layer.insert_chat_history(
|
||||||
user_question=question,
|
user_question=question,
|
||||||
|
|||||||
61
fam-core/src/fam_core/img_proxy.py
Normal file
61
fam-core/src/fam_core/img_proxy.py
Normal file
@@ -0,0 +1,61 @@
|
|||||||
|
"""
|
||||||
|
Img-Proxy - NAS 端图片代理(图源在 Oracle,计算全部在 Oracle)
|
||||||
|
|
||||||
|
浏览器不直连 Oracle(避免公网暴露 5000 端口与 token 外泄),
|
||||||
|
而是访问 NAS fam-core 的 /api/proxy/*,由 NAS 出网到 Oracle 拉取 jpeg 回传。
|
||||||
|
|
||||||
|
Oracle 侧已有磁盘缓存 / VLM 人物定位裁剪,NAS 端仅透传,不做图像计算。
|
||||||
|
"""
|
||||||
|
from flask import Blueprint, Response, request
|
||||||
|
|
||||||
|
from .oracle_sync import get_sync
|
||||||
|
from .logger import setup_logger
|
||||||
|
|
||||||
|
logger = setup_logger('fam-core.img_proxy')
|
||||||
|
|
||||||
|
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
|
||||||
|
|
||||||
|
|
||||||
|
@img_bp.route('/api/proxy/frame', methods=['GET'])
|
||||||
|
def proxy_frame():
|
||||||
|
"""事件缩略帧: video_id + ts -> jpeg"""
|
||||||
|
data = _fetch('/api/oracle/frame', {
|
||||||
|
'video_id': request.args.get('video_id', type=int) or 0,
|
||||||
|
'ts': request.args.get('ts', ''),
|
||||||
|
'w': request.args.get('w', 400, type=int),
|
||||||
|
})
|
||||||
|
if data is None:
|
||||||
|
return '', 404
|
||||||
|
return Response(data, mimetype='image/jpeg',
|
||||||
|
headers={'Cache-Control': 'public, max-age=86400'})
|
||||||
|
|
||||||
|
|
||||||
|
@img_bp.route('/api/proxy/avatar', methods=['GET'])
|
||||||
|
def proxy_avatar():
|
||||||
|
"""人物头像: label/canonical -> jpeg"""
|
||||||
|
data = _fetch('/api/oracle/avatar', {
|
||||||
|
'label': request.args.get('label', ''),
|
||||||
|
'w': request.args.get('w', 160, type=int),
|
||||||
|
})
|
||||||
|
if data is None:
|
||||||
|
return '', 404
|
||||||
|
return Response(data, mimetype='image/jpeg',
|
||||||
|
headers={'Cache-Control': 'public, max-age=86400'})
|
||||||
@@ -1,10 +1,9 @@
|
|||||||
"""
|
"""
|
||||||
日志工具 - 统一格式,带 task_id 作为 trace_id
|
日志工具 - 统一格式(stdout + 文件双写)
|
||||||
"""
|
"""
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
from datetime import datetime
|
|
||||||
|
|
||||||
# 日志目录:fam-core/logs/(相对 src 的上一级),失败则退化为仅 stdout
|
# 日志目录:fam-core/logs/(相对 src 的上一级),失败则退化为仅 stdout
|
||||||
_LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), 'logs')
|
_LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), 'logs')
|
||||||
@@ -33,12 +32,3 @@ def setup_logger(name='fam-core', level=logging.INFO):
|
|||||||
except OSError:
|
except OSError:
|
||||||
pass # 目录不可写时退化为仅 stdout
|
pass # 目录不可写时退化为仅 stdout
|
||||||
return logger
|
return logger
|
||||||
|
|
||||||
|
|
||||||
def log_task(logger, task_id, stage, message, level=logging.INFO, duration_ms=None):
|
|
||||||
"""带 task_id 的结构化日志"""
|
|
||||||
parts = [f"[task_id={task_id}]", stage]
|
|
||||||
if duration_ms is not None:
|
|
||||||
parts.append(f"done in {duration_ms}ms")
|
|
||||||
parts.append(message)
|
|
||||||
logger.log(level, ' '.join(parts))
|
|
||||||
|
|||||||
7
fam-core/tests/conftest.py
Normal file
7
fam-core/tests/conftest.py
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
import os
|
||||||
|
import sys
|
||||||
|
|
||||||
|
# 让测试能直接 `from fam_core.xxx import yyy`,无需先 pip install -e .
|
||||||
|
_SRC = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), 'src')
|
||||||
|
if _SRC not in sys.path:
|
||||||
|
sys.path.insert(0, _SRC)
|
||||||
55
fam-core/tests/test_chat_handler.py
Normal file
55
fam-core/tests/test_chat_handler.py
Normal file
@@ -0,0 +1,55 @@
|
|||||||
|
from fam_core.chat_handler.chat_handler import _format_events
|
||||||
|
|
||||||
|
|
||||||
|
def test_multiple_persons_and_attention():
|
||||||
|
rows = [{
|
||||||
|
"ts": "2026-08-21 16:11:25", "camera_name": "客厅",
|
||||||
|
"description": "两人在客厅玩耍",
|
||||||
|
"person_list_json": '["汤圆", "爷爷"]',
|
||||||
|
"is_attention_event": True,
|
||||||
|
}]
|
||||||
|
out = _format_events(rows)
|
||||||
|
assert out == "[2026-08-21 16:11 客厅] 汤圆,爷爷: 两人在客厅玩耍 [关注事件]"
|
||||||
|
|
||||||
|
|
||||||
|
def test_no_person_defaults_to_wu_ren():
|
||||||
|
rows = [{
|
||||||
|
"ts": "2026-08-21 16:11:25", "camera_name": "客厅",
|
||||||
|
"description": "空房间", "person_list_json": '[]',
|
||||||
|
"is_attention_event": False,
|
||||||
|
}]
|
||||||
|
out = _format_events(rows)
|
||||||
|
assert "无人:" in out
|
||||||
|
assert "[关注事件]" not in out
|
||||||
|
|
||||||
|
|
||||||
|
def test_person_list_already_decoded_list():
|
||||||
|
"""pymysql 某些驱动/字段类型可能已经把 JSON 解码成 list,不是字符串。"""
|
||||||
|
rows = [{
|
||||||
|
"ts": "2026-08-21 16:11:25", "camera_name": "客厅",
|
||||||
|
"description": "在客厅", "person_list_json": ["汤圆"],
|
||||||
|
"is_attention_event": False,
|
||||||
|
}]
|
||||||
|
out = _format_events(rows)
|
||||||
|
assert "汤圆:" in out
|
||||||
|
|
||||||
|
|
||||||
|
def test_malformed_person_list_json_degrades_gracefully():
|
||||||
|
rows = [{
|
||||||
|
"ts": "2026-08-21 16:11:25", "camera_name": "客厅",
|
||||||
|
"description": "在客厅", "person_list_json": "not valid json",
|
||||||
|
"is_attention_event": False,
|
||||||
|
}]
|
||||||
|
out = _format_events(rows)
|
||||||
|
assert "无人:" in out
|
||||||
|
|
||||||
|
|
||||||
|
def test_multiple_rows_joined_by_newline():
|
||||||
|
rows = [
|
||||||
|
{"ts": "2026-08-21 16:11:25", "camera_name": "客厅",
|
||||||
|
"description": "第一条", "person_list_json": '["汤圆"]', "is_attention_event": False},
|
||||||
|
{"ts": "2026-08-21 16:12:00", "camera_name": "客厅",
|
||||||
|
"description": "第二条", "person_list_json": '["汤圆"]', "is_attention_event": False},
|
||||||
|
]
|
||||||
|
out = _format_events(rows)
|
||||||
|
assert len(out.split("\n")) == 2
|
||||||
25
fam-core/tests/test_config_loader.py
Normal file
25
fam-core/tests/test_config_loader.py
Normal file
@@ -0,0 +1,25 @@
|
|||||||
|
import os
|
||||||
|
|
||||||
|
from fam_core.config_loader import _resolve_env_vars
|
||||||
|
|
||||||
|
|
||||||
|
def test_resolves_string_env_var(monkeypatch):
|
||||||
|
monkeypatch.setenv("FAM_CORE_TEST_VAR", "hello")
|
||||||
|
assert _resolve_env_vars("${FAM_CORE_TEST_VAR}") == "hello"
|
||||||
|
|
||||||
|
|
||||||
|
def test_unset_env_var_left_as_literal():
|
||||||
|
assert _resolve_env_vars("${SOME_TOTALLY_UNSET_VAR_XYZ}") == "${SOME_TOTALLY_UNSET_VAR_XYZ}"
|
||||||
|
|
||||||
|
|
||||||
|
def test_resolves_nested_dict(monkeypatch):
|
||||||
|
monkeypatch.setenv("FAM_CORE_TEST_VAR", "secret")
|
||||||
|
data = {"a": {"b": "${FAM_CORE_TEST_VAR}"}, "c": ["${FAM_CORE_TEST_VAR}", "plain"]}
|
||||||
|
out = _resolve_env_vars(data)
|
||||||
|
assert out == {"a": {"b": "secret"}, "c": ["secret", "plain"]}
|
||||||
|
|
||||||
|
|
||||||
|
def test_non_string_values_passed_through():
|
||||||
|
assert _resolve_env_vars(42) == 42
|
||||||
|
assert _resolve_env_vars(True) is True
|
||||||
|
assert _resolve_env_vars(None) is None
|
||||||
Reference in New Issue
Block a user