重构路由
This commit is contained in:
@@ -0,0 +1,127 @@
|
||||
# services/task_query_service.py
|
||||
"""任务状态查询服务 — 从 task_router.py 中的持久化任务组装逻辑抽取"""
|
||||
|
||||
from typing import Optional, Dict, Any, List
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from services.storage_integration_rustfs import StorageIntegrationService
|
||||
from services.redis_task_manager import redis_task_manager
|
||||
from models.database import ProcessingTask, STPFile, MeshData
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class TaskQueryService:
|
||||
"""任务状态查询与视图组装"""
|
||||
|
||||
@staticmethod
|
||||
async def get_task_view(db_session: AsyncSession, task_id: str) -> Optional[Dict[str, Any]]:
|
||||
"""
|
||||
获取任务视图 — 优先返回 Redis 缓存,否则从 PostgreSQL + RustFS 组装
|
||||
|
||||
Returns:
|
||||
任务视图字典,如果任务不存在返回 None
|
||||
"""
|
||||
# 1. Redis/内存任务(进行中的任务)
|
||||
task = await redis_task_manager.get_task(task_id)
|
||||
if task:
|
||||
logger.info(f"返回缓存任务状态:{task_id} - {task.get('status')}")
|
||||
return task
|
||||
|
||||
# 2. 持久化任务(已完成/失败,或服务重启后的任务)
|
||||
storage_service = StorageIntegrationService()
|
||||
|
||||
# 查询任务和文件元数据
|
||||
result = await db_session.execute(
|
||||
select(ProcessingTask, STPFile)
|
||||
.join(STPFile, ProcessingTask.stp_file_id == STPFile.id)
|
||||
.where(ProcessingTask.task_id == task_id)
|
||||
)
|
||||
row = result.first()
|
||||
if not row:
|
||||
return None
|
||||
|
||||
processing_task, stp_file = row
|
||||
|
||||
# 从 RustFS 取几何 / 型腔 / 网格详细 JSON
|
||||
try:
|
||||
file_with_data = await storage_service.get_stp_file_with_data(
|
||||
db_session, stp_file_id=stp_file.id
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(f"获取文件数据失败: {e}")
|
||||
file_with_data = {}
|
||||
|
||||
# 解析 geometry_json
|
||||
geometry_json = TaskQueryService._extract_geometry_json(file_with_data)
|
||||
cavity_json: Optional[Dict[str, Any]] = file_with_data.get("mold_cavity_data")
|
||||
features_json: List[Dict[str, Any]] = file_with_data.get("features", [])
|
||||
recommendations_json: List[Dict[str, Any]] = file_with_data.get("recommendations", [])
|
||||
|
||||
# 组装网格摘要
|
||||
mesh_summary = await TaskQueryService._get_mesh_summary(db_session, stp_file.id)
|
||||
|
||||
# 构造与内存任务兼容的任务视图
|
||||
task_view = {
|
||||
"task_id": processing_task.task_id,
|
||||
"status": processing_task.status,
|
||||
"filename": stp_file.original_filename if stp_file else "",
|
||||
"file_path": stp_file.file_path or "",
|
||||
"file_size": stp_file.file_size if stp_file else 0,
|
||||
"upload_time": processing_task.created_time.isoformat()
|
||||
if processing_task.created_time
|
||||
else "",
|
||||
"completed_at": processing_task.completed_time.isoformat()
|
||||
if processing_task.completed_time
|
||||
else "",
|
||||
"geometry_data": geometry_json,
|
||||
"key_info": cavity_json,
|
||||
"cavity_data": cavity_json,
|
||||
"mesh_summary": mesh_summary,
|
||||
"analysis_result": {
|
||||
"geometry_data": geometry_json,
|
||||
"detected_features": features_json,
|
||||
"design_recommendations": recommendations_json,
|
||||
"quality_metrics": {
|
||||
"volume_utilization": file_with_data.get("analysis_metrics", {}).get("volume_utilization", 0),
|
||||
"topology_complexity": file_with_data.get("analysis_metrics", {}).get("topology_complexity", 0),
|
||||
"wall_uniformity": file_with_data.get("analysis_metrics", {}).get("wall_uniformity", 0)
|
||||
},
|
||||
"analysis_summary": file_with_data.get("analysis_metrics", {}).get("analysis_summary", "分析完成")
|
||||
} if geometry_json or features_json or recommendations_json else None,
|
||||
"error": processing_task.error_message or stp_file.error_message or None,
|
||||
}
|
||||
|
||||
logger.info(f"返回持久化任务状态: {task_id} - {processing_task.status}")
|
||||
return task_view
|
||||
|
||||
@staticmethod
|
||||
def _extract_geometry_json(file_with_data: dict) -> Optional[Dict[str, Any]]:
|
||||
"""从 file_with_data 中提取 geometry_json"""
|
||||
if not file_with_data.get("geometry_data"):
|
||||
return None
|
||||
geo_raw = file_with_data["geometry_data"]
|
||||
if isinstance(geo_raw, dict):
|
||||
if "geometry_data" in geo_raw:
|
||||
return geo_raw["geometry_data"]
|
||||
return geo_raw
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
async def _get_mesh_summary(db_session: AsyncSession, stp_file_id: int) -> Optional[Dict[str, Any]]:
|
||||
"""从数据库查询网格摘要"""
|
||||
mesh_record = await db_session.execute(
|
||||
select(MeshData).where(MeshData.stp_file_id == stp_file_id)
|
||||
)
|
||||
mesh_record = mesh_record.scalar_one_or_none()
|
||||
if mesh_record:
|
||||
return {
|
||||
"vertex_count": mesh_record.vertex_count,
|
||||
"face_count": mesh_record.face_count,
|
||||
"point_count": mesh_record.point_count,
|
||||
"quality": mesh_record.quality,
|
||||
}
|
||||
return None
|
||||
Reference in New Issue
Block a user