feat: Webhook 入库/出库后动态生成扫码任务节点与操作日志
- mom-inbound/mom-outbound 在标记已入库/已出库后,追加'扫码入库/扫码出库'主线任务 - 新任务 parent_task_id 取最后一个主线任务,task_type=TRANSFER 保证画在主干道 - 生成 TaskRecord 操作日志,前端流转树最底部长出节点 - 幂等:已存在同名任务则不重复插入
This commit is contained in:
@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user