"""外部系统回调 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 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 _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 把误出库的设备物理回滚到仓库 → 本接口强制执行特权回滚。 - 公司归属: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"} 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 实例 @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)。 - 公司归属: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"} # ── 按 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}