feat(fam-edge): 新增 DiskGuard 磁盘空间守护,剩余空间不足自动清理旧素材

背景:Oracle 磁盘曾经被写满(gdrive_videos 持续下载新素材、旧文件迟迟没
清理),触发 rclone 的一个安全机制——同步遇到 IO 错误就整体拒绝执行删除,
形成"越满越删不掉,越删不掉越满"的死循环,最终连新视频都下载不了。这次
排查+手动清理已经解决了当次故障,但需要一道独立于 rclone 同步之外的兜底,
防止再次悄悄写满没人发现。

- oracle_db.py 新增 get_oldest_purgeable_material():只挑最旧的、已完成
  分割阶段(status='done')的整段素材(非 motion_ 前缀),绝不碰运动片段
  (事件时间轴/人物头像依赖它)和还在处理中的素材
- disk_guard.py 新增 DiskGuard 后台线程:5 分钟检查一次,剩余空间 <10GB
  触发清理,删到 15GB 水位为止(留缓冲避免刚清完又立刻触发),复用已有的
  delete_video() 完成实际删除
- app.py 启动这个后台服务;/api/oracle/activity 新增 disk 字段(实时剩余
  空间 + 最近一次清理动作),供服务状态页展示
- 新增 9 个单元测试,全部通过(139/139)

已部署 Oracle 验证:DiskGuard 正常启动,配置生效。

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
ericwyuan
2026-08-28 12:08:04 +08:00
parent 3ec79de911
commit ef56ae5662
7 changed files with 320 additions and 1 deletions

View File

@@ -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],

View File

@@ -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)

View File

@@ -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")

View File

@@ -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 = ''):
"""文件校验不通过(损坏/非视频等),标记 invalidproducer 不再重试。"""
now = _now_iso()