37c0364e7e
Replace all print()/stderr logging with Python logging module using logger = logging.getLogger(__name__) pattern for consistent log levels and formatting. Extract score conversion to shared score_utils module. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
201 lines
6.7 KiB
Python
201 lines
6.7 KiB
Python
"""
|
|
文件监控服务
|
|
"""
|
|
import os
|
|
import time
|
|
import logging
|
|
import threading
|
|
from pathlib import Path
|
|
from typing import Optional, Callable
|
|
from watchdog.observers import Observer
|
|
from watchdog.events import FileSystemEventHandler, FileCreatedEvent, FileModifiedEvent, FileDeletedEvent
|
|
|
|
from ..core.config import get_settings
|
|
from ..core.database import get_db
|
|
from .knowledge_base_service import KnowledgeBaseService
|
|
|
|
logger = logging.getLogger(__name__)
|
|
settings = get_settings()
|
|
|
|
|
|
class KnowledgeBaseHandler(FileSystemEventHandler):
|
|
"""知识库文件监控处理器"""
|
|
|
|
def __init__(self, knowledge_base_service: KnowledgeBaseService):
|
|
self.kb_service = knowledge_base_service
|
|
self.knowledge_base_dir = Path(settings.knowledge_base_dir)
|
|
self.supported_extensions = settings.allowed_extensions
|
|
|
|
def on_created(self, event):
|
|
"""处理文件创建事件"""
|
|
if not event.is_directory and self._is_supported_file(event.src_path):
|
|
logger.info(f"检测到新文件: {event.src_path}")
|
|
self._process_file_async(event.src_path, "created")
|
|
|
|
def on_modified(self, event):
|
|
"""处理文件修改事件"""
|
|
if not event.is_directory and self._is_supported_file(event.src_path):
|
|
logger.info(f"检测到文件修改: {event.src_path}")
|
|
self._process_file_async(event.src_path, "modified")
|
|
|
|
def on_deleted(self, event):
|
|
"""处理文件删除事件"""
|
|
if not event.is_directory and self._is_supported_file(event.src_path):
|
|
logger.info(f"检测到文件删除: {event.src_path}")
|
|
self._handle_file_deletion(event.src_path)
|
|
|
|
def _is_supported_file(self, file_path: str) -> bool:
|
|
"""检查是否为支持的文件类型"""
|
|
return Path(file_path).suffix.lower() in self.supported_extensions
|
|
|
|
def _process_file_async(self, file_path: str, event_type: str):
|
|
"""异步处理文件"""
|
|
def process():
|
|
try:
|
|
# 等待文件写入完成
|
|
time.sleep(1)
|
|
|
|
result = self.kb_service.process_file(file_path)
|
|
logger.info(f"文件处理结果 ({event_type}): {result}")
|
|
|
|
except Exception as e:
|
|
logger.error(f"处理文件失败: {file_path}, 错误: {str(e)}")
|
|
|
|
# 在后台线程中处理
|
|
thread = threading.Thread(target=process)
|
|
thread.daemon = True
|
|
thread.start()
|
|
|
|
def _handle_file_deletion(self, file_path: str):
|
|
"""处理文件删除"""
|
|
try:
|
|
# 查找并删除对应的数据库记录
|
|
from ..models.document import Document
|
|
from sqlalchemy import and_
|
|
|
|
document = self.kb_service.db.query(Document).filter(
|
|
and_(
|
|
Document.file_path == file_path,
|
|
Document.source_type == "knowledge_base"
|
|
)
|
|
).first()
|
|
|
|
if document:
|
|
result = self.kb_service.delete_document(document.id)
|
|
logger.info(f"文件删除处理结果: {result}")
|
|
else:
|
|
logger.info(f"未找到对应的数据库记录: {file_path}")
|
|
|
|
except Exception as e:
|
|
logger.error(f"处理文件删除失败: {file_path}, 错误: {str(e)}")
|
|
|
|
|
|
class FileWatcherService:
|
|
"""文件监控服务"""
|
|
|
|
def __init__(self):
|
|
self.observer: Optional[Observer] = None
|
|
self.knowledge_base_dir = Path(settings.knowledge_base_dir)
|
|
self.is_running = False
|
|
|
|
def start(self):
|
|
"""启动文件监控"""
|
|
if self.is_running:
|
|
logger.info("文件监控服务已在运行")
|
|
return
|
|
|
|
if not settings.enable_file_watcher:
|
|
logger.info("文件监控服务已禁用")
|
|
return
|
|
|
|
if not self.knowledge_base_dir.exists():
|
|
logger.info(f"知识库目录不存在: {self.knowledge_base_dir}")
|
|
return
|
|
|
|
try:
|
|
# 创建数据库会话
|
|
db = next(get_db())
|
|
kb_service = KnowledgeBaseService(db)
|
|
|
|
# 创建事件处理器
|
|
event_handler = KnowledgeBaseHandler(kb_service)
|
|
|
|
# 创建观察者
|
|
self.observer = Observer()
|
|
self.observer.schedule(
|
|
event_handler,
|
|
str(self.knowledge_base_dir),
|
|
recursive=True # 递归监控子目录
|
|
)
|
|
|
|
# 启动观察者
|
|
self.observer.start()
|
|
self.is_running = True
|
|
|
|
logger.info(f"文件监控服务已启动,监控目录: {self.knowledge_base_dir}")
|
|
|
|
# 确保子目录对应的系统知识库存在
|
|
logger.info("确保系统知识库与目录同步...")
|
|
kb_service.ensure_system_knowledge_bases()
|
|
|
|
# 执行初始扫描
|
|
logger.info("执行初始知识库扫描...")
|
|
scan_result = kb_service.scan_directory()
|
|
logger.info(f"初始扫描结果: {scan_result}")
|
|
|
|
except Exception as e:
|
|
logger.error(f"启动文件监控服务失败: {str(e)}")
|
|
self.is_running = False
|
|
|
|
def stop(self):
|
|
"""停止文件监控"""
|
|
if self.observer and self.is_running:
|
|
self.observer.stop()
|
|
self.observer.join()
|
|
self.is_running = False
|
|
logger.info("文件监控服务已停止")
|
|
|
|
def is_active(self) -> bool:
|
|
"""检查监控服务是否活跃"""
|
|
return self.is_running and self.observer and self.observer.is_alive()
|
|
|
|
def get_status(self) -> dict:
|
|
"""获取监控服务状态"""
|
|
return {
|
|
"is_running": self.is_running,
|
|
"is_active": self.is_active(),
|
|
"knowledge_base_dir": str(self.knowledge_base_dir),
|
|
"directory_exists": self.knowledge_base_dir.exists(),
|
|
"enable_file_watcher": settings.enable_file_watcher
|
|
}
|
|
|
|
|
|
# 全局文件监控服务实例
|
|
file_watcher_service: Optional[FileWatcherService] = None
|
|
|
|
|
|
def get_file_watcher_service() -> FileWatcherService:
|
|
"""获取文件监控服务实例"""
|
|
global file_watcher_service
|
|
if file_watcher_service is None:
|
|
file_watcher_service = FileWatcherService()
|
|
return file_watcher_service
|
|
|
|
|
|
def start_file_watcher():
|
|
"""启动文件监控服务"""
|
|
watcher = get_file_watcher_service()
|
|
watcher.start()
|
|
|
|
|
|
def stop_file_watcher():
|
|
"""停止文件监控服务"""
|
|
watcher = get_file_watcher_service()
|
|
watcher.stop()
|
|
|
|
|
|
|
|
|
|
|
|
|