diff --git a/backend/app/api/v1/endpoints/webhooks.py b/backend/app/api/v1/endpoints/webhooks.py index 6fa4f6a..8437c3f 100644 --- a/backend/app/api/v1/endpoints/webhooks.py +++ b/backend/app/api/v1/endpoints/webhooks.py @@ -15,7 +15,7 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.core.config import settings from app.core.database import get_db from app.models.product import Product -from app.models.task import Task +from app.models.task import Task, TaskRecord from app.models.task_log import TaskLog router = APIRouter(prefix="/external/webhooks", tags=["外部回调"]) @@ -113,12 +113,85 @@ async def mom_inbound_webhook( )) changed = True + # ── 动态生成"扫码入库"主线任务节点 + 操作日志(流转树最底部长出入库节点) ── + if await _append_warehouse_task(db, product, "扫码入库", "通过 MOM 系统扫码入库完成"): + changed = True + if changed: await db.commit() return {"ok": True, "matched": True, "serial_number": product.serial_number} +async def _append_warehouse_task( + db: AsyncSession, + product: Product, + task_name: str, + record_remark: str, +) -> bool: + """在流转树主干道最底部追加一个主线任务节点(扫码入库/扫码出库)+ 操作日志。 + + 逻辑: + - 找该产品最后一个主线任务(created_at 最晚且为主线)的 id 作为 parent_task_id, + 保证树状主干连贯; + - 插入 task_type='TRANSFER' 的主线任务(现有主线枚举 → is_main=True,画在中央主干道), + 状态直接 COMPLETED; + - 生成一条 TaskRecord 供前端"查看操作日志"展示。 + + 返回是否新增了节点(供调用方置 changed=True 触发提交)。 + """ + from uuid import uuid4 + from datetime import datetime as _dt + + # 0) 幂等:该产品若已有同名主线任务(扫码入库/扫码出库),则不重复插入 + existing = ( + await db.execute( + select(Task.id).where( + Task.product_id == product.id, + Task.task_name == task_name, + ).limit(1) + ) + ).scalars().first() + if existing: + return False + + # 1) 最后一个主线任务(created_at 最晚,且无父任务 或 task_type 为主线枚举) + last_main_task = ( + await db.execute( + select(Task) + .where( + Task.product_id == product.id, + or_( + Task.parent_task_id.is_(None), + Task.task_type.in_(["TRANSFER", "RECOVERY"]), + ), + ) + .order_by(Task.created_at.desc()) + .limit(1) + ) + ).scalars().first() + + # 2) 插入"扫码入库/扫码出库"主线任务 + new_task = Task( + product_id=product.id, + parent_task_id=last_main_task.id if last_main_task else None, + task_name=task_name, + assignee_id="MOM 仓储系统", + status="COMPLETED", + task_type="TRANSFER", # 现有主线枚举 → is_main=True,画在中央主干道 + completed_at=_dt.now(), + remark=record_remark, + ) + db.add(new_task) + await db.flush() # 生成 new_task.id + + # 3) 生成操作日志(TaskRecord),供前端"查看操作日志"有真实数据 + db.add(TaskRecord( + task_id=new_task.id, + remark=record_remark, + )) + return True + class MomOutboundPayload(BaseModel): """MOM 仓储系统发货出库的回调载荷""" serial_number: str | None = None # 产品 16 位身份证(优先匹配) @@ -211,6 +284,10 @@ async def mom_outbound_webhook( )) changed = True + # ── 动态生成"扫码出库"主线任务节点 + 操作日志(流转树最底部长出出库节点) ── + if await _append_warehouse_task(db, product, "扫码出库", "通过 MOM 系统扫码出库完成"): + changed = True + if changed: await db.commit()