sample_adapters.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324
  1. """三通道共用的对象目录、canonical 样例装载与适配器(任务书 P1-A / P1-F)。
  2. - 样例来源:doc/db/mdp/mock_api/samples/*.json(复用不迁移);无样例对象用 BUILTIN_SAMPLES。
  3. - 所有业务键统一附加运行编号前缀 SIM-{runId}-,不与真实业务键碰撞。
  4. - adapt_inbound_rows 按 MdpInboundFieldMapper.IdentityFields 的必填/可选字段做别名映射,
  5. 缺值字段进 missing_required,由调用方在预览中显示,不得静默造假业务数据。
  6. """
  7. from __future__ import annotations
  8. import json
  9. import re
  10. import uuid
  11. from datetime import datetime, timezone, timedelta
  12. from pathlib import Path
  13. SIM_ROOT = Path(__file__).resolve().parents[1]
  14. REPO_ROOT = SIM_ROOT.parents[1]
  15. MOCK_API_DIR = REPO_ROOT / "doc" / "db" / "mdp" / "mock_api"
  16. CATALOG_DIR = SIM_ROOT / "catalog"
  17. RUN_ID_RE = re.compile(r"^\d{14}-[0-9a-f]{6}$")
  18. BIZ_KEY_PREFIX = "SIM-"
  19. # 通用字段别名(小写比较):把既有样例字段归一为 sourceUpdatedAt
  20. _UPDATED_AT_ALIASES = ("updateTime", "sourceUpdatedAt", "inspectStartTime", "WorkDate", "RctDate", "Date", "sxrq", "effDate")
  21. def new_run_id(now: datetime | None = None) -> str:
  22. """运行编号 SIM-{yyyyMMddHHmmss}-{shortId} 中的时间戳与短 id 部分。"""
  23. now = now or datetime.now(timezone(timedelta(hours=8)))
  24. return f"{now.strftime('%Y%m%d%H%M%S')}-{uuid.uuid4().hex[:6]}"
  25. def biz_prefix(run_id: str) -> str:
  26. return f"{BIZ_KEY_PREFIX}{run_id}-"
  27. def now_iso() -> str:
  28. return datetime.now(timezone(timedelta(hours=8))).strftime("%Y-%m-%dT%H:%M:%S+08:00")
  29. # ---------------------------------------------------------------------------
  30. # 对象目录
  31. # ---------------------------------------------------------------------------
  32. _CATALOG_CACHE: dict | None = None
  33. def load_catalog() -> dict:
  34. global _CATALOG_CACHE
  35. if _CATALOG_CACHE is None:
  36. with (CATALOG_DIR / "entities.json").open("r", encoding="utf-8") as f:
  37. _CATALOG_CACHE = json.load(f)
  38. return _CATALOG_CACHE
  39. def load_scenarios() -> dict:
  40. with (CATALOG_DIR / "scenarios.json").open("r", encoding="utf-8") as f:
  41. return json.load(f)
  42. def find_object(object_code: str) -> dict | None:
  43. code = (object_code or "").strip().upper()
  44. for obj in load_catalog()["objects"]:
  45. if obj["objectCode"] == code:
  46. return obj
  47. return None
  48. # ---------------------------------------------------------------------------
  49. # canonical 样例
  50. # ---------------------------------------------------------------------------
  51. BUILTIN_SAMPLES: dict[str, list[dict]] = {
  52. "CUSTOMER": [
  53. {"bizKey": "CUST-001", "Cust": "CUST-001", "Name": "模拟客户001", "Domain": "SIM", "Status": "1"}
  54. ],
  55. "LOCATION": [
  56. {"bizKey": "LOC-A01", "Location": "SIM-A-01", "Descr": "模拟库位A01", "Domain": "SIM", "Typed": "WH"}
  57. ],
  58. "EMPLOYEE_HEADCOUNT": [
  59. {"bizKey": "EMP-0001", "Domain": "SIM", "Employee": "EMP-0001", "Name": "模拟员工001",
  60. "Department": "生产部", "Position": "操作工", "EmploymentStatus": "在岗"}
  61. ],
  62. "INVENTORY_OPENING_BALANCE": [
  63. {"bizKey": "OPEN-INV-001", "idid": "SIM-OPEN-INV-001", "code": "RM-CT-001", "sl": 100,
  64. "slzx": 100, "shyn": "mock", "shtime": "2026-09-01T00:00:00", "date0": "2026-09-01"}
  65. ],
  66. "FINISHED_ONHAND": [
  67. {"bizKey": "FG-ONHAND-001", "code": "FG-CT-001", "sl": 50, "date0": "2026-09-30", "location": "FG-01"}
  68. ],
  69. "FINISHED_OPENING_BALANCE": [
  70. {"bizKey": "FG-OPEN-001", "code": "FG-CT-001", "sl": 20, "date0": "2026-09-01", "location": "FG-01"}
  71. ],
  72. }
  73. def load_canonical_rows(obj: dict) -> list[dict]:
  74. """读取对象 canonical 样例;sample 为 null 时使用内置最小样例。"""
  75. sample = obj.get("sample")
  76. if not sample:
  77. rows = BUILTIN_SAMPLES.get(obj["objectCode"])
  78. if rows is None:
  79. raise FileNotFoundError(f"object {obj['objectCode']} has no sample and no builtin rows")
  80. return [dict(r) for r in rows]
  81. path = REPO_ROOT / sample
  82. if not path.exists():
  83. raise FileNotFoundError(f"sample not found: {sample}")
  84. with path.open("r", encoding="utf-8") as f:
  85. data = json.load(f)
  86. if not isinstance(data, list):
  87. raise ValueError(f"sample must be a JSON array: {sample}")
  88. return [dict(r) for r in data]
  89. def _extract_updated_at(row: dict) -> str:
  90. for alias in _UPDATED_AT_ALIASES:
  91. for key in row:
  92. if key.lower() == alias.lower() and row[key]:
  93. return str(row[key])
  94. return now_iso()
  95. def _prefix_biz_key(row: dict, run_id: str) -> str:
  96. """给业务键加 SIM-{runId}- 前缀,返回新 bizKey。兼容 bizKey / Id / id 三种载体。"""
  97. prefix = biz_prefix(run_id)
  98. for key in ("bizKey", "Id", "id"):
  99. if key in row and row[key] is not None:
  100. raw = str(row[key])
  101. row[key] = raw if raw.startswith(prefix) else prefix + raw
  102. return str(row[key])
  103. row["bizKey"] = prefix + uuid.uuid4().hex[:8]
  104. return row["bizKey"]
  105. # ---------------------------------------------------------------------------
  106. # 通道适配器
  107. # ---------------------------------------------------------------------------
  108. def adapt_db_rows(obj: dict, run_id: str, rows: list[dict] | None = None) -> list[dict]:
  109. """DB_SYNC:源表行 = canonical 行 + SIM 前缀业务键 + sourceUpdatedAt 增量列。"""
  110. rows = [dict(r) for r in (rows if rows is not None else load_canonical_rows(obj))]
  111. out = []
  112. for row in rows:
  113. _prefix_biz_key(row, run_id)
  114. row.setdefault("Domain", "SIM")
  115. row["sourceUpdatedAt"] = _extract_updated_at(row)
  116. out.append(row)
  117. return out
  118. def adapt_api_rows(obj: dict, run_id: str, rows: list[dict] | None = None) -> list[dict]:
  119. """API_PULL:与 DB_SYNC 同语义,另保证 bizKey 存在且按 bizKey 排序(cursor 语义)。"""
  120. rows = adapt_db_rows(obj, run_id, rows)
  121. for row in rows:
  122. row["bizKey"] = str(row.get("bizKey") or row.get("Id") or row.get("id"))
  123. rows.sort(key=lambda r: str(r["bizKey"]))
  124. return rows
  125. # entityCode -> 契约字段(与 MdpInboundFieldMapper.IdentityFields 一致);("字段", 必填, 别名列表)
  126. # 版本语义:C# 端只有显式 v2 才使用 IdentityFields(中立必填集合);v1 特化仅存在于
  127. # S5_INVENTORY_TXN 与 MDM_EMPLOYEE_HEADCOUNT。第一阶段目录统一登记 contractVersion=v2,
  128. # 因此本表按 IdentityFields 生成样例,避免“按 v1 造数、按 v2 校验”的口径错配。
  129. INBOUND_SPECS: dict[str, list[tuple[str, bool, list[str]]]] = {
  130. "MDM_ITEM": [("ItemNum", True, ["itemCode", "item_number", "number", "ItemNum", "wlbm"]),
  131. ("Descr", True, ["itemName", "descr", "wlmc"]),
  132. ("Drawing", False, []), ("UM", False, ["uom"]), ("ItemType", False, []),
  133. ("Status", False, ["status"]), ("Domain", False, ["Domain"])],
  134. "MDM_CUSTOMER": [("Cust", True, ["customerNo", "Cust"]), ("Name", True, ["customerName", "Name"]),
  135. ("Domain", False, ["Domain"]), ("Status", False, ["status"])],
  136. "MDM_SUPPLIER": [("Supp", True, ["supplierCode", "Supp", "supplier_number"]),
  137. ("Name", True, ["supplierName", "Name"]),
  138. ("Domain", False, ["Domain"]), ("Status", False, ["status"])],
  139. "MDM_LOCATION": [("Location", True, ["Location", "location"]), ("Descr", True, ["Descr", "descr", "Address"]),
  140. ("Domain", False, ["Domain"]), ("Typed", False, [])],
  141. "MDM_SOURCE_LIST": [("Supp", True, ["supplier_number", "Supp", "supplierCode"]),
  142. ("ItemNum", True, ["number", "itemCode", "ItemNum"]),
  143. ("Domain", False, ["Domain"]), ("Status", False, ["status"])],
  144. "MDM_EMPLOYEE_HEADCOUNT": [("Domain", True, ["Domain"]), ("Employee", True, ["Employee"]),
  145. ("Position", True, ["Position"]), ("EmploymentStatus", True, ["EmploymentStatus"]),
  146. ("Name", False, ["Name"]), ("Department", False, ["Department"])],
  147. "S1_SALES_ORDER_ENTRY": [("bill_no", True, ["bill_no", "orderNo"]), ("entry_seq", True, ["entry_seq", "Line"]),
  148. ("seorder_id", True, ["seorder_id", "id", "Id"]), ("qty", True, ["qty", "Qty"]),
  149. ("item_number", False, ["item_code", "itemCode", "item_number", "ItemNum"]),
  150. ("item_name", False, ["itemName"]), ("plan_date", False, ["plan_date"]),
  151. ("date", False, ["date"]), ("progress", False, []), ("urgent", False, [])],
  152. "S1_REQUIREMENT_EXAMINE_RESULT": [("Id", True, ["Id", "id"]), ("bill_no", True, ["bill_no"]),
  153. ("morder_no", True, ["morder_no"]), ("sentry_id", False, ["sentry_id"]),
  154. ("create_time", False, ["create_time"]), ("IsDeleted", False, [])],
  155. "S1_REQUIREMENT_EXAMINE_DETAIL": [("Id", True, ["Id", "id"]), ("examine_id", True, ["examine_id"]),
  156. ("item_number", True, ["item_number"]), ("lack_qty", True, ["lack_qty", "qty"]),
  157. ("level", True, ["level"]), ("num", False, ["num"]),
  158. ("item_name", False, ["item_name"]), ("needCount", False, []),
  159. ("qty", False, ["qty"]), ("is_use", False, [])],
  160. "S2_WORK_ORDER_SCHEDULE": [("WorkOrder", True, ["WorkOrd", "work_order", "workOrderNo", "OrderNo"]),
  161. ("ItemCode", True, ["ItemNum", "itemCode", "item_number", "componentItem"]),
  162. ("QtyOrdered", True, ["Qty", "qty", "QtyReq", "sl"]),
  163. ("DocType", False, []), ("DueDate", False, ["WorkDate", "DueDate"]),
  164. ("QtyCompleted", False, []), ("Status", False, ["Status", "status", "ztid"]),
  165. ("ProdLine", False, ["ProdLine"])],
  166. "S3_PURCHASE_ORDER": [("PoNo", True, ["poNo", "PurOrd", "OrdNbr"]), ("PoLine", True, ["poLine", "OrdLine", "Line"]),
  167. ("ItemCode", True, ["itemCode", "ItemNum"]), ("OrderQty", True, ["qty", "Qty", "QtyOrded"]),
  168. ("SupplierCode", False, ["supplierCode", "Supp"]), ("DueDate", False, []),
  169. ("OrderDate", False, []), ("Status", False, ["status", "Status"])],
  170. "S3_PURCHASE_RECEIPT": [("Domain", True, ["Domain"]), ("Receiver", True, ["Receiver", "RctNbr"]),
  171. ("Line", True, ["Line", "OrdLine"]), ("ItemCode", True, ["ItemNum", "itemCode"]),
  172. ("QtyReceived", True, ["QtyReceived", "qty"]), ("PoNo", False, ["OrdNbr", "poNo"]),
  173. ("SupplierCode", False, ["Supp", "supplierCode"]), ("ReceiptDate", False, ["RctDate"])],
  174. "S4_SHIPMENT": [("ShipmentNo", True, ["shipmentNo", "shddh"]), ("PoNo", True, ["poNo"]),
  175. ("ItemCode", True, ["itemCode"]), ("ShipQty", True, ["qty", "Qty"]),
  176. ("PoLine", False, ["poLine"]), ("SupplierCode", False, ["supplierCode"]),
  177. ("ShipDate", False, ["ShipDate"])],
  178. "S4_IQC": [("PoNo", True, ["poNo"]), ("PoLine", True, ["poLine"]), ("ItemCode", True, ["itemCode"]),
  179. ("ReceiptQty", True, ["qty", "ReceiptQty"]), ("SupplierCode", False, ["supplierCode"]),
  180. ("DefectQty", False, []), ("QcResult", False, ["result", "QcResult"]),
  181. ("ReceiptDate", False, ["ReceiptDate"])],
  182. "S4_RETURN": [("PoNo", True, ["poNo", "lydjbh"]), ("PoLine", True, ["poLine", "bbh"]),
  183. ("ItemCode", True, ["item_code", "wlbm"]), ("ReturnQty", True, ["qty", "ybl"]),
  184. ("SupplierCode", False, ["supplierCode"]), ("ReturnReason", False, ["bz"]),
  185. ("ReturnStatus", False, ["status"])],
  186. "S4_SHORTAGE": [("WorkOrder", True, ["work_order", "WorkOrd"]),
  187. ("ItemCode", True, ["component_item_code", "ItemNum"]),
  188. ("ShortageQty", True, ["shortage_qty", "lack_qty"]), ("SupplierCode", False, []),
  189. ("RiskLevel", False, []), ("NeedDate", False, [])],
  190. "S5_WORK_ORDER_BOM": [("Domain", True, ["Domain"]), ("OrderNo", True, ["orderNo", "WorkOrd"]),
  191. ("ItemCode", True, ["componentItem", "ItemNum", "itemCode"]),
  192. ("QtyRequired", True, ["qtyPer", "QtyReq", "qty"]),
  193. ("LineNo", False, ["Line"]), ("Unit", False, ["UM", "uom"])],
  194. "S5_INVENTORY_TXN": [("Id", True, ["Id", "id", "bizKey"]), ("BizDocType", True, ["bizDocType", "transType"]),
  195. ("ItemCode", False, ["itemNum", "code"]), ("Qty", False, ["qtyChange", "sl"]),
  196. ("ApprovedFlag", False, ["approvedFlag"]), ("VoidFlag", False, ["voidFlag"]),
  197. ("SummaryFlag", False, ["summaryFlag"]), ("DocDate", False, ["effDate", "jhdate"]),
  198. ("TransTime", False, ["addtime"]), ("Domain", False, ["Domain"]),
  199. ("TransType", False, ["transType"]), ("Location", False, ["location"]),
  200. ("LotSerial", False, ["lotSerial"]), ("CreateUser", False, ["createUser"]),
  201. ("EndBalance", False, ["endBalance"]), ("BeginBalance", False, []),
  202. ("DocQty", False, []), ("LineClosedFlag", False, [])],
  203. "S5_INVENTORY_OPENING_BALANCE": [("Id", True, ["Id", "idid", "bizKey"]), ("idid", False, ["idid"]),
  204. ("code", False, ["code", "ItemNum"]), ("sl", False, ["sl"]),
  205. ("slzx", False, ["slzx"]), ("shyn", False, ["shyn"]),
  206. ("shtime", False, ["shtime"]), ("date0", False, ["date0"])],
  207. "S5_INVENTORY_BALANCE_MONTHLY": [("PeriodYm", True, ["periodYm"]),
  208. ("AvgBalanceAmount", True, ["avgBalanceAmount"]),
  209. ("CategoryCode", False, []), ("WarehouseCode", False, ["warehouseCode"]),
  210. ("ItemCode", False, ["itemCode"]), ("IssueCostAmount", False, ["issueCostAmount"])],
  211. "S6_WORK_ORDER_LINE": [("Domain", True, ["Domain"]), ("OrderNo", True, ["WorkOrd", "orderNo"]),
  212. ("ItemCode", True, ["ItemNum", "itemCode"]), ("QtyPlanned", True, ["QtyReq", "qty", "Qty"]),
  213. ("QtyCompleted", False, []), ("PlanFinishDate", False, ["WorkDate"]),
  214. ("ReleaseTime", False, []), ("ClosedFlag", False, []), ("VoidFlag", False, [])],
  215. "S6_REPORT_TXN": [("Domain", True, ["Domain"]), ("ReportId", True, ["Id", "noid", "ReportId"]),
  216. ("WorkOrderNo", True, ["noid", "WorkOrd", "workOrderNo"]),
  217. ("ReportDate", True, ["kgdate", "ReportDate"]), ("ReportQty", True, ["sl", "qty", "ReportQty"]),
  218. ("StartWorkDate", False, ["kgdate"])],
  219. "S7_FQC_TASK_TXN": [("BillNo", True, ["billNo", "BillNo"]),
  220. ("ProductionOrderNo", True, ["productionOrderNo", "scph"]),
  221. ("MaterialCode", True, ["materialCode", "wlbm"]), ("Qty", True, ["qty"]),
  222. ("SalesOrderNo", False, ["salesOrderNo"]), ("Domain", False, ["Domain"])],
  223. "S7_SALES_ORDER_LINE": [("Domain", True, ["Domain"]), ("OrderNo", True, ["bill_no", "orderNo", "OrdNbr"]),
  224. ("LineNo", True, ["entry_seq", "Line"]), ("ItemCode", True, ["item_code", "itemCode", "ItemNum"]),
  225. ("QtyPlanned", True, ["qty", "Qty"]), ("PlanFinishDate", True, ["plan_date", "WorkDate"]),
  226. ("QtyCompleted", False, []), ("ReleaseTime", False, []),
  227. ("ClosedFlag", False, []), ("VoidFlag", False, [])],
  228. "S7_FINISHED_ONHAND": [("Id", True, ["Id", "bizKey"]), ("code", False, ["code", "ItemNum"]),
  229. ("sl", False, ["sl", "QtyRec"]), ("date0", False, ["date0", "Date"]),
  230. ("location", False, ["location", "LocationTo"])],
  231. "S7_FINISHED_OPENING_BALANCE": [("Id", True, ["Id", "bizKey"]), ("code", False, ["code"]),
  232. ("sl", False, ["sl"]), ("date0", False, ["date0"]),
  233. ("location", False, ["location"])],
  234. }
  235. def _pick(row_lower: dict[str, object], aliases: list[str]) -> object | None:
  236. for alias in aliases:
  237. if alias.lower() in row_lower and row_lower[alias.lower()] not in (None, ""):
  238. return row_lower[alias.lower()]
  239. return None
  240. def adapt_inbound_rows(obj: dict, run_id: str, rows: list[dict] | None = None,
  241. domain: str = "SIM") -> dict:
  242. """API_INBOUND:canonical 行 → 契约字段行。返回 rows / missing_required / biz_keys。"""
  243. inbound = obj.get("apiInbound") or {}
  244. entity_code = inbound.get("entityCode")
  245. if not inbound.get("supported") or not entity_code:
  246. raise ValueError(f"object {obj['objectCode']} does not support API_INBOUND")
  247. spec = INBOUND_SPECS.get(entity_code)
  248. if spec is None:
  249. raise ValueError(f"no inbound spec for {entity_code}")
  250. rows = [dict(r) for r in (rows if rows is not None else load_canonical_rows(obj))]
  251. out: list[dict] = []
  252. missing: set[str] = set()
  253. biz_keys: list[str] = []
  254. for row in rows:
  255. biz_key = _prefix_biz_key(row, run_id)
  256. biz_keys.append(biz_key)
  257. row_lower = {str(k).lower(): v for k, v in row.items()}
  258. mapped: dict[str, object] = {}
  259. for field, required, aliases in spec:
  260. value = _pick(row_lower, [field, *aliases])
  261. if value is None and field == "Domain":
  262. value = domain
  263. if value is None and required:
  264. missing.add(field)
  265. continue
  266. if value is not None:
  267. mapped[field] = value
  268. # Id 类必填键允许退化为 SIM 业务键
  269. if "Id" in {f for f, r, _ in spec if r} and "Id" not in mapped:
  270. mapped["Id"] = biz_key
  271. mapped["sourceUpdatedAt"] = _extract_updated_at(row)
  272. out.append(mapped)
  273. return {"rows": out, "missing_required": sorted(missing), "biz_keys": biz_keys}
  274. def inbound_envelope(rows: list[dict], snapshot_id: str | None = None, seq: int | None = None) -> dict:
  275. """推数信封,与 push_demo.py / MdpInboundReceiveService 契约一致。"""
  276. data: dict = {"list": rows}
  277. if snapshot_id:
  278. data["snapshotId"] = snapshot_id
  279. if seq is not None:
  280. data["seq"] = seq
  281. return {"data": data}