diff --git a/fam-core/config/config.yaml b/fam-core/config/config.yaml index fd7e708..2347964 100644 --- a/fam-core/config/config.yaml +++ b/fam-core/config/config.yaml @@ -33,26 +33,35 @@ chat_handler: qa_url: "http://129.146.203.203:5000/api/edge/chat/ask" timeout: 120 -# 运动监测通知服务(2026-08-22 重构为 Webhook 驱动,不轮询) -# SS 配置 Webhook 把运动事件 POST 到 NAS /api/ss/webhook,本服务映射字段后 -# 推送甲骨文 FAM-Edge(单向 NAS -> Oracle,甲骨文不再反向访问 NAS)。 -# SS Webhook 不提供 event_id / camera_id,由本服务合成 event_id、把摄像头名映射成 -# camera_id(启动一次性从 SS 拉取 + 以下静态映射兜底),并把 EVENT_TIME 解析成 epoch。 +# 运动监测通知服务(2026-08-22 定稿:轮询主路径) +# NAS 本机轮询 SS EventCenter.Event.List(真实 event_id/start_time/duration), +# 增量推送到甲骨文 FAM-Edge(单向 NAS -> Oracle,甲骨文不再反向访问 NAS)。 +# Webhook(/api/ss/webhook)保留为可选低延迟补充,SS 行动规则未配置则不触发。 # SS 凭据走环境变量(${DSM_ACCOUNT}/${DSM_PASSWORD}),由启动脚本 source 的 .env 提供。 motion_notifier: enabled: true - poll_enabled: false # 关闭轮询,改用 SS Webhook 推送 - dsm_host: "192.168.50.64" # Surveillance Station 所在地址(NAS 本机,用于一次性拉摄像头映射) + poll_enabled: true # 轮询主路径(默认开启) + dsm_host: "192.168.50.64" # Surveillance Station 所在地址(NAS 本机) dsm_port: 5000 dsm_account: "${DSM_ACCOUNT}" dsm_password: "${DSM_PASSWORD}" - # 摄像头名 -> camera_id 静态映射(优先);启动时还会从 SS 拉一份补全 + camera_ids: [2] # 轮询关注的摄像头(Generic_ONVIF-001) + # 摄像头名 -> camera_id 静态映射(Webhook 路径用;启动时还会从 SS 拉一份补全) camera_name_to_id: "Generic_ONVIF-001": 2 oracle_base_url: "http://129.146.203.203:5000" # 与 oracle_sync.base_url 一致 oracle_token: "${ORACLE_SYNC_TOKEN}" # 与 oracle_sync.token 一致 - timeout_sec: 10 # 单次 SS 请求超时(摄像头映射拉取用) - # 以下为可选轮询参数(poll_enabled=true 时才生效,当前默认关闭) - poll_interval_sec: 60 - poll_window_hours: 2 - batch_size: 100 + timeout_sec: 10 # 单次 SS 请求超时 + # 心跳:跟轮询 SS 无关,只是定期空 POST 一下甲骨文的 /api/ss/motion,证明 + # NAS->Oracle 这条推送链路本身还活着(enabled=true 就跑,不受 poll_enabled 影响)。 + # 甲骨文侧 dsm_motion_prefilter.max_heartbeat_age_sec(默认 900s)据此判断"无运动" + # 结论是否可信——这个心跳间隔要明显小于那个阈值,否则会被误判成链路已死。 + 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 diff --git a/fam-core/src/fam_core/motion_bp.py b/fam-core/src/fam_core/motion_bp.py index 0307410..690094a 100644 --- a/fam-core/src/fam_core/motion_bp.py +++ b/fam-core/src/fam_core/motion_bp.py @@ -1,9 +1,11 @@ """ -运动监测 Webhook 接收(fam-core) +运动监测 Webhook 接收(fam-core,可选补充,非主路径) -群晖 Surveillance Station 配置「行動規則 -> 事件=偵測到動作 -> 動作=Webhook」, -把运动事件实时 POST 到本端点,本端点即时映射字段并推送到甲骨文 FAM-Edge。 -数据源完全来自 SS Webhook(不轮询 SS),数据方向 NAS -> Oracle 单向。 +主数据源是 MotionNotifier 轮询 SS EventCenter.Event.List(真实 +event_id/start_time/duration)。本端点是可选的低延迟补充:若在群晖 +Surveillance Station 配置「行動規則 -> 事件=偵測到動作 -> 動作=Webhook」, +SS 会把运动事件实时 POST 到本端点,本端点映射字段后走同一个 +push_events_to_oracle 推送到甲骨文 FAM-Edge(数据方向 NAS -> Oracle 单向)。 SS Webhook 仅提供模板变量(无 event_id / camera_id 数字字段): %EVENT_TIME% -> event_time (本地时间字符串) diff --git a/fam-core/src/fam_core/motion_notifier/motion_notifier.py b/fam-core/src/fam_core/motion_notifier/motion_notifier.py index 137502d..468d74c 100644 --- a/fam-core/src/fam_core/motion_notifier/motion_notifier.py +++ b/fam-core/src/fam_core/motion_notifier/motion_notifier.py @@ -1,21 +1,27 @@ """ -MotionNotifier - NAS 端运动监测通知服务(Webhook 驱动,2026-08-22 重构) +MotionNotifier - NAS 端运动监测通知服务(轮询主路径,2026-08-22 定型) 职责: - 作为 SS Webhook 的接收侧客户端,把群晖 Surveillance Station 推送来的运动侦测 - 事件(SYNO.SurveillanceStation.EventCenter.Event 的字段,event_type=10 即运动) - 映射后主动 POST 推送到甲骨文 FAM-Edge 的 /api/ss/motion 接口。 + 在 NAS 本机**轮询**群晖 Surveillance Station 的运动侦测事件 + (SYNO.SurveillanceStation.EventCenter.Event method=List,参数名下划线风格 + camera_ids/start_time/end_time,event_type=10 即运动),增量推送到甲骨文 + FAM-Edge 的 /api/ss/motion 接口。轮询能拿到真实 event_id/start_time/duration。 数据方向(关键约束): NAS -> Oracle,单向。甲骨文不再反向访问 NAS。 - 原 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 本地预过滤使用。 -轮询已禁用(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 re @@ -59,8 +65,15 @@ class MotionNotifier: 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)) - # 是否启用轮询:用户要求不轮询,改由 SS Webhook 推送(默认关闭) - self.poll_enabled = bool(cfg.get('poll_enabled', False)) + # 是否启用轮询(默认开启):轮询 EventCenter.Event.List 为主数据源,拿到 + # 真实 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 路径用摄像头名映射 self.camera_ids = cfg.get('camera_ids', [2]) # 摄像头名 -> camera_id 静态映射(优先),启动时再从 SS 拉一份补全 @@ -77,6 +90,10 @@ class MotionNotifier: 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): @@ -118,7 +135,7 @@ class MotionNotifier: "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}, + "limit": 1000, "_sid": self._sid}, timeout=self.timeout) data = resp.json() except Exception as e: @@ -312,20 +329,66 @@ class MotionNotifier: 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: 取当前 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()) 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}") + 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()) @@ -339,10 +402,15 @@ class MotionNotifier: 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] - self.push_events_to_oracle(batch) - self._last_event_id = max(int(e.get('id', 0)) for e in new) + stored = self.push_events_to_oracle(batch) + 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)) def _run(self): @@ -371,9 +439,16 @@ class MotionNotifier: 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 推送客户端 + 摄像头名映射使用") + "Webhook 推送客户端 + 心跳 + 摄像头名映射使用") return if self._running: return @@ -389,9 +464,15 @@ class MotionNotifier: 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,