Files
rag/api/sync_routes.py
lacerate551 6d5a4b5b5e feat: 性能插桩、缓存清除端点、LLM top_p 与 JSON 强制输出
- chat_routes.py: 新增 8 阶段性能计时(stages_ms)+ SSE 输出
- sync_routes.py: 新增 POST /cache/clear 缓存清除端点
- engine.py: LLM 调用增加 top_p 参数
- intent_analyzer.py: response_format 加配置开关(INTENT_RESPONSE_FORMAT)
2026-06-21 20:25:12 +08:00

363 lines
10 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
知识库同步 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, jsonify, current_app
import logging
logger = logging.getLogger(__name__)
from auth.gateway import require_gateway_auth
from core.status_codes import SYNC_SUCCESS, SYNC_ERROR, INTERNAL_ERROR
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(
error="SERVICE_UNAVAILABLE",
error_code=INTERNAL_ERROR,
message="同步服务未启用",
http_status=503
)
return service, None
# ==================== 同步 API ====================
@sync_bp.route('/sync', methods=['POST'])
@require_gateway_auth
def trigger_sync() -> Tuple[Any, int]:
"""
手动触发知识库同步
扫描文档目录,检测变更并执行向量化处理。
请求体 (可选):
{
"collection": "向量库名称", // 可选,不传则同步所有
"full_sync": false // 是否全量同步
}
Returns:
成功: {"success": true, "data": {"result": {...}}}
失败: {"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
try:
result = service.sync_now()
return success_response(
data={"result": result.to_dict() if hasattr(result, 'to_dict') else result},
status_code=SYNC_SUCCESS,
message="同步完成"
)
except Exception as e:
logger.error(f"同步操作异常: {e}")
return error_response(
error="SYNC_ERROR",
error_code=SYNC_ERROR,
message="操作失败",
http_status=500
)
@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 jsonify({
"status": "failed",
"status_code": INTERNAL_ERROR,
"enabled": False,
"message": "同步服务未启用"
})
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 jsonify(status)
except Exception as e:
logger.error(f"状态查询异常: {e}")
return jsonify({
"status": "failed",
"status_code": INTERNAL_ERROR,
"enabled": True,
"error": "操作失败"
})
@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 jsonify({"history": history})
except Exception as e:
logger.error(f"同步操作异常: {e}")
return error_response(
error="SYNC_ERROR",
error_code=SYNC_ERROR,
message="操作失败",
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 jsonify({"changes": changes})
except Exception as e:
logger.error(f"同步操作异常: {e}")
return error_response(
error="SYNC_ERROR",
error_code=SYNC_ERROR,
message="操作失败",
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 jsonify({"status": "success", "status_code": SYNC_SUCCESS, "message": "文件监控已在运行"})
if hasattr(service, 'start'):
success = service.start()
if success:
return jsonify({"status": "success", "status_code": SYNC_SUCCESS, "message": "文件监控已启动"})
else:
return error_response(
error="SYNC_ERROR",
error_code=SYNC_ERROR,
message="启动文件监控失败",
http_status=500
)
else:
return jsonify({"status": "success", "message": "文件监控功能不可用"})
except Exception as e:
logger.error(f"同步操作异常: {e}")
return error_response(
error="SYNC_ERROR",
error_code=SYNC_ERROR,
message="操作失败",
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 jsonify({"status": "success", "status_code": SYNC_SUCCESS, "message": "文件监控已停止"})
except Exception as e:
logger.error(f"同步操作异常: {e}")
return error_response(
error="SYNC_ERROR",
error_code=SYNC_ERROR,
message="操作失败",
http_status=500
)
# ==================== 缓存管理 API ====================
@sync_bp.route('/cache/clear', methods=['POST'])
@require_gateway_auth
def clear_all_caches() -> Tuple[Any, int]:
"""
清除所有缓存层Query / Embedding / Rerank / Semantic / Intent 精确缓存)
用于测试和调试,不影响向量库数据。
Returns:
{"success": true, "message": "缓存已清除", "cleared": [...]}
"""
cleared = []
# 1. 精确匹配缓存Query / Embedding / Rerank
try:
from core.cache import get_cache_manager
cache = get_cache_manager()
if cache:
cache.clear_all()
cleared.append("query_cache")
cleared.append("embedding_cache")
cleared.append("rerank_cache")
except Exception as e:
logger.warning(f"清除精确缓存失败: {e}")
# 2. 语义缓存FAISS
try:
from core.semantic_cache import get_semantic_cache
sc = get_semantic_cache()
if sc:
sc.clear()
cleared.append("semantic_cache")
except Exception as e:
logger.warning(f"清除语义缓存失败: {e}")
logger.info(f"缓存清除完成: {cleared}")
return jsonify({
"success": True,
"message": f"已清除 {len(cleared)} 层缓存",
"cleared": cleared
})