Files
track/backend/app/api/v1/endpoints/webhooks.py
duxingchen edd43fec29 fix(webhook): MomOutboundPayload 补 outbound_type,与 LICA 实例解析一致
MOM 一直在载荷里发 outbound_type(出库类型:SALES 销售 / USE 领用 /
PRODUCTION 生产),本实例此前没声明该字段,被 Pydantic 静默丢弃。

当前没有任何代码读它,补上不影响行为;目的是让两个实例对同一载荷的解析结果
一致 —— 否则将来谁写了读这个字段的代码,会在 LICA 拿到值、在本实例拿到 None,
而且这种不一致是静默的,不会报错。
2026-09-22 13:46:11 +08:00

498 lines
23 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""外部系统回调 Webhook — Track 作为接收方
MOM 仓储系统确认接收产品入库后,回调本接口,将 Track 中该产品的状态
真正标记为"已入库闭环"(更新宏观状态 + 记录 task_logs 证明仓库已接收)。
同一条入站通道还承担【撤回出库】的强制回滚:MOM 把误点出库的设备物理
回滚到仓库时,Track 必须被动跟随 MOM 的权威物理状态(详见
_mom_inbound_revoke 上方的特权通道说明)。
"""
from __future__ import annotations
from datetime import datetime
from fastapi import APIRouter, Depends, Header, HTTPException, Request
from pydantic import BaseModel
from sqlalchemy import or_, select
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.config import settings
from app.core.database import get_db
from app.core.lifecycle import sync_product_status
from app.models.product import Product
from app.models.task import Task, TaskRecord
from app.models.task_log import TaskLog
router = APIRouter(prefix="/external/webhooks", tags=["外部回调"])
class MomInboundPayload(BaseModel):
"""MOM 仓储系统确认接收入库 / 撤回出库的回调载荷"""
serial_number: str | None = None # 产品 16 位身份证(可空,优先匹配)
sku: str | None = None # 规格型号 spec_model(serial 缺失时的兜底匹配)
operator: str | None = None # 入库操作人(写入 task_logs.operator_id)
inbound_time: datetime | None = None # 入库确认时间
# ↓ MOM 侧一直在发、此前被 Pydantic 静默丢弃的字段。撤回信号靠它们识别。
event: str | None = None # 事件名,如 inbound.created / outbound.revoked
action: str | None = None # 显式动作指令,如 revoke_outbound
source_table: str | None = None # stock_product / stock_semi
company_name: str | None = None # 目标公司(IRIS / LICA),MOM 据此分流到不同 Track 实例
# 「撤回出库」信号词 —— 只在 action / event 里做子串匹配。
# MOM 侧的字段命名尚未冻结,故刻意宽松:revoke_outbound / outbound.revoked /
# rollback_outbound 都能命中,避免因对方改个词就整条链路失联。
_OUTBOUND_REVOKE_TOKENS = ("revoke", "rollback", "revert", "cancel")
# ── 公司归属分流 ──────────────────────────────────────────────────────────
# MOM 现在会在载荷里带 company_name,同一套物理库可能同时向多个 Track 实例
# (IRIS / LICA)回调。本实例服务的是 IRIS,故只放行 IRIS 与空白值。
#
# ⚠️ 判定刻意做成「只排除已知的外来公司」,而非「白名单只认 IRIS」:
# MOM 在无法确定公司归属时会回落到扁平配置,该配置指向本实例 —— 这类
# 消息的 company_name 会是空 / 缺失。若此处按白名单把空白也拒掉,它们
# 就彻底丢了:MOM 那边已收到 200、认为投递成功,不会再重推。
# 同理,未见过的新值(不是 IRIS 也不是 LICA)也一律照常处理。
_FOREIGN_COMPANIES = {"LICA"}
def _is_foreign_company(company_name: str | None) -> bool:
"""载荷是否属于本实例不该处理的其它公司。
返回 True 表示应原样忽略(仍回 200,避免 MOM 反复重推)。
"""
# 大小写 / 首尾空白都容忍:MOM 侧常量书写方式未必冻结,误判的代价是
# 一条消息被错误地当成本公司处理(有唯一匹配约束,最坏是 matched=False)。
return (company_name or "").strip().upper() in _FOREIGN_COMPANIES
def _attribute_audit_to_mom_operator(request: Request, operator: str | None) -> None:
"""把外部回调归因到 MOM 侧的实际操作人。
外部回调走 X-API-Key 鉴权、没有 JWT,所以 JWT 依赖不执行,
审计中间件读到的 request.state.audit_user 永远是空 ——
操作审计里就出现一堆没有归属的「外部系统对接」记录。
但 MOM 载荷里本来就带着实际操作人(operator,即 MOM 侧扫码的那位),
写进 request.state 即可让审计归因到人。
⚠️ 必须在 X-API-Key 校验【之后】调用:密钥不对说明载荷本身就不可信,
此时把 operator 写进审计等于允许伪造人。
"""
who = (operator or "").strip()
if not who:
# 取不到操作人时留一个明确的系统标记,而不是继续显示「未认证」——
# 「MOM系统」至少说明这是一次机器回调,不是"一个匿名的人"。
request.state.audit_user = "MOM系统"
return
request.state.audit_user = who
try:
# 尽力而为:查不到中文名也不影响审计(前端会回退显示账号)
from app.services.mom_cache import get_display_names
request.state.audit_display_name = get_display_names([who]).get(who) or ""
except Exception: # noqa: BLE001 —— 姓名解析失败绝不能影响回调处理
pass
def _is_outbound_revoke(payload: MomInboundPayload) -> bool:
"""payload 是否携带**显式**的撤回出库信号。
注意:返回 False 不代表「不是撤回」——MOM 也可能不加任何标记、直接以
常规 inbound.created 重推。那种隐式信号由调用方用「产品此刻是否处于
已出库」兜底判定(见 mom_inbound_webhook 里的 was_outbound)。
"""
for raw in (payload.action, payload.event):
token = (raw or "").strip().lower()
if token and any(word in token for word in _OUTBOUND_REVOKE_TOKENS):
return True
return False
async def _pick_warehouse_log_task(db: AsyncSession, product: Product) -> Task | None:
"""挑一条挂日志的任务:优先「在库」任务,其次该产品最新任务,都没有则 None。"""
task = (
await db.execute(
select(Task)
.where(Task.product_id == product.id, Task.task_name.ilike("%在库%"))
.order_by(Task.created_at.desc())
.limit(1)
)
).scalars().first()
if task is None:
task = (
await db.execute(
select(Task)
.where(Task.product_id == product.id)
.order_by(Task.created_at.desc())
.limit(1)
)
).scalars().first()
return task
async def _match_inbound_product(
db: AsyncSession, payload: MomInboundPayload, *, allow_outbound: bool,
) -> Product | None:
"""按 serial_number(优先)或 sku 匹配产品。
allow_outbound=False:只认「当前挂在虚拟仓库池」的产品(常规入库的既有语义)。
allow_outbound=True :额外放行「已出库」产品 —— 出库回调会把 current_location_id
置为 None,若仍用原条件,撤回信号必然失配并静默 return matched=False,
造成 MOM 认为货已回库、Track 却永远停在「已出库」的数据脑裂。
"""
location_cond = Product.current_location_id == "virtual_warehouse"
where_cond = (
or_(
location_cond,
Product.overall_status == "已出库",
Product.status == "OUTBOUND",
)
if allow_outbound
else location_cond
)
if payload.serial_number:
return (
await db.execute(
select(Product).where(
or_(
Product.serial_number == payload.serial_number,
Product.external_serial == payload.serial_number,
),
where_cond,
)
)
).scalar_one_or_none()
if payload.sku:
return (
await db.execute(
select(Product)
.where(Product.spec_model == payload.sku, where_cond)
.order_by(Product.created_at.desc())
)
).scalars().first()
return None
@router.post("/mom-inbound")
async def mom_inbound_webhook(
payload: MomInboundPayload,
request: Request,
x_api_key: str | None = Header(default=None, alias="X-API-Key"),
db: AsyncSession = Depends(get_db),
) -> dict:
"""MOM 确认接收入库 / 撤回出库后回调本接口。
- 鉴权:Header X-API-Key 必须等于环境变量 TRACK_WEBHOOK_KEY。
- 常规入库:用 serial_number(优先)或 sku 匹配「当前位于 virtual_warehouse」
的产品,命中则标记"已实收"(overall_status=已入库 + 记录 task_logs)。
- 撤回出库:MOM 把误出库的设备物理回滚到仓库 → 本接口强制执行特权回滚。
- 公司归属:company_name 明确写着其它公司(LICA)时原样忽略;空白 / 缺失
一律照常处理(见 _is_foreign_company 的说明)。
- 未命中返回 200(MOM 可能操作了非 Track 生产的物料,直接忽略)。
"""
# ── 鉴权 ──
if not settings.TRACK_WEBHOOK_KEY or x_api_key != settings.TRACK_WEBHOOK_KEY:
raise HTTPException(status_code=401, detail="Unauthorized: invalid X-API-Key")
# ── 公司归属:不是本实例的消息原样忽略(仍回 200,避免 MOM 当作失败而重推) ──
# ⚠️ 键名 reason / 取值 "ignored_company" 与 LICA 实例(~/track-lica)保持一致:
# MOM 侧不解析它,但排查时两边日志对着看,字段名不一致会白白浪费时间。
if _is_foreign_company(payload.company_name):
return {"ok": True, "matched": False, "reason": "ignored_company"}
# 归因到 MOM 侧实际扫码的人(必须在鉴权通过之后,见函数注释)
_attribute_audit_to_mom_operator(request, payload.operator)
explicit_revoke = _is_outbound_revoke(payload)
# ── 匹配产品 ──
# 常规入库保持严格匹配;撤回(显式标记,或带 serial 可精确定位)才放宽到已出库产品。
# 刻意不给 sku 兜底也无条件放宽:同型号可能有多台,放宽后可能误标到别的设备。
product = await _match_inbound_product(
db, payload, allow_outbound=explicit_revoke or bool(payload.serial_number),
)
# ── 未命中:可能是非 Track 生产的物料,直接忽略 ──
if product is None:
return {"ok": True, "matched": False}
# 隐式撤回:payload 没带任何标记,但产品此刻正处于「已出库」。
# 对一台已发货的设备来说,任何入库回调都只能意味着「货回来了」。
was_outbound = (
(product.overall_status or "").strip() == "已出库"
or (product.status or "").strip().upper() == "OUTBOUND"
)
is_revoke = explicit_revoke or was_outbound
# ══════════════════════════════════════════════════════════════════════
# ★ 特权通道 — MOM 的物理状态同步优先级最高,强制覆写、不受任何内部守卫约束
#
# 与 task_service.py 的【绝对物理终态保护】(PHYSICAL_TERMINAL_OVERALL,
# task_service.py:115-127) 方向刻意相反:那套保护约束的是「车间内部流转
# 不许用工序名抹掉物理终态」;而本接口是物理事实的**权威来源**——MOM 说
# 货已回到仓库,Track 必须无条件跟随。
#
# ⚠️ 后续维护者:不要在此处添加 _is_physical_terminal / 状态互斥 / 仅当
# 状态为 X 才允许覆写 之类的校验。那会让设备永远卡在「已出库」,
# 与 MOM 账面对不上——正是本次要消灭的数据脑裂。
# ══════════════════════════════════════════════════════════════════════
changed = False
# 1) 宏观状态强制覆写为「已入库」(撤回时从「已出库」拉回)
if product.overall_status != "已入库":
product.overall_status = "已入库"
changed = True
# 2) 物理位置强制回滚到虚拟仓库池(出库回调曾把它置为 None)
if product.current_location_id != "virtual_warehouse":
product.current_location_id = "virtual_warehouse"
changed = True
# 3) 双字段同步:lifecycle.py 约定凡改写 overall_status 必调一次。
# (原实现在这里硬编码 product.status="ARCHIVED",绕过了约定,一并纠正)
# ⚠️ 必须把 status 的变化也计入 changed:否则当 overall_status / location
# 本来就已经正确时,这一处纠偏会因为 changed 保持 False 而永远不提交。
prev_status = product.status
sync_product_status(product)
if product.status != prev_status:
changed = True
# ── 记录日志(优先"在库"任务,其次该产品最新任务;无任务则仅更新状态) ──
log_task = await _pick_warehouse_log_task(db, product)
if log_task is not None:
if is_revoke:
signal = payload.action or payload.event or "inbound.created(隐式)"
remark = (
f"MOM 撤回出库 → 强制回滚:宏观状态已入库、"
f"位置已回到 virtual_warehouse(信号: {signal})"
)
action_type = "warehouse_outbound_revoked"
else:
time_str = payload.inbound_time.isoformat() if payload.inbound_time else "—"
remark = f"MOM 仓储系统确认接收入库(inbound_time: {time_str})"
action_type = "warehouse_inbound"
db.add(TaskLog(
task_id=log_task.id,
operator_id=(payload.operator or "virtual_warehouse")[:64],
action_type=action_type,
remark=remark,
))
changed = True
# ── 动态生成主线任务节点 + 操作日志(流转树最底部长出节点) ──
if is_revoke:
# 撤回必须留痕:否则流转树末节点仍是「扫码出库」,而产品徽标已是
# 「已入库」,这种可见的自相矛盾会让车间不敢信这套数据。
#
# ⚠️ 节点名里的「(重新入库)」不是装饰,是 [必须保留] 的契约:
# product_service.py:166-174 的 _has_warehouse_task() 用**子串**判定
# 仓库节点("在库" in task_name or "入库" in task_name)。而
# 「撤回出库」四个字里只有"出库"、不含"入库",会让它判定为"无仓库任务",
# 进而给 location==virtual_warehouse 的产品注入一个假的「已完成 /
# 待仓库扫码」虚拟节点(product_service.py:237-242 的情况 A)——
# 该设备明明已入库且在仓库里,树尾却显示待收货。
# 补上「重新入库」后关键字命中,虚拟节点不再注入。
appended = await _append_warehouse_task(
db, product, "撤回出库(重新入库)", "MOM 撤回出库,设备已物理回滚至仓库",
)
else:
appended = await _append_warehouse_task(
db, product, "扫码入库", "通过 MOM 系统扫码入库完成",
)
if appended:
changed = True
if changed:
await db.commit()
return {
"ok": True,
"matched": True,
"serial_number": product.serial_number,
"revoked": is_revoke,
}
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 显式置 None:不能填非 UUID 字符串,否则前端解析头像/用户信息报错导致节点跳过渲染
assignee_id=None,
status="COMPLETED",
task_type="WAREHOUSE", # 仓储任务类型(is_main 判断已兼容 WAREHOUSE → 画在中央主干道)
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 位身份证(优先匹配)
sku: str | None = None # 规格型号 spec_model(serial 缺失时的兜底匹配)
operator: str | None = None # 出库操作人(写入 task_logs.operator_id)
outbound_time: datetime | None = None # 出库时间
company_name: str | None = None # 目标公司(IRIS / LICA),MOM 据此分流到不同 Track 实例
# ↓ 2026-09 新增:MOM 一直在发、此前被 Pydantic 静默丢弃。当前无人读取,
# 先接住是为了与 LICA 实例(~/track-lica)对同一载荷的解析结果保持一致 ——
# 否则将来谁写了读这个字段的代码,会在 LICA 拿到值、在本实例拿到 None。
outbound_type: str | None = None # SALES / USE / PRODUCTION
@router.post("/mom-outbound")
async def mom_outbound_webhook(
payload: MomOutboundPayload,
request: Request,
x_api_key: str | None = Header(default=None, alias="X-API-Key"),
db: AsyncSession = Depends(get_db),
) -> dict:
"""MOM 仓储系统发货出库后回调本接口,将 Track 产品标记为"已出库"。
- 鉴权:Header X-API-Key 必须等于环境变量 TRACK_WEBHOOK_KEY。
- 用 serial_number(优先)或 sku 匹配"在仓库/已入库"的产品;
命中则标记"已出库"(overall_status=已出库 + status=OUTBOUND + 记录 task_logs)。
- 公司归属:company_name 明确写着其它公司(LICA)时原样忽略;空白 / 缺失
一律照常处理(见 _is_foreign_company 的说明)。
- 未命中返回 200(MOM 出库的可能是非 Track 生产的物料,直接忽略)。
"""
# ── 鉴权 ──
if not settings.TRACK_WEBHOOK_KEY or x_api_key != settings.TRACK_WEBHOOK_KEY:
raise HTTPException(status_code=401, detail="Unauthorized: invalid X-API-Key")
# ── 公司归属:不是本实例的消息原样忽略(仍回 200,避免 MOM 当作失败而重推) ──
# ⚠️ 键名 reason / 取值 "ignored_company" 与 LICA 实例(~/track-lica)保持一致:
# MOM 侧不解析它,但排查时两边日志对着看,字段名不一致会白白浪费时间。
if _is_foreign_company(payload.company_name):
return {"ok": True, "matched": False, "reason": "ignored_company"}
# 归因到 MOM 侧实际出库的人(必须在鉴权通过之后,见函数注释)
_attribute_audit_to_mom_operator(request, payload.operator)
# ── 按 serial_number(优先)或 sku 匹配"在仓库/已入库"的产品 ──
product = None
where_cond = or_(
Product.current_location_id == "virtual_warehouse",
Product.overall_status.in_(["已入库", "在库"]),
)
if payload.serial_number:
product = (
await db.execute(
select(Product).where(
or_(
Product.serial_number == payload.serial_number,
Product.external_serial == payload.serial_number,
),
where_cond,
)
)
).scalars().first()
elif payload.sku:
product = (
await db.execute(
select(Product)
.where(Product.spec_model == payload.sku, where_cond)
.order_by(Product.created_at.desc())
)
).scalars().first()
# ── 未命中:可能出库的是非 Track 生产的物料,直接忽略 ──
if product is None:
return {"ok": True, "matched": False}
# ── 标记"已出库" ──
changed = False
if product.overall_status != "已出库":
product.overall_status = "已出库"
product.status = "OUTBOUND"
changed = True
# 🚚 同步出清厂内位置(与 _recalc_product_location 的「货发走就离场」对齐):
# 本回调的匹配条件之一就是 current_location_id == "virtual_warehouse",
# 即设备此刻还挂在仓库池里。既然 MOM 已确认发货,就不该再显示为厂内仓库/工位。
# 置空后产品列表的「当前位置」显示为「—」。
if product.current_location_id is not None:
product.current_location_id = None
changed = True
# 记录出库日志(优先"在库"任务,其次该产品最新任务)
outbound_task = await _pick_warehouse_log_task(db, product)
if outbound_task is not None:
time_str = payload.outbound_time.isoformat() if payload.outbound_time else "—"
db.add(TaskLog(
task_id=outbound_task.id,
operator_id=(payload.operator or "virtual_warehouse")[:64],
action_type="warehouse_outbound",
remark=f"MOM 仓储系统发货出库(outbound_time: {time_str})",
))
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}