"""三通道共用的对象目录、canonical 样例装载与适配器(任务书 P1-A / P1-F)。 - 样例来源:doc/db/mdp/mock_api/samples/*.json(复用不迁移);无样例对象用 BUILTIN_SAMPLES。 - 所有业务键统一附加运行编号前缀 SIM-{runId}-,不与真实业务键碰撞。 - adapt_inbound_rows 按 MdpInboundFieldMapper.IdentityFields 的必填/可选字段做别名映射, 缺值字段进 missing_required,由调用方在预览中显示,不得静默造假业务数据。 """ from __future__ import annotations import json import re import uuid from datetime import datetime, timezone, timedelta from pathlib import Path SIM_ROOT = Path(__file__).resolve().parents[1] REPO_ROOT = SIM_ROOT.parents[1] MOCK_API_DIR = REPO_ROOT / "doc" / "db" / "mdp" / "mock_api" CATALOG_DIR = SIM_ROOT / "catalog" RUN_ID_RE = re.compile(r"^\d{14}-[0-9a-f]{6}$") BIZ_KEY_PREFIX = "SIM-" # 通用字段别名(小写比较):把既有样例字段归一为 sourceUpdatedAt _UPDATED_AT_ALIASES = ("updateTime", "sourceUpdatedAt", "inspectStartTime", "WorkDate", "RctDate", "Date", "sxrq", "effDate") def new_run_id(now: datetime | None = None) -> str: """运行编号 SIM-{yyyyMMddHHmmss}-{shortId} 中的时间戳与短 id 部分。""" now = now or datetime.now(timezone(timedelta(hours=8))) return f"{now.strftime('%Y%m%d%H%M%S')}-{uuid.uuid4().hex[:6]}" def biz_prefix(run_id: str) -> str: return f"{BIZ_KEY_PREFIX}{run_id}-" def now_iso() -> str: return datetime.now(timezone(timedelta(hours=8))).strftime("%Y-%m-%dT%H:%M:%S+08:00") # --------------------------------------------------------------------------- # 对象目录 # --------------------------------------------------------------------------- _CATALOG_CACHE: dict | None = None def load_catalog() -> dict: global _CATALOG_CACHE if _CATALOG_CACHE is None: with (CATALOG_DIR / "entities.json").open("r", encoding="utf-8") as f: _CATALOG_CACHE = json.load(f) return _CATALOG_CACHE def load_scenarios() -> dict: with (CATALOG_DIR / "scenarios.json").open("r", encoding="utf-8") as f: return json.load(f) def find_object(object_code: str) -> dict | None: code = (object_code or "").strip().upper() for obj in load_catalog()["objects"]: if obj["objectCode"] == code: return obj return None # --------------------------------------------------------------------------- # canonical 样例 # --------------------------------------------------------------------------- BUILTIN_SAMPLES: dict[str, list[dict]] = { "CUSTOMER": [ {"bizKey": "CUST-001", "Cust": "CUST-001", "Name": "模拟客户001", "Domain": "SIM", "Status": "1"} ], "LOCATION": [ {"bizKey": "LOC-A01", "Location": "SIM-A-01", "Descr": "模拟库位A01", "Domain": "SIM", "Typed": "WH"} ], "EMPLOYEE_HEADCOUNT": [ {"bizKey": "EMP-0001", "Domain": "SIM", "Employee": "EMP-0001", "Name": "模拟员工001", "Department": "生产部", "Position": "操作工", "EmploymentStatus": "在岗"} ], "INVENTORY_OPENING_BALANCE": [ {"bizKey": "OPEN-INV-001", "idid": "SIM-OPEN-INV-001", "code": "RM-CT-001", "sl": 100, "slzx": 100, "shyn": "mock", "shtime": "2026-09-01T00:00:00", "date0": "2026-09-01"} ], "FINISHED_ONHAND": [ {"bizKey": "FG-ONHAND-001", "code": "FG-CT-001", "sl": 50, "date0": "2026-09-30", "location": "FG-01"} ], "FINISHED_OPENING_BALANCE": [ {"bizKey": "FG-OPEN-001", "code": "FG-CT-001", "sl": 20, "date0": "2026-09-01", "location": "FG-01"} ], } def load_canonical_rows(obj: dict) -> list[dict]: """读取对象 canonical 样例;sample 为 null 时使用内置最小样例。""" sample = obj.get("sample") if not sample: rows = BUILTIN_SAMPLES.get(obj["objectCode"]) if rows is None: raise FileNotFoundError(f"object {obj['objectCode']} has no sample and no builtin rows") return [dict(r) for r in rows] path = REPO_ROOT / sample if not path.exists(): raise FileNotFoundError(f"sample not found: {sample}") with path.open("r", encoding="utf-8") as f: data = json.load(f) if not isinstance(data, list): raise ValueError(f"sample must be a JSON array: {sample}") return [dict(r) for r in data] def _extract_updated_at(row: dict) -> str: for alias in _UPDATED_AT_ALIASES: for key in row: if key.lower() == alias.lower() and row[key]: return str(row[key]) return now_iso() def _prefix_biz_key(row: dict, run_id: str) -> str: """给业务键加 SIM-{runId}- 前缀,返回新 bizKey。兼容 bizKey / Id / id 三种载体。""" prefix = biz_prefix(run_id) for key in ("bizKey", "Id", "id"): if key in row and row[key] is not None: raw = str(row[key]) row[key] = raw if raw.startswith(prefix) else prefix + raw return str(row[key]) row["bizKey"] = prefix + uuid.uuid4().hex[:8] return row["bizKey"] # --------------------------------------------------------------------------- # 通道适配器 # --------------------------------------------------------------------------- def adapt_db_rows(obj: dict, run_id: str, rows: list[dict] | None = None) -> list[dict]: """DB_SYNC:源表行 = canonical 行 + SIM 前缀业务键 + sourceUpdatedAt 增量列。""" rows = [dict(r) for r in (rows if rows is not None else load_canonical_rows(obj))] out = [] for row in rows: _prefix_biz_key(row, run_id) row.setdefault("Domain", "SIM") row["sourceUpdatedAt"] = _extract_updated_at(row) out.append(row) return out def adapt_api_rows(obj: dict, run_id: str, rows: list[dict] | None = None) -> list[dict]: """API_PULL:与 DB_SYNC 同语义,另保证 bizKey 存在且按 bizKey 排序(cursor 语义)。""" rows = adapt_db_rows(obj, run_id, rows) for row in rows: row["bizKey"] = str(row.get("bizKey") or row.get("Id") or row.get("id")) rows.sort(key=lambda r: str(r["bizKey"])) return rows # entityCode -> 契约字段(与 MdpInboundFieldMapper.IdentityFields 一致);("字段", 必填, 别名列表) # 版本语义:C# 端只有显式 v2 才使用 IdentityFields(中立必填集合);v1 特化仅存在于 # S5_INVENTORY_TXN 与 MDM_EMPLOYEE_HEADCOUNT。第一阶段目录统一登记 contractVersion=v2, # 因此本表按 IdentityFields 生成样例,避免“按 v1 造数、按 v2 校验”的口径错配。 INBOUND_SPECS: dict[str, list[tuple[str, bool, list[str]]]] = { "MDM_ITEM": [("ItemNum", True, ["itemCode", "item_number", "number", "ItemNum", "wlbm"]), ("Descr", True, ["itemName", "descr", "wlmc"]), ("Drawing", False, []), ("UM", False, ["uom"]), ("ItemType", False, []), ("Status", False, ["status"]), ("Domain", False, ["Domain"])], "MDM_CUSTOMER": [("Cust", True, ["customerNo", "Cust"]), ("Name", True, ["customerName", "Name"]), ("Domain", False, ["Domain"]), ("Status", False, ["status"])], "MDM_SUPPLIER": [("Supp", True, ["supplierCode", "Supp", "supplier_number"]), ("Name", True, ["supplierName", "Name"]), ("Domain", False, ["Domain"]), ("Status", False, ["status"])], "MDM_LOCATION": [("Location", True, ["Location", "location"]), ("Descr", True, ["Descr", "descr", "Address"]), ("Domain", False, ["Domain"]), ("Typed", False, [])], "MDM_SOURCE_LIST": [("Supp", True, ["supplier_number", "Supp", "supplierCode"]), ("ItemNum", True, ["number", "itemCode", "ItemNum"]), ("Domain", False, ["Domain"]), ("Status", False, ["status"])], "MDM_EMPLOYEE_HEADCOUNT": [("Domain", True, ["Domain"]), ("Employee", True, ["Employee"]), ("Position", True, ["Position"]), ("EmploymentStatus", True, ["EmploymentStatus"]), ("Name", False, ["Name"]), ("Department", False, ["Department"])], "S1_SALES_ORDER_ENTRY": [("bill_no", True, ["bill_no", "orderNo"]), ("entry_seq", True, ["entry_seq", "Line"]), ("seorder_id", True, ["seorder_id", "id", "Id"]), ("qty", True, ["qty", "Qty"]), ("item_number", False, ["item_code", "itemCode", "item_number", "ItemNum"]), ("item_name", False, ["itemName"]), ("plan_date", False, ["plan_date"]), ("date", False, ["date"]), ("progress", False, []), ("urgent", False, [])], "S1_REQUIREMENT_EXAMINE_RESULT": [("Id", True, ["Id", "id"]), ("bill_no", True, ["bill_no"]), ("morder_no", True, ["morder_no"]), ("sentry_id", False, ["sentry_id"]), ("create_time", False, ["create_time"]), ("IsDeleted", False, [])], "S1_REQUIREMENT_EXAMINE_DETAIL": [("Id", True, ["Id", "id"]), ("examine_id", True, ["examine_id"]), ("item_number", True, ["item_number"]), ("lack_qty", True, ["lack_qty", "qty"]), ("level", True, ["level"]), ("num", False, ["num"]), ("item_name", False, ["item_name"]), ("needCount", False, []), ("qty", False, ["qty"]), ("is_use", False, [])], "S2_WORK_ORDER_SCHEDULE": [("WorkOrder", True, ["WorkOrd", "work_order", "workOrderNo", "OrderNo"]), ("ItemCode", True, ["ItemNum", "itemCode", "item_number", "componentItem"]), ("QtyOrdered", True, ["Qty", "qty", "QtyReq", "sl"]), ("DocType", False, []), ("DueDate", False, ["WorkDate", "DueDate"]), ("QtyCompleted", False, []), ("Status", False, ["Status", "status", "ztid"]), ("ProdLine", False, ["ProdLine"])], "S3_PURCHASE_ORDER": [("PoNo", True, ["poNo", "PurOrd", "OrdNbr"]), ("PoLine", True, ["poLine", "OrdLine", "Line"]), ("ItemCode", True, ["itemCode", "ItemNum"]), ("OrderQty", True, ["qty", "Qty", "QtyOrded"]), ("SupplierCode", False, ["supplierCode", "Supp"]), ("DueDate", False, []), ("OrderDate", False, []), ("Status", False, ["status", "Status"])], "S3_PURCHASE_RECEIPT": [("Domain", True, ["Domain"]), ("Receiver", True, ["Receiver", "RctNbr"]), ("Line", True, ["Line", "OrdLine"]), ("ItemCode", True, ["ItemNum", "itemCode"]), ("QtyReceived", True, ["QtyReceived", "qty"]), ("PoNo", False, ["OrdNbr", "poNo"]), ("SupplierCode", False, ["Supp", "supplierCode"]), ("ReceiptDate", False, ["RctDate"])], "S4_SHIPMENT": [("ShipmentNo", True, ["shipmentNo", "shddh"]), ("PoNo", True, ["poNo"]), ("ItemCode", True, ["itemCode"]), ("ShipQty", True, ["qty", "Qty"]), ("PoLine", False, ["poLine"]), ("SupplierCode", False, ["supplierCode"]), ("ShipDate", False, ["ShipDate"])], "S4_IQC": [("PoNo", True, ["poNo"]), ("PoLine", True, ["poLine"]), ("ItemCode", True, ["itemCode"]), ("ReceiptQty", True, ["qty", "ReceiptQty"]), ("SupplierCode", False, ["supplierCode"]), ("DefectQty", False, []), ("QcResult", False, ["result", "QcResult"]), ("ReceiptDate", False, ["ReceiptDate"])], "S4_RETURN": [("PoNo", True, ["poNo", "lydjbh"]), ("PoLine", True, ["poLine", "bbh"]), ("ItemCode", True, ["item_code", "wlbm"]), ("ReturnQty", True, ["qty", "ybl"]), ("SupplierCode", False, ["supplierCode"]), ("ReturnReason", False, ["bz"]), ("ReturnStatus", False, ["status"])], "S4_SHORTAGE": [("WorkOrder", True, ["work_order", "WorkOrd"]), ("ItemCode", True, ["component_item_code", "ItemNum"]), ("ShortageQty", True, ["shortage_qty", "lack_qty"]), ("SupplierCode", False, []), ("RiskLevel", False, []), ("NeedDate", False, [])], "S5_WORK_ORDER_BOM": [("Domain", True, ["Domain"]), ("OrderNo", True, ["orderNo", "WorkOrd"]), ("ItemCode", True, ["componentItem", "ItemNum", "itemCode"]), ("QtyRequired", True, ["qtyPer", "QtyReq", "qty"]), ("LineNo", False, ["Line"]), ("Unit", False, ["UM", "uom"])], "S5_INVENTORY_TXN": [("Id", True, ["Id", "id", "bizKey"]), ("BizDocType", True, ["bizDocType", "transType"]), ("ItemCode", False, ["itemNum", "code"]), ("Qty", False, ["qtyChange", "sl"]), ("ApprovedFlag", False, ["approvedFlag"]), ("VoidFlag", False, ["voidFlag"]), ("SummaryFlag", False, ["summaryFlag"]), ("DocDate", False, ["effDate", "jhdate"]), ("TransTime", False, ["addtime"]), ("Domain", False, ["Domain"]), ("TransType", False, ["transType"]), ("Location", False, ["location"]), ("LotSerial", False, ["lotSerial"]), ("CreateUser", False, ["createUser"]), ("EndBalance", False, ["endBalance"]), ("BeginBalance", False, []), ("DocQty", False, []), ("LineClosedFlag", False, [])], "S5_INVENTORY_OPENING_BALANCE": [("Id", True, ["Id", "idid", "bizKey"]), ("idid", False, ["idid"]), ("code", False, ["code", "ItemNum"]), ("sl", False, ["sl"]), ("slzx", False, ["slzx"]), ("shyn", False, ["shyn"]), ("shtime", False, ["shtime"]), ("date0", False, ["date0"])], "S5_INVENTORY_BALANCE_MONTHLY": [("PeriodYm", True, ["periodYm"]), ("AvgBalanceAmount", True, ["avgBalanceAmount"]), ("CategoryCode", False, []), ("WarehouseCode", False, ["warehouseCode"]), ("ItemCode", False, ["itemCode"]), ("IssueCostAmount", False, ["issueCostAmount"])], "S6_WORK_ORDER_LINE": [("Domain", True, ["Domain"]), ("OrderNo", True, ["WorkOrd", "orderNo"]), ("ItemCode", True, ["ItemNum", "itemCode"]), ("QtyPlanned", True, ["QtyReq", "qty", "Qty"]), ("QtyCompleted", False, []), ("PlanFinishDate", False, ["WorkDate"]), ("ReleaseTime", False, []), ("ClosedFlag", False, []), ("VoidFlag", False, [])], "S6_REPORT_TXN": [("Domain", True, ["Domain"]), ("ReportId", True, ["Id", "noid", "ReportId"]), ("WorkOrderNo", True, ["noid", "WorkOrd", "workOrderNo"]), ("ReportDate", True, ["kgdate", "ReportDate"]), ("ReportQty", True, ["sl", "qty", "ReportQty"]), ("StartWorkDate", False, ["kgdate"])], "S7_FQC_TASK_TXN": [("BillNo", True, ["billNo", "BillNo"]), ("ProductionOrderNo", True, ["productionOrderNo", "scph"]), ("MaterialCode", True, ["materialCode", "wlbm"]), ("Qty", True, ["qty"]), ("SalesOrderNo", False, ["salesOrderNo"]), ("Domain", False, ["Domain"])], "S7_SALES_ORDER_LINE": [("Domain", True, ["Domain"]), ("OrderNo", True, ["bill_no", "orderNo", "OrdNbr"]), ("LineNo", True, ["entry_seq", "Line"]), ("ItemCode", True, ["item_code", "itemCode", "ItemNum"]), ("QtyPlanned", True, ["qty", "Qty"]), ("PlanFinishDate", True, ["plan_date", "WorkDate"]), ("QtyCompleted", False, []), ("ReleaseTime", False, []), ("ClosedFlag", False, []), ("VoidFlag", False, [])], "S7_FINISHED_ONHAND": [("Id", True, ["Id", "bizKey"]), ("code", False, ["code", "ItemNum"]), ("sl", False, ["sl", "QtyRec"]), ("date0", False, ["date0", "Date"]), ("location", False, ["location", "LocationTo"])], "S7_FINISHED_OPENING_BALANCE": [("Id", True, ["Id", "bizKey"]), ("code", False, ["code"]), ("sl", False, ["sl"]), ("date0", False, ["date0"]), ("location", False, ["location"])], } def _pick(row_lower: dict[str, object], aliases: list[str]) -> object | None: for alias in aliases: if alias.lower() in row_lower and row_lower[alias.lower()] not in (None, ""): return row_lower[alias.lower()] return None def adapt_inbound_rows(obj: dict, run_id: str, rows: list[dict] | None = None, domain: str = "SIM") -> dict: """API_INBOUND:canonical 行 → 契约字段行。返回 rows / missing_required / biz_keys。""" inbound = obj.get("apiInbound") or {} entity_code = inbound.get("entityCode") if not inbound.get("supported") or not entity_code: raise ValueError(f"object {obj['objectCode']} does not support API_INBOUND") spec = INBOUND_SPECS.get(entity_code) if spec is None: raise ValueError(f"no inbound spec for {entity_code}") rows = [dict(r) for r in (rows if rows is not None else load_canonical_rows(obj))] out: list[dict] = [] missing: set[str] = set() biz_keys: list[str] = [] for row in rows: biz_key = _prefix_biz_key(row, run_id) biz_keys.append(biz_key) row_lower = {str(k).lower(): v for k, v in row.items()} mapped: dict[str, object] = {} for field, required, aliases in spec: value = _pick(row_lower, [field, *aliases]) if value is None and field == "Domain": value = domain if value is None and required: missing.add(field) continue if value is not None: mapped[field] = value # Id 类必填键允许退化为 SIM 业务键 if "Id" in {f for f, r, _ in spec if r} and "Id" not in mapped: mapped["Id"] = biz_key mapped["sourceUpdatedAt"] = _extract_updated_at(row) out.append(mapped) return {"rows": out, "missing_required": sorted(missing), "biz_keys": biz_keys} def inbound_envelope(rows: list[dict], snapshot_id: str | None = None, seq: int | None = None) -> dict: """推数信封,与 push_demo.py / MdpInboundReceiveService 契约一致。""" data: dict = {"list": rows} if snapshot_id: data["snapshotId"] = snapshot_id if seq is not None: data["seq"] = seq return {"data": data}