Compare commits
2 Commits
0d8e62607c
...
01bbb39750
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
01bbb39750 | ||
|
|
38dea6ba2e |
@@ -32,3 +32,21 @@ chat_handler:
|
|||||||
# 智能问答统一走 FAM-Edge 编排端点(Gemini → NVIDIA → 本地 Ollama 兜底)
|
# 智能问答统一走 FAM-Edge 编排端点(Gemini → NVIDIA → 本地 Ollama 兜底)
|
||||||
qa_url: "http://129.146.203.203:5000/api/edge/chat/ask"
|
qa_url: "http://129.146.203.203:5000/api/edge/chat/ask"
|
||||||
timeout: 120
|
timeout: 120
|
||||||
|
|
||||||
|
# 运动监测通知服务(2026-08-22 新增)
|
||||||
|
# 在 NAS 本机轮询群晖 Surveillance Station 的运动侦测事件,把所需信息主动 POST
|
||||||
|
# 推送到甲骨文 FAM-Edge(单向 NAS -> Oracle,甲骨文不再反向访问 NAS)。
|
||||||
|
# SS 凭据走环境变量(${DSM_ACCOUNT}/${DSM_PASSWORD}),由启动脚本 source 的 .env 提供。
|
||||||
|
motion_notifier:
|
||||||
|
enabled: true
|
||||||
|
dsm_host: "192.168.50.64" # Surveillance Station 所在地址(NAS 本机)
|
||||||
|
dsm_port: 5000
|
||||||
|
dsm_account: "${DSM_ACCOUNT}"
|
||||||
|
dsm_password: "${DSM_PASSWORD}"
|
||||||
|
camera_ids: [2] # 关注的摄像头(Generic_ONVIF-001 的 camera_id)
|
||||||
|
oracle_base_url: "http://129.146.203.203:5000" # 与 oracle_sync.base_url 一致
|
||||||
|
oracle_token: "${ORACLE_SYNC_TOKEN}" # 与 oracle_sync.token 一致
|
||||||
|
poll_interval_sec: 60 # 轮询间隔
|
||||||
|
poll_window_hours: 2 # 每轮回看窗口(小时),覆盖轮询间隔内的新事件
|
||||||
|
batch_size: 100 # 单批推送上限
|
||||||
|
timeout_sec: 10 # 单次 SS 请求超时
|
||||||
|
|||||||
@@ -17,9 +17,11 @@ sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
|||||||
from .config_loader import load_config
|
from .config_loader import load_config
|
||||||
from .logger import setup_logger
|
from .logger import setup_logger
|
||||||
from .oracle_sync import get_sync
|
from .oracle_sync import get_sync
|
||||||
|
from .motion_notifier.motion_notifier import get_motion_notifier
|
||||||
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
|
from .img_proxy import img_bp
|
||||||
|
from .motion_bp import motion_bp
|
||||||
from .ui_api import ui_bp
|
from .ui_api import ui_bp
|
||||||
from .static_app import static_bp
|
from .static_app import static_bp
|
||||||
|
|
||||||
@@ -32,6 +34,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.register_blueprint(img_bp)
|
||||||
|
app.register_blueprint(motion_bp)
|
||||||
app.register_blueprint(ui_bp)
|
app.register_blueprint(ui_bp)
|
||||||
app.register_blueprint(static_bp)
|
app.register_blueprint(static_bp)
|
||||||
|
|
||||||
@@ -56,6 +59,15 @@ try:
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Oracle-Sync 启动失败: {e}")
|
logger.error(f"Oracle-Sync 启动失败: {e}")
|
||||||
|
|
||||||
|
# 初始化运动监测通知服务(NAS 轮询 SS 事件 -> 推送甲骨文;enabled 才真正启动线程)
|
||||||
|
_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'])
|
@app.route('/api/status', methods=['GET'])
|
||||||
def status():
|
def status():
|
||||||
@@ -63,6 +75,7 @@ def status():
|
|||||||
return jsonify({
|
return jsonify({
|
||||||
"service": "fam-core",
|
"service": "fam-core",
|
||||||
"sync": _sync.status() if _sync else {"running": False, "error": "未初始化"},
|
"sync": _sync.status() if _sync else {"running": False, "error": "未初始化"},
|
||||||
|
"motion": _notifier.status() if _notifier else {"running": False, "error": "未初始化"},
|
||||||
}), 200
|
}), 200
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -413,6 +413,35 @@ def set_sync_cursor(value: str):
|
|||||||
conn.close()
|
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()
|
||||||
|
|
||||||
|
|
||||||
# ============================================================
|
# ============================================================
|
||||||
# 统计
|
# 统计
|
||||||
# ============================================================
|
# ============================================================
|
||||||
|
|||||||
53
fam-core/src/fam_core/motion_bp.py
Normal file
53
fam-core/src/fam_core/motion_bp.py
Normal file
@@ -0,0 +1,53 @@
|
|||||||
|
"""
|
||||||
|
运动监测 Webhook 接收(fam-core)
|
||||||
|
|
||||||
|
群晖 Surveillance Station 可配置 HTTP 推送(Webhook),将事件实时 POST 到本端点,
|
||||||
|
本端点即时转发到甲骨文 FAM-Edge。与 MotionNotifier 轮询互补,提供更低的事件延迟。
|
||||||
|
|
||||||
|
SS Webhook 的实际 payload 格式随套件版本而异,本路由做宽松解析:
|
||||||
|
- 接受 JSON 或表单;
|
||||||
|
- 兼容 {events:[...]} / {data:[...]} / 单事件对象 / 裸数组;
|
||||||
|
- 字段名兼容 id/event_id、thumbnail_url/thumbnail_dir。
|
||||||
|
无法识别的字段会原样透传,甲骨文端按已知字段落库(未知字段忽略)。
|
||||||
|
"""
|
||||||
|
from flask import Blueprint, request, 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__)
|
||||||
|
|
||||||
|
|
||||||
|
def _coerce_events(payload) -> list:
|
||||||
|
"""从各种可能的 SS payload 形态中提取事件列表。"""
|
||||||
|
if isinstance(payload, list):
|
||||||
|
return payload
|
||||||
|
if isinstance(payload, dict):
|
||||||
|
for key in ('events', 'data', 'event', 'items'):
|
||||||
|
v = payload.get(key)
|
||||||
|
if isinstance(v, list):
|
||||||
|
return v
|
||||||
|
# 单事件对象
|
||||||
|
return [payload]
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
@motion_bp.route('/api/ss/webhook', methods=['POST'])
|
||||||
|
def ss_webhook():
|
||||||
|
data = request.get_json(silent=True)
|
||||||
|
if not data:
|
||||||
|
data = request.form.to_dict() or None
|
||||||
|
if not data:
|
||||||
|
return jsonify({"error": "Invalid payload"}), 400
|
||||||
|
events = _coerce_events(data)
|
||||||
|
if not events:
|
||||||
|
return jsonify({"error": "no events found in payload"}), 400
|
||||||
|
pushed = get_motion_notifier().push_events_to_oracle(events)
|
||||||
|
return jsonify({"status": "ok", "received": len(events), "pushed": pushed}), 200
|
||||||
|
|
||||||
|
|
||||||
|
@motion_bp.route('/api/ss/status', methods=['GET'])
|
||||||
|
def ss_status():
|
||||||
|
return jsonify(get_motion_notifier().status()), 200
|
||||||
4
fam-core/src/fam_core/motion_notifier/__init__.py
Normal file
4
fam-core/src/fam_core/motion_notifier/__init__.py
Normal file
@@ -0,0 +1,4 @@
|
|||||||
|
"""MotionNotifier 包:NAS 端运动监测通知服务。"""
|
||||||
|
from .motion_notifier import MotionNotifier, get_motion_notifier
|
||||||
|
|
||||||
|
__all__ = ["MotionNotifier", "get_motion_notifier"]
|
||||||
260
fam-core/src/fam_core/motion_notifier/motion_notifier.py
Normal file
260
fam-core/src/fam_core/motion_notifier/motion_notifier.py
Normal file
@@ -0,0 +1,260 @@
|
|||||||
|
"""
|
||||||
|
MotionNotifier - NAS 端运动监测通知服务(新架构,2026-08-22)
|
||||||
|
|
||||||
|
职责:
|
||||||
|
在 NAS 本机轮询群晖 Surveillance Station 的运动侦测事件
|
||||||
|
(SYNO.SurveillanceStation.EventCenter.Event,参数名下划线风格 camera_ids/
|
||||||
|
start_time/end_time,event_type=10 即运动),把"需要的信息"主动 POST 推送到
|
||||||
|
甲骨文 FAM-Edge 的 /api/ss/motion 接口。
|
||||||
|
|
||||||
|
数据方向(关键约束): NAS -> Oracle,单向。甲骨文不再反向访问 NAS。
|
||||||
|
- 原 fam-edge 的 dsm_motion_client(甲骨文主动查 SS API)已停用;
|
||||||
|
- 改由本服务在 NAS 上拉取 SS 事件,推送后由甲骨文本地落库 ss_motion_events,
|
||||||
|
供 video_processor 本地预过滤使用。
|
||||||
|
|
||||||
|
两种推送触发(都汇聚到 push_events_to_oracle):
|
||||||
|
1. 轮询(默认开): 每 poll_interval_sec 拉一次新事件,增量推送到 Oracle。
|
||||||
|
2. Webhook(可选): fam-core 暴露 POST /api/ss/webhook,群晖 SS 配置 HTTP 推送后
|
||||||
|
可实时转发(见 fam_core.motion_bp)。
|
||||||
|
"""
|
||||||
|
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,Webhook 路由经此获取)。"""
|
||||||
|
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.203.203: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))
|
||||||
|
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
|
||||||
|
self._last_poll_at = None
|
||||||
|
self._last_error = None
|
||||||
|
self._pushed_total = 0
|
||||||
|
|
||||||
|
@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": 100, "_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_id(CMS 多机场景的服务器 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 _init_cursor(self):
|
||||||
|
"""启动时初始化 last_event_id: 取当前 SS 最大事件 id,只推增量(不回灌历史)。"""
|
||||||
|
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:
|
||||||
|
saved = db_layer.get_motion_cursor()
|
||||||
|
self._last_event_id = int(saved) if saved else 0
|
||||||
|
logger.info(f"运动通知游标初始化: 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 # 查询失败,下一轮重试
|
||||||
|
# 只保留 id 大于游标的新事件(SS 事件 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)))
|
||||||
|
for i in range(0, len(new), self.batch_size):
|
||||||
|
batch = new[i:i + self.batch_size]
|
||||||
|
self.push_events_to_oracle(batch)
|
||||||
|
self._last_event_id = max(int(e.get('id', 0)) for e in new)
|
||||||
|
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
|
||||||
|
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)
|
||||||
|
|
||||||
|
def status(self) -> dict:
|
||||||
|
return {
|
||||||
|
"running": self.is_alive(),
|
||||||
|
"enabled": self.enabled,
|
||||||
|
"last_event_id": self._last_event_id,
|
||||||
|
"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}",
|
||||||
|
}
|
||||||
@@ -47,22 +47,17 @@ video_processing:
|
|||||||
# 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed
|
# 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed
|
||||||
vision_order: ["gemini", "nvidia"]
|
vision_order: ["gemini", "nvidia"]
|
||||||
|
|
||||||
# DSM 运动侦测预过滤(2026-08-22 新增):分析前先查一下群晖 Surveillance Station
|
# 运动侦测预过滤(2026-08-22 重构):
|
||||||
# 自己记录的运动侦测事件(SYNO.SurveillanceStation.EventCenter.Event,未公开文档的
|
# - 甲骨文【不再反向访问 NAS】。原 dsm_motion_client(甲骨文主动查 SS API)已停用,
|
||||||
# 内部接口,参数名是下划线风格 camera_ids/start_time/end_time),这段时间窗口一条
|
# 改由 NAS 端 fam-core 的 MotionNotifier 轮询/接收 SS 事件后,主动 POST 推送到
|
||||||
# 运动事件都没有就跳过云端分析(标记 done,compute_provider=skipped_no_motion),
|
# 甲骨文 /api/ss/motion,落库 ss_motion_events。
|
||||||
# 省掉长期无人时段白白消耗的 Gemini/NVIDIA 配额。
|
# - video_processor 分析前调用 db.has_motion_in_range_local()(本地运动事件表)做
|
||||||
# 账号密码走 .env(DSM_ACCOUNT/DSM_PASSWORD),不明文入库;查询失败/未配置一律
|
# 预过滤:窗口内无运动事件则跳过云端分析(compute_provider=skipped_no_motion)。
|
||||||
# fail-open(照常送云端分析),不会因为这层可选优化漏检真实事件。
|
# - 本地运动事件表为空(冷启动/尚未收到推送)一律 fail-open(照常分析),不漏检。
|
||||||
|
# 此 block 仅保留 min_motion_seconds 语义参考;host/account 等反向访问字段已弃用。
|
||||||
dsm_motion_prefilter:
|
dsm_motion_prefilter:
|
||||||
enabled: true
|
enabled: false # 已停用:甲骨文不再主动访问 NAS SS
|
||||||
host: "192.168.50.64"
|
|
||||||
port: 5000
|
|
||||||
account: "${DSM_ACCOUNT}"
|
|
||||||
password: "${DSM_PASSWORD}"
|
|
||||||
camera_id: 2 # Surveillance Station 里 Generic_ONVIF-001 的 camera_id
|
|
||||||
min_motion_seconds: 0 # 窗口内运动事件总时长需 ≥ 此值才算"有运动"(0=有事件就算)
|
min_motion_seconds: 0 # 窗口内运动事件总时长需 ≥ 此值才算"有运动"(0=有事件就算)
|
||||||
timeout_sec: 10
|
|
||||||
|
|
||||||
# 智能问答降级链(与视频分析独立):Gemini -> NVIDIA -> 本地 Ollama
|
# 智能问答降级链(与视频分析独立):Gemini -> NVIDIA -> 本地 Ollama
|
||||||
models:
|
models:
|
||||||
|
|||||||
@@ -197,6 +197,36 @@ def oracle_avatar():
|
|||||||
return Response(data, mimetype='image/jpeg')
|
return Response(data, mimetype='image/jpeg')
|
||||||
|
|
||||||
|
|
||||||
|
@api_bp.route('/api/ss/motion', methods=['POST'])
|
||||||
|
def ss_motion():
|
||||||
|
"""接收 NAS 推送的运动侦测事件(单向:NAS -> Oracle)。
|
||||||
|
|
||||||
|
请求体: {"token": "...", "events": [
|
||||||
|
{"event_id": 25349, "camera_id": 2, "event_type": 10,
|
||||||
|
"start_time": 1787318023, "duration": 149, "thumbnail_url": "107471,12075"}
|
||||||
|
]}
|
||||||
|
start_time/duration 为 Unix epoch(与 SS 同源,时区无关)。
|
||||||
|
落库 ss_motion_events(按 event_id 幂等),供 video_processor 本地预过滤使用。
|
||||||
|
"""
|
||||||
|
if not _check_token():
|
||||||
|
return jsonify({"error": "unauthorized"}), 401
|
||||||
|
data = request.get_json(silent=True)
|
||||||
|
if not data or 'events' not in data:
|
||||||
|
return jsonify({"error": "缺少 events"}), 400
|
||||||
|
events = data.get('events') or []
|
||||||
|
if not isinstance(events, list):
|
||||||
|
return jsonify({"error": "events 必须是数组"}), 400
|
||||||
|
try:
|
||||||
|
stored = state.get_db().record_motion_events(events)
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"ss_motion 落库异常: {e}")
|
||||||
|
return jsonify({"error": str(e)}), 500
|
||||||
|
if stored:
|
||||||
|
state.get_db().record_activity(
|
||||||
|
'motion', 'push', f"接收 NAS 运动事件 {stored} 条")
|
||||||
|
return jsonify({"status": "ok", "received": len(events), "stored": stored}), 200
|
||||||
|
|
||||||
|
|
||||||
@api_bp.route('/health', methods=['GET'])
|
@api_bp.route('/health', methods=['GET'])
|
||||||
def health():
|
def health():
|
||||||
"""健康检查:DB 连通性 + producer/consumer 线程存活状态。
|
"""健康检查:DB 连通性 + producer/consumer 线程存活状态。
|
||||||
|
|||||||
@@ -1,5 +1,15 @@
|
|||||||
"""
|
"""
|
||||||
DsmMotionClient - 查询群晖 Surveillance Station 的运动侦测事件(预过滤用)
|
[已弃用 / DEPRECATED] 本模块自 2026-08-22 起不再被调用。
|
||||||
|
|
||||||
|
架构约束:甲骨文 FAM-Edge 不得反向访问 NAS。原 video_processor 在此处主动查询
|
||||||
|
NAS 的 Surveillance Station(甲骨文 -> NAS),违反该约束,已停用。
|
||||||
|
|
||||||
|
替代方案:由 NAS 端 fam-core 的 MotionNotifier 轮询/接收 SS 事件后,主动 POST
|
||||||
|
推送到甲骨文的 /api/ss/motion,落库 ss_motion_events;video_processor 改用
|
||||||
|
db.has_motion_in_range_local() 做本地预过滤。本文件保留仅作参考,请勿再实例化。
|
||||||
|
|
||||||
|
---
|
||||||
|
DsmMotionClient(旧实现,仅作历史参考) - 查询群晖 Surveillance Station 的运动侦测事件(预过滤用)
|
||||||
|
|
||||||
背景: fam-edge 原来不管这段 30 分钟录像里有没有人走动,一律整段送云端 VLM 分析,
|
背景: fam-edge 原来不管这段 30 分钟录像里有没有人走动,一律整段送云端 VLM 分析,
|
||||||
配额/耗时都花在长期无人的空转时段上。DSM 的 Surveillance Station 用
|
配额/耗时都花在长期无人的空转时段上。DSM 的 Surveillance Station 用
|
||||||
|
|||||||
@@ -110,10 +110,21 @@ class OracleDB:
|
|||||||
detail TEXT,
|
detail TEXT,
|
||||||
ts TEXT
|
ts TEXT
|
||||||
);
|
);
|
||||||
|
CREATE TABLE IF NOT EXISTS ss_motion_events (
|
||||||
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
|
event_id INTEGER UNIQUE,
|
||||||
|
camera_id INTEGER,
|
||||||
|
event_type INTEGER,
|
||||||
|
start_time INTEGER,
|
||||||
|
duration INTEGER,
|
||||||
|
thumbnail_url TEXT,
|
||||||
|
received_at TEXT
|
||||||
|
);
|
||||||
CREATE INDEX IF NOT EXISTS idx_videos_updated ON videos(updated_at);
|
CREATE INDEX IF NOT EXISTS idx_videos_updated ON videos(updated_at);
|
||||||
CREATE INDEX IF NOT EXISTS idx_events_video ON events(video_id);
|
CREATE INDEX IF NOT EXISTS idx_events_video ON events(video_id);
|
||||||
CREATE INDEX IF NOT EXISTS idx_model_calls_created ON model_calls(created_at);
|
CREATE INDEX IF NOT EXISTS idx_model_calls_created ON model_calls(created_at);
|
||||||
CREATE INDEX IF NOT EXISTS idx_activity_ts ON service_activity(ts);
|
CREATE INDEX IF NOT EXISTS idx_activity_ts ON service_activity(ts);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_motion_window ON ss_motion_events(start_time, event_type);
|
||||||
""")
|
""")
|
||||||
# 兼容旧库:补 retry_count / file_valid / media 等列(生产-消费队列用)
|
# 兼容旧库:补 retry_count / file_valid / media 等列(生产-消费队列用)
|
||||||
cols = [r[1] for r in c.execute("PRAGMA table_info(videos)").fetchall()]
|
cols = [r[1] for r in c.execute("PRAGMA table_info(videos)").fetchall()]
|
||||||
@@ -193,6 +204,71 @@ class OracleDB:
|
|||||||
"ORDER BY id DESC LIMIT ?", (int(limit),)).fetchall()
|
"ORDER BY id DESC LIMIT ?", (int(limit),)).fetchall()
|
||||||
return [dict(r) for r in rows]
|
return [dict(r) for r in rows]
|
||||||
|
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
# 运动侦测事件(NAS 推送,单向:NAS -> Oracle,不再反向访问 NAS)
|
||||||
|
# 由 fam-core 的 MotionNotifier 轮询/接收 SS 事件后 POST 到 /api/ss/motion。
|
||||||
|
# ------------------------------------------------------------------
|
||||||
|
def record_motion_events(self, events: List[Dict]) -> int:
|
||||||
|
"""批量 upsert NAS 推送来的运动事件(按 event_id 幂等)。返回成功条数。"""
|
||||||
|
n = 0
|
||||||
|
now = _now_iso()
|
||||||
|
for e in (events or []):
|
||||||
|
eid = e.get('event_id')
|
||||||
|
if eid is None:
|
||||||
|
continue
|
||||||
|
self._conn.execute(
|
||||||
|
"""INSERT INTO ss_motion_events
|
||||||
|
(event_id, camera_id, event_type, start_time, duration,
|
||||||
|
thumbnail_url, received_at)
|
||||||
|
VALUES (?,?,?,?,?,?,?)
|
||||||
|
ON CONFLICT(event_id) DO UPDATE SET
|
||||||
|
camera_id=excluded.camera_id,
|
||||||
|
event_type=excluded.event_type,
|
||||||
|
start_time=excluded.start_time,
|
||||||
|
duration=excluded.duration,
|
||||||
|
thumbnail_url=excluded.thumbnail_url,
|
||||||
|
received_at=excluded.received_at""",
|
||||||
|
(int(eid), e.get('camera_id'), e.get('event_type'),
|
||||||
|
e.get('start_time'), e.get('duration'),
|
||||||
|
e.get('thumbnail_url'), now))
|
||||||
|
n += 1
|
||||||
|
if n:
|
||||||
|
self._conn.commit()
|
||||||
|
return n
|
||||||
|
|
||||||
|
def get_recent_motion_events(self, limit: int = 50) -> List[Dict]:
|
||||||
|
"""最近运动事件(按 start_time 倒序)。"""
|
||||||
|
rows = self._conn.execute(
|
||||||
|
"SELECT id, event_id, camera_id, event_type, start_time, duration, "
|
||||||
|
"thumbnail_url, received_at FROM ss_motion_events "
|
||||||
|
"ORDER BY start_time DESC LIMIT ?", (int(limit),)).fetchall()
|
||||||
|
return [dict(r) for r in rows]
|
||||||
|
|
||||||
|
def has_motion_in_range_local(self, start_ts: int, end_ts: int,
|
||||||
|
camera_id: int = None) -> Optional[bool]:
|
||||||
|
"""本地运动预过滤:判断 [start_ts, end_ts] 窗口内是否存在运动事件。
|
||||||
|
|
||||||
|
替代原 dsm_motion_client 反向访问 NAS 的做法。返回:
|
||||||
|
- None : 本地运动事件表为空(冷启动,尚未收到 NAS 推送),调用方
|
||||||
|
必须 fail-open(照常送云端分析),不能当作"无运动"跳过。
|
||||||
|
- True/False : 窗口内确有/确无运动事件。
|
||||||
|
start_ts/end_ts 为 Unix epoch(与 SS 事件 start_time 同源,时区无关)。
|
||||||
|
"""
|
||||||
|
total = self._conn.execute(
|
||||||
|
"SELECT COUNT(*) c FROM ss_motion_events").fetchone()['c']
|
||||||
|
if total == 0:
|
||||||
|
return None
|
||||||
|
sql = ("SELECT COUNT(*) c FROM ss_motion_events "
|
||||||
|
"WHERE event_type = 10 "
|
||||||
|
"AND start_time <= ? "
|
||||||
|
"AND (start_time + COALESCE(duration,0)) >= ?")
|
||||||
|
params = [end_ts, start_ts]
|
||||||
|
if camera_id is not None:
|
||||||
|
sql += " AND camera_id = ?"
|
||||||
|
params.append(camera_id)
|
||||||
|
cnt = self._conn.execute(sql, params).fetchone()['c']
|
||||||
|
return cnt > 0
|
||||||
|
|
||||||
def get_queue_status(self) -> Dict:
|
def get_queue_status(self) -> Dict:
|
||||||
"""实时队列/处理状态(前端服务状态卡用)。"""
|
"""实时队列/处理状态(前端服务状态卡用)。"""
|
||||||
total = self._conn.execute("SELECT COUNT(*) c FROM videos").fetchone()['c']
|
total = self._conn.execute("SELECT COUNT(*) c FROM videos").fetchone()['c']
|
||||||
|
|||||||
@@ -22,7 +22,6 @@ from .logger import setup_logger
|
|||||||
from .config_loader import load_config
|
from .config_loader import load_config
|
||||||
from .model_adapters.adapter_factory import build_adapters
|
from .model_adapters.adapter_factory import build_adapters
|
||||||
from .model_adapters.base_adapter import BaseModelAdapter
|
from .model_adapters.base_adapter import BaseModelAdapter
|
||||||
from .dsm_motion_client import DsmMotionClient
|
|
||||||
from . import oracle_db
|
from . import oracle_db
|
||||||
|
|
||||||
logger = setup_logger('fam-edge.video_processor')
|
logger = setup_logger('fam-edge.video_processor')
|
||||||
@@ -194,7 +193,6 @@ class VideoProcessor:
|
|||||||
adapters = build_adapters(self.config.get('models', []))
|
adapters = build_adapters(self.config.get('models', []))
|
||||||
self.vision_adapters: Dict[str, BaseModelAdapter] = {
|
self.vision_adapters: Dict[str, BaseModelAdapter] = {
|
||||||
a.provider_name: a for a in adapters if a.get_role() == 'vision'}
|
a.provider_name: a for a in adapters if a.get_role() == 'vision'}
|
||||||
self.dsm_motion = DsmMotionClient(self.config.get('dsm_motion_prefilter', {}))
|
|
||||||
|
|
||||||
def _ordered_vision_adapters(self) -> List[BaseModelAdapter]:
|
def _ordered_vision_adapters(self) -> List[BaseModelAdapter]:
|
||||||
ordered = []
|
ordered = []
|
||||||
@@ -238,18 +236,23 @@ class VideoProcessor:
|
|||||||
if event_start:
|
if event_start:
|
||||||
self.db.set_event_start_time(video_id, event_start)
|
self.db.set_event_start_time(video_id, event_start)
|
||||||
|
|
||||||
# DSM 运动侦测预过滤:这段时间窗口里群晖自己记录的运动事件一条都没有,
|
# 运动侦测预过滤(v2: 改用 NAS 推送来的本地运动事件,不再反向访问 NAS):
|
||||||
# 就跳过云端分析(省配额)。查询本身失败/未配置一律 fail-open(照常分析),
|
# 这段时间窗口里没有任何运动事件,就跳过云端分析(省配额)。
|
||||||
|
# 本地运动事件表为空(冷启动,尚未收到 NAS 推送)一律 fail-open(照常分析),
|
||||||
# 绝不能因为这层可选优化漏检真实事件。
|
# 绝不能因为这层可选优化漏检真实事件。
|
||||||
duration_sec = (vmeta or {}).get('duration_sec') if self.file_validate else None
|
duration_sec = (vmeta or {}).get('duration_sec') if self.file_validate else None
|
||||||
if event_start and duration_sec:
|
if event_start and duration_sec:
|
||||||
try:
|
try:
|
||||||
start_dt = datetime.strptime(event_start, '%Y-%m-%d %H:%M:%S')
|
# event_start 是北京时间墙钟,转成与 SS 事件同源的 Unix epoch
|
||||||
has_motion = self.dsm_motion.has_motion_in_range(start_dt, duration_sec)
|
start_dt = datetime.strptime(event_start, '%Y-%m-%d %H:%M:%S').replace(
|
||||||
|
tzinfo=timezone(timedelta(hours=8)))
|
||||||
|
start_ts = int(start_dt.timestamp())
|
||||||
|
end_ts = start_ts + int(duration_sec)
|
||||||
|
has_motion = self.db.has_motion_in_range_local(start_ts, end_ts)
|
||||||
except ValueError:
|
except ValueError:
|
||||||
has_motion = None
|
has_motion = None
|
||||||
if has_motion is False:
|
if has_motion is False:
|
||||||
logger.info(f"[video_id={video_id}] DSM 运动预过滤:该时段无运动,跳过云端分析")
|
logger.info(f"[video_id={video_id}] 运动预过滤:该时段无运动,跳过云端分析")
|
||||||
self.db.mark_video_processed(
|
self.db.mark_video_processed(
|
||||||
video_id, "(自动跳过:该时段未检测到运动)", [], [], 'skipped_no_motion')
|
video_id, "(自动跳过:该时段未检测到运动)", [], [], 'skipped_no_motion')
|
||||||
return True
|
return True
|
||||||
|
|||||||
Reference in New Issue
Block a user