Files
rag/api/sync_routes.py
lacerate551 183a57e7f1 refactor(api): 统一响应格式迁移 + 异步任务系统 + 状态码体系完善
- 将全部路由文件(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 迁移计划
2026-06-05 22:56:00 +08:00

342 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, 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
)