""" 知识库同步 API 本模块提供文档同步服务的管理接口,包括: - 手动触发同步 - 同步状态查询 - 文件监控控制 路由列表: POST /sync : 手动触发同步 GET /sync/status : 获取同步状态 GET /sync/history : 同步历史记录 GET /sync/changes : 变更日志 POST /sync/start : 启动文件监控 POST /sync/stop : 停止文件监控 架构说明: - 同步服务负责检测文档变更并自动向量化 - 订阅通知功能由后端服务负责 - 权限验证由后端网关完成 同步流程: 1. 扫描文档目录 2. 计算文件哈希,与历史记录对比 3. 检测新增、修改、删除的文件 4. 解析变更文件并更新向量库 Example: curl -X POST http://localhost:5001/sync \\ -H "Authorization: Bearer mock-token-admin" """ from typing import Optional, Tuple, Any from flask import Blueprint, request, current_app import logging logger = logging.getLogger(__name__) from auth.gateway import require_gateway_auth from core.status_codes import SUCCESS, SYNC_SUCCESS, SYNC_ERROR, INTERNAL_ERROR, SERVICE_UNAVAILABLE from api.response_utils import success_response, error_response sync_bp = Blueprint('sync', __name__) def _get_sync_service() -> Optional[Any]: """ 获取同步服务实例 Returns: 同步服务实例,不可用时返回 None """ try: return current_app.config.get('SYNC_SERVICE') except Exception as e: logger.debug(f"获取同步服务失败: {e}") return None def _require_sync_service() -> Tuple[Optional[Any], Optional[Tuple]]: """ 检查同步服务是否可用 Returns: 元组 (sync_service, error_response): - 成功时 error_response 为 None - 失败时 sync_service 为 None """ service = _get_sync_service() if not service: return None, error_response( "SERVICE_UNAVAILABLE", INTERNAL_ERROR, "同步服务未启用", http_status=503 ) return service, None # ==================== 同步 API ==================== @sync_bp.route('/sync', methods=['POST']) @require_gateway_auth def trigger_sync() -> Tuple[Any, int]: """ 手动触发知识库同步(异步任务) 立即返回 task_id,后台线程执行同步。 客户端通过 GET /tasks/ 轮询进度。 请求体 (可选): { "collection": "向量库名称", // 可选,不传则同步所有 "full_sync": false // 是否全量同步 } Returns: 成功: {"success": true, "data": {"task_id": "xxx", "message": "同步任务已启动"}} 失败: {"error": "...", "error_code": "..."} Example: curl -X POST http://localhost:5001/sync \\ -H "Authorization: Bearer mock-token-admin" """ service, err = _require_sync_service() if err: return err from core.task_registry import get_registry registry = get_registry() # 检查是否有正在运行的同步任务 running = registry.list_tasks(status='running', task_type='sync', limit=1) if running: return error_response( "TASK_RUNNING", SYNC_ERROR, f"同步任务正在执行中 (task_id: {running[0].id}),请等待完成", http_status=409 ) # 创建异步任务 task = registry.create_task('sync', '文档同步') def _do_sync(task, sync_service): """后台执行同步""" # 注册进度回调 processed = [0] def on_change(change): processed[0] += 1 registry.update_progress( task.id, current=processed[0], stage='处理文件', message=f"已处理: {change.document_name if hasattr(change, 'document_name') else change.document_id}" ) old_callback = sync_service.on_change_callback sync_service.on_change_callback = on_change try: registry.update_progress(task.id, stage='扫描文档', message='正在检测变更...') result = sync_service.sync_now() result_dict = result.to_dict() if hasattr(result, 'to_dict') else { 'documents_processed': result.documents_processed, 'documents_added': result.documents_added, 'documents_modified': result.documents_modified, 'documents_deleted': result.documents_deleted, 'errors': result.errors, } return result_dict finally: sync_service.on_change_callback = old_callback registry.start_task(task.id, _do_sync, service) return success_response( data={ 'task_id': task.id, 'message': '同步任务已启动,通过 GET /tasks/' + task.id + ' 查询进度' }, status_code=SYNC_SUCCESS, message="同步任务已启动" ) @sync_bp.route('/sync/status', methods=['GET']) @require_gateway_auth def get_sync_status() -> Tuple[Any, int]: """ 获取同步状态 返回同步服务的当前状态,包括: - 是否启用 - 文件监控是否运行 - 最后同步时间 - 跟踪的文档数量 Returns: { "enabled": bool, "monitoring": bool, "last_sync": "ISO 8601", "documents_tracked": N } """ service, err = _require_sync_service() if err: return error_response("SERVICE_UNAVAILABLE", SERVICE_UNAVAILABLE, "同步服务未启用", http_status=503, enabled=False) try: # 获取状态信息 status = { "enabled": True, "monitoring": service.is_running() if hasattr(service, 'is_running') else False, "last_sync": None, "documents_tracked": 0 } # 尝试获取更多状态信息 if hasattr(service, 'get_status'): status.update(service.get_status()) return success_response(data=status) except Exception as e: logger.error(f"状态查询异常: {e}") return error_response("INTERNAL_ERROR", INTERNAL_ERROR, "操作失败", http_status=500, enabled=True) @sync_bp.route('/sync/history', methods=['GET']) @require_gateway_auth def get_sync_history() -> Tuple[Any, int]: """ 获取同步历史 返回最近的同步操作记录。 查询参数: limit: 返回数量限制(默认 20) Returns: {"history": [...]} """ service, err = _require_sync_service() if err: return err limit = request.args.get('limit', 20, type=int) try: history = service.get_sync_history(limit=limit) if hasattr(service, 'get_sync_history') else [] return success_response(data={"history": history}) except Exception as e: logger.error(f"同步操作异常: {e}") return error_response( "SYNC_ERROR", SYNC_ERROR, "操作失败", http_status=500 ) @sync_bp.route('/sync/changes', methods=['GET']) @require_gateway_auth def get_change_logs() -> Tuple[Any, int]: """ 获取变更日志 返回文档变更的详细记录,包括新增、修改、删除等操作。 查询参数: limit: 返回数量限制(默认 50) collection: 过滤指定向量库(可选) Returns: {"changes": [...]} """ service, err = _require_sync_service() if err: return err limit = request.args.get('limit', 50, type=int) collection = request.args.get('collection') try: changes = service.get_change_logs(limit=limit, collection=collection) if hasattr(service, 'get_change_logs') else [] return success_response(data={"changes": changes}) except Exception as e: logger.error(f"同步操作异常: {e}") return error_response( "SYNC_ERROR", SYNC_ERROR, "操作失败", http_status=500 ) @sync_bp.route('/sync/start', methods=['POST']) @require_gateway_auth def start_sync_monitor() -> Tuple[Any, int]: """ 启动文件监控 启用实时文件监控,自动检测文档变更并触发同步。 适用于需要实时更新的场景。 Returns: {"status": "success", "message": "文件监控已启动"} Note: 文件监控会持续运行,直到调用 /sync/stop 或服务重启 """ service, err = _require_sync_service() if err: return err try: if hasattr(service, 'is_running') and service.is_running(): return success_response(status_code=SYNC_SUCCESS, message="文件监控已在运行") if hasattr(service, 'start'): success = service.start() if success: return success_response(status_code=SYNC_SUCCESS, message="文件监控已启动") else: return error_response( "SYNC_ERROR", SYNC_ERROR, "启动文件监控失败", http_status=500 ) else: return success_response(message="文件监控功能不可用") except Exception as e: logger.error(f"同步操作异常: {e}") return error_response( "SYNC_ERROR", SYNC_ERROR, "操作失败", http_status=500 ) @sync_bp.route('/sync/stop', methods=['POST']) @require_gateway_auth def stop_sync_monitor() -> Tuple[Any, int]: """ 停止文件监控 停止实时文件监控服务。已同步的数据保持不变。 Returns: {"status": "success", "message": "文件监控已停止"} """ service, err = _require_sync_service() if err: return err try: if hasattr(service, 'stop'): service.stop() return success_response(status_code=SYNC_SUCCESS, message="文件监控已停止") except Exception as e: logger.error(f"同步操作异常: {e}") return error_response( "SYNC_ERROR", SYNC_ERROR, "操作失败", http_status=500 )