[重构] 运动监测改为 Webhook 驱动(去轮询) - 关闭 MotionNotifier 轮询,NAS 端合成稳定 event_id、摄像头名映射 camera_id、EVENT_TIME 解析为 epoch,SS Webhook 直接推送甲骨文

This commit is contained in:
ericwyuan
2026-08-22 10:53:03 +08:00
parent 01bbb39750
commit ca6d425b19
4 changed files with 211 additions and 34 deletions

View File

@@ -33,20 +33,26 @@ chat_handler:
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 新增 # 运动监测通知服务2026-08-22 重构为 Webhook 驱动,不轮询
# 在 NAS 本机轮询群晖 Surveillance Station 的运动侦测事件,把所需信息主动 POST # SS 配置 Webhook 把运动事件 POST 到 NAS /api/ss/webhook本服务映射字段后
# 推送甲骨文 FAM-Edge单向 NAS -> Oracle甲骨文不再反向访问 NAS # 推送甲骨文 FAM-Edge单向 NAS -> Oracle甲骨文不再反向访问 NAS
# SS Webhook 不提供 event_id / camera_id由本服务合成 event_id、把摄像头名映射成
# camera_id启动一次性从 SS 拉取 + 以下静态映射兜底),并把 EVENT_TIME 解析成 epoch。
# SS 凭据走环境变量(${DSM_ACCOUNT}/${DSM_PASSWORD}),由启动脚本 source 的 .env 提供。 # SS 凭据走环境变量(${DSM_ACCOUNT}/${DSM_PASSWORD}),由启动脚本 source 的 .env 提供。
motion_notifier: motion_notifier:
enabled: true enabled: true
dsm_host: "192.168.50.64" # Surveillance Station 所在地址NAS 本机) poll_enabled: false # 关闭轮询,改用 SS Webhook 推送
dsm_host: "192.168.50.64" # Surveillance Station 所在地址NAS 本机,用于一次性拉摄像头映射)
dsm_port: 5000 dsm_port: 5000
dsm_account: "${DSM_ACCOUNT}" dsm_account: "${DSM_ACCOUNT}"
dsm_password: "${DSM_PASSWORD}" dsm_password: "${DSM_PASSWORD}"
camera_ids: [2] # 关注的摄像头Generic_ONVIF-001 的 camera_id # 摄像头名 -> camera_id 静态映射(优先);启动时还会从 SS 拉一份补全
camera_name_to_id:
"Generic_ONVIF-001": 2
oracle_base_url: "http://129.146.203.203:5000" # 与 oracle_sync.base_url 一致 oracle_base_url: "http://129.146.203.203:5000" # 与 oracle_sync.base_url 一致
oracle_token: "${ORACLE_SYNC_TOKEN}" # 与 oracle_sync.token 一致 oracle_token: "${ORACLE_SYNC_TOKEN}" # 与 oracle_sync.token 一致
poll_interval_sec: 60 # 轮询间隔 timeout_sec: 10 # 单次 SS 请求超时(摄像头映射拉取用)
poll_window_hours: 2 # 每轮回看窗口(小时),覆盖轮询间隔内的新事件 # 以下为可选轮询参数poll_enabled=true 时才生效,当前默认关闭)
batch_size: 100 # 单批推送上限 poll_interval_sec: 60
timeout_sec: 10 # 单次 SS 请求超时 poll_window_hours: 2
batch_size: 100

View File

@@ -59,11 +59,17 @@ try:
except Exception as e: except Exception as e:
logger.error(f"Oracle-Sync 启动失败: {e}") logger.error(f"Oracle-Sync 启动失败: {e}")
# 初始化运动监测通知服务(NAS 轮询 SS 事件 -> 推送甲骨文enabled 才真正启动线程 # 初始化运动监测通知服务(Webhook 驱动,不轮询;提供推送客户端 + 摄像头名映射
_notifier = None _notifier = None
try: try:
_notifier = get_motion_notifier() _notifier = get_motion_notifier()
_notifier.start() _notifier.start()
# 启动时一次性从 SS 拉取摄像头名->id 映射非轮询失败仅告警Webhook 仍可运行
try:
_notifier.refresh_camera_map()
logger.info("MotionNotifier 摄像头映射已加载")
except Exception as e:
logger.warning(f"摄像头映射加载失败Webhook 仍可运行camera_id 可能为空): {e}")
logger.info("MotionNotifier 已初始化") logger.info("MotionNotifier 已初始化")
except Exception as e: except Exception as e:
logger.error(f"MotionNotifier 初始化失败: {e}") logger.error(f"MotionNotifier 初始化失败: {e}")

View File

@@ -1,14 +1,25 @@
""" """
运动监测 Webhook 接收fam-core 运动监测 Webhook 接收fam-core
群晖 Surveillance Station 配置 HTTP 推送Webhook将事件实时 POST 到本端点 群晖 Surveillance Station 配置「行動規則 -> 事件=偵測到動作 -> 動作=Webhook」
本端点即时转发到甲骨文 FAM-Edge。与 MotionNotifier 轮询互补,提供更低的事件延迟 把运动事件实时 POST 到本端点,本端点即时映射字段并推送到甲骨文 FAM-Edge
数据源完全来自 SS Webhook不轮询 SS数据方向 NAS -> Oracle 单向。
SS Webhook 的实际 payload 格式随套件版本而异,本路由做宽松解析: SS Webhook 仅提供模板变量(无 event_id / camera_id 数字字段):
- 接受 JSON 或表单; %EVENT_TIME% -> event_time (本地时间字符串)
- 兼容 {events:[...]} / {data:[...]} / 单事件对象 / 裸数组; %DEVICE_NAME% -> device_name (摄像头名,需映射到 camera_id)
- 字段名兼容 id/event_id、thumbnail_url/thumbnail_dir。 %EVENT_NAME% -> event_name
无法识别的字段会原样透传,甲骨文端按已知字段落库(未知字段忽略)。 %SERVER_NAME% -> server_name
%THUMBNAIL_URL% -> thumbnail_url
因此 SS 端配置 Webhook 时,请按以下「参数名 -> 模板变量」一一添加(参数名必须一致):
event_time = %EVENT_TIME%
device_name = %DEVICE_NAME%
event_name = %EVENT_NAME%
server_name = %SERVER_NAME%
thumbnail_url = %THUMBNAIL_URL%
本端点接受 JSON 或表单;支持单条对象 / {events:[...]} / 裸数组,字段名已做兼容。
""" """
from flask import Blueprint, request, jsonify from flask import Blueprint, request, jsonify
@@ -21,7 +32,7 @@ motion_bp = Blueprint('motion_bp', __name__)
def _coerce_events(payload) -> list: def _coerce_events(payload) -> list:
"""各种可能的 SS payload 形态中提取事件列表。""" """ SS Webhook 各种形态中提取事件列表。"""
if isinstance(payload, list): if isinstance(payload, list):
return payload return payload
if isinstance(payload, dict): if isinstance(payload, dict):
@@ -29,7 +40,7 @@ def _coerce_events(payload) -> list:
v = payload.get(key) v = payload.get(key)
if isinstance(v, list): if isinstance(v, list):
return v return v
# 单事件对象 # 单事件对象(表单字段即顶层 key
return [payload] return [payload]
return [] return []
@@ -41,11 +52,21 @@ def ss_webhook():
data = request.form.to_dict() or None data = request.form.to_dict() or None
if not data: if not data:
return jsonify({"error": "Invalid payload"}), 400 return jsonify({"error": "Invalid payload"}), 400
events = _coerce_events(data) raw_events = _coerce_events(data)
if not events: if not raw_events:
return jsonify({"error": "no events found in payload"}), 400 return jsonify({"error": "no events found in payload"}), 400
pushed = get_motion_notifier().push_events_to_oracle(events) notifier = get_motion_notifier()
return jsonify({"status": "ok", "received": len(events), "pushed": pushed}), 200 norm = []
for raw in raw_events:
ev = notifier.build_event_from_webhook(raw)
if ev:
norm.append(ev)
if not norm:
return jsonify({"error": "no mappable events",
"received": len(raw_events)}), 400
pushed = notifier.push_events_to_oracle(norm)
return jsonify({"status": "ok", "received": len(raw_events),
"mapped": len(norm), "pushed": pushed}), 200
@motion_bp.route('/api/ss/status', methods=['GET']) @motion_bp.route('/api/ss/status', methods=['GET'])

View File

@@ -1,25 +1,26 @@
""" """
MotionNotifier - NAS 端运动监测通知服务(新架构2026-08-22 MotionNotifier - NAS 端运动监测通知服务(Webhook 驱动2026-08-22 重构
职责: 职责:
在 NAS 本机轮询群晖 Surveillance Station 的运动侦测事件 作为 SS Webhook 的接收侧客户端,把群晖 Surveillance Station 推送来的运动侦测
SYNO.SurveillanceStation.EventCenter.Event,参数名下划线风格 camera_ids/ 事件SYNO.SurveillanceStation.EventCenter.Event 的字段event_type=10 即运动)
start_time/end_timeevent_type=10 即运动),把"需要的信息"主动 POST 推送到 映射后主动 POST 推送到甲骨文 FAM-Edge 的 /api/ss/motion 接口。
甲骨文 FAM-Edge 的 /api/ss/motion 接口。
数据方向(关键约束): NAS -> Oracle单向。甲骨文不再反向访问 NAS。 数据方向(关键约束): NAS -> Oracle单向。甲骨文不再反向访问 NAS。
- 原 fam-edge 的 dsm_motion_client甲骨文主动查 SS API已停用 - 原 fam-edge 的 dsm_motion_client甲骨文主动查 SS API已停用
- 改由本服务在 NAS 上拉取 SS 事件,推送后由甲骨文本地落库 ss_motion_events - SS Webhook 不提供 event_id / camera_id 数字字段,本服务负责:
供 video_processor 本地预过滤使用。 * 合成稳定 event_iddevice_name+event_time+thumbnail_url 哈希,幂等去重)
* 把 %DEVICE_NAME% 经摄像头名->id 映射解析为 camera_id
* 把 %EVENT_TIME% 解析为 Unix epoch
- 推送后由甲骨文本地落库 ss_motion_events供 video_processor 本地预过滤使用。
两种推送触发(都汇聚到 push_events_to_oracle: 轮询已禁用poll_enabled=false用户要求不轮询仅保留轮询代码路径作为可选能力。
1. 轮询(默认开): 每 poll_interval_sec 拉一次新事件,增量推送到 Oracle 摄像头名->id 映射在启动时一次性从 SS 拉取补全(非轮询),并以配置 camera_name_to_id 兜底
2. Webhook可选: fam-core 暴露 POST /api/ss/webhook群晖 SS 配置 HTTP 推送后
可实时转发(见 fam_core.motion_bp
""" """
import os import os
import re import re
import time import time
import hashlib
import threading import threading
from datetime import datetime, timezone from datetime import datetime, timezone
@@ -58,6 +59,14 @@ class MotionNotifier:
self.poll_window_hours = int(cfg.get('poll_window_hours', 2)) self.poll_window_hours = int(cfg.get('poll_window_hours', 2))
self.batch_size = int(cfg.get('batch_size', 100)) self.batch_size = int(cfg.get('batch_size', 100))
self.timeout = int(cfg.get('timeout_sec', 10)) self.timeout = int(cfg.get('timeout_sec', 10))
# 是否启用轮询:用户要求不轮询,改由 SS Webhook 推送(默认关闭)
self.poll_enabled = bool(cfg.get('poll_enabled', False))
# 摄像头名 -> camera_id 静态映射(优先),启动时再从 SS 拉一份补全
self.camera_name_to_id = {
str(k): int(v) for k, v in (cfg.get('camera_name_to_id') or {}).items()
}
self._ss_name_to_id = {}
self._camera_loaded = False
self._base = f"http://{self.dsm_host}:{self.dsm_port}/webapi" self._base = f"http://{self.dsm_host}:{self.dsm_port}/webapi"
self._sid = None self._sid = None
self._running = False self._running = False
@@ -129,6 +138,134 @@ class MotionNotifier:
return events return events
return None return None
# ------------------------------------------------------------------
# 摄像头名 -> camera_id 映射(一次性从 SS 拉取,供 Webhook 使用,非轮询)
# ------------------------------------------------------------------
def refresh_camera_map(self):
"""登录 SS 拉取摄像头列表,构建 name->id 映射(启动时调用一次)。
失败时仅告警,不影响 Webhook 运行(配置兜底 camera_name_to_id 仍可生效)。
"""
try:
sid = self._login()
if not sid:
logger.warning("摄像头映射刷新失败SS 登录失败")
return
resp = requests.get(
f"{self._base}/entry.cgi",
params={"api": "SYNO.SurveillanceStation.Camera", "version": 1,
"method": "List", "_sid": sid},
timeout=self.timeout)
data = resp.json()
if not data.get('success'):
logger.warning(f"摄像头映射刷新失败:{data.get('error')}")
return
for c in (data.get('data') or {}).get('cameras', []):
name = c.get('name')
cid = c.get('id')
if name and cid is not None:
self._ss_name_to_id[str(name)] = int(cid)
self._camera_loaded = True
logger.info(f"摄像头映射已刷新:{self._ss_name_to_id}")
except Exception as e:
logger.warning(f"摄像头映射刷新异常:{e}")
def resolve_camera_id(self, device_name):
"""把 SS Webhook 的 %DEVICE_NAME%(摄像头名)解析为 camera_idint|None"""
if not device_name:
return None
name = str(device_name).strip()
for mapping in (self.camera_name_to_id, self._ss_name_to_id):
if name in mapping:
return mapping[name]
low = name.lower()
for k, v in mapping.items():
if k.lower() == low:
return v
return None
@staticmethod
def parse_ss_time_to_epoch(s):
"""把 SS Webhook 的 %EVENT_TIME%(本地时间字符串)解析为 Unix epoch。
兼容 '2026-08-22T10:35:00' / '2026-08-22T10:35' / 带空格 / 带时区偏移。
无时区信息时按 NAS 本地时区CST, +8解释得到绝对 epoch
与甲骨文 ss_motion_events.start_time同为绝对 epoch一致。
"""
if not s:
return None
s = str(s).strip()
# 带 Z / +08:00 偏移:交给 fromisoformat
try:
dt = datetime.fromisoformat(s.replace('Z', '+00:00'))
if dt.tzinfo is None:
dt = dt.replace(tzinfo=datetime.now().astimezone().tzinfo)
return int(dt.timestamp())
except ValueError:
pass
# 无偏移的多种格式兜底
base = s.replace('T', ' ')
for fmt in ('%Y-%m-%d %H:%M:%S', '%Y-%m-%d %H:%M',
'%Y/%m/%d %H:%M:%S', '%Y/%m/%d %H:%M'):
try:
dt = datetime.strptime(base, fmt)
dt = dt.replace(tzinfo=datetime.now().astimezone().tzinfo)
return int(dt.timestamp())
except ValueError:
continue
logger.warning(f"无法解析 SS 时间:{s}")
return None
@staticmethod
def _synth_event_id(device_name, event_time, thumbnail_url):
"""合成稳定的 event_id用于幂等去重
含 thumbnail_url 使同分钟内的多次真实事件可区分SS 重试(相同 payload
则哈希相同 -> 甲骨文 UNIQUE 约束幂等覆盖。
"""
raw = f"{device_name}|{event_time}|{thumbnail_url or ''}"
return int(hashlib.sha256(raw.encode('utf-8')).hexdigest()[:12], 16)
def build_event_from_webhook(self, raw):
"""把一条 SS Webhook 原始字段映射为甲骨文所需的归一化运动事件。
SS Webhook 不提供 event_id / camera_id 数字字段,此处:
- event_id : 由 device_name+event_time+thumbnail_url 合成(稳定、可去重)
- camera_id : 由 device_name 经摄像头映射解析(解析失败则 None
- event_type : 运动固定 10
- start_time : EVENT_TIME 解析为 epoch
- duration : Webhook 无此字段,置 0
返回 dict 或 None字段不足无法构造
"""
if not isinstance(raw, dict):
return None
device_name = (raw.get('device_name') or raw.get('camera') or
raw.get('device') or raw.get('cam_name'))
event_time = raw.get('event_time') or raw.get('time')
thumbnail_url = (raw.get('thumbnail_url') or raw.get('thumbnail') or
raw.get('thumb'))
if not device_name or not event_time:
logger.warning(f"Webhook 事件缺字段 device_name/event_time: {raw}")
return None
start_ts = self.parse_ss_time_to_epoch(event_time)
if start_ts is None:
# 解析失败:退化为当前时间,宁可误报也不漏报
start_ts = int(datetime.now(timezone.utc).timestamp())
logger.warning(f"Webhook 事件时间解析失败,退化为当前时间: {event_time}")
camera_id = self.resolve_camera_id(device_name)
if camera_id is None and not self._camera_loaded:
# 懒加载一次摄像头映射后重试
self.refresh_camera_map()
camera_id = self.resolve_camera_id(device_name)
return {
"event_id": self._synth_event_id(device_name, event_time, thumbnail_url),
"camera_id": camera_id,
"event_type": 10,
"start_time": start_ts,
"duration": 0,
"thumbnail_url": thumbnail_url,
}
# ------------------------------------------------------------------ # ------------------------------------------------------------------
# 推送到甲骨文 # 推送到甲骨文
# ------------------------------------------------------------------ # ------------------------------------------------------------------
@@ -232,6 +369,10 @@ class MotionNotifier:
if not self.enabled: if not self.enabled:
logger.info("MotionNotifier 未启用motion_notifier.enabled=false") logger.info("MotionNotifier 未启用motion_notifier.enabled=false")
return return
if not self.poll_enabled:
logger.info("MotionNotifier 轮询已禁用poll_enabled=false仅作为 "
"Webhook 推送客户端 + 摄像头名映射使用")
return
if self._running: if self._running:
return return
self._running = True self._running = True
@@ -251,6 +392,9 @@ class MotionNotifier:
return { return {
"running": self.is_alive(), "running": self.is_alive(),
"enabled": self.enabled, "enabled": self.enabled,
"poll_enabled": self.poll_enabled,
"camera_map": {**self.camera_name_to_id, **self._ss_name_to_id},
"camera_loaded": self._camera_loaded,
"last_event_id": self._last_event_id, "last_event_id": self._last_event_id,
"last_poll_at": self._last_poll_at.isoformat() if self._last_poll_at else None, "last_poll_at": self._last_poll_at.isoformat() if self._last_poll_at else None,
"last_error": self._last_error, "last_error": self._last_error,