diff --git a/fam-edge/config/config.yaml b/fam-edge/config/config.yaml index e839d04..363f604 100644 --- a/fam-edge/config/config.yaml +++ b/fam-edge/config/config.yaml @@ -71,6 +71,19 @@ motion_segment: min_duration_sec: 1 # 短于该时长的事件不分割 unfinished_grace_sec: 10 # start_time+duration 距当前 ≤ 该秒视为"已结束"容差 +# 磁盘空间守护(2026-08-28 新增):Oracle 磁盘曾经被写满,触发 rclone 遇到 +# IO 错误就拒绝删除的保护机制,形成"越满越删不掉"的死循环,导致新视频下载 +# 和运动片段分割全部失败。这里加一道独立于 rclone 同步之外的兜底:不管上游 +# 同步是否正常,剩余空间跌破 min_free_gb 就主动清理最旧的、已完成分割阶段 +# 的整段素材(gdrive_videos 里的原始录像,不碰 motion_clips 里的运动片段)。 +disk_guard: + enabled: true + check_interval_sec: 300 # 5 分钟检查一次 + min_free_gb: 10 # 剩余空间低于此值触发清理 + target_free_gb: 15 # 清理到这个水位就停(留缓冲,避免刚清完又立刻再触发) + watch_path: "/opt/fam-edge" # 检查这个路径所在磁盘分区的剩余空间 + max_delete_per_round: 50 # 单轮最多清理几个文件,防止候选异常多时一次删太多 + # 闭集人物识别(2026-08-22 新增,家里固定 4 人:爷爷/爸爸/媳妇/汤圆): # 原来靠大模型自己编的"人物A/B/C"临时 uid + 文字特征描述跨视频合并,验证下来 # 不可靠(用户原话"现在的识别全是错的")。人脸向量方案也验证过,家庭监控这种 diff --git a/fam-edge/src/fam_edge/api_gateway/api_gateway.py b/fam-edge/src/fam_edge/api_gateway/api_gateway.py index 9ea378f..cb07491 100644 --- a/fam-edge/src/fam_edge/api_gateway/api_gateway.py +++ b/fam-edge/src/fam_edge/api_gateway/api_gateway.py @@ -261,10 +261,20 @@ def activity(): segment['clips_files'] = 0 # 事件↔片段一致性对账(数量可验证) segment['consistency'] = db.get_segment_consistency() + # 磁盘空间实时用量 + DiskGuard 最近一次清理动作(2026-08-28 新增) + disk_info = {"free_gb": None, "last_activity": _last_activity('disk_guard')} + try: + import shutil as _shutil + from ..config_loader import load_config as _load_config + watch_path = _load_config().get('disk_guard', {}).get('watch_path', '/opt/fam-edge') + disk_info["free_gb"] = round(_shutil.disk_usage(watch_path).free / (1024 ** 3), 1) + except Exception as e: + logger.warning(f"磁盘空间查询异常: {e}") return jsonify({ "queue": q_status, "db": db.get_queue_status(), "segment": segment, + "disk": disk_info, "rclone": _last_activity('rclone'), "person": _last_activity('person'), "model_calls": [dict(m) for m in model_calls], diff --git a/fam-edge/src/fam_edge/app.py b/fam-edge/src/fam_edge/app.py index 3d4136f..03129e1 100644 --- a/fam-edge/src/fam_edge/app.py +++ b/fam-edge/src/fam_edge/app.py @@ -5,6 +5,7 @@ FAM-Edge 主应用 - Flask 单进程(新架构 v2.1) - API-Gateway(同步拉取 / 命名校正 / 智能问答) - VideoQueue(生产-消费队列:同步落地新视频入队,独立消费者云端分析,模型超时 ×2) - PersonService(人物汇总合并,定时) + - DiskGuard(磁盘空间守护,定时检查,剩余空间过低时清理最旧的已完成素材) """ import os import sys @@ -18,6 +19,7 @@ from .api_gateway.api_gateway import api_bp from . import state from .video_queue import VideoQueue from .person_service import PersonService +from .disk_guard import DiskGuard logger = setup_logger('fam-edge.app') @@ -31,9 +33,10 @@ def root(): "mode": "drive-sync + whole-video analysis (producer-consumer queue)"}), 200 -# 启动生产-消费队列 + 人物服务 +# 启动生产-消费队列 + 人物服务 + 磁盘空间守护 _queue = None _person = None +_disk_guard = None try: db = state.get_db() _queue = VideoQueue(db) @@ -44,6 +47,9 @@ try: _person = PersonService(db) _person.start() logger.info("PersonService 已启动") + + _disk_guard = DiskGuard(db) + _disk_guard.start() except Exception as e: logger.error(f"后台服务启动失败: {e}", exc_info=True) diff --git a/fam-edge/src/fam_edge/disk_guard.py b/fam-edge/src/fam_edge/disk_guard.py new file mode 100644 index 0000000..914de30 --- /dev/null +++ b/fam-edge/src/fam_edge/disk_guard.py @@ -0,0 +1,114 @@ +""" +DiskGuard - 磁盘空间守护(2026-08-28 新增) + +背景:Oracle 机器磁盘曾经被写满(gdrive_videos 持续下载新素材、旧文件迟迟 +没被清理),触发了 rclone 的一个安全机制——同步过程中一旦遇到 IO 错误(写 +不进去)就整体拒绝执行删除,导致"越满越删不掉,越删不掉越满"的死循环, +最终连新视频都下载不了。这个模块是一道独立于 rclone 同步逻辑之外的兜底: +不管上游同步是否正常,只要本机磁盘剩余空间跌破阈值,就主动清理最旧的、已 +经处理完成的整段素材文件,把空间抢回来。 + +清理对象:仅限"整段素材"(gdrive_videos 里落地的原始录像文件,filename +不是 motion_ 前缀)且 status='done'(已完成分割阶段,不会再被 Video-Queue +重新捡起)——不碰运动片段(motion_clips 里的文件,是独立的分析产物,事件 +时间轴展示、人物头像裁剪都依赖它,删了会导致图片丢失)、不碰还在 +pending/processing 中的素材(避免删掉还没来得及处理的数据)。 +""" +import shutil +import threading +import time +import os + +from .logger import setup_logger +from .config_loader import load_config + +logger = setup_logger('fam-edge.disk_guard') + + +class DiskGuard: + def __init__(self, db): + self.db = db + cfg = load_config().get('disk_guard', {}) + self.enabled = bool(cfg.get('enabled', True)) + self.check_interval_sec = int(cfg.get('check_interval_sec', 300)) + self.min_free_gb = float(cfg.get('min_free_gb', 10)) + # 清理到这个水位就停:留出缓冲,避免刚清理完又立刻因为新文件写入 + # 跌破阈值、频繁触发清理循环 + self.target_free_gb = float(cfg.get('target_free_gb', 15)) + self.watch_path = cfg.get('watch_path', '/opt/fam-edge') + # 单轮最多清理几个文件:防止候选异常多时一次性删太多,先清一部分 + # 观察效果,下一轮检查再继续(check_interval_sec 后很快就会再跑一次) + self.max_delete_per_round = int(cfg.get('max_delete_per_round', 50)) + self._running = False + self._thread = None + + def start(self): + if not self.enabled: + logger.info("DiskGuard 未启用") + return + self._running = True + self._thread = threading.Thread(target=self._run, daemon=True, name='disk-guard') + self._thread.start() + logger.info( + f"DiskGuard 已启动(阈值 {self.min_free_gb}GB," + f"目标水位 {self.target_free_gb}GB,检查间隔 {self.check_interval_sec}s)") + + def stop(self): + self._running = False + if self._thread: + self._thread.join(timeout=5) + + def is_alive(self) -> bool: + return self._thread is not None and self._thread.is_alive() + + def _free_gb(self) -> float: + return shutil.disk_usage(self.watch_path).free / (1024 ** 3) + + def _run(self): + while self._running: + try: + self.check_once() + except Exception as e: + logger.error(f"DiskGuard 检查异常: {e}", exc_info=True) + for _ in range(self.check_interval_sec): + if not self._running: + return + time.sleep(1) + + def check_once(self): + """检查一次磁盘空间,不足则清理最旧的已完成素材直到恢复到目标水位。 + + 供后台循环调用,也可单独调用做一次性检查(比如手动触发/测试)。 + """ + free_gb = self._free_gb() + if free_gb >= self.min_free_gb: + return + logger.warning(f"磁盘剩余 {free_gb:.1f}GB 低于阈值 {self.min_free_gb}GB,开始清理旧素材") + self.db.record_activity( + 'disk_guard', 'low_space', + f"剩余 {free_gb:.1f}GB 低于阈值 {self.min_free_gb}GB,开始清理") + + deleted = 0 + freed_bytes = 0 + while deleted < self.max_delete_per_round and self._free_gb() < self.target_free_gb: + candidate = self.db.get_oldest_purgeable_material() + if not candidate: + logger.warning("磁盘空间仍然紧张,但已经没有可清理的素材了") + self.db.record_activity( + 'disk_guard', 'no_candidate', + f"剩余 {self._free_gb():.1f}GB 仍不足,且没有可清理的素材") + break + video_id, local_path = candidate['id'], candidate.get('local_path') + size = os.path.getsize(local_path) if local_path and os.path.isfile(local_path) else 0 + self.db.delete_video(video_id) + deleted += 1 + freed_bytes += size + logger.info(f"DiskGuard 清理素材 video_id={video_id}(约 {size/1024/1024:.0f}MB)") + + if deleted: + free_gb = self._free_gb() + self.db.record_activity( + 'disk_guard', 'cleaned', + f"清理 {deleted} 个素材,释放约 {freed_bytes/1024**3:.1f}GB," + f"当前剩余 {free_gb:.1f}GB") + logger.info(f"DiskGuard 本轮清理完成:{deleted} 个文件,当前剩余 {free_gb:.1f}GB") diff --git a/fam-edge/src/fam_edge/oracle_db.py b/fam-edge/src/fam_edge/oracle_db.py index b175685..de9b2bd 100644 --- a/fam-edge/src/fam_edge/oracle_db.py +++ b/fam-edge/src/fam_edge/oracle_db.py @@ -567,6 +567,23 @@ class OracleDB: logger.warning(f"删除视频文件失败 {local_path}: {e}") return local_path or '' + def get_oldest_purgeable_material(self) -> Optional[Dict]: + """磁盘空间紧张时的清理候选:最旧的、已完成分割阶段的整段素材 + (filename 不是 motion_ 前缀)。按 id 升序取第一个(id 越小越早入库)。 + + 只挑 status='done' 的——素材一旦完成分割就不会再被 Video-Queue 重新 + 捡起,删掉它的本地文件不影响任何功能(运动片段是独立文件,事件时间 + 轴/人物头像只依赖 motion_clips 里的片段,不依赖原始整段素材);绝不 + 碰 pending/processing 中的,避免删掉还没来得及处理的数据。 + """ + row = self._conn.execute( + "SELECT id, local_path FROM videos " + "WHERE status='done' AND local_path IS NOT NULL AND local_path != '' " + "AND filename NOT LIKE 'motion_%' " + "ORDER BY id ASC LIMIT 1" + ).fetchone() + return dict(row) if row else None + def mark_video_invalid(self, video_id: int, error: str = ''): """文件校验不通过(损坏/非视频等),标记 invalid,producer 不再重试。""" now = _now_iso() diff --git a/fam-edge/tests/test_disk_guard.py b/fam-edge/tests/test_disk_guard.py new file mode 100644 index 0000000..fdec0f1 --- /dev/null +++ b/fam-edge/tests/test_disk_guard.py @@ -0,0 +1,120 @@ +import pytest + +from fam_edge.disk_guard import DiskGuard + + +def _cfg(**overrides): + base = { + "enabled": True, + "min_free_gb": 10, + "target_free_gb": 15, + "watch_path": "/opt/fam-edge", + "max_delete_per_round": 50, + "check_interval_sec": 300, + } + base.update(overrides) + return base + + +class _FakeDB: + def __init__(self, candidates=None): + self._candidates = list(candidates or []) + self.deleted_ids = [] + self.activities = [] + + def get_oldest_purgeable_material(self): + if not self._candidates: + return None + return self._candidates.pop(0) + + def delete_video(self, video_id): + self.deleted_ids.append(video_id) + return None + + def record_activity(self, service, action, detail=''): + self.activities.append((service, action, detail)) + + +def _guard(monkeypatch, db, free_gb_sequence, **cfg_overrides): + """free_gb_sequence: 每次调用 _free_gb() 依次返回的值列表(最后一个值 + 耗尽后保持不变),用来模拟"清理一个文件后空间逐步恢复"的过程。""" + monkeypatch.setattr( + "fam_edge.disk_guard.load_config", + lambda: {"disk_guard": _cfg(**cfg_overrides)}) + guard = DiskGuard(db) + seq = list(free_gb_sequence) + + def fake_free_gb(): + if len(seq) > 1: + return seq.pop(0) + return seq[0] + + monkeypatch.setattr(guard, "_free_gb", fake_free_gb) + return guard + + +def test_check_once_does_nothing_when_space_sufficient(monkeypatch): + db = _FakeDB() + guard = _guard(monkeypatch, db, [20.0]) + + guard.check_once() + + assert db.deleted_ids == [] + assert db.activities == [] + + +def test_check_once_cleans_until_target_reached(monkeypatch): + """核心诉求: 低于 min_free_gb 触发清理,一直清到 target_free_gb 为止, + 不是清一个就停(否则马上又会跌破阈值,频繁触发)。""" + db = _FakeDB(candidates=[ + {"id": 1, "local_path": "/tmp/a.mp4"}, + {"id": 2, "local_path": "/tmp/b.mp4"}, + {"id": 3, "local_path": "/tmp/c.mp4"}, + ]) + # 初始 8GB(< min_free_gb=10),每删一个恢复到 8/12/16GB(16 >= target=15 时停) + guard = _guard(monkeypatch, db, [8.0, 8.0, 12.0, 16.0]) + + guard.check_once() + + assert db.deleted_ids == [1, 2] + actions = [a[1] for a in db.activities] + assert "low_space" in actions + assert "cleaned" in actions + + +def test_check_once_stops_when_no_candidate_left(monkeypatch): + """核心诉求: 候选清空了但空间依然不足,不能死循环,要停下来并记录一条 + "没有可清理素材"的警告,让人能在服务状态页看到这个异常情况——即便如此, + 已经发生的清理动作本身也要记录(能看到确实清过、释放了多少),不因为 + 没完全达标就把 cleaned 记录吞掉。""" + db = _FakeDB(candidates=[{"id": 1, "local_path": "/tmp/a.mp4"}]) + guard = _guard(monkeypatch, db, [8.0, 8.0, 9.0]) # 删完仅 1 个后仍然 <15GB + + guard.check_once() + + assert db.deleted_ids == [1] + actions = [a[1] for a in db.activities] + assert "no_candidate" in actions + assert "cleaned" in actions + + +def test_check_once_respects_max_delete_per_round(monkeypatch): + """核心诉求: 单轮清理有上限,防止候选异常多时一次性删太多——下一轮检查 + 很快就会再触发,没必要在一轮里清空所有候选。""" + candidates = [{"id": i, "local_path": f"/tmp/{i}.mp4"} for i in range(1, 6)] + db = _FakeDB(candidates=candidates) + # 空间一直卡在 8GB 不涨(模拟每个文件都很小,删多少都到不了 target) + guard = _guard(monkeypatch, db, [8.0], max_delete_per_round=3) + + guard.check_once() + + assert db.deleted_ids == [1, 2, 3] + + +def test_disabled_guard_does_not_start(monkeypatch): + db = _FakeDB() + guard = _guard(monkeypatch, db, [1.0], enabled=False) + + guard.start() + + assert guard.is_alive() is False diff --git a/fam-edge/tests/test_oracle_db.py b/fam-edge/tests/test_oracle_db.py index a4fb008..d93c947 100644 --- a/fam-edge/tests/test_oracle_db.py +++ b/fam-edge/tests/test_oracle_db.py @@ -407,6 +407,45 @@ def test_delete_video_nonexistent_returns_none(tmp_path): assert db.delete_video(99999) is None +def test_get_oldest_purgeable_material_none_when_empty(tmp_path): + db = _db(tmp_path) + assert db.get_oldest_purgeable_material() is None + + +def test_get_oldest_purgeable_material_ignores_motion_clips(tmp_path): + """核心诉求: 运动片段(motion_ 前缀)是独立的分析产物,事件时间轴/人物 + 头像都依赖它,磁盘清理绝不能碰它,只能清理原始整段素材。""" + db = _db(tmp_path) + vid = db.ensure_video("motion_1_1000.mp4", "/tmp/motion_1_1000.mp4", + event_start_time="2026-08-22 10:00:00") + db.mark_video_processed(vid, "摘要", [], [], "gemini") + assert db.get_oldest_purgeable_material() is None + + +def test_get_oldest_purgeable_material_ignores_non_done_status(tmp_path): + """核心诉求: 还在 pending/processing 的素材不能被清理,避免删掉还没 + 来得及处理的数据。""" + db = _db(tmp_path) + db.ensure_video("Generic_ONVIF-001-20260815-000000.mp4", + "/tmp/Generic_ONVIF-001-20260815-000000.mp4") + assert db.get_oldest_purgeable_material() is None + + +def test_get_oldest_purgeable_material_returns_oldest_done_material(tmp_path): + db = _db(tmp_path) + vid1 = db.ensure_video("Generic_ONVIF-001-20260815-000000.mp4", + "/tmp/Generic_ONVIF-001-20260815-000000.mp4") + db.mark_video_processed(vid1, "(整段素材已分割 0 段运动片段)", [], [], 'motion_segment') + vid2 = db.ensure_video("Generic_ONVIF-001-20260816-000000.mp4", + "/tmp/Generic_ONVIF-001-20260816-000000.mp4") + db.mark_video_processed(vid2, "(整段素材已分割 0 段运动片段)", [], [], 'motion_segment') + + candidate = db.get_oldest_purgeable_material() + + assert candidate["id"] == vid1 + assert candidate["local_path"] == "/tmp/Generic_ONVIF-001-20260815-000000.mp4" + + def test_delete_video_does_not_touch_ss_motion_events(tmp_path): """核心诉求: ss_motion_events 是运动侦测源事件,跟切出来的视频片段生命周期 独立,删视频不该连带删掉源事件(否则分割逻辑的幂等判断会被破坏)。"""