diff --git a/fam-edge/config/config.yaml b/fam-edge/config/config.yaml index 10cfbed..f64deed 100644 --- a/fam-edge/config/config.yaml +++ b/fam-edge/config/config.yaml @@ -47,6 +47,23 @@ video_processing: # 降级顺序:先 gemini 整视频,失败再 nvidia 整视频;两者都失败 -> 标记 failed vision_order: ["gemini", "nvidia"] +# DSM 运动侦测预过滤(2026-08-22 新增):分析前先查一下群晖 Surveillance Station +# 自己记录的运动侦测事件(SYNO.SurveillanceStation.EventCenter.Event,未公开文档的 +# 内部接口,参数名是下划线风格 camera_ids/start_time/end_time),这段时间窗口一条 +# 运动事件都没有就跳过云端分析(标记 done,compute_provider=skipped_no_motion), +# 省掉长期无人时段白白消耗的 Gemini/NVIDIA 配额。 +# 账号密码走 .env(DSM_ACCOUNT/DSM_PASSWORD),不明文入库;查询失败/未配置一律 +# fail-open(照常送云端分析),不会因为这层可选优化漏检真实事件。 +dsm_motion_prefilter: + enabled: true + 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=有事件就算) + timeout_sec: 10 + # 智能问答降级链(与视频分析独立):Gemini -> NVIDIA -> 本地 Ollama models: - provider: "gemini" diff --git a/fam-edge/src/fam_edge/dsm_motion_client.py b/fam-edge/src/fam_edge/dsm_motion_client.py new file mode 100644 index 0000000..9de698c --- /dev/null +++ b/fam-edge/src/fam_edge/dsm_motion_client.py @@ -0,0 +1,125 @@ +""" +DsmMotionClient - 查询群晖 Surveillance Station 的运动侦测事件(预过滤用) + +背景: fam-edge 原来不管这段 30 分钟录像里有没有人走动,一律整段送云端 VLM 分析, +配额/耗时都花在长期无人的空转时段上。DSM 的 Surveillance Station 用 +SYNO.SurveillanceStation.EventCenter.Event 这个未公开文档的内部 API 记录了摄像 +头真实的运动侦测窗口(start_time + duration,event_type=10 即运动),比自己在本地 +用 ffmpeg 帧差分更准,也不需要额外算力。 + +用法: process_video() 分析前,用视频的 [event_start, event_start+duration] 时间窗 +查一次这个接口——窗口内一条运动事件都没有,就跳过云端分析(标记 done, +compute_provider='skipped_no_motion'),有任何一条就正常送去分析。 + +失败即放行(fail-open): 网络错误/认证失败/账号未授权(401)等任何异常都视为 +"无法判断",返回 None,调用方必须当作"照常分析"处理,而不是当作"跳过"处理—— +宁可多花一次配额,也不能因为这个可选的省钱层漏检真实事件。 + +未公开文档的内部 API,字段名/权限模型可能随群晖套件升级变化,出问题时表现为 +每次都 fail-open(不再省钱),不会导致漏检。 +""" +import time +from datetime import datetime +from typing import Optional + +import requests + +from .logger import setup_logger + +logger = setup_logger('fam-edge.dsm_motion_client') + + +class DsmMotionClient: + def __init__(self, config: dict): + self.enabled = bool(config.get('enabled', False)) + self.host = config.get('host', '') + self.port = config.get('port', 5000) + self.account = self._resolve(config.get('account', '')) + self.password = self._resolve(config.get('password', '')) + self.camera_id = config.get('camera_id') + self.min_motion_seconds = float(config.get('min_motion_seconds', 0)) + self.timeout = config.get('timeout_sec', 10) + self._sid: Optional[str] = None + self._base = f"http://{self.host}:{self.port}/webapi" + + @staticmethod + def _resolve(raw: str) -> str: + import os + if raw.startswith('${') and raw.endswith('}'): + return os.environ.get(raw[2:-1], '') + return raw + + def _login(self) -> Optional[str]: + try: + resp = requests.get( + f"{self._base}/auth.cgi", + params={ + "api": "SYNO.API.Auth", "version": 6, "method": "login", + "account": self.account, "passwd": self.password, + "session": "SurveillanceStation", "format": "sid", + }, + timeout=self.timeout, + ) + data = resp.json() + if data.get('success'): + return data['data']['sid'] + logger.warning(f"DSM 登录失败: {data.get('error')}") + except Exception as e: + logger.warning(f"DSM 登录异常: {e}") + return None + + def has_motion_in_range(self, start_dt: datetime, duration_sec: float) -> Optional[bool]: + """查询 [start_dt, start_dt+duration_sec] 窗口内是否有运动事件。 + + 返回 True/False 为确定结果;返回 None 表示查询本身失败,调用方必须 + fail-open(当作"有运动"处理,照常送云端分析),不能当作"无运动"跳过。 + """ + if not self.enabled or not self.host or not self.account or not self.password \ + or self.camera_id is None: + return None + if self._sid is None: + self._sid = self._login() + if self._sid is None: + return None + start_ts = int(start_dt.timestamp()) + end_ts = int(start_dt.timestamp() + max(duration_sec, 0)) + for attempt in range(2): # 第 2 次仅用于 session 过期后重新登录重试一次 + try: + resp = requests.get( + f"{self._base}/entry.cgi", + params={ + "api": "SYNO.SurveillanceStation.EventCenter.Event", + "version": 1, "method": "List", + "camera_ids": str(self.camera_id), + "start_time": start_ts, "end_time": end_ts, + "limit": 50, "_sid": self._sid, + }, + timeout=self.timeout, + ) + data = resp.json() + except Exception as e: + logger.warning(f"DSM 运动事件查询异常: {e}") + return None + if not data.get('success'): + err_code = (data.get('error') or {}).get('code') + if err_code in (106, 107, 119) and attempt == 0: + # session 过期/被顶掉,重新登录重试一次 + self._sid = self._login() + if self._sid is None: + return None + continue + logger.warning(f"DSM 运动事件查询失败: {data.get('error')}") + return None + # data 按 ds_id(群晖多主机 CMS 场景的服务器 id,单机固定是 "0")分组, + # 不是按 camera_id 分组——这里把所有 ds_id 下的事件拉平, + # 用事件自带的 camera_id 字段再确认一次(双保险,服务端已按 camera_ids 过滤)。 + events = [ + e for group in (data.get('data') or {}).values() + for e in (group or []) + if e.get('camera_id') == self.camera_id + ] + if not events: + return False + total_motion = sum(float(e.get('duration', 0) or 0) for e in events) + return total_motion >= self.min_motion_seconds + return None diff --git a/fam-edge/src/fam_edge/video_processor.py b/fam-edge/src/fam_edge/video_processor.py index e1a862d..30855f2 100644 --- a/fam-edge/src/fam_edge/video_processor.py +++ b/fam-edge/src/fam_edge/video_processor.py @@ -22,6 +22,7 @@ from .logger import setup_logger from .config_loader import load_config from .model_adapters.adapter_factory import build_adapters from .model_adapters.base_adapter import BaseModelAdapter +from .dsm_motion_client import DsmMotionClient from . import oracle_db logger = setup_logger('fam-edge.video_processor') @@ -193,6 +194,7 @@ class VideoProcessor: adapters = build_adapters(self.config.get('models', [])) self.vision_adapters: Dict[str, BaseModelAdapter] = { 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]: ordered = [] @@ -236,6 +238,22 @@ class VideoProcessor: if event_start: self.db.set_event_start_time(video_id, event_start) + # DSM 运动侦测预过滤:这段时间窗口里群晖自己记录的运动事件一条都没有, + # 就跳过云端分析(省配额)。查询本身失败/未配置一律 fail-open(照常分析), + # 绝不能因为这层可选优化漏检真实事件。 + duration_sec = (vmeta or {}).get('duration_sec') if self.file_validate else None + if event_start and duration_sec: + try: + start_dt = datetime.strptime(event_start, '%Y-%m-%d %H:%M:%S') + has_motion = self.dsm_motion.has_motion_in_range(start_dt, duration_sec) + except ValueError: + has_motion = None + if has_motion is False: + logger.info(f"[video_id={video_id}] DSM 运动预过滤:该时段无运动,跳过云端分析") + self.db.mark_video_processed( + video_id, "(自动跳过:该时段未检测到运动)", [], [], 'skipped_no_motion') + return True + known = self.db.get_known_members_context() logger.info(f"[video_id={video_id}] 开始整视频分析: {filename} " f"(event_start={event_start}, known_members={'有' if known else '无'})") diff --git a/fam-edge/tests/test_dsm_motion_client.py b/fam-edge/tests/test_dsm_motion_client.py new file mode 100644 index 0000000..f6f8ac2 --- /dev/null +++ b/fam-edge/tests/test_dsm_motion_client.py @@ -0,0 +1,157 @@ +from datetime import datetime + +from fam_edge.dsm_motion_client import DsmMotionClient + + +def _cfg(**overrides): + base = { + "enabled": True, + "host": "192.168.50.64", + "port": 5000, + "account": "ericwyuan", + "password": "secret", + "camera_id": 2, + "min_motion_seconds": 0, + "timeout_sec": 5, + } + base.update(overrides) + return base + + +class _FakeResp: + def __init__(self, payload): + self._payload = payload + + def json(self): + return self._payload + + +def _events_payload(events, ds_id="0"): + return {"success": True, "data": {ds_id: events}} + + +def test_disabled_config_always_returns_none(): + c = DsmMotionClient(_cfg(enabled=False)) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is None + + +def test_missing_credentials_returns_none(): + c = DsmMotionClient(_cfg(account="", password="")) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is None + + +def test_login_failure_returns_none(monkeypatch): + def fake_get(url, params=None, timeout=None): + return _FakeResp({"success": False, "error": {"code": 400}}) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg()) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is None + + +def test_login_network_exception_returns_none(monkeypatch): + def fake_get(url, params=None, timeout=None): + raise ConnectionError("boom") + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg()) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is None + + +def test_no_motion_events_returns_false(monkeypatch): + calls = {"n": 0} + + def fake_get(url, params=None, timeout=None): + calls["n"] += 1 + if params.get("api") == "SYNO.API.Auth": + return _FakeResp({"success": True, "data": {"sid": "abc"}}) + return _FakeResp(_events_payload([])) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg()) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is False + + +def test_motion_event_present_returns_true(monkeypatch): + def fake_get(url, params=None, timeout=None): + if params.get("api") == "SYNO.API.Auth": + return _FakeResp({"success": True, "data": {"sid": "abc"}}) + return _FakeResp(_events_payload( + [{"camera_id": 2, "event_type": 10, "start_time": 1000, "duration": 12}])) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg()) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is True + + +def test_events_grouped_by_ds_id_not_camera_id(monkeypatch): + """核心场景: 响应按 ds_id 分组(单机固定 "0"),不是按 camera_id 分组—— + 错把 camera_id 当 dict key 去取会永远拿到空列表。""" + def fake_get(url, params=None, timeout=None): + if params.get("api") == "SYNO.API.Auth": + return _FakeResp({"success": True, "data": {"sid": "abc"}}) + return _FakeResp(_events_payload( + [{"camera_id": 2, "event_type": 10, "start_time": 1000, "duration": 5}], ds_id="0")) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg(camera_id=2)) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is True + + +def test_events_for_other_camera_ignored(monkeypatch): + def fake_get(url, params=None, timeout=None): + if params.get("api") == "SYNO.API.Auth": + return _FakeResp({"success": True, "data": {"sid": "abc"}}) + return _FakeResp(_events_payload( + [{"camera_id": 99, "event_type": 10, "start_time": 1000, "duration": 999}])) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg(camera_id=2)) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is False + + +def test_total_motion_below_threshold_is_treated_as_no_motion(monkeypatch): + def fake_get(url, params=None, timeout=None): + if params.get("api") == "SYNO.API.Auth": + return _FakeResp({"success": True, "data": {"sid": "abc"}}) + return _FakeResp(_events_payload( + [{"camera_id": 2, "event_type": 10, "start_time": 1000, "duration": 1}])) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg(min_motion_seconds=3)) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is False + + +def test_total_motion_meeting_threshold_counts_as_motion(monkeypatch): + def fake_get(url, params=None, timeout=None): + if params.get("api") == "SYNO.API.Auth": + return _FakeResp({"success": True, "data": {"sid": "abc"}}) + return _FakeResp(_events_payload( + [{"camera_id": 2, "event_type": 10, "start_time": 1000, "duration": 2}, + {"camera_id": 2, "event_type": 10, "start_time": 1100, "duration": 2}])) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg(min_motion_seconds=3)) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is True + + +def test_expired_session_relogins_and_retries_once(monkeypatch): + """session 过期(错误码 106)时应该重新登录一次再查,而不是直接放弃。""" + state = {"logins": 0, "queries": 0} + + def fake_get(url, params=None, timeout=None): + if params.get("api") == "SYNO.API.Auth": + state["logins"] += 1 + return _FakeResp({"success": True, "data": {"sid": f"sid{state['logins']}"}}) + state["queries"] += 1 + if state["queries"] == 1: + return _FakeResp({"success": False, "error": {"code": 106}}) + return _FakeResp(_events_payload( + [{"camera_id": 2, "event_type": 10, "start_time": 1000, "duration": 5}])) + monkeypatch.setattr("fam_edge.dsm_motion_client.requests.get", fake_get) + c = DsmMotionClient(_cfg()) + assert c.has_motion_in_range(datetime(2026, 8, 22), 1800) is True + assert state["logins"] == 2 + assert state["queries"] == 2 + + +def test_env_var_credentials_resolved(): + import os + os.environ["TEST_DSM_ACCOUNT_XYZ"] = "realaccount" + try: + c = DsmMotionClient(_cfg(account="${TEST_DSM_ACCOUNT_XYZ}")) + assert c.account == "realaccount" + finally: + del os.environ["TEST_DSM_ACCOUNT_XYZ"]