- 将全部路由文件(12个)的 jsonify 响应迁移至 success_response/error_response 统一格式 - 修复 sync_routes.py error_response 参数错误(P0) - 新增异步任务系统:task_registry + task_routes - 新增状态码:TASK_NOT_FOUND(4014)、TASK_CONFLICT(4015)、REINDEX_ERROR(5040) - 修正 task_routes/exam_pkg 中语义不匹配的状态码 - 更新 curl 测试手册、后端对接规范文档 - 添加缓存性能报告和 Redis 迁移计划
342 lines
10 KiB
Python
342 lines
10 KiB
Python
"""
|
||
知识库同步 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/<task_id> 轮询进度。
|
||
|
||
请求体 (可选):
|
||
{
|
||
"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
|
||
)
|