mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-22 00:32:50 +08:00
fix(monitor): 目录监控自愈、快照语义修正与覆盖保护闭环 (#6210)
This commit is contained in:
14
app/monitor/__init__.py
Normal file
14
app/monitor/__init__.py
Normal file
@@ -0,0 +1,14 @@
|
||||
"""
|
||||
目录监控包。
|
||||
|
||||
- watcher.py 本地目录监控线程(watchfiles)
|
||||
- syslimits.py 系统限制探测与监控模式决策
|
||||
- snapshot.py 远程快照存取与比对
|
||||
- dispatcher.py 监控事件到整理链的分发
|
||||
- poller.py 远程目录轮询监控
|
||||
- monitor.py Monitor 门面:装配、生命周期与健康检查
|
||||
"""
|
||||
from app.monitor.watcher import DirectoryChangeEvent, LocalDirectoryWatcher
|
||||
from app.monitor.monitor import Monitor
|
||||
|
||||
__all__ = ["DirectoryChangeEvent", "LocalDirectoryWatcher", "Monitor"]
|
||||
210
app/monitor/dispatcher.py
Normal file
210
app/monitor/dispatcher.py
Normal file
@@ -0,0 +1,210 @@
|
||||
import re
|
||||
import traceback
|
||||
from pathlib import Path
|
||||
from threading import Lock
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from app.chain.transfer import TransferChain
|
||||
from app.core.cache import TTLCache
|
||||
from app.core.config import settings
|
||||
from app.db.transferhistory_oper import TransferHistoryOper
|
||||
from app.log import logger
|
||||
from app.schemas import FileItem
|
||||
|
||||
|
||||
class TransferDispatcher:
|
||||
"""
|
||||
将监控事件分发到整理链:候选判定、TTL 去重、整理历史查重与整理触发。
|
||||
"""
|
||||
# 历史查询失败待重试队列上限,防止长时间故障期间无限增长
|
||||
MAX_PENDING_RETRIES = 1000
|
||||
# 单个文件的最大重试次数(按健康检查周期计,60 次约 1 小时)
|
||||
MAX_RETRY_ATTEMPTS = 60
|
||||
|
||||
def __init__(self, all_exts: Optional[List[str]] = None, cache: Optional[Any] = None):
|
||||
"""
|
||||
初始化整理分发器。
|
||||
:param all_exts: 监控的文件扩展名,默认取系统配置
|
||||
:param cache: 去重缓存,默认使用 10 秒 TTL 缓存
|
||||
"""
|
||||
self.all_exts = all_exts if all_exts is not None else (
|
||||
settings.RMT_MEDIAEXT + settings.RMT_SUBEXT + settings.RMT_AUDIOEXT)
|
||||
self._cache = cache if cache is not None else TTLCache(region="monitor", maxsize=1024, ttl=10)
|
||||
self._lock = Lock()
|
||||
# 历史查询失败待重试的文件
|
||||
self._pending_retries: Dict[str, Dict[str, Any]] = {}
|
||||
self._pending_guard = Lock()
|
||||
|
||||
@staticmethod
|
||||
def _is_bluray_sub(_path: Path) -> bool:
|
||||
"""
|
||||
判断是否蓝光原盘目录内的媒体流文件。
|
||||
"""
|
||||
return True if re.search(r"BDMV[/\\]STREAM", _path.as_posix(), re.IGNORECASE) else False
|
||||
|
||||
@staticmethod
|
||||
def _get_bluray_dir(_path: Path) -> Optional[Path]:
|
||||
"""
|
||||
获取蓝光原盘BDMV目录的上级目录。
|
||||
"""
|
||||
for p in _path.parents:
|
||||
if p.name == "BDMV":
|
||||
return p.parent
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def _has_suffix_in(file_path: Path, extensions: List[str]) -> bool:
|
||||
"""
|
||||
判断路径后缀是否命中给定扩展名列表。
|
||||
"""
|
||||
if not file_path.suffix:
|
||||
return False
|
||||
return file_path.suffix.casefold() in {ext.casefold() for ext in extensions}
|
||||
|
||||
def is_transfer_candidate_path(self, file_path: Path) -> bool:
|
||||
"""
|
||||
判断监控事件路径是否需要进入整理链。
|
||||
"""
|
||||
if self._has_suffix_in(file_path, settings.DOWNLOAD_TMPEXT):
|
||||
return False
|
||||
return self._has_suffix_in(file_path, self.all_exts)
|
||||
|
||||
@staticmethod
|
||||
def _build_transfer_src_path(event_path: Path, is_bluray_folder: bool) -> str:
|
||||
"""
|
||||
生成整理记录使用的源路径。
|
||||
"""
|
||||
if is_bluray_folder:
|
||||
return f"{event_path.as_posix()}/"
|
||||
return event_path.as_posix()
|
||||
|
||||
@staticmethod
|
||||
def _has_transfer_history(storage: str, src_path: str) -> Optional[bool]:
|
||||
"""
|
||||
判断源文件是否已经存在整理记录。
|
||||
:return: True/False 查询成功,None 查询失败
|
||||
"""
|
||||
try:
|
||||
return bool(TransferHistoryOper().get_by_src(src_path, storage=storage))
|
||||
except Exception as err:
|
||||
logger.error(f"查询整理历史失败: {src_path} - {err}")
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def _pending_key(storage: str, event_path: Path) -> str:
|
||||
"""
|
||||
生成待重试文件的唯一键。
|
||||
"""
|
||||
return f"{storage}:{Path(event_path).as_posix()}"
|
||||
|
||||
def _register_pending(self, storage: str, event_path: Path, file_size: float = None):
|
||||
"""
|
||||
登记历史查询失败的文件待重试,重复失败累计次数,超限后放弃。
|
||||
:param storage: 存储
|
||||
:param event_path: 原始事件路径
|
||||
:param file_size: 文件大小
|
||||
"""
|
||||
key = self._pending_key(storage, event_path)
|
||||
with self._pending_guard:
|
||||
entry = self._pending_retries.get(key)
|
||||
if entry:
|
||||
entry["attempts"] += 1
|
||||
if entry["attempts"] >= self.MAX_RETRY_ATTEMPTS:
|
||||
self._pending_retries.pop(key, None)
|
||||
logger.error(f"整理历史查询持续失败,已放弃重试: {key}")
|
||||
return
|
||||
if len(self._pending_retries) >= self.MAX_PENDING_RETRIES:
|
||||
logger.error(f"整理重试队列已满,丢弃: {key}")
|
||||
return
|
||||
self._pending_retries[key] = {
|
||||
"storage": storage,
|
||||
"event_path": event_path,
|
||||
"file_size": file_size,
|
||||
"attempts": 1
|
||||
}
|
||||
logger.warn(f"整理历史查询失败,已登记待重试: {key}")
|
||||
|
||||
def _discard_pending(self, storage: str, event_path: Path):
|
||||
"""
|
||||
历史查询已得到确定结果,移除待重试登记。
|
||||
:param storage: 存储
|
||||
:param event_path: 原始事件路径
|
||||
"""
|
||||
with self._pending_guard:
|
||||
self._pending_retries.pop(self._pending_key(storage, event_path), None)
|
||||
|
||||
def retry_pending(self):
|
||||
"""
|
||||
重试历史查询失败的文件,由健康检查周期驱动。
|
||||
成功或得到确定结果的条目在 handle_file 内部自动移除。
|
||||
"""
|
||||
with self._pending_guard:
|
||||
items = list(self._pending_retries.values())
|
||||
for item in items:
|
||||
logger.info(f"重试整理: {item['storage']}:{item['event_path']}")
|
||||
self.handle_file(storage=item["storage"], event_path=item["event_path"],
|
||||
file_size=item["file_size"])
|
||||
|
||||
def handle_file(self, storage: str, event_path: Path, file_size: float = None) -> bool:
|
||||
"""
|
||||
整理一个文件。
|
||||
:param storage: 存储
|
||||
:param event_path: 事件文件路径
|
||||
:param file_size: 文件大小
|
||||
:return: 是否进入整理链
|
||||
"""
|
||||
with self._lock:
|
||||
# 登记重试用原始事件路径,蓝光目录解析在重试时重新执行
|
||||
origin_path = event_path
|
||||
is_bluray_folder = False
|
||||
# 蓝光原盘文件处理
|
||||
if self._is_bluray_sub(event_path):
|
||||
event_path = self._get_bluray_dir(event_path)
|
||||
if not event_path:
|
||||
return False
|
||||
is_bluray_folder = True
|
||||
elif not self.is_transfer_candidate_path(event_path):
|
||||
return False
|
||||
|
||||
# TTL缓存控重
|
||||
if self._cache.get(str(event_path)):
|
||||
return False
|
||||
self._cache[str(event_path)] = True
|
||||
|
||||
src_path = self._build_transfer_src_path(
|
||||
event_path=event_path,
|
||||
is_bluray_folder=is_bluray_folder,
|
||||
)
|
||||
has_transfer_history = self._has_transfer_history(
|
||||
storage=storage,
|
||||
src_path=src_path,
|
||||
)
|
||||
if has_transfer_history is None:
|
||||
# 查询失败是暂时故障,登记待重试(由健康检查周期驱动),不能永久跳过
|
||||
self._register_pending(storage=storage, event_path=origin_path, file_size=file_size)
|
||||
return False
|
||||
self._discard_pending(storage=storage, event_path=origin_path)
|
||||
if has_transfer_history:
|
||||
return False
|
||||
|
||||
try:
|
||||
if is_bluray_folder:
|
||||
logger.info(f"开始整理蓝光原盘: {event_path}")
|
||||
else:
|
||||
logger.info(f"开始整理文件: {event_path}")
|
||||
# 开始整理
|
||||
TransferChain().do_transfer(
|
||||
fileitem=FileItem(
|
||||
storage=storage,
|
||||
path=src_path,
|
||||
type="file" if not is_bluray_folder else "dir",
|
||||
name=event_path.name,
|
||||
basename=event_path.stem,
|
||||
extension=event_path.suffix[1:],
|
||||
size=file_size
|
||||
)
|
||||
)
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error("目录监控整理文件发生错误:%s - %s" % (str(e), traceback.format_exc()))
|
||||
return False
|
||||
508
app/monitor/monitor.py
Normal file
508
app/monitor/monitor.py
Normal file
@@ -0,0 +1,508 @@
|
||||
import traceback
|
||||
from pathlib import Path
|
||||
from threading import Lock
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from apscheduler.schedulers.background import BackgroundScheduler
|
||||
|
||||
from app.core.config import settings
|
||||
from app.helper.directory import DirectoryHelper
|
||||
from app.helper.message import MessageHelper
|
||||
from app.log import logger
|
||||
from app.monitor.dispatcher import TransferDispatcher
|
||||
from app.monitor.poller import RemotePoller
|
||||
from app.monitor.snapshot import SnapshotStore
|
||||
from app.monitor.syslimits import decide_monitor_mode, get_system_optimization_tips
|
||||
from app.monitor.watcher import LocalDirectoryWatcher
|
||||
from app.schemas.types import SystemConfigKey
|
||||
from app.utils.mixins import ConfigReloadMixin
|
||||
from app.utils.singleton import SingletonClass
|
||||
from app.utils.system import SystemUtils
|
||||
|
||||
|
||||
class Monitor(ConfigReloadMixin, metaclass=SingletonClass):
|
||||
"""
|
||||
目录监控门面,单例模式:装配本地/远程监控、维护生命周期与健康检查。
|
||||
"""
|
||||
CONFIG_WATCH = {SystemConfigKey.Directories.value}
|
||||
# 目录监控健康检查间隔(秒)
|
||||
WATCHDOG_INTERVAL = 60
|
||||
# 连续多少个健康检查周期无新增重启后才宣告恢复,避免反复崩溃时告警刷屏
|
||||
RECOVERY_STABLE_CYCLES = 5
|
||||
|
||||
def __init__(self):
|
||||
super().__init__()
|
||||
# 本地目录监控服务
|
||||
self._watchers = []
|
||||
# 本地目录监控列表读写锁
|
||||
self._watcher_lock = Lock()
|
||||
# 启动失败待重试的本地监控配置
|
||||
self._pending_locals: List[Dict[str, Any]] = []
|
||||
# 已告警的监控目录,避免重复推送
|
||||
self._alerted_paths: set = set()
|
||||
# 各监控目录已告警过的自动重启次数
|
||||
self._restart_marks: Dict[str, int] = {}
|
||||
# 各监控目录连续稳定的健康检查周期数
|
||||
self._stable_cycles: Dict[str, int] = {}
|
||||
# 定时服务
|
||||
self._scheduler = None
|
||||
# 整理分发器
|
||||
self._dispatcher = TransferDispatcher()
|
||||
# 快照存储
|
||||
self._store = SnapshotStore()
|
||||
# 远程轮询监控
|
||||
self._poller = RemotePoller(store=self._store, dispatcher=self._dispatcher,
|
||||
alert_cb=self.__poller_alert)
|
||||
# 启动目录监控和文件整理
|
||||
self.init()
|
||||
|
||||
def on_config_changed(self):
|
||||
self.init()
|
||||
|
||||
def get_reload_name(self):
|
||||
return "目录监控"
|
||||
|
||||
def save_snapshot(self, storage: str, snapshot: Dict, file_count: int = 0,
|
||||
last_snapshot_time: Optional[float] = None):
|
||||
"""
|
||||
保存快照到文件缓存。
|
||||
"""
|
||||
self._store.save(storage, snapshot, file_count=file_count, last_snapshot_time=last_snapshot_time)
|
||||
|
||||
def load_snapshot(self, storage: str) -> Optional[Dict]:
|
||||
"""
|
||||
从文件缓存加载快照。
|
||||
"""
|
||||
return self._store.load(storage)
|
||||
|
||||
def reset_snapshot(self, storage: str) -> bool:
|
||||
"""
|
||||
重置快照,强制下次扫描时重新建立基准。
|
||||
"""
|
||||
return self._store.reset(storage)
|
||||
|
||||
def force_full_scan(self, storage: str, mon_path: Path) -> bool:
|
||||
"""
|
||||
强制全量扫描并处理所有文件(包括已存在的文件)。
|
||||
"""
|
||||
return self._poller.force_full_scan(storage=storage, mon_path=mon_path)
|
||||
|
||||
@staticmethod
|
||||
def adjust_monitor_interval(file_count: int) -> int:
|
||||
"""
|
||||
根据文件数量动态调整监控间隔。
|
||||
"""
|
||||
return SnapshotStore.adjust_interval(file_count)
|
||||
|
||||
@staticmethod
|
||||
def compare_snapshots(old_snapshot: Dict, new_snapshot: Dict) -> Dict[str, List]:
|
||||
"""
|
||||
比对快照,找出变化的文件。
|
||||
"""
|
||||
return SnapshotStore.compare(old_snapshot, new_snapshot)
|
||||
|
||||
def init(self):
|
||||
"""
|
||||
启动监控
|
||||
"""
|
||||
# 停止现有任务
|
||||
self.stop()
|
||||
|
||||
# 读取目录配置
|
||||
monitor_dirs = DirectoryHelper().get_download_dirs()
|
||||
if not monitor_dirs:
|
||||
logger.info("未找到任何目录监控配置")
|
||||
return
|
||||
|
||||
messagehelper = MessageHelper()
|
||||
|
||||
# 先筛出有效的监控配置,再按下载目录去重,避免非监控配置顶掉监控配置
|
||||
valid_dirs = []
|
||||
for mon_dir in monitor_dirs:
|
||||
if not mon_dir.library_path:
|
||||
logger.warn(f"跳过监控配置 {mon_dir.download_path}:未设置媒体库目录")
|
||||
continue
|
||||
if mon_dir.monitor_type != "monitor":
|
||||
logger.debug(f"跳过监控配置 {mon_dir.download_path}:监控类型为 {mon_dir.monitor_type}")
|
||||
continue
|
||||
valid_dirs.append(mon_dir)
|
||||
|
||||
deduped: Dict[str, Any] = {}
|
||||
for mon_dir in valid_dirs:
|
||||
key = f"{mon_dir.storage}_{mon_dir.download_path}"
|
||||
if key in deduped:
|
||||
logger.warn(f"监控配置重复,忽略后一条: {mon_dir.download_path}"
|
||||
f"(媒体库 {mon_dir.library_path})")
|
||||
continue
|
||||
deduped[key] = mon_dir
|
||||
monitor_dirs = list(deduped.values())
|
||||
logger.info(f"找到 {len(monitor_dirs)} 个目录监控配置")
|
||||
|
||||
# 启动定时服务进程
|
||||
self._scheduler = BackgroundScheduler(timezone=settings.TZ)
|
||||
|
||||
mon_storages: Dict[str, List[Path]] = {}
|
||||
# 本地监控启动结果计数,用于输出真实的启动总结
|
||||
local_started = 0
|
||||
local_failed = 0
|
||||
for mon_dir in monitor_dirs:
|
||||
# 检查媒体库目录是不是下载目录的子目录
|
||||
mon_path = Path(mon_dir.download_path)
|
||||
target_path = Path(mon_dir.library_path)
|
||||
if target_path.is_relative_to(mon_path):
|
||||
logger.warn(f"{target_path} 是监控目录 {mon_path} 的子目录,无法监控!")
|
||||
messagehelper.put(f"{target_path} 是监控目录 {mon_path} 的子目录,无法监控", title="目录监控")
|
||||
continue
|
||||
|
||||
# 启动监控
|
||||
if mon_dir.storage == "local":
|
||||
if self.__start_local_monitor(mon_path=mon_path, monitor_mode=mon_dir.monitor_mode):
|
||||
local_started += 1
|
||||
else:
|
||||
local_failed += 1
|
||||
else:
|
||||
mon_storages.setdefault(mon_dir.storage, []).append(mon_path)
|
||||
|
||||
for storage, paths in mon_storages.items():
|
||||
# 远程目录监控 - 使用智能间隔
|
||||
# 先尝试加载已有快照获取文件数量
|
||||
snapshot_data = self._store.load(storage)
|
||||
file_count = snapshot_data.get('file_count', 0) if snapshot_data else 0
|
||||
interval = SnapshotStore.adjust_interval(file_count)
|
||||
for path in paths:
|
||||
logger.info(f"正在启动远程目录监控: {path} [{storage}]")
|
||||
logger.info("*** 重要提示:远程目录监控只处理新增和修改的文件,不会处理监控启动前已存在的文件 ***")
|
||||
logger.info(f"预估文件数量: {file_count}, 监控间隔: {interval}分钟")
|
||||
|
||||
self._scheduler.add_job(
|
||||
self.polling_observer,
|
||||
'interval',
|
||||
minutes=interval,
|
||||
kwargs={
|
||||
'storage': storage,
|
||||
'mon_paths': paths
|
||||
},
|
||||
id=f"monitor_{storage}",
|
||||
replace_existing=True
|
||||
)
|
||||
logger.info(f"✓ 远程目录监控已启动: [间隔: {interval}分钟]")
|
||||
|
||||
# 监控健康检查:重建异常监控线程、重试启动失败目录、重试历史查询失败的文件
|
||||
if local_started or local_failed or mon_storages:
|
||||
self._scheduler.add_job(
|
||||
self.watchdog,
|
||||
'interval',
|
||||
seconds=self.WATCHDOG_INTERVAL,
|
||||
id="monitor_watchdog",
|
||||
replace_existing=True
|
||||
)
|
||||
logger.info(f"✓ 目录监控健康检查已启动: [间隔: {self.WATCHDOG_INTERVAL}秒]")
|
||||
|
||||
# 启动定时服务
|
||||
if self._scheduler.get_jobs():
|
||||
self._scheduler.print_jobs()
|
||||
self._scheduler.start()
|
||||
logger.info("定时监控服务已启动")
|
||||
|
||||
# 输出监控总结,报告实际启动成功数而不是配置数
|
||||
remote_count = sum(len(paths) for paths in mon_storages.values())
|
||||
summary = f"目录监控启动完成: 本地监控 {local_started} 个成功"
|
||||
if local_failed:
|
||||
summary += f"、{local_failed} 个失败(将自动退避重试)"
|
||||
summary += f",远程监控 {remote_count} 个"
|
||||
if local_failed:
|
||||
logger.warn(summary)
|
||||
else:
|
||||
logger.info(summary)
|
||||
|
||||
def __start_local_monitor(self, mon_path: Path, monitor_mode: str) -> bool:
|
||||
"""
|
||||
启动单个本地目录监控,失败时登记待重试。
|
||||
:param mon_path: 监控目录
|
||||
:param monitor_mode: 配置的监控模式
|
||||
:return: 是否启动成功
|
||||
"""
|
||||
logger.info(f"正在启动本地目录监控: {mon_path}")
|
||||
logger.info("*** 重要提示:目录监控只处理新增和修改的文件,不会处理监控启动前已存在的文件 ***")
|
||||
|
||||
try:
|
||||
# 检查是否需要使用轮询模式(兼容模式/网络存储不做启动期目录遍历)
|
||||
use_polling, reason, limits, file_count = decide_monitor_mode(mon_path, monitor_mode)
|
||||
logger.info(f"监控模式决策: {reason}")
|
||||
|
||||
mode_name = "兼容模式(轮询)" if use_polling else "快速模式"
|
||||
logger.info(f"使用{mode_name}监控 {mon_path}")
|
||||
if file_count is not None:
|
||||
logger.info(f"监控目录 {mon_path} 包含约 {file_count} 个文件")
|
||||
if not use_polling and limits:
|
||||
if limits['warnings']:
|
||||
for warning in limits['warnings']:
|
||||
logger.warn(f"系统限制警告: {warning}")
|
||||
if limits['max_user_watches'] > 0 and file_count is not None:
|
||||
usage_percent = (file_count / limits['max_user_watches']) * 100
|
||||
logger.info(
|
||||
f"系统监控资源使用率: {usage_percent:.1f}% ({file_count}/{limits['max_user_watches']})")
|
||||
|
||||
# 网络/FUSE 挂载轮询降频,减少监控自身对挂载后端的持续 stat 压力
|
||||
poll_delay_ms = None
|
||||
if use_polling and SystemUtils.is_network_filesystem(mon_path):
|
||||
poll_delay_ms = LocalDirectoryWatcher.POLL_DELAY_NETWORK_MS
|
||||
logger.info(f"检测到网络文件系统,轮询扫描间隔调整为 {poll_delay_ms}ms: {mon_path}")
|
||||
|
||||
watcher = LocalDirectoryWatcher(
|
||||
mon_path=mon_path,
|
||||
callback=self,
|
||||
force_polling=True if use_polling else None,
|
||||
poll_delay_ms=poll_delay_ms
|
||||
)
|
||||
# 启动成功后再登记,避免失败的监控残留在列表中
|
||||
watcher.start()
|
||||
with self._watcher_lock:
|
||||
self._watchers.append(watcher)
|
||||
self._pending_locals = [
|
||||
pending for pending in self._pending_locals
|
||||
if pending["mon_path"] != mon_path
|
||||
]
|
||||
self.__clear_alert(mon_path, f"本地目录监控已恢复: {mon_path} [{mode_name}]")
|
||||
|
||||
logger.info(f"✓ 本地目录监控已启动: {mon_path} [{mode_name}]")
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
self.__handle_start_failure(mon_path=mon_path, monitor_mode=monitor_mode, err=e)
|
||||
return False
|
||||
|
||||
def __handle_start_failure(self, mon_path: Path, monitor_mode: str, err: Exception):
|
||||
"""
|
||||
处理本地目录监控启动失败,登记待重试并按需告警。
|
||||
:param mon_path: 监控目录
|
||||
:param monitor_mode: 配置的监控模式
|
||||
:param err: 启动异常
|
||||
"""
|
||||
err_msg = str(err)
|
||||
logger.error(f"启动本地目录监控失败: {mon_path}")
|
||||
logger.error(f"错误详情: {err_msg}")
|
||||
|
||||
if "inotify" in err_msg.lower():
|
||||
logger.error("inotify 相关错误,这通常是由于系统监控数量限制导致的")
|
||||
logger.error("解决方案:")
|
||||
for tip in get_system_optimization_tips():
|
||||
logger.error(f" {tip}")
|
||||
logger.error("执行上述命令后重启 MoviePilot")
|
||||
elif "permission" in err_msg.lower():
|
||||
logger.error("权限错误,请检查 MoviePilot 是否有足够的权限访问监控目录")
|
||||
elif isinstance(err, (FileNotFoundError, NotADirectoryError)):
|
||||
logger.error("监控目录当前不可用,网络存储/FUSE 挂载可能尚未就绪,将自动重试")
|
||||
elif monitor_mode != "compatibility":
|
||||
logger.error("建议尝试使用兼容模式进行监控")
|
||||
|
||||
with self._watcher_lock:
|
||||
if all(pending["mon_path"] != mon_path for pending in self._pending_locals):
|
||||
self._pending_locals.append({
|
||||
"mon_path": mon_path,
|
||||
"monitor_mode": monitor_mode
|
||||
})
|
||||
self.__send_alert(mon_path,
|
||||
f"启动本地目录监控失败: {mon_path}\n错误: {err_msg}\n"
|
||||
f"将自动退避重试")
|
||||
|
||||
def watchdog(self):
|
||||
"""
|
||||
目录监控健康检查:重建崩溃或静默失效的监控线程,并重试启动失败的监控目录。
|
||||
"""
|
||||
try:
|
||||
self.__check_watchers()
|
||||
self.__retry_pending_locals()
|
||||
self._dispatcher.retry_pending()
|
||||
except Exception as e:
|
||||
logger.error(f"目录监控健康检查出现错误:{e}\n{traceback.format_exc()}")
|
||||
|
||||
def __check_watchers(self):
|
||||
"""
|
||||
检查本地目录监控线程状态,异常时重建。
|
||||
"""
|
||||
with self._watcher_lock:
|
||||
watchers = list(self._watchers)
|
||||
for watcher in watchers:
|
||||
key = str(watcher.watch_path)
|
||||
if watcher.is_stalled():
|
||||
reason = f"监控循环超过 {LocalDirectoryWatcher.STALL_TIMEOUT} 秒无任何活动,判定为静默失效"
|
||||
elif not watcher.is_alive():
|
||||
reason = "监控线程已退出"
|
||||
else:
|
||||
# 线程已自愈,但崩溃过就要告警,避免自动重启把故障变成新的静默
|
||||
if watcher.restart_count > self._restart_marks.get(key, 0):
|
||||
self._restart_marks[key] = watcher.restart_count
|
||||
self._stable_cycles[key] = 0
|
||||
self.__send_alert(watcher.watch_path,
|
||||
f"目录监控发生错误并已自动重启"
|
||||
f"(累计 {watcher.restart_count} 次): {watcher.watch_path}")
|
||||
else:
|
||||
# 稳定满恢复窗口才宣告恢复,避免反复崩溃时告警/恢复消息来回刷屏
|
||||
self._stable_cycles[key] = self._stable_cycles.get(key, 0) + 1
|
||||
if self._stable_cycles[key] >= self.RECOVERY_STABLE_CYCLES:
|
||||
self.__clear_alert(watcher.watch_path, f"目录监控已恢复正常: {watcher.watch_path}")
|
||||
continue
|
||||
logger.error(f"目录监控异常: {watcher.watch_path} - {reason},正在重建监控线程 ...")
|
||||
self.__send_alert(watcher.watch_path,
|
||||
f"目录监控异常: {watcher.watch_path}\n原因: {reason}\n正在自动重建监控")
|
||||
self.__rebuild_watcher(watcher)
|
||||
|
||||
def __rebuild_watcher(self, watcher: LocalDirectoryWatcher):
|
||||
"""
|
||||
重建一个本地目录监控线程。
|
||||
:param watcher: 需要重建的监控
|
||||
"""
|
||||
# 卡死的线程阻塞在底层调用中无法强制回收,只能请求停止后由守护线程自然退出
|
||||
watcher.stop()
|
||||
new_watcher = LocalDirectoryWatcher(
|
||||
mon_path=watcher.watch_path,
|
||||
callback=self,
|
||||
force_polling=watcher.force_polling,
|
||||
poll_delay_ms=watcher.poll_delay_ms
|
||||
)
|
||||
try:
|
||||
new_watcher.start()
|
||||
except Exception as e:
|
||||
logger.error(f"重建目录监控失败: {watcher.watch_path} - {e}")
|
||||
with self._watcher_lock:
|
||||
self._watchers = [item for item in self._watchers if item is not watcher]
|
||||
if all(pending["mon_path"] != watcher.watch_path for pending in self._pending_locals):
|
||||
self._pending_locals.append({
|
||||
"mon_path": watcher.watch_path,
|
||||
# 重建沿用原监控模式,force_polling 为 True 即兼容模式
|
||||
"monitor_mode": "compatibility" if watcher.force_polling else "fast"
|
||||
})
|
||||
return
|
||||
with self._watcher_lock:
|
||||
self._watchers = [new_watcher if item is watcher else item for item in self._watchers]
|
||||
# 新监控的重启计数从零开始,同步重置告警基准
|
||||
self._restart_marks.pop(str(watcher.watch_path), None)
|
||||
self._stable_cycles.pop(str(watcher.watch_path), None)
|
||||
logger.info(f"✓ 目录监控已重建: {watcher.watch_path}")
|
||||
self.__clear_alert(watcher.watch_path, f"目录监控已自动恢复: {watcher.watch_path}")
|
||||
|
||||
def __retry_pending_locals(self):
|
||||
"""
|
||||
重试启动失败的本地目录监控,给网络存储/FUSE 挂载留出就绪时间。
|
||||
"""
|
||||
with self._watcher_lock:
|
||||
pending = list(self._pending_locals)
|
||||
for item in pending:
|
||||
# 失败次数越多重试间隔越长(按健康检查周期数退避),长时间故障时不刷屏
|
||||
if item.get("skip_cycles", 0) > 0:
|
||||
item["skip_cycles"] -= 1
|
||||
continue
|
||||
logger.info(f"重试启动本地目录监控: {item['mon_path']}")
|
||||
if not self.__start_local_monitor(mon_path=item["mon_path"], monitor_mode=item["monitor_mode"]):
|
||||
item["attempts"] = item.get("attempts", 0) + 1
|
||||
item["skip_cycles"] = min(item["attempts"], 10)
|
||||
|
||||
def __send_alert(self, mon_path: Path, message: str):
|
||||
"""
|
||||
推送目录监控异常告警,同一目录仅在状态变化时推送一次。
|
||||
:param mon_path: 监控目录
|
||||
:param message: 告警内容
|
||||
"""
|
||||
key = str(mon_path)
|
||||
with self._watcher_lock:
|
||||
if key in self._alerted_paths:
|
||||
return
|
||||
self._alerted_paths.add(key)
|
||||
MessageHelper().put(message, title="目录监控")
|
||||
|
||||
@staticmethod
|
||||
def __poller_alert(storage: str, message: str):
|
||||
"""
|
||||
远程轮询监控告警回调,复用消息渠道推送。
|
||||
:param storage: 存储名称
|
||||
:param message: 告警内容
|
||||
"""
|
||||
logger.warn(f"[{storage}] {message}")
|
||||
MessageHelper().put(message, title="目录监控")
|
||||
|
||||
def __clear_alert(self, mon_path: Path, message: str):
|
||||
"""
|
||||
清除目录监控异常告警状态,并在此前告警过时推送恢复消息。
|
||||
:param mon_path: 监控目录
|
||||
:param message: 恢复内容
|
||||
"""
|
||||
key = str(mon_path)
|
||||
with self._watcher_lock:
|
||||
if key not in self._alerted_paths:
|
||||
return
|
||||
self._alerted_paths.discard(key)
|
||||
logger.info(message)
|
||||
MessageHelper().put(message, title="目录监控")
|
||||
|
||||
def polling_observer(self, storage: str, mon_paths: List[Path]):
|
||||
"""
|
||||
轮询监控:执行一轮快照并按结果动态调整监控间隔。
|
||||
"""
|
||||
file_count = self._poller.poll(storage=storage, mon_paths=mon_paths)
|
||||
if file_count is None or not self._scheduler:
|
||||
return
|
||||
# 动态调整监控间隔
|
||||
new_interval = SnapshotStore.adjust_interval(file_count)
|
||||
try:
|
||||
current_job = self._scheduler.get_job(f"monitor_{storage}")
|
||||
if current_job and current_job.trigger.interval.total_seconds() / 60 != new_interval:
|
||||
self._scheduler.modify_job(
|
||||
f"monitor_{storage}",
|
||||
trigger='interval',
|
||||
minutes=new_interval
|
||||
)
|
||||
logger.info(f"{storage} 监控间隔已调整为 {new_interval} 分钟")
|
||||
except Exception as e:
|
||||
logger.error(f"调整监控间隔失败: {storage} - {e}")
|
||||
|
||||
def event_handler(self, event, text: str, event_path: str, file_size: float = None):
|
||||
"""
|
||||
处理文件变化。
|
||||
:param event: 事件
|
||||
:param text: 事件描述
|
||||
:param event_path: 事件文件路径
|
||||
:param file_size: 文件大小
|
||||
"""
|
||||
if event.is_directory:
|
||||
return
|
||||
if not self._dispatcher.is_transfer_candidate_path(Path(event_path)):
|
||||
return
|
||||
# 整理文件
|
||||
self._dispatcher.handle_file(storage="local", event_path=Path(event_path), file_size=file_size)
|
||||
|
||||
def stop(self):
|
||||
"""
|
||||
退出监控
|
||||
"""
|
||||
# 先停定时服务,避免健康检查在停止过程中重建监控线程
|
||||
if self._scheduler:
|
||||
self._scheduler.remove_all_jobs()
|
||||
if self._scheduler.running:
|
||||
try:
|
||||
self._scheduler.shutdown()
|
||||
logger.info("定时监控服务已停止")
|
||||
except Exception as e:
|
||||
logger.error(f"停止定时服务出现了错误:{e}")
|
||||
self._scheduler = None
|
||||
with self._watcher_lock:
|
||||
watchers = self._watchers
|
||||
self._watchers = []
|
||||
self._pending_locals = []
|
||||
self._alerted_paths = set()
|
||||
self._restart_marks = {}
|
||||
self._stable_cycles = {}
|
||||
if watchers:
|
||||
logger.info("正在停止本地目录监控服务...")
|
||||
for watcher in watchers:
|
||||
try:
|
||||
watcher.stop()
|
||||
watcher.join(timeout=5)
|
||||
if watcher.is_alive():
|
||||
logger.warning(f"本地目录监控线程在5秒内未能停止: {watcher.watch_path}")
|
||||
else:
|
||||
logger.debug(f"已停止本地目录监控服务: {watcher.watch_path}")
|
||||
except Exception as e:
|
||||
logger.error(f"停止目录监控服务出现了错误:{e}")
|
||||
logger.info("本地目录监控服务已停止")
|
||||
# 缓存与快照存储是共享后端的代理,生命周期由应用全局管理,这里不再关闭
|
||||
228
app/monitor/poller.py
Normal file
228
app/monitor/poller.py
Normal file
@@ -0,0 +1,228 @@
|
||||
import traceback
|
||||
from pathlib import Path
|
||||
from threading import Lock
|
||||
from typing import Callable, Dict, List, Optional
|
||||
|
||||
from app.chain.storage import StorageChain
|
||||
from app.log import logger
|
||||
from app.monitor.dispatcher import TransferDispatcher
|
||||
from app.monitor.snapshot import SnapshotStore
|
||||
|
||||
|
||||
class RemotePoller:
|
||||
"""
|
||||
远程目录轮询监控:快照、比对并分发变化文件。
|
||||
"""
|
||||
# 同一存储连续异常达到该次数后推送告警
|
||||
FAILURE_ALERT_THRESHOLD = 3
|
||||
|
||||
def __init__(self, store: SnapshotStore, dispatcher: TransferDispatcher,
|
||||
alert_cb: Optional[Callable[[str, str], None]] = None):
|
||||
"""
|
||||
初始化远程轮询监控。
|
||||
:param store: 快照存储
|
||||
:param dispatcher: 整理分发器
|
||||
:param alert_cb: 告警回调 (storage, message)
|
||||
"""
|
||||
self._store = store
|
||||
self._dispatcher = dispatcher
|
||||
self._alert_cb = alert_cb
|
||||
# 快照锁按存储隔离,避免一个慢存储阻塞其他存储的轮询
|
||||
self._locks: Dict[str, Lock] = {}
|
||||
self._locks_guard = Lock()
|
||||
# 各存储连续异常次数
|
||||
self._failure_counts: Dict[str, int] = {}
|
||||
|
||||
def _get_lock(self, storage: str) -> Lock:
|
||||
"""
|
||||
获取指定存储的快照锁。
|
||||
:param storage: 存储名称
|
||||
:return: 快照锁
|
||||
"""
|
||||
with self._locks_guard:
|
||||
return self._locks.setdefault(storage, Lock())
|
||||
|
||||
def _note_failure(self, storage: str, reason: str):
|
||||
"""
|
||||
记录一次轮询异常,连续异常达到阈值时推送告警。
|
||||
:param storage: 存储名称
|
||||
:param reason: 异常原因
|
||||
"""
|
||||
count = self._failure_counts.get(storage, 0) + 1
|
||||
self._failure_counts[storage] = count
|
||||
logger.warn(f"远程目录监控异常(连续第 {count} 次): {storage} - {reason}")
|
||||
if count == self.FAILURE_ALERT_THRESHOLD and self._alert_cb:
|
||||
self._alert_cb(storage,
|
||||
f"远程目录监控连续 {count} 次异常: {storage}\n原因: {reason}\n将继续按周期重试")
|
||||
|
||||
def _note_success(self, storage: str):
|
||||
"""
|
||||
记录一次轮询成功,此前告警过时推送恢复消息。
|
||||
:param storage: 存储名称
|
||||
"""
|
||||
if self._failure_counts.get(storage, 0) >= self.FAILURE_ALERT_THRESHOLD and self._alert_cb:
|
||||
self._alert_cb(storage, f"远程目录监控已恢复: {storage}")
|
||||
self._failure_counts[storage] = 0
|
||||
|
||||
def poll(self, storage: str, mon_paths: List[Path]) -> Optional[int]:
|
||||
"""
|
||||
执行一轮轮询监控。
|
||||
:param storage: 存储名称
|
||||
:param mon_paths: 监控路径列表
|
||||
:return: 基线文件数量,本轮无有效结果时返回 None
|
||||
"""
|
||||
monitor_scope = ",".join(str(mon_path) for mon_path in mon_paths) or "未配置路径"
|
||||
with self._get_lock(storage):
|
||||
try:
|
||||
# 加载上次快照数据,读取失败不能当作首次快照,否则会丢弃已有基线
|
||||
old_snapshot_data, load_ok = self._store.load_checked(storage)
|
||||
if not load_ok:
|
||||
self._note_failure(storage, "读取快照基线失败,跳过本轮")
|
||||
return None
|
||||
old_snapshot = old_snapshot_data.get('snapshot', {}) if old_snapshot_data else {}
|
||||
last_snapshot_time = old_snapshot_data.get('timestamp', 0) if old_snapshot_data else 0
|
||||
is_first_snapshot = old_snapshot_data is None
|
||||
|
||||
new_snapshot = {}
|
||||
failed_paths = []
|
||||
for mon_path in mon_paths:
|
||||
logger.debug(f"开始对 {storage}:{mon_path} 进行快照...")
|
||||
|
||||
# 生成新快照(增量模式)
|
||||
snapshot = StorageChain().snapshot_storage(
|
||||
storage=storage,
|
||||
path=mon_path,
|
||||
last_snapshot_time=last_snapshot_time
|
||||
)
|
||||
|
||||
if snapshot is None:
|
||||
failed_paths.append(str(mon_path))
|
||||
logger.warn(f"获取 {storage}:{mon_path} 快照失败")
|
||||
continue
|
||||
new_snapshot.update(snapshot)
|
||||
logger.info(f"{storage}:{mon_path} 快照完成,发现 {len(snapshot)} 个文件")
|
||||
|
||||
if failed_paths and (is_first_snapshot or len(failed_paths) == len(mon_paths)):
|
||||
# 首次基线必须完整建立;全部路径失败时本轮没有有效数据,均不落盘
|
||||
self._note_failure(storage, f"快照失败: {','.join(failed_paths)}")
|
||||
return None
|
||||
|
||||
# 增量快照只包含变化子树,必须与基线合并才是完整视图;
|
||||
# 直接把增量当基线会导致下一轮把未扫到的旧文件全部误判为新增
|
||||
merged_snapshot = {**old_snapshot, **new_snapshot}
|
||||
file_count = len(merged_snapshot)
|
||||
|
||||
if not is_first_snapshot:
|
||||
self._handle_changes(storage, old_snapshot, new_snapshot)
|
||||
else:
|
||||
logger.info(f"{storage} 首次快照完成,共 {file_count} 个文件")
|
||||
logger.info("*** 首次快照仅建立基准,不会处理现有文件。后续监控将处理新增和修改的文件 ***")
|
||||
|
||||
# 保存合并后的基线
|
||||
if not self._store.save(storage, merged_snapshot, file_count, last_snapshot_time):
|
||||
self._note_failure(storage, "保存快照基线失败")
|
||||
return None
|
||||
|
||||
if failed_paths:
|
||||
# 部分路径失败:成功路径已合并,失败路径保留旧基线,下轮重试
|
||||
self._note_failure(storage, f"部分路径快照失败: {','.join(failed_paths)}")
|
||||
else:
|
||||
self._note_success(storage)
|
||||
return file_count
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"轮询监控 {storage}:{monitor_scope} 出现错误:{e}\n{traceback.format_exc()}")
|
||||
self._note_failure(storage, str(e))
|
||||
return None
|
||||
|
||||
def _handle_changes(self, storage: str, old_snapshot: dict, new_snapshot: dict):
|
||||
"""
|
||||
比对快照并把变化文件送入整理链。
|
||||
:param storage: 存储名称
|
||||
:param old_snapshot: 旧基线
|
||||
:param new_snapshot: 本轮增量快照
|
||||
"""
|
||||
changes = SnapshotStore.compare(old_snapshot, new_snapshot)
|
||||
added_files = [
|
||||
file_path
|
||||
for file_path in changes['added']
|
||||
if self._dispatcher.is_transfer_candidate_path(Path(file_path))
|
||||
]
|
||||
modified_files = [
|
||||
file_path
|
||||
for file_path in changes['modified']
|
||||
if self._dispatcher.is_transfer_candidate_path(Path(file_path))
|
||||
]
|
||||
|
||||
# 处理新增文件
|
||||
handled_added_count = 0
|
||||
for new_file in added_files:
|
||||
file_info = new_snapshot.get(new_file, {})
|
||||
file_size = file_info.get('size', 0) if isinstance(file_info, dict) else file_info
|
||||
if self._dispatcher.handle_file(storage=storage, event_path=Path(new_file), file_size=file_size):
|
||||
handled_added_count += 1
|
||||
|
||||
# 处理修改文件
|
||||
handled_modified_count = 0
|
||||
for modified_file in modified_files:
|
||||
file_info = new_snapshot.get(modified_file, {})
|
||||
file_size = file_info.get('size', 0) if isinstance(file_info, dict) else file_info
|
||||
if self._dispatcher.handle_file(storage=storage, event_path=Path(modified_file), file_size=file_size):
|
||||
handled_modified_count += 1
|
||||
|
||||
if handled_added_count or handled_modified_count:
|
||||
logger.info(f"{storage} 发现 {handled_added_count} 个新增文件,{handled_modified_count} 个修改文件")
|
||||
else:
|
||||
logger.debug(f"{storage} 无文件变化")
|
||||
|
||||
def force_full_scan(self, storage: str, mon_path: Path) -> bool:
|
||||
"""
|
||||
强制全量扫描并处理所有文件(包括已存在的文件)。
|
||||
:param storage: 存储名称
|
||||
:param mon_path: 监控路径
|
||||
:return: 是否成功
|
||||
"""
|
||||
try:
|
||||
logger.info(f"开始强制全量扫描: {storage}:{mon_path}")
|
||||
|
||||
# 生成快照
|
||||
new_snapshot = StorageChain().snapshot_storage(
|
||||
storage=storage,
|
||||
path=mon_path,
|
||||
last_snapshot_time=0 # 全量扫描,不使用增量
|
||||
)
|
||||
|
||||
if new_snapshot is None:
|
||||
logger.warn(f"获取 {storage}:{mon_path} 快照失败")
|
||||
return False
|
||||
|
||||
file_count = len(new_snapshot)
|
||||
logger.info(f"{storage}:{mon_path} 全量扫描完成,发现 {file_count} 个文件")
|
||||
|
||||
# 处理所有文件
|
||||
processed_count = 0
|
||||
for file_path, file_info in new_snapshot.items():
|
||||
try:
|
||||
if not self._dispatcher.is_transfer_candidate_path(Path(file_path)):
|
||||
continue
|
||||
file_size = file_info.get('size', 0) if isinstance(file_info, dict) else file_info
|
||||
if self._dispatcher.handle_file(storage=storage, event_path=Path(file_path),
|
||||
file_size=file_size):
|
||||
processed_count += 1
|
||||
except Exception as e:
|
||||
logger.error(f"处理文件 {file_path} 失败: {e}")
|
||||
continue
|
||||
|
||||
logger.info(f"{storage}:{mon_path} 全量扫描完成,共处理 {processed_count}/{file_count} 个文件")
|
||||
|
||||
# 全量扫描覆盖单个路径,与已有基线合并后落盘,避免覆盖其他监控路径的基线
|
||||
old_snapshot_data, load_ok = self._store.load_checked(storage)
|
||||
old_snapshot = old_snapshot_data.get('snapshot', {}) if (load_ok and old_snapshot_data) else {}
|
||||
merged_snapshot = {**old_snapshot, **new_snapshot}
|
||||
self._store.save(storage, merged_snapshot, len(merged_snapshot))
|
||||
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"强制全量扫描失败: {storage}:{mon_path} - {e}")
|
||||
return False
|
||||
148
app/monitor/snapshot.py
Normal file
148
app/monitor/snapshot.py
Normal file
@@ -0,0 +1,148 @@
|
||||
import json
|
||||
import time
|
||||
from typing import Dict, List, Optional, Tuple
|
||||
|
||||
from app.core.cache import FileCache
|
||||
from app.core.config import settings
|
||||
from app.log import logger
|
||||
|
||||
|
||||
class SnapshotStore:
|
||||
"""
|
||||
远程目录监控快照的存取与比对。
|
||||
"""
|
||||
|
||||
def __init__(self, cache: Optional[FileCache] = None):
|
||||
"""
|
||||
初始化快照存储。
|
||||
:param cache: 快照文件缓存,默认使用 CACHE_PATH/snapshots
|
||||
"""
|
||||
self._cache = cache if cache is not None else FileCache(base=settings.CACHE_PATH / "snapshots")
|
||||
|
||||
def save(self, storage: str, snapshot: Dict, file_count: int = 0,
|
||||
last_snapshot_time: Optional[float] = None) -> bool:
|
||||
"""
|
||||
保存快照到文件缓存。
|
||||
:param storage: 存储名称
|
||||
:param snapshot: 快照数据
|
||||
:param file_count: 文件数量,用于调整监控间隔
|
||||
:param last_snapshot_time: 上次快照时间戳
|
||||
:return: 是否保存成功
|
||||
"""
|
||||
try:
|
||||
snapshot_time = max((item.get('modify_time', 0) for item in snapshot.values()), default=None)
|
||||
if snapshot_time is None:
|
||||
snapshot_time = last_snapshot_time or time.time()
|
||||
snapshot_data = {
|
||||
'timestamp': snapshot_time,
|
||||
'file_count': file_count,
|
||||
'snapshot': snapshot
|
||||
}
|
||||
cache_key = f"{storage}_snapshot"
|
||||
snapshot_json = json.dumps(snapshot_data, ensure_ascii=False, indent=2)
|
||||
self._cache.set(cache_key, snapshot_json.encode('utf-8'), region="snapshots")
|
||||
logger.debug(f"快照已保存到缓存: {storage}")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"保存快照失败: {e}")
|
||||
return False
|
||||
|
||||
def load_checked(self, storage: str) -> Tuple[Optional[Dict], bool]:
|
||||
"""
|
||||
从文件缓存加载快照,并区分「快照不存在」与「读取失败」。
|
||||
读取失败时不能当作首次快照处理,否则会静默丢弃已有基线。
|
||||
:param storage: 存储名称
|
||||
:return: (快照数据或None, 是否读取成功)
|
||||
"""
|
||||
try:
|
||||
cache_key = f"{storage}_snapshot"
|
||||
snapshot_data = self._cache.get(cache_key, region="snapshots")
|
||||
if snapshot_data:
|
||||
data = json.loads(snapshot_data.decode('utf-8'))
|
||||
logger.debug(f"成功加载快照: {storage}, 包含 {len(data.get('snapshot', {}))} 个文件")
|
||||
return data, True
|
||||
logger.debug(f"快照文件不存在: {storage}")
|
||||
return None, True
|
||||
except Exception as e:
|
||||
logger.error(f"加载快照失败: {e}")
|
||||
return None, False
|
||||
|
||||
def load(self, storage: str) -> Optional[Dict]:
|
||||
"""
|
||||
从文件缓存加载快照。
|
||||
:param storage: 存储名称
|
||||
:return: 快照数据或None
|
||||
"""
|
||||
data, _ = self.load_checked(storage)
|
||||
return data
|
||||
|
||||
def reset(self, storage: str) -> bool:
|
||||
"""
|
||||
重置快照,强制下次扫描时重新建立基准。
|
||||
:param storage: 存储名称
|
||||
:return: 是否成功
|
||||
"""
|
||||
try:
|
||||
cache_key = f"{storage}_snapshot"
|
||||
if self._cache.exists(cache_key, region="snapshots"):
|
||||
self._cache.delete(cache_key, region="snapshots")
|
||||
logger.info(f"快照已重置: {storage}")
|
||||
return True
|
||||
logger.debug(f"快照文件不存在,无需重置: {storage}")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"重置快照失败: {storage} - {e}")
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
def compare(old_snapshot: Dict, new_snapshot: Dict) -> Dict[str, List]:
|
||||
"""
|
||||
比对快照,找出变化的文件(只处理新增和修改,不处理删除)。
|
||||
:param old_snapshot: 旧快照
|
||||
:param new_snapshot: 新快照
|
||||
:return: 变化信息
|
||||
"""
|
||||
changes = {
|
||||
'added': [],
|
||||
'modified': []
|
||||
}
|
||||
|
||||
old_files = set(old_snapshot.keys())
|
||||
new_files = set(new_snapshot.keys())
|
||||
|
||||
# 新增文件
|
||||
changes['added'] = list(new_files - old_files)
|
||||
|
||||
# 修改文件(大小或时间变化)
|
||||
for file_path in old_files & new_files:
|
||||
old_info = old_snapshot[file_path]
|
||||
new_info = new_snapshot[file_path]
|
||||
|
||||
# 检查文件大小变化
|
||||
old_size = old_info.get('size', 0) if isinstance(old_info, dict) else old_info
|
||||
new_size = new_info.get('size', 0) if isinstance(new_info, dict) else new_info
|
||||
|
||||
# 检查修改时间变化(如果有的话)
|
||||
old_time = old_info.get('modify_time', 0) if isinstance(old_info, dict) else 0
|
||||
new_time = new_info.get('modify_time', 0) if isinstance(new_info, dict) else 0
|
||||
|
||||
if old_size != new_size or (old_time and new_time and old_time != new_time):
|
||||
changes['modified'].append(file_path)
|
||||
|
||||
return changes
|
||||
|
||||
@staticmethod
|
||||
def adjust_interval(file_count: int) -> int:
|
||||
"""
|
||||
根据文件数量动态调整监控间隔。
|
||||
:param file_count: 文件数量
|
||||
:return: 监控间隔(分钟)
|
||||
"""
|
||||
if file_count < 100:
|
||||
return 5 # 5分钟
|
||||
elif file_count < 500:
|
||||
return 10 # 10分钟
|
||||
elif file_count < 1000:
|
||||
return 15 # 15分钟
|
||||
else:
|
||||
return 30 # 30分钟
|
||||
134
app/monitor/syslimits.py
Normal file
134
app/monitor/syslimits.py
Normal file
@@ -0,0 +1,134 @@
|
||||
import os
|
||||
import platform
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
from app.log import logger
|
||||
from app.utils.system import SystemUtils
|
||||
|
||||
|
||||
def count_directory_entries(directory: Path, max_check: int = 10000) -> Tuple[int, int]:
|
||||
"""
|
||||
统计目录下的文件与子目录数量(用于检测是否超过系统限制)。
|
||||
:param directory: 目录路径
|
||||
:param max_check: 最大检查文件数量,避免长时间阻塞
|
||||
:return: (文件数量, 目录数量)
|
||||
"""
|
||||
file_count = 0
|
||||
dir_count = 0
|
||||
try:
|
||||
for _, dirs, files in os.walk(str(directory)):
|
||||
file_count += len(files)
|
||||
dir_count += len(dirs)
|
||||
if file_count > max_check:
|
||||
break
|
||||
except Exception as err:
|
||||
logger.debug(f"统计目录规模失败: {err}")
|
||||
return file_count, dir_count
|
||||
|
||||
|
||||
def count_directory_files(directory: Path, max_check: int = 10000) -> int:
|
||||
"""
|
||||
统计目录下的文件数量。
|
||||
:param directory: 目录路径
|
||||
:param max_check: 最大检查数量,避免长时间阻塞
|
||||
:return: 文件数量
|
||||
"""
|
||||
file_count, _ = count_directory_entries(directory, max_check=max_check)
|
||||
return file_count
|
||||
|
||||
|
||||
def check_system_limits() -> Dict[str, Any]:
|
||||
"""
|
||||
检查系统监控相关限制。
|
||||
:return: 系统限制信息
|
||||
"""
|
||||
limits = {
|
||||
'max_user_watches': 0,
|
||||
'max_user_instances': 0,
|
||||
'warnings': []
|
||||
}
|
||||
|
||||
try:
|
||||
if platform.system() == 'Linux':
|
||||
# 检查 inotify 限制
|
||||
try:
|
||||
with open('/proc/sys/fs/inotify/max_user_watches', 'r', encoding='utf-8', errors='replace') as f:
|
||||
limits['max_user_watches'] = int(f.read().strip())
|
||||
except Exception as e:
|
||||
logger.debug(f"读取 inotify 限制失败: {e}")
|
||||
limits['max_user_watches'] = 8192 # 默认值
|
||||
|
||||
try:
|
||||
with open('/proc/sys/fs/inotify/max_user_instances', 'r', encoding='utf-8', errors='replace') as f:
|
||||
limits['max_user_instances'] = int(f.read().strip())
|
||||
except Exception as e:
|
||||
logger.debug(f"读取 inotify 实例限制失败: {e}")
|
||||
except Exception as e:
|
||||
limits['warnings'].append(f"检查系统限制时出错: {e}")
|
||||
|
||||
return limits
|
||||
|
||||
|
||||
def get_system_optimization_tips() -> List[str]:
|
||||
"""
|
||||
获取系统优化建议。
|
||||
:return: 优化建议列表
|
||||
"""
|
||||
tips = []
|
||||
system = platform.system()
|
||||
|
||||
if system == 'Linux':
|
||||
tips.extend([
|
||||
"增加 inotify 监控数量限制:",
|
||||
"echo fs.inotify.max_user_watches=524288 | sudo tee -a /etc/sysctl.conf",
|
||||
"echo fs.inotify.max_user_instances=524288 | sudo tee -a /etc/sysctl.conf",
|
||||
"sudo sysctl -p",
|
||||
"",
|
||||
"如果在Docker中运行,请在宿主机上执行以上命令"
|
||||
])
|
||||
elif system == 'Darwin':
|
||||
tips.extend([
|
||||
"macOS 系统优化建议:",
|
||||
"sudo sysctl kern.maxfiles=65536",
|
||||
"sudo sysctl kern.maxfilesperproc=32768",
|
||||
"ulimit -n 32768"
|
||||
])
|
||||
elif system == 'Windows':
|
||||
tips.extend([
|
||||
"Windows 系统优化建议:",
|
||||
"1. 关闭不必要的实时保护软件对监控目录的扫描",
|
||||
"2. 将监控目录添加到Windows Defender排除列表",
|
||||
"3. 确保有足够的可用内存"
|
||||
])
|
||||
|
||||
return tips
|
||||
|
||||
|
||||
def decide_monitor_mode(directory: Path,
|
||||
monitor_mode: str) -> Tuple[bool, str, Optional[Dict[str, Any]], Optional[int]]:
|
||||
"""
|
||||
决策监控模式。兼容模式与网络文件系统直接短路,只有快速模式候选才统计
|
||||
目录规模与系统限制,避免启动期对网络挂载做无谓的全量遍历。
|
||||
|
||||
inotify 的 max_user_watches 按监视点(目录)计数,因此用目录数量而不是
|
||||
文件数量与上限比较。
|
||||
|
||||
:param directory: 监控目录
|
||||
:param monitor_mode: 配置的监控模式
|
||||
:return: (是否使用轮询, 原因, 系统限制信息或None, 文件数量或None)
|
||||
"""
|
||||
if monitor_mode == "compatibility":
|
||||
return True, "用户配置为兼容模式", None, None
|
||||
|
||||
# 检查网络文件系统
|
||||
if SystemUtils.is_network_filesystem(directory):
|
||||
return True, "检测到网络文件系统,建议使用兼容模式", None, None
|
||||
|
||||
limits = check_system_limits()
|
||||
file_count, dir_count = count_directory_entries(directory)
|
||||
max_watches = limits.get('max_user_watches')
|
||||
if max_watches and dir_count > max_watches * 0.8:
|
||||
return (True, f"目录数量({dir_count})接近 inotify 监控上限({max_watches})",
|
||||
limits, file_count)
|
||||
return False, "使用快速模式", limits, file_count
|
||||
303
app/monitor/watcher.py
Normal file
303
app/monitor/watcher.py
Normal file
@@ -0,0 +1,303 @@
|
||||
import threading
|
||||
import time
|
||||
import traceback
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
from watchfiles import Change, DefaultFilter, watch
|
||||
|
||||
from app.log import logger
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class DirectoryChangeEvent:
|
||||
"""
|
||||
目录文件变化事件,隔离底层 watchfiles 事件结构。
|
||||
"""
|
||||
change_type: Change
|
||||
src_path: str
|
||||
is_directory: bool
|
||||
|
||||
|
||||
class LocalDirectoryWatcher:
|
||||
"""
|
||||
基于 watchfiles 的本地目录监控线程。
|
||||
"""
|
||||
_HANDLE_CHANGES = {Change.added, Change.modified}
|
||||
# 监控循环异常退出后的重启退避秒数,网络存储/FUSE 挂载抖动通常是暂时的
|
||||
RESTART_BACKOFF = (5, 15, 30, 60, 120, 300)
|
||||
# 单次监控循环存活超过该秒数视为已恢复,重置退避
|
||||
HEALTHY_UPTIME = 60
|
||||
# 超过该秒数监控循环没有任何活动,判定为静默失效
|
||||
STALL_TIMEOUT = 600
|
||||
# 轮询模式目录扫描间隔(毫秒):本地磁盘用 watchfiles 默认值
|
||||
POLL_DELAY_LOCAL_MS = 300
|
||||
# 网络/FUSE 挂载轮询降频,减少监控自身对挂载后端的持续 stat 压力
|
||||
POLL_DELAY_NETWORK_MS = 5000
|
||||
|
||||
def __init__(self, mon_path: Path, callback: Any, force_polling: Optional[bool] = None,
|
||||
poll_delay_ms: Optional[int] = None):
|
||||
"""
|
||||
初始化本地目录监控。
|
||||
:param mon_path: 监控目录
|
||||
:param callback: 目录变化回调对象
|
||||
:param force_polling: 是否强制使用轮询模式,None 表示由 watchfiles 自动选择
|
||||
:param poll_delay_ms: 轮询模式目录扫描间隔(毫秒),仅轮询时生效
|
||||
"""
|
||||
self._watch_path = mon_path
|
||||
self._callback = callback
|
||||
self._force_polling = force_polling
|
||||
self._poll_delay_ms = poll_delay_ms or self.POLL_DELAY_LOCAL_MS
|
||||
self._stop_event = threading.Event()
|
||||
self._thread: Optional[threading.Thread] = None
|
||||
self._watch_filter = DefaultFilter()
|
||||
# 最近一次监控循环活动时间(monotonic),用于检测静默失效
|
||||
self._last_activity: float = 0.0
|
||||
# 累计自动重启次数
|
||||
self._restart_count: int = 0
|
||||
|
||||
@property
|
||||
def watch_path(self) -> Path:
|
||||
"""
|
||||
获取监控目录。
|
||||
:return: 监控目录
|
||||
"""
|
||||
return self._watch_path
|
||||
|
||||
@property
|
||||
def force_polling(self) -> Optional[bool]:
|
||||
"""
|
||||
获取监控模式配置,重建监控线程时沿用。
|
||||
:return: 是否强制轮询
|
||||
"""
|
||||
return self._force_polling
|
||||
|
||||
@property
|
||||
def restart_count(self) -> int:
|
||||
"""
|
||||
获取累计自动重启次数。
|
||||
:return: 自动重启次数
|
||||
"""
|
||||
return self._restart_count
|
||||
|
||||
@property
|
||||
def poll_delay_ms(self) -> int:
|
||||
"""
|
||||
获取轮询模式目录扫描间隔(毫秒),重建监控线程时沿用。
|
||||
:return: 扫描间隔
|
||||
"""
|
||||
return self._poll_delay_ms
|
||||
|
||||
def start(self):
|
||||
"""
|
||||
启动本地目录监控线程。
|
||||
"""
|
||||
if not self._watch_path.exists():
|
||||
raise FileNotFoundError(f"监控目录不存在: {self._watch_path}")
|
||||
if not self._watch_path.is_dir():
|
||||
raise NotADirectoryError(f"监控路径不是目录: {self._watch_path}")
|
||||
if self.is_alive():
|
||||
logger.info(f"本地目录监控已在运行中: {self._watch_path}")
|
||||
return
|
||||
self._stop_event.clear()
|
||||
self._mark_activity()
|
||||
self._thread = threading.Thread(
|
||||
target=self._run,
|
||||
name=f"MoviePilot-DirectoryWatcher-{self._watch_path.name}",
|
||||
daemon=True
|
||||
)
|
||||
self._thread.start()
|
||||
|
||||
def stop(self):
|
||||
"""
|
||||
请求停止本地目录监控线程。
|
||||
"""
|
||||
self._stop_event.set()
|
||||
|
||||
def join(self, timeout: Optional[float] = None):
|
||||
"""
|
||||
等待本地目录监控线程退出。
|
||||
:param timeout: 最长等待秒数
|
||||
"""
|
||||
if self._thread:
|
||||
self._thread.join(timeout=timeout)
|
||||
|
||||
def is_alive(self) -> bool:
|
||||
"""
|
||||
判断监控线程是否仍在运行。
|
||||
:return: 线程存活状态
|
||||
"""
|
||||
return bool(self._thread and self._thread.is_alive())
|
||||
|
||||
def is_stalled(self) -> bool:
|
||||
"""
|
||||
判断监控线程是否已静默失效(线程存活但监控循环长时间无任何活动)。
|
||||
:return: 是否静默失效
|
||||
"""
|
||||
if self._stop_event.is_set() or not self.is_alive():
|
||||
return False
|
||||
if not self._last_activity:
|
||||
return False
|
||||
return (time.monotonic() - self._last_activity) > self.STALL_TIMEOUT
|
||||
|
||||
def _mark_activity(self):
|
||||
"""
|
||||
记录一次监控循环活动时间,作为静默失效检测的心跳。
|
||||
"""
|
||||
self._last_activity = time.monotonic()
|
||||
|
||||
def _run(self):
|
||||
"""
|
||||
运行 watchfiles 主循环,异常时退避重启,避免一次故障导致监控永久停摆。
|
||||
"""
|
||||
# 快速模式失败后降级为轮询,降级后的失败一律走退避重启
|
||||
force_polling = self._force_polling
|
||||
attempt = 0
|
||||
while not self._stop_event.is_set():
|
||||
started_at = time.monotonic()
|
||||
try:
|
||||
self._mark_activity()
|
||||
self._run_watch(force_polling=force_polling)
|
||||
# 正常返回表示收到停止信号
|
||||
return
|
||||
except Exception as err:
|
||||
if self._stop_event.is_set():
|
||||
return
|
||||
# 崩溃堆栈按 ERROR 级输出,生产环境 LOG_LEVEL=ERROR 时也能落盘
|
||||
logger.error(f"本地目录监控异常堆栈: {self._watch_path}\n{traceback.format_exc()}")
|
||||
if force_polling is not True:
|
||||
logger.warn(f"快速模式监控 {self._watch_path} 失败,将自动切换到兼容模式: {err}")
|
||||
force_polling = True
|
||||
continue
|
||||
if time.monotonic() - started_at >= self.HEALTHY_UPTIME:
|
||||
# 上一轮监控已稳定运行过,重新从最短间隔开始退避
|
||||
attempt = 0
|
||||
delay = self.RESTART_BACKOFF[min(attempt, len(self.RESTART_BACKOFF) - 1)]
|
||||
attempt += 1
|
||||
self._restart_count += 1
|
||||
logger.error(f"本地目录监控发生错误,{delay} 秒后自动重启"
|
||||
f"(累计第 {self._restart_count} 次): {self._watch_path} - {err}")
|
||||
if self._stop_event.wait(timeout=delay):
|
||||
return
|
||||
|
||||
def _run_watch(self, force_polling: Optional[bool]):
|
||||
"""
|
||||
执行一次 watchfiles 监控循环。
|
||||
:param force_polling: 是否强制轮询
|
||||
"""
|
||||
for changes in watch(
|
||||
str(self._watch_path),
|
||||
watch_filter=self._watch_filter,
|
||||
stop_event=self._stop_event,
|
||||
rust_timeout=1000,
|
||||
yield_on_timeout=True,
|
||||
force_polling=force_polling,
|
||||
poll_delay_ms=self._poll_delay_ms,
|
||||
recursive=True,
|
||||
ignore_permission_denied=True):
|
||||
self._mark_activity()
|
||||
if self._stop_event.is_set():
|
||||
break
|
||||
if not changes:
|
||||
continue
|
||||
self._handle_changes(changes)
|
||||
self._mark_activity()
|
||||
|
||||
def _handle_changes(self, changes: set[tuple[Change, str]]):
|
||||
"""
|
||||
将 watchfiles 原始变更转换为目录监控事件。
|
||||
:param changes: watchfiles 返回的变更集合
|
||||
"""
|
||||
changes = self._expand_added_directories(changes)
|
||||
for change_type, path_str in sorted(changes, key=lambda item: item[1]):
|
||||
# 批量整理可能持续较久,逐个文件刷新心跳,避免被误判为静默失效
|
||||
self._mark_activity()
|
||||
if change_type not in self._HANDLE_CHANGES:
|
||||
continue
|
||||
event_path = Path(path_str)
|
||||
event = self._build_event(change_type=change_type, event_path=event_path)
|
||||
if not event or event.is_directory:
|
||||
continue
|
||||
file_size = self._get_file_size(event_path)
|
||||
if file_size is None:
|
||||
continue
|
||||
text = self._change_text(change_type)
|
||||
try:
|
||||
self._callback.event_handler(
|
||||
event=event,
|
||||
text=text,
|
||||
event_path=path_str,
|
||||
file_size=file_size
|
||||
)
|
||||
except Exception as err:
|
||||
logger.error(f"处理本地目录监控事件失败: {path_str} - {err}")
|
||||
|
||||
def _expand_added_directories(self, changes: set[tuple[Change, str]]) -> set[tuple[Change, str]]:
|
||||
"""
|
||||
将整体移入监控范围的新增目录展开为内部文件事件。
|
||||
:param changes: watchfiles 返回的变更集合
|
||||
:return: 包含目录内新增文件的变更集合
|
||||
"""
|
||||
expanded_changes = set(changes)
|
||||
for change_type, path_str in changes:
|
||||
if change_type != Change.added:
|
||||
continue
|
||||
event_path = Path(path_str)
|
||||
try:
|
||||
if not event_path.is_dir():
|
||||
continue
|
||||
for nested_path in event_path.rglob("*"):
|
||||
if not nested_path.is_file():
|
||||
continue
|
||||
nested_path_str = nested_path.as_posix()
|
||||
if self._watch_filter(Change.added, nested_path_str):
|
||||
expanded_changes.add((Change.added, nested_path_str))
|
||||
except OSError as err:
|
||||
logger.debug(f"扫描新增目录失败: {event_path} - {err}")
|
||||
return expanded_changes
|
||||
|
||||
@staticmethod
|
||||
def _build_event(change_type: Change, event_path: Path) -> Optional[DirectoryChangeEvent]:
|
||||
"""
|
||||
构建目录变化事件,路径已不存在时忽略。
|
||||
:param change_type: watchfiles 变化类型
|
||||
:param event_path: 变化路径
|
||||
:return: 目录变化事件
|
||||
"""
|
||||
try:
|
||||
is_directory = event_path.is_dir()
|
||||
except OSError as err:
|
||||
logger.debug(f"读取目录监控事件路径失败: {event_path} - {err}")
|
||||
return None
|
||||
if not event_path.exists():
|
||||
return None
|
||||
return DirectoryChangeEvent(
|
||||
change_type=change_type,
|
||||
src_path=event_path.as_posix(),
|
||||
is_directory=is_directory
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _get_file_size(event_path: Path) -> Optional[int]:
|
||||
"""
|
||||
读取事件文件大小,文件已消失时返回 None。
|
||||
:param event_path: 事件文件路径
|
||||
:return: 文件大小
|
||||
"""
|
||||
try:
|
||||
return event_path.stat().st_size
|
||||
except OSError as err:
|
||||
logger.debug(f"读取目录监控文件大小失败: {event_path} - {err}")
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def _change_text(change_type: Change) -> str:
|
||||
"""
|
||||
转换 watchfiles 事件类型为日志文案。
|
||||
:param change_type: watchfiles 变化类型
|
||||
:return: 事件描述
|
||||
"""
|
||||
if change_type == Change.modified:
|
||||
return "修改"
|
||||
return "新增"
|
||||
Reference in New Issue
Block a user