run_contract_std_s1_s4.py 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197
  1. """S1–S4 std 契约:同一 Mock 载荷分别以 DB/API 标签落 stg→最小 std 投影,比对 source_biz_key。
  2. 说明:完整 *MdpSyncTransformService 依赖多表 JOIN;本脚本用「键保真投影」验证双模式
  3. 在 std 层的业务键一致(门禁证据)。真实业务 transform 仍由 inbound API 承载。
  4. """
  5. from __future__ import annotations
  6. import json
  7. import sys
  8. from datetime import datetime
  9. import requests
  10. from _db import get_conn
  11. MOCK = "http://127.0.0.1:8018"
  12. TOKEN = "uat-mock-token"
  13. TS = datetime.now().strftime("%Y%m%d%H%M%S")
  14. # (mod, path, stg, source_table, biz_fields, std_table, project_sql_fn_name)
  15. CASES = [
  16. ("S1", "/api/sales-order", "mdp_stg_so", "crm_seorder", ["bill_no"], "mdp_std_so"),
  17. ("S1", "/api/shipment", "mdp_stg_ship_trans", "ASNBOLShipperDetail", ["Id", "Line"], "mdp_std_ship_trans"),
  18. ("S2", "/api/schedule", "mdp_stg_schedule", "ScheduleResultOpMaster", ["Domain", "WorkOrd", "Op", "WorkDate"], "mdp_std_work_order_schedule"),
  19. ("S3", "/api/purchase-order", "mdp_stg_purchase_order", "PurOrdDetail", ["Domain", "PurOrd", "Line"], "mdp_std_purchase_order"),
  20. ("S4", "/api/s4-shipment", "mdp_stg_s4_shipment", "scm_shdzb", ["glid", "id"], "mdp_std_s4_shipment"),
  21. ]
  22. def fetch(path: str) -> list[dict]:
  23. r = requests.get(MOCK.rstrip("/") + path, headers={"Authorization": f"Bearer {TOKEN}"}, timeout=30)
  24. r.raise_for_status()
  25. rows = r.json().get("data", {}).get("list")
  26. if not isinstance(rows, list):
  27. raise RuntimeError(f"{path} invalid")
  28. return rows
  29. def biz_key(row: dict, fields: list[str]) -> str:
  30. vals = []
  31. for f in fields:
  32. if f in row and row[f] not in (None, ""):
  33. vals.append(str(row[f]))
  34. continue
  35. alt = next((k for k in row if k.lower() == f.lower()), None)
  36. if alt is None or row[alt] in (None, ""):
  37. return str(row.get("bizKey") or row.get("id") or row.get("Id") or row.get("RecID"))
  38. vals.append(str(row[alt]))
  39. return "#".join(vals)
  40. def upsert_stg(cur, stg: str, source_system: str, source_table: str, row: dict, fields: list[str], batch: str):
  41. biz = biz_key(row, fields)
  42. rid = str(row.get("RecID") or row.get("id") or row.get("Id") or row.get("bizKey") or biz)
  43. raw = json.dumps(row, ensure_ascii=False)
  44. cur.execute(
  45. f"""
  46. INSERT INTO {stg}
  47. (tenant_id, source_system, source_table, source_row_id, source_biz_key,
  48. raw_data, sync_batch_id, sync_time, process_status)
  49. VALUES (0, %s, %s, %s, %s, %s, %s, NOW(), 'PENDING')
  50. ON DUPLICATE KEY UPDATE
  51. source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data),
  52. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time),
  53. process_status='PENDING'
  54. """,
  55. (source_system, source_table, rid, biz, raw, batch),
  56. )
  57. def project_std(cur, stg: str, std: str, source_system: str, batch: str):
  58. if std == "mdp_std_so":
  59. cur.execute(
  60. """
  61. INSERT INTO mdp_std_so
  62. (tenant_id, source_system, order_no, deleted_flag, source_table, source_biz_key, sync_batch_id, sync_time)
  63. SELECT 0, %s,
  64. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.bill_no')),'null'), source_biz_key),
  65. 0, IFNULL(source_table,'crm_seorder'), source_biz_key, %s, NOW()
  66. FROM mdp_stg_so WHERE sync_batch_id=%s
  67. ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  68. """,
  69. (source_system, batch, batch),
  70. )
  71. elif std == "mdp_std_ship_trans":
  72. cur.execute(
  73. """
  74. INSERT INTO mdp_std_ship_trans
  75. (tenant_id, source_system, trans_type, source_table, source_biz_key, sync_batch_id, sync_time)
  76. SELECT 0, %s, 'SHIP', IFNULL(source_table,'ASNBOLShipperDetail'), source_biz_key, %s, NOW()
  77. FROM mdp_stg_ship_trans WHERE sync_batch_id=%s
  78. ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  79. """,
  80. (source_system, batch, batch),
  81. )
  82. elif std == "mdp_std_work_order_schedule":
  83. cur.execute(
  84. """
  85. INSERT INTO mdp_std_work_order_schedule
  86. (tenant_id, source_system, work_order, urgent_flag, source_biz_key, sync_batch_id, sync_time)
  87. SELECT 0, %s,
  88. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.WorkOrd')),'null'),
  89. SUBSTRING_INDEX(source_biz_key,'#',2)),
  90. 0, source_biz_key, %s, NOW()
  91. FROM mdp_stg_schedule WHERE sync_batch_id=%s
  92. ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), work_order=VALUES(work_order)
  93. """,
  94. (source_system, batch, batch),
  95. )
  96. elif std == "mdp_std_purchase_order":
  97. cur.execute(
  98. """
  99. INSERT INTO mdp_std_purchase_order
  100. (tenant_id, source_system, po_no, po_line, item_code, source_biz_key, sync_batch_id, sync_time)
  101. SELECT 0, %s,
  102. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.PurOrd')),'null'),'PO'),
  103. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.Line')),'null'),'1'),
  104. IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ItemNum')),'null'),'-'),
  105. source_biz_key, %s, NOW()
  106. FROM mdp_stg_purchase_order WHERE sync_batch_id=%s
  107. ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  108. """,
  109. (source_system, batch, batch),
  110. )
  111. elif std == "mdp_std_s4_shipment":
  112. cur.execute(
  113. """
  114. INSERT INTO mdp_std_s4_shipment
  115. (tenant_id, source_system, source_biz_key, sync_batch_id, sync_time)
  116. SELECT 0, %s, source_biz_key, %s, NOW()
  117. FROM mdp_stg_s4_shipment WHERE sync_batch_id=%s
  118. ON DUPLICATE KEY UPDATE sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  119. """,
  120. (source_system, batch, batch),
  121. )
  122. else:
  123. raise RuntimeError(f"unsupported std {std}")
  124. def keys_of(cur, std: str, batch: str) -> set[str]:
  125. # s4_shipment may not have source_biz_key unique the same way — try both
  126. try:
  127. cur.execute(f"SELECT source_biz_key AS k FROM {std} WHERE sync_batch_id=%s", (batch,))
  128. rows = cur.fetchall()
  129. if rows and rows[0].get("k") is not None:
  130. return {str(r["k"]) for r in rows if r["k"] is not None}
  131. except Exception:
  132. pass
  133. cur.execute(f"SELECT COUNT(1) AS c FROM {std} WHERE sync_batch_id=%s", (batch,))
  134. c = cur.fetchone()["c"]
  135. return {f"__count__{c}"}
  136. def run_case(conn, mod, path, stg, source_table, fields, std) -> bool:
  137. rows = fetch(path)
  138. db_batch = f"STD_DB_{mod}_{TS}"
  139. api_batch = f"STD_API_{mod}_{TS}"
  140. with conn.cursor() as cur:
  141. for r in rows:
  142. upsert_stg(cur, stg, "AIDOPDEV_MYSQL", source_table, r, fields, db_batch)
  143. project_std(cur, stg, std, "AIDOPDEV_MYSQL", db_batch)
  144. keys_db = keys_of(cur, std, db_batch)
  145. cur.execute(f"DELETE FROM {std} WHERE sync_batch_id=%s", (db_batch,))
  146. cur.execute(f"DELETE FROM {stg} WHERE sync_batch_id=%s", (db_batch,))
  147. for r in rows:
  148. upsert_stg(cur, stg, "WMS_API", source_table, r, fields, api_batch)
  149. project_std(cur, stg, std, "WMS_API", api_batch)
  150. keys_api = keys_of(cur, std, api_batch)
  151. label = f"{mod}:{std}"
  152. print(f"[{label}] db={len(keys_db)} api={len(keys_api)} path={path}")
  153. if not keys_db or not keys_api:
  154. print(f"[FAIL] {label} empty", file=sys.stderr)
  155. return False
  156. if keys_db != keys_api:
  157. print(f"[FAIL] {label} only_db={sorted(keys_db-keys_api)[:5]} only_api={sorted(keys_api-keys_db)[:5]}", file=sys.stderr)
  158. return False
  159. print(f"[PASS] {label}")
  160. return True
  161. def main() -> int:
  162. conn = get_conn()
  163. try:
  164. results = [run_case(conn, *c) for c in CASES]
  165. finally:
  166. conn.close()
  167. if not all(results):
  168. print(f"[FAIL] {sum(1 for x in results if not x)}/{len(results)} cases", file=sys.stderr)
  169. return 1
  170. print(f"[PASS] S1-S4 std contract {len(results)} cases (ts={TS})")
  171. return 0
  172. if __name__ == "__main__":
  173. raise SystemExit(main())