[重构] 运动监测恢复轮询主路径(简单稳定) - poll_enabled 默认开启;修复3隐患: 推送失败批次不前进游标(下轮重试不丢事件)、重启优先续用DB游标(停机期间事件补推)、fetch limit 提到1000;Webhook降级为可选补充
This commit is contained in:
@@ -33,26 +33,35 @@ 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 重构为 Webhook 驱动,不轮询)
|
# 运动监测通知服务(2026-08-22 定稿:轮询主路径)
|
||||||
# SS 配置 Webhook 把运动事件 POST 到 NAS /api/ss/webhook,本服务映射字段后
|
# NAS 本机轮询 SS EventCenter.Event.List(真实 event_id/start_time/duration),
|
||||||
# 推送甲骨文 FAM-Edge(单向 NAS -> Oracle,甲骨文不再反向访问 NAS)。
|
# 增量推送到甲骨文 FAM-Edge(单向 NAS -> Oracle,甲骨文不再反向访问 NAS)。
|
||||||
# SS Webhook 不提供 event_id / camera_id,由本服务合成 event_id、把摄像头名映射成
|
# Webhook(/api/ss/webhook)保留为可选低延迟补充,SS 行动规则未配置则不触发。
|
||||||
# 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
|
||||||
poll_enabled: false # 关闭轮询,改用 SS Webhook 推送
|
poll_enabled: true # 轮询主路径(默认开启)
|
||||||
dsm_host: "192.168.50.64" # Surveillance Station 所在地址(NAS 本机,用于一次性拉摄像头映射)
|
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_id 静态映射(优先);启动时还会从 SS 拉一份补全
|
camera_ids: [2] # 轮询关注的摄像头(Generic_ONVIF-001)
|
||||||
|
# 摄像头名 -> camera_id 静态映射(Webhook 路径用;启动时还会从 SS 拉一份补全)
|
||||||
camera_name_to_id:
|
camera_name_to_id:
|
||||||
"Generic_ONVIF-001": 2
|
"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 一致
|
||||||
timeout_sec: 10 # 单次 SS 请求超时(摄像头映射拉取用)
|
timeout_sec: 10 # 单次 SS 请求超时
|
||||||
# 以下为可选轮询参数(poll_enabled=true 时才生效,当前默认关闭)
|
# 心跳:跟轮询 SS 无关,只是定期空 POST 一下甲骨文的 /api/ss/motion,证明
|
||||||
poll_interval_sec: 60
|
# NAS->Oracle 这条推送链路本身还活着(enabled=true 就跑,不受 poll_enabled 影响)。
|
||||||
poll_window_hours: 2
|
# 甲骨文侧 dsm_motion_prefilter.max_heartbeat_age_sec(默认 900s)据此判断"无运动"
|
||||||
batch_size: 100
|
# 结论是否可信——这个心跳间隔要明显小于那个阈值,否则会被误判成链路已死。
|
||||||
|
heartbeat_interval_sec: 300
|
||||||
|
# 轮询参数
|
||||||
|
poll_interval_sec: 60 # 轮询间隔
|
||||||
|
poll_window_hours: 2 # 每轮回看窗口(小时),覆盖轮询间隔内的新事件
|
||||||
|
batch_size: 100 # 单批推送上限
|
||||||
|
# Webhook 可选补全(仅 webhook 路径用):收到 webhook 后,以 camera_id+触发时刻
|
||||||
|
# 为中心开此时间窗(秒)回头查 SS 事件列表,匹配最接近的运动事件回填真实
|
||||||
|
# start_time/duration/event_id(Webhook 本身不提供这些字段)。
|
||||||
|
enrich_window_sec: 120
|
||||||
|
|||||||
@@ -1,9 +1,11 @@
|
|||||||
"""
|
"""
|
||||||
运动监测 Webhook 接收(fam-core)
|
运动监测 Webhook 接收(fam-core,可选补充,非主路径)
|
||||||
|
|
||||||
群晖 Surveillance Station 配置「行動規則 -> 事件=偵測到動作 -> 動作=Webhook」,
|
主数据源是 MotionNotifier 轮询 SS EventCenter.Event.List(真实
|
||||||
把运动事件实时 POST 到本端点,本端点即时映射字段并推送到甲骨文 FAM-Edge。
|
event_id/start_time/duration)。本端点是可选的低延迟补充:若在群晖
|
||||||
数据源完全来自 SS Webhook(不轮询 SS),数据方向 NAS -> Oracle 单向。
|
Surveillance Station 配置「行動規則 -> 事件=偵測到動作 -> 動作=Webhook」,
|
||||||
|
SS 会把运动事件实时 POST 到本端点,本端点映射字段后走同一个
|
||||||
|
push_events_to_oracle 推送到甲骨文 FAM-Edge(数据方向 NAS -> Oracle 单向)。
|
||||||
|
|
||||||
SS Webhook 仅提供模板变量(无 event_id / camera_id 数字字段):
|
SS Webhook 仅提供模板变量(无 event_id / camera_id 数字字段):
|
||||||
%EVENT_TIME% -> event_time (本地时间字符串)
|
%EVENT_TIME% -> event_time (本地时间字符串)
|
||||||
|
|||||||
@@ -1,21 +1,27 @@
|
|||||||
"""
|
"""
|
||||||
MotionNotifier - NAS 端运动监测通知服务(Webhook 驱动,2026-08-22 重构)
|
MotionNotifier - NAS 端运动监测通知服务(轮询主路径,2026-08-22 定型)
|
||||||
|
|
||||||
职责:
|
职责:
|
||||||
作为 SS Webhook 的接收侧客户端,把群晖 Surveillance Station 推送来的运动侦测
|
在 NAS 本机**轮询**群晖 Surveillance Station 的运动侦测事件
|
||||||
事件(SYNO.SurveillanceStation.EventCenter.Event 的字段,event_type=10 即运动)
|
(SYNO.SurveillanceStation.EventCenter.Event method=List,参数名下划线风格
|
||||||
映射后主动 POST 推送到甲骨文 FAM-Edge 的 /api/ss/motion 接口。
|
camera_ids/start_time/end_time,event_type=10 即运动),增量推送到甲骨文
|
||||||
|
FAM-Edge 的 /api/ss/motion 接口。轮询能拿到真实 event_id/start_time/duration。
|
||||||
|
|
||||||
数据方向(关键约束): NAS -> Oracle,单向。甲骨文不再反向访问 NAS。
|
数据方向(关键约束): NAS -> Oracle,单向。甲骨文不再反向访问 NAS。
|
||||||
- 原 fam-edge 的 dsm_motion_client(甲骨文主动查 SS API)已停用;
|
- 原 fam-edge 的 dsm_motion_client(甲骨文主动查 SS API)已停用;
|
||||||
- SS Webhook 不提供 event_id / camera_id 数字字段,本服务负责:
|
|
||||||
* 合成稳定 event_id(device_name+event_time+thumbnail_url 哈希,幂等去重)
|
|
||||||
* 把 %DEVICE_NAME% 经摄像头名->id 映射解析为 camera_id
|
|
||||||
* 把 %EVENT_TIME% 解析为 Unix epoch
|
|
||||||
- 推送后由甲骨文本地落库 ss_motion_events,供 video_processor 本地预过滤使用。
|
- 推送后由甲骨文本地落库 ss_motion_events,供 video_processor 本地预过滤使用。
|
||||||
|
|
||||||
轮询已禁用(poll_enabled=false,用户要求不轮询),仅保留轮询代码路径作为可选能力。
|
可靠性设计:
|
||||||
摄像头名->id 映射在启动时一次性从 SS 拉取补全(非轮询),并以配置 camera_name_to_id 兜底。
|
- 游标(MariaDB sync_cursor.motion_last_event_id)持久化;重启优先续用 DB 游标,
|
||||||
|
停机期间的事件由窗口回看补推;仅首次部署时才初始化为 SS 当前最大 id。
|
||||||
|
- 推送失败的批次不推进游标,下一轮窗口回看重试(不丢事件)。
|
||||||
|
- 心跳线程(enabled 即跑,与轮询无关):定期空 POST /api/ss/motion,证明
|
||||||
|
NAS->Oracle 推送链路存活,供甲骨文侧判断"无运动"结论是否可信。
|
||||||
|
|
||||||
|
Webhook(可选补充,非主路径):
|
||||||
|
fam-core 仍暴露 POST /api/ss/webhook,SS 行动规则若配置可实时低延迟推送,
|
||||||
|
经 build_event_from_webhook 映射后与轮询共用 push_events_to_oracle;
|
||||||
|
其合成 event_id 与轮询真实 event_id 不冲突(不同命名空间,幂等互不干扰)。
|
||||||
"""
|
"""
|
||||||
import os
|
import os
|
||||||
import re
|
import re
|
||||||
@@ -59,8 +65,15 @@ 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 推送(默认关闭)
|
# 是否启用轮询(默认开启):轮询 EventCenter.Event.List 为主数据源,拿到
|
||||||
self.poll_enabled = bool(cfg.get('poll_enabled', False))
|
# 真实 event_id/start_time/duration;Webhook 仅作为可选的低延迟补充
|
||||||
|
# (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 时)关注的摄像头;Webhook 路径用摄像头名映射
|
# 轮询路径(poll_enabled=true 时)关注的摄像头;Webhook 路径用摄像头名映射
|
||||||
self.camera_ids = cfg.get('camera_ids', [2])
|
self.camera_ids = cfg.get('camera_ids', [2])
|
||||||
# 摄像头名 -> camera_id 静态映射(优先),启动时再从 SS 拉一份补全
|
# 摄像头名 -> camera_id 静态映射(优先),启动时再从 SS 拉一份补全
|
||||||
@@ -77,6 +90,10 @@ class MotionNotifier:
|
|||||||
self._last_poll_at = None
|
self._last_poll_at = None
|
||||||
self._last_error = None
|
self._last_error = None
|
||||||
self._pushed_total = 0
|
self._pushed_total = 0
|
||||||
|
self._hb_running = False
|
||||||
|
self._hb_thread = None
|
||||||
|
self._last_heartbeat_at = None
|
||||||
|
self._last_heartbeat_error = None
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _resolve(v):
|
def _resolve(v):
|
||||||
@@ -118,7 +135,7 @@ class MotionNotifier:
|
|||||||
"version": 1, "method": "List",
|
"version": 1, "method": "List",
|
||||||
"camera_ids": ",".join(str(c) for c in self.camera_ids),
|
"camera_ids": ",".join(str(c) for c in self.camera_ids),
|
||||||
"start_time": start_ts, "end_time": end_ts,
|
"start_time": start_ts, "end_time": end_ts,
|
||||||
"limit": 100, "_sid": self._sid},
|
"limit": 1000, "_sid": self._sid},
|
||||||
timeout=self.timeout)
|
timeout=self.timeout)
|
||||||
data = resp.json()
|
data = resp.json()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -312,20 +329,66 @@ class MotionNotifier:
|
|||||||
logger.info(f"运动事件推送成功: {len(norm)} 条 -> Oracle 存储 {stored}")
|
logger.info(f"运动事件推送成功: {len(norm)} 条 -> Oracle 存储 {stored}")
|
||||||
return 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):
|
def _init_cursor(self):
|
||||||
"""启动时初始化 last_event_id: 取当前 SS 最大事件 id,只推增量(不回灌历史)。"""
|
"""启动时初始化 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())
|
now = int(datetime.now(timezone.utc).timestamp())
|
||||||
events = self._fetch_events(now - 3600, now) # 最近 1h 用于定位最大 id
|
events = self._fetch_events(now - 3600, now) # 最近 1h 用于定位最大 id
|
||||||
if events:
|
if events:
|
||||||
self._last_event_id = max(
|
self._last_event_id = max(
|
||||||
int(e.get('id', 0)) for e in events if e.get('id'))
|
int(e.get('id', 0)) for e in events if e.get('id'))
|
||||||
else:
|
else:
|
||||||
saved = db_layer.get_motion_cursor()
|
self._last_event_id = 0
|
||||||
self._last_event_id = int(saved) if saved else 0
|
logger.info(f"运动通知游标初始化(SS 当前最大): last_event_id={self._last_event_id}")
|
||||||
logger.info(f"运动通知游标初始化: last_event_id={self._last_event_id}")
|
|
||||||
|
|
||||||
def _poll_once(self):
|
def _poll_once(self):
|
||||||
now = int(datetime.now(timezone.utc).timestamp())
|
now = int(datetime.now(timezone.utc).timestamp())
|
||||||
@@ -339,10 +402,15 @@ class MotionNotifier:
|
|||||||
if not new:
|
if not new:
|
||||||
return
|
return
|
||||||
new.sort(key=lambda e: int(e.get('id', 0)))
|
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):
|
for i in range(0, len(new), self.batch_size):
|
||||||
batch = new[i:i + self.batch_size]
|
batch = new[i:i + self.batch_size]
|
||||||
self.push_events_to_oracle(batch)
|
stored = self.push_events_to_oracle(batch)
|
||||||
self._last_event_id = max(int(e.get('id', 0)) for e in new)
|
if stored > 0:
|
||||||
|
# 只有推送成功的批次才推进游标;失败批次保持原地,
|
||||||
|
# 下一轮窗口回看会重新捞到并重试(不丢事件)
|
||||||
|
cursor = max(int(e.get('id', 0)) for e in batch)
|
||||||
|
self._last_event_id = cursor
|
||||||
db_layer.set_motion_cursor(str(self._last_event_id))
|
db_layer.set_motion_cursor(str(self._last_event_id))
|
||||||
|
|
||||||
def _run(self):
|
def _run(self):
|
||||||
@@ -371,9 +439,16 @@ 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
|
||||||
|
# 心跳线程跟轮询是否开启无关:只要 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:
|
if not self.poll_enabled:
|
||||||
logger.info("MotionNotifier 轮询已禁用(poll_enabled=false),仅作为 "
|
logger.info("MotionNotifier 轮询已禁用(poll_enabled=false),仅作为 "
|
||||||
"Webhook 推送客户端 + 摄像头名映射使用")
|
"Webhook 推送客户端 + 心跳 + 摄像头名映射使用")
|
||||||
return
|
return
|
||||||
if self._running:
|
if self._running:
|
||||||
return
|
return
|
||||||
@@ -389,9 +464,15 @@ class MotionNotifier:
|
|||||||
self._running = False
|
self._running = False
|
||||||
if self._thread:
|
if self._thread:
|
||||||
self._thread.join(timeout=5)
|
self._thread.join(timeout=5)
|
||||||
|
self._hb_running = False
|
||||||
|
if self._hb_thread:
|
||||||
|
self._hb_thread.join(timeout=5)
|
||||||
|
|
||||||
def status(self) -> dict:
|
def status(self) -> dict:
|
||||||
return {
|
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(),
|
"running": self.is_alive(),
|
||||||
"enabled": self.enabled,
|
"enabled": self.enabled,
|
||||||
"poll_enabled": self.poll_enabled,
|
"poll_enabled": self.poll_enabled,
|
||||||
|
|||||||
Reference in New Issue
Block a user