Files
track-LICA/backend/app/api/v1/endpoints/webhooks.py
duxingchen a68b2bbca0 feat(webhook): 部门校验 —— 只认 company_name == "LICA"
MOM 现在同时对接 IRIS 与 LICA 两个 Track 实例,按载荷里的 company_name 分流。
本实例采取**严格白名单**:
  · company_name == "LICA"(strip 后)→ 正常处理
  · 空白 / 缺失 / "IRIS" / 未知值      → 忽略

⚠️ 与 IRIS 实例的策略**刻意相反**,别"顺手统一成一样":
   IRIS 对空白值要放行 —— MOM 判定不出公司时会回落到指向 IRIS 的扁平配置,
   不收就彻底丢了。
   LICA 没有兜底角色,空白值只可能来自「MOM 没判定出公司」,那本就该由 IRIS 兜。
   **宁可漏,不可误收** —— 误收会把别的部门的设备状态改掉,那是数据污染,
   比漏一条通知严重得多。

实现:
  · MomInboundPayload / MomOutboundPayload 补 company_name 字段
    (不补的话会被 Pydantic 静默丢弃,校验无从谈起)
  · 新增 _belongs_to_this_org(),在**鉴权之后、匹配产品之前**拦截
  · 拦截时返回 200 + matched=False + reason=org_mismatch —— 与「未命中」保持
    同一契约,避免 MOM 侧把它当成故障去重试
  · 顺带补上 outbound 一直在发、但此前被丢弃的 outbound_type 字段

实测:
  · 7 种 company_name 取值全部符合预期
    ("LICA" / "LICA " 通过;"IRIS" / "" / 缺失 / null / "UNKNOWN" 全部忽略)
  · 无 X-API-Key 仍返回 401(鉴权没有被绕过)
  · 真实闭环:LICA 入库回调 → 产品「待仓库收货」→「已入库」
  · 对照:对同一产品发 IRIS 的回调 → 状态纹丝不动 
  · 测试数据已还原为原始值
2026-09-22 13:09:22 +08:00

462 lines
21 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 上方的特权通道说明)。
═══ 部门校验(本实例 = LICA═══
MOM 现在同时对接 IRIS 与 LICA 两个 Track 实例,按载荷里的 `company_name`
分流。本实例采取**严格白名单**:只有 `company_name == "LICA"` 才处理,
空白 / 缺失 / "IRIS" / 未知值一律忽略。
这与 IRIS 实例的策略**刻意相反** —— IRIS 对空白值要放行,因为 MOM 判定不出
公司时会回落到指向 IRIS 的扁平配置,不收就彻底丢了;而 LICA 没有兜底角色,
空白值本就不该由它承担。**宁可漏,不可误收**:误收会把别的部门的设备
状态改掉,那是数据污染,比漏一条通知严重得多。
"""
from __future__ import annotations
from datetime import datetime
from fastapi import APIRouter, Depends, Header, HTTPException
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_modelserial 缺失时的兜底匹配)
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
# ↓ 2026-09 新增MOM 用它区分这条业务属于哪个部门,本实例只认 "LICA"。
# 不加这个字段的话会被 Pydantic 静默丢弃,部门校验就无从谈起。
company_name: str | None = None # IRIS / LICA
def _belongs_to_this_org(company_name: str | None) -> bool:
"""这条回调是不是发给本部门LICA的 —— **严格白名单**。
只有明确写着 "LICA" 才返回 True空白 / 缺失 / "IRIS" / 未知部门名
一律 False忽略
⚠️ 与 IRIS 实例的策略刻意相反,别"顺手统一"
IRIS 对空白值要放行 —— MOM 判定不出公司时会回落到指向 IRIS 的扁平
配置,不收就彻底丢了。
LICA 没有兜底角色,空白值只可能来自"MOM 没判定出公司",那本就该由
IRIS 兜。**宁可漏,不可误收**:误收会把别的部门的设备状态改掉。
"""
return (company_name or "").strip() == settings.ORG_DEPARTMENT
# 「撤回出库」信号词 —— 只在 action / event 里做子串匹配。
# MOM 侧的字段命名尚未冻结故刻意宽松revoke_outbound / outbound.revoked /
# rollback_outbound 都能命中,避免因对方改个词就整条链路失联。
_OUTBOUND_REVOKE_TOKENS = ("revoke", "rollback", "revert", "cancel")
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,
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 把误出库的设备物理回滚到仓库 → 本接口强制执行特权回滚。
- 未命中返回 200MOM 可能操作了非 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")
# ── 部门校验(纵深防御)──
# MOM 侧已按 company_name 分流,这里再挡一道:万一 MOM 路由写错、或有人拿
# 旧配置直接打这个接口,也不会把别的部门的设备状态改掉。
# 返回 200 + matched=False 而非 4xx —— 与「未命中」保持同一契约,
# 避免 MOM 侧把它当成故障去重试。
if not _belongs_to_this_org(payload.company_name):
return {"ok": True, "matched": False, "reason": "org_mismatch"}
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_modelserial 缺失时的兜底匹配)
operator: str | None = None # 出库操作人(写入 task_logs.operator_id
outbound_time: datetime | None = None # 出库时间
# ↓ 2026-09 新增:部门路由键,本实例只认 "LICA"(见 _belongs_to_this_org
company_name: str | None = None # IRIS / LICA
outbound_type: str | None = None # SALES / PRODUCTIONMOM 一直在发,此前被丢弃)
@router.post("/mom-outbound")
async def mom_outbound_webhook(
payload: MomOutboundPayload,
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
- 未命中返回 200MOM 出库的可能是非 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")
# ── 部门校验(纵深防御)—— 同 mom_inbound理由见该处注释
if not _belongs_to_this_org(payload.company_name):
return {"ok": True, "matched": False, "reason": "org_mismatch"}
# ── 按 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}