D17 批 1:Human-in-Loop 老师傅经验反馈(数据 + 权限 + 写入 API)
新增 experience_feedback 表(32 表迁移,alembic head b7d1f4a92c3e),
老师傅对系统推荐方案给出"采纳 / 调整 / 拒绝"反馈,按"产品指纹 +
工艺参数"为键跨任务匹配,下次同指纹产品分析自动消费。
变更内容:
- src/moldinsight/models/experience_feedback.py(new)ORM:Base 单点来源、
跨模块裸 FK(user_id / processing_task_id / stp_file_id)、不建 ORM
relationship;fingerprint JSON 列存跨任务匹配键
- src/moldinsight/models/__init__.py 导出 ExperienceFeedback
- migrations/versions/b7d1f4a92c3e_add_experience_feedback.py(new)32 表
迁移;fingerprint 列 PG 下加 GIN 索引(jsonb_path_query 支持)
- src/shared/database/init_db.py 加 3 个权限码(view_experience_feedback /
feedback_experience_hint / manage_experience_feedback)+ 新角色
process_engineer;admin 角色 permissions 同步补齐;init_permissions /
init_roles 改为按 code 比对(新增保留已有 id,避免 FK 引用失效)——
修复既有 DB 启动期漏掉新权限的幂等 bug
- src/moldinsight/services/experience_feedback_service.py(new)service:
compute_fingerprint 分桶(bbox_aspect / volume_bucket / face_bucket /
undercut_class / material_family / is_foam)/ record_feedback(D9 边界:
service.flush + 路由 commit;D17 衰减:同 stp_file_id 整体续期 90 天 TTL,
无 celery beat 依赖)/ list_hints_for_task / resolve_for_process_params
- src/moldinsight/api/experience_feedback_router.py(new)路由:Pydantic
模型写在路由文件内(项目硬规则);POST /api/tasks/{task_id}/experience-feedback
+ GET /api/tasks/{task_id}/experience-hints;归属 TaskQueryService.ensure_task_access
+ User.has_permission 全仓首次调用点
- src/moldinsight/api/__init__.py ROUTE_MODULES 注册新路由
- tests/test_model_ownership.py EXPECTED_TABLES 加 experience_feedback
(31→32)
- tests/test_experience_feedback_fingerprint.py(new)分桶参数化覆盖
bbox / volume / face / undercut / material / is_foam 各边界值
- tests/test_experience_feedback_router.py(new)API 契约 9 例
(401/403/422/200 路径 + 衰减续期 + 任务归属校验 + ORM 注册收口)
- docs/STATUS.md 顶部加 2026-09-23 批 1 日志条目
- docs/TECH_DEBT.md D17 加批 1 已完成描述 + 剩余工作清单
- docs/API_CONTRACT.md §3.2 加 D17 端点表格
测试基线:185 passed, 9 skipped(净增 59 测试)。
Co-Authored-By: Claude Code <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,394 @@
|
||||
"""老师傅经验反馈服务:D17 Human-in-Loop 闭环。
|
||||
|
||||
三层使用:
|
||||
1. `record_feedback` — 路由层 POST 调用;写入 ExperienceFeedback;同 stp_file_id
|
||||
整体续期(D17 衰减机制);只 flush,由路由 commit(D9 边界)。
|
||||
2. `list_hints_for_task` — 路由层 GET 调用;按 stp_file_id + material_family +
|
||||
is_foam 锚定,聚合返回前端 ResultView 用 hints 摘要。
|
||||
3. `resolve_for_process_params` — ProcessingService 调用;返回 OCC worker payload
|
||||
用的 hints dict,OCC 子进程透传给 MultiSchemeMoldPlanner。
|
||||
"""
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Dict, Any, List, Optional
|
||||
|
||||
from sqlalchemy import select, update, or_
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from shared.models.identity import User, Role, UserRole
|
||||
from shared.utils.logger import get_logger
|
||||
from moldinsight.models import (
|
||||
ExperienceFeedback,
|
||||
ProcessingTask,
|
||||
STPFile,
|
||||
GeometryData,
|
||||
MoldCavityData,
|
||||
)
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
FEEDBACK_TTL_DAYS = 90
|
||||
|
||||
|
||||
def compute_fingerprint(
|
||||
geometry_data: Optional[dict],
|
||||
material_name: str,
|
||||
is_foam: bool,
|
||||
) -> Dict[str, str]:
|
||||
"""计算产品指纹(跨任务匹配键)。
|
||||
|
||||
分桶策略见 plan §5.1:
|
||||
- bbox_aspect: cube/compact/slab/elongated/long_bar
|
||||
- volume_bucket: xs/s/m/l/xl(mm³ → cm³)
|
||||
- face_bucket: simple/normal/complex/dense
|
||||
- undercut_class: none/mild/moderate/heavy
|
||||
- material_family: foam/abs/pp/pa/other(粗粒度,避免材料名变化导致匹配失效)
|
||||
- is_foam: "true"/"false"
|
||||
"""
|
||||
geo = geometry_data or {}
|
||||
dims = geo.get("bounding_box", {}).get("dimensions") or [0, 0, 0]
|
||||
sorted_dims = sorted(dims or [0, 0, 0])
|
||||
if sorted_dims[0] <= 0:
|
||||
ratio = 1.0
|
||||
else:
|
||||
ratio = sorted_dims[1] / sorted_dims[0]
|
||||
|
||||
if ratio < 1.0:
|
||||
bbox_aspect = "cube"
|
||||
elif ratio < 1.5:
|
||||
bbox_aspect = "compact"
|
||||
elif ratio < 3.0:
|
||||
bbox_aspect = "slab"
|
||||
elif ratio < 6.0:
|
||||
bbox_aspect = "elongated"
|
||||
else:
|
||||
bbox_aspect = "long_bar"
|
||||
|
||||
volume_cm3 = (geo.get("volume", 0) or 0) / 1000.0
|
||||
if volume_cm3 < 10:
|
||||
volume_bucket = "xs"
|
||||
elif volume_cm3 < 100:
|
||||
volume_bucket = "s"
|
||||
elif volume_cm3 < 500:
|
||||
volume_bucket = "m"
|
||||
elif volume_cm3 < 2000:
|
||||
volume_bucket = "l"
|
||||
else:
|
||||
volume_bucket = "xl"
|
||||
|
||||
face_count = (
|
||||
geo.get("topology_faces", 0)
|
||||
or geo.get("topology", {}).get("faces", 0)
|
||||
or 0
|
||||
)
|
||||
if face_count < 100:
|
||||
face_bucket = "simple"
|
||||
elif face_count < 500:
|
||||
face_bucket = "normal"
|
||||
elif face_count < 2000:
|
||||
face_bucket = "complex"
|
||||
else:
|
||||
face_bucket = "dense"
|
||||
|
||||
undercut_count = geo.get("undercut_count", 0) or 0
|
||||
if undercut_count == 0:
|
||||
undercut_class = "none"
|
||||
elif undercut_count < 4:
|
||||
undercut_class = "mild"
|
||||
elif undercut_count < 9:
|
||||
undercut_class = "moderate"
|
||||
else:
|
||||
undercut_class = "heavy"
|
||||
|
||||
mat_lower = (material_name or "").lower()
|
||||
if "al" in mat_lower and "si" in mat_lower:
|
||||
material_family = "foam"
|
||||
elif "abs" in mat_lower:
|
||||
material_family = "abs"
|
||||
elif "pp" in mat_lower:
|
||||
material_family = "pp"
|
||||
elif "pa" in mat_lower:
|
||||
material_family = "pa"
|
||||
else:
|
||||
material_family = "other"
|
||||
|
||||
return {
|
||||
"bbox_aspect": bbox_aspect,
|
||||
"volume_bucket": volume_bucket,
|
||||
"face_bucket": face_bucket,
|
||||
"undercut_class": undercut_class,
|
||||
"material_family": material_family,
|
||||
"is_foam": "true" if is_foam else "false",
|
||||
}
|
||||
|
||||
|
||||
class ExperienceFeedbackService:
|
||||
"""老师傅经验反馈:写入 / 同指纹 hints 摘要 / OCC worker 用 hints 解析。"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
# 显式无依赖:与 processing_service / task_storage_service 范式一致
|
||||
pass
|
||||
|
||||
async def record_feedback(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
*,
|
||||
task_id: str,
|
||||
scheme_id: str,
|
||||
feedback_status: str,
|
||||
feedback_reason: Optional[str],
|
||||
adjust_suggestion: Optional[str],
|
||||
user: User,
|
||||
confidence_at_submit: Optional[float] = None,
|
||||
score_at_submit: Optional[float] = None,
|
||||
process_params_snapshot: Optional[Dict[str, Any]] = None,
|
||||
) -> ExperienceFeedback:
|
||||
"""写入方案级反馈。仅 flush,由路由 commit(D9 边界)。
|
||||
|
||||
同时同 stp_file_id 整体续期(expires_at = now() + 90d),
|
||||
"写入即消费"的语义保证新反馈立刻进入有效集。
|
||||
"""
|
||||
if feedback_status not in ("adopted", "adjust", "rejected"):
|
||||
raise ValueError(f"feedback_status 非法: {feedback_status}")
|
||||
|
||||
# 1. 找 task + stp_file
|
||||
row = await session.execute(
|
||||
select(ProcessingTask, STPFile)
|
||||
.join(STPFile, ProcessingTask.stp_file_id == STPFile.id)
|
||||
.where(ProcessingTask.task_id == task_id)
|
||||
)
|
||||
row = row.first()
|
||||
if not row:
|
||||
raise ValueError(f"任务不存在: {task_id}")
|
||||
processing_task, stp_file = row
|
||||
|
||||
# 2. 解析 material_name + is_foam_material
|
||||
params = processing_task.parameters or {}
|
||||
material_name = str(params.get("material") or "ABS")
|
||||
is_foam_material = bool(params.get("is_foam_material", False))
|
||||
|
||||
# 3. 取 geometry_data / cavity_data 计算 fingerprint
|
||||
geometry_summary: Dict[str, Any] = {}
|
||||
geo_row = await session.execute(
|
||||
select(GeometryData).where(GeometryData.stp_file_id == stp_file.id)
|
||||
)
|
||||
geo = geo_row.scalar_one_or_none()
|
||||
if geo is not None:
|
||||
geometry_summary = {
|
||||
"volume": geo.volume,
|
||||
"bounding_box": {
|
||||
"min": geo.bounding_box_min,
|
||||
"max": geo.bounding_box_max,
|
||||
},
|
||||
"topology_faces": geo.topology_faces,
|
||||
}
|
||||
|
||||
cavity_row = await session.execute(
|
||||
select(MoldCavityData).where(MoldCavityData.stp_file_id == stp_file.id)
|
||||
)
|
||||
cavity = cavity_row.scalar_one_or_none()
|
||||
if cavity is not None:
|
||||
ki = cavity.cavity_key_info or {}
|
||||
side_actions = (
|
||||
ki.get("quality_considerations") if isinstance(ki, dict) else None
|
||||
) or {}
|
||||
if isinstance(side_actions, dict):
|
||||
geometry_summary["undercut_count"] = side_actions.get("undercut_count", 0) or 0
|
||||
|
||||
fingerprint = compute_fingerprint(geometry_summary, material_name, is_foam_material)
|
||||
|
||||
# 4. 找 scheme 的 axis + method(从 cavity_key_info.candidate_schemes)
|
||||
scheme_axis = "Z"
|
||||
scheme_method: Optional[str] = None
|
||||
if cavity is not None:
|
||||
ki = cavity.cavity_key_info or {}
|
||||
candidate_schemes = ki.get("candidate_schemes") if isinstance(ki, dict) else None
|
||||
if isinstance(candidate_schemes, list):
|
||||
for cs in candidate_schemes:
|
||||
if isinstance(cs, dict) and cs.get("scheme_id") == scheme_id:
|
||||
scheme_axis = (
|
||||
cs.get("axis")
|
||||
or cs.get("parting_axis")
|
||||
or "Z"
|
||||
)
|
||||
scheme_method = cs.get("method") or cs.get("scheme_method")
|
||||
break
|
||||
|
||||
# 5. 取用户角色(显式 JOIN 查询,避免 user.roles 在跨 session 下 lazy load 失败)
|
||||
role_codes = await self._fetch_user_role_codes(session, user.id)
|
||||
if getattr(user, "is_superuser", False):
|
||||
role_code = "admin"
|
||||
elif role_codes:
|
||||
role_code = role_codes[0]
|
||||
else:
|
||||
role_code = "user"
|
||||
|
||||
# 6. 写 ExperienceFeedback
|
||||
now = datetime.now(timezone.utc)
|
||||
new_ttl = now + timedelta(days=FEEDBACK_TTL_DAYS)
|
||||
feedback = ExperienceFeedback(
|
||||
processing_task_id=processing_task.id,
|
||||
stp_file_id=stp_file.id,
|
||||
scheme_id=scheme_id,
|
||||
scheme_axis=str(scheme_axis)[:1],
|
||||
scheme_method=scheme_method,
|
||||
feedback_status=feedback_status,
|
||||
feedback_reason=feedback_reason,
|
||||
adjust_suggestion=adjust_suggestion,
|
||||
process_params_snapshot=process_params_snapshot,
|
||||
fingerprint=fingerprint,
|
||||
confidence_at_submit=confidence_at_submit,
|
||||
score_at_submit=score_at_submit,
|
||||
user_id=user.id,
|
||||
role_code=role_code,
|
||||
expires_at=new_ttl,
|
||||
)
|
||||
session.add(feedback)
|
||||
await session.flush()
|
||||
|
||||
# 7. 同 stp_file_id 整体续期(D17 衰减:仅刷新过期 / NULL 行)
|
||||
await session.execute(
|
||||
update(ExperienceFeedback)
|
||||
.where(
|
||||
ExperienceFeedback.stp_file_id == stp_file.id,
|
||||
or_(
|
||||
ExperienceFeedback.expires_at.is_(None),
|
||||
ExperienceFeedback.expires_at < now,
|
||||
),
|
||||
)
|
||||
.values(expires_at=new_ttl)
|
||||
)
|
||||
await session.flush()
|
||||
return feedback
|
||||
|
||||
async def list_hints_for_task(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
*,
|
||||
stp_file_id: int,
|
||||
material_name: str,
|
||||
is_foam: bool,
|
||||
limit: int = 10,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""同指纹历史采纳摘要,给前端 ResultView 用。
|
||||
|
||||
排除 expires_at < now() 的过期反馈;按 material_family + is_foam 锚定;
|
||||
按 scheme_axis 聚合(adopted/rejected/adjust 计数 + 加权 confidence)。
|
||||
"""
|
||||
now = datetime.now(timezone.utc)
|
||||
is_foam_str = "true" if is_foam else "false"
|
||||
mat_lower = (material_name or "").lower()
|
||||
if "al" in mat_lower and "si" in mat_lower:
|
||||
material_family = "foam"
|
||||
elif "abs" in mat_lower:
|
||||
material_family = "abs"
|
||||
elif "pp" in mat_lower:
|
||||
material_family = "pp"
|
||||
elif "pa" in mat_lower:
|
||||
material_family = "pa"
|
||||
else:
|
||||
material_family = "other"
|
||||
|
||||
rows = await session.execute(
|
||||
select(ExperienceFeedback)
|
||||
.where(
|
||||
ExperienceFeedback.stp_file_id == stp_file_id,
|
||||
or_(
|
||||
ExperienceFeedback.expires_at.is_(None),
|
||||
ExperienceFeedback.expires_at > now,
|
||||
),
|
||||
)
|
||||
.order_by(ExperienceFeedback.created_at.desc())
|
||||
.limit(limit * 4)
|
||||
)
|
||||
feedbacks = rows.scalars().all()
|
||||
|
||||
axis_summary: Dict[str, Dict[str, Any]] = {}
|
||||
for fb in feedbacks:
|
||||
fp = fb.fingerprint or {}
|
||||
# 锚定:material_family + is_foam 必须一致
|
||||
if fp.get("material_family") != material_family:
|
||||
continue
|
||||
if fp.get("is_foam") != is_foam_str:
|
||||
continue
|
||||
axis = fb.scheme_axis or "Z"
|
||||
summary = axis_summary.setdefault(axis, {
|
||||
"scheme_axis": axis,
|
||||
"adopted_count": 0,
|
||||
"rejected_count": 0,
|
||||
"adjust_count": 0,
|
||||
"sample_count": 0,
|
||||
})
|
||||
summary["sample_count"] += 1
|
||||
if fb.feedback_status == "adopted":
|
||||
summary["adopted_count"] += 1
|
||||
elif fb.feedback_status == "rejected":
|
||||
summary["rejected_count"] += 1
|
||||
elif fb.feedback_status == "adjust":
|
||||
summary["adjust_count"] += 1
|
||||
|
||||
result: List[Dict[str, Any]] = []
|
||||
for axis, s in axis_summary.items():
|
||||
total = s["adopted_count"] + s["rejected_count"] + s["adjust_count"]
|
||||
if total == 0:
|
||||
continue
|
||||
confidence = (s["adopted_count"] - s["rejected_count"]) / max(total, 1)
|
||||
confidence = max(-1.0, min(1.0, confidence))
|
||||
weight = max(0.0, confidence) # weight 仅正向上有效(不"扣分"老算法)
|
||||
s["confidence"] = round(confidence, 3)
|
||||
s["weight"] = round(weight, 3)
|
||||
result.append(s)
|
||||
|
||||
result.sort(key=lambda x: (-x["weight"], -x["sample_count"]))
|
||||
return result[:limit]
|
||||
|
||||
@staticmethod
|
||||
async def _fetch_user_role_codes(session: AsyncSession, user_id: int) -> List[str]:
|
||||
"""显式 JOIN 拿用户角色 codes,避免 user.roles 在跨 session 下 detached lazy load 失败。
|
||||
|
||||
测试场景下 user 是从一个 session 取出传到另一个 session,访问 user.roles 会触发
|
||||
DetachedInstanceError;生产场景下也以显式查询更稳(不依赖 ORM relationship 配置)。
|
||||
"""
|
||||
rows = await session.execute(
|
||||
select(Role.code)
|
||||
.join(UserRole, UserRole.role_id == Role.id)
|
||||
.where(UserRole.user_id == user_id)
|
||||
)
|
||||
return [row[0] for row in rows.all()]
|
||||
|
||||
async def resolve_for_process_params(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
*,
|
||||
task_id: str,
|
||||
process_params: Dict[str, Any],
|
||||
bucket_hint: Optional[Dict[str, str]] = None,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""返回 OCC worker payload 用的 hints。
|
||||
|
||||
按 fingerprint bucket 找同指纹最近 N 条采纳,每条
|
||||
{scheme_axis, weight, sample_count, summary},传给 PartingSchemeScorer
|
||||
加成和 PartingCandidateGenerator 优先级加成。
|
||||
"""
|
||||
row = await session.execute(
|
||||
select(ProcessingTask, STPFile)
|
||||
.join(STPFile, ProcessingTask.stp_file_id == STPFile.id)
|
||||
.where(ProcessingTask.task_id == task_id)
|
||||
)
|
||||
row = row.first()
|
||||
if not row:
|
||||
return []
|
||||
processing_task, stp_file = row
|
||||
params = processing_task.parameters or {}
|
||||
material_name = str(params.get("material") or "ABS")
|
||||
is_foam = bool(params.get("is_foam_material", False))
|
||||
|
||||
return await self.list_hints_for_task(
|
||||
session,
|
||||
stp_file_id=stp_file.id,
|
||||
material_name=material_name,
|
||||
is_foam=is_foam,
|
||||
)
|
||||
|
||||
|
||||
# 模块级单例(与 processing_service / task_storage_service 范式一致)
|
||||
experience_feedback_service = ExperienceFeedbackService()
|
||||
Reference in New Issue
Block a user