196 lines
7.1 KiB
Python
196 lines
7.1 KiB
Python
"""AI 任务管理路由;创建、查询、日志、确认和取消均经过统一队列。"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from typing import Any, Dict, Optional
|
|
from uuid import uuid4
|
|
|
|
from fastapi import APIRouter, Header, HTTPException, Query
|
|
from fastapi.responses import JSONResponse
|
|
|
|
from models.ai_tasks import AITaskCancelRequest, AITaskCreateRequest
|
|
from services.ai_tasks import AITaskError, ai_task_service
|
|
from services.workbench import ControlPlaneStorageUnavailable
|
|
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
def _error(exc: AITaskError) -> HTTPException:
|
|
return HTTPException(
|
|
status_code=exc.status_code,
|
|
detail={
|
|
"code": exc.code,
|
|
"message": exc.message,
|
|
"trace_id": uuid4().hex,
|
|
"retryable": exc.code in {"device_offline", "idempotency_in_progress"},
|
|
"details": exc.details,
|
|
},
|
|
)
|
|
|
|
|
|
def _actor(value: Optional[str]) -> str:
|
|
return str(value or "ai").strip() or "ai"
|
|
|
|
|
|
def _storage_error(exc: Exception, *, data: Optional[dict] = None) -> JSONResponse:
|
|
"""存储没接上时也给出可追踪回执,明确没有向设备下发。"""
|
|
trace_id = uuid4().hex
|
|
return JSONResponse(
|
|
status_code=503,
|
|
content={
|
|
"code": 503,
|
|
"success": False,
|
|
"error_code": "control_plane_storage_unavailable",
|
|
"error_message": str(exc),
|
|
"trace_id": trace_id,
|
|
"channel_used": "control-plane",
|
|
"raw_rpc_receipt": None,
|
|
"readback": {
|
|
"verified": False,
|
|
"write_performed": False,
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"source": "ai_task_storage_guard",
|
|
"reason": "storage_unavailable_before_device_dispatch",
|
|
},
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"data": data or {"items": [], "total": 0, "data_state": "storage_unavailable"},
|
|
},
|
|
)
|
|
|
|
|
|
def _task_read_error(exc: AITaskError, *, task_id: str) -> JSONResponse:
|
|
"""任务只读失败也明确说明没有向手机下发或写入。"""
|
|
trace_id = uuid4().hex
|
|
return JSONResponse(
|
|
status_code=exc.status_code,
|
|
content={
|
|
"code": exc.status_code,
|
|
"success": False,
|
|
"error_code": exc.code,
|
|
"error_message": exc.message,
|
|
"trace_id": trace_id,
|
|
"channel_used": "control-plane",
|
|
"raw_rpc_receipt": None,
|
|
"readback": {
|
|
"verified": True,
|
|
"write_performed": False,
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"source": "control_plane_store",
|
|
"reason": "task_not_found_before_device_dispatch" if exc.code == "resource_not_found" else "task_read_failed_before_device_dispatch",
|
|
},
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"data": {"task_id": task_id, "found": False},
|
|
},
|
|
)
|
|
|
|
|
|
@router.post("/ai/tasks", tags=["AI任务"], responses={403: {"description": "AI任务权限不足"}, 422: {"description": "参数不完整"}, 503: {"description": "AI任务存储未就绪"}})
|
|
async def create_ai_task(
|
|
request: AITaskCreateRequest,
|
|
idempotency_key: Optional[str] = Header(default=None, alias="Idempotency-Key"),
|
|
x_actor: Optional[str] = Header(default=None),
|
|
) -> Dict[str, Any]:
|
|
try:
|
|
payload = request.dict() if hasattr(request, "dict") else request.model_dump()
|
|
return {"code": 200, "data": await ai_task_service.create_task(payload, actor=_actor(x_actor), idempotency_key=idempotency_key or "")}
|
|
except AITaskError as exc:
|
|
raise _error(exc) from exc
|
|
except ControlPlaneStorageUnavailable as exc:
|
|
return _storage_error(exc)
|
|
|
|
|
|
@router.get("/ai/tasks", tags=["AI任务记录"], responses={422: {"description": "查询条件不合法"}, 503: {"description": "AI任务存储未就绪"}})
|
|
async def list_ai_tasks(
|
|
status: Optional[str] = Query(default=None),
|
|
limit: int = Query(default=50, ge=1, le=100),
|
|
offset: int = Query(default=0, ge=0),
|
|
) -> Dict[str, Any]:
|
|
try:
|
|
data = await ai_task_service.list_tasks(status=status, limit=limit, offset=offset)
|
|
return {
|
|
"code": 200,
|
|
"success": True,
|
|
"trace_id": uuid4().hex,
|
|
"channel_used": "control-plane",
|
|
"raw_rpc_receipt": None,
|
|
"readback": {"verified": True, "write_performed": False, "source": "control_plane_store"},
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"data": data,
|
|
}
|
|
except AITaskError as exc:
|
|
raise _error(exc) from exc
|
|
except ControlPlaneStorageUnavailable as exc:
|
|
return _storage_error(exc)
|
|
|
|
|
|
@router.get("/ai/tasks/{task_id}", tags=["AI任务"], responses={503: {"description": "AI任务存储未就绪"}})
|
|
async def get_ai_task(task_id: str) -> Dict[str, Any]:
|
|
try:
|
|
return {
|
|
"code": 200,
|
|
"success": True,
|
|
"trace_id": uuid4().hex,
|
|
"channel_used": "control-plane",
|
|
"raw_rpc_receipt": None,
|
|
"readback": {
|
|
"verified": True,
|
|
"write_performed": False,
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"source": "control_plane_store",
|
|
},
|
|
"dispatch_count": 0,
|
|
"downlink_count": 0,
|
|
"data": await ai_task_service.get_task(task_id),
|
|
}
|
|
except AITaskError as exc:
|
|
return _task_read_error(exc, task_id=task_id)
|
|
except ControlPlaneStorageUnavailable as exc:
|
|
return _storage_error(exc)
|
|
|
|
|
|
@router.get("/ai/tasks/{task_id}/logs", tags=["AI任务日志"])
|
|
async def get_ai_task_logs(task_id: str) -> Dict[str, Any]:
|
|
try:
|
|
return {"code": 200, "data": {"items": await ai_task_service.list_logs(task_id)}}
|
|
except AITaskError as exc:
|
|
raise _error(exc) from exc
|
|
except ControlPlaneStorageUnavailable as exc:
|
|
return _storage_error(exc)
|
|
|
|
|
|
@router.get("/ai/tasks/{task_id}/audit", tags=["AI任务审计"])
|
|
async def get_ai_task_audit(task_id: str) -> Dict[str, Any]:
|
|
try:
|
|
return {"code": 200, "data": {"items": await ai_task_service.list_audit(task_id)}}
|
|
except AITaskError as exc:
|
|
raise _error(exc) from exc
|
|
except ControlPlaneStorageUnavailable as exc:
|
|
return _storage_error(exc)
|
|
|
|
|
|
@router.post("/ai/tasks/{task_id}/confirm", tags=["AI任务"])
|
|
async def confirm_ai_task(task_id: str, x_actor: Optional[str] = Header(default=None)) -> Dict[str, Any]:
|
|
try:
|
|
return {"code": 200, "data": await ai_task_service.confirm_task(task_id, actor=_actor(x_actor))}
|
|
except AITaskError as exc:
|
|
raise _error(exc) from exc
|
|
|
|
|
|
@router.post("/ai/tasks/{task_id}/cancel", tags=["AI任务"])
|
|
async def cancel_ai_task(
|
|
task_id: str,
|
|
request: AITaskCancelRequest,
|
|
x_actor: Optional[str] = Header(default=None),
|
|
) -> Dict[str, Any]:
|
|
try:
|
|
return {"code": 200, "data": await ai_task_service.cancel_task(task_id, actor=_actor(x_actor), reason=request.reason)}
|
|
except AITaskError as exc:
|
|
raise _error(exc) from exc
|