"""S7 成品入库(NBR/WOI) 端到端:Mock API → mdp_stg_production_receipt → mdp_std_production_receipt。 用法: set MOCK_TOKEN=uat-mock-token python e2e_api_to_std_s7.py """ from __future__ import annotations import json import sys from datetime import datetime import requests from _db import get_conn MOCK = "http://127.0.0.1:8018" TOKEN = "uat-mock-token" BATCH = f"CT_API_E2E_S7_{datetime.now():%Y%m%d%H%M%S}" def fetch(path: str) -> list[dict]: r = requests.get( MOCK.rstrip("/") + path, headers={"Authorization": f"Bearer {TOKEN}"}, timeout=30, ) r.raise_for_status() rows = r.json().get("data", {}).get("list") if not isinstance(rows, list): raise RuntimeError(f"{path} data.list invalid") return rows def upsert_stg(cur, source_table: str, row: dict, biz_expr: list[str]): vals = [] for f in biz_expr: if f not in row or row[f] is None or row[f] == "": vals = None break vals.append(str(row[f])) biz = "#".join(vals) if vals else str(row.get("bizKey") or row.get("RecID") or row.get("id")) rid = str(row.get("RecID") or row.get("bizKey") or biz) raw = json.dumps(row, ensure_ascii=False) cur.execute( """ INSERT INTO mdp_stg_production_receipt (tenant_id, source_system, source_table, source_row_id, source_biz_key, raw_data, sync_batch_id, sync_time, process_status, create_time) VALUES (0, 'WMS_API', %s, %s, %s, %s, %s, NOW(), 'PENDING', NOW()) ON DUPLICATE KEY UPDATE source_row_id=VALUES(source_row_id), raw_data=VALUES(raw_data), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), process_status='PENDING', update_time=NOW() """, (source_table, rid, biz, raw, BATCH), ) def transform_std(cur): cur.execute( """ INSERT INTO mdp_std_production_receipt (tenant_id, factory_id, source_system, domain, master_rec_id, detail_rec_id, nbr, line, receipt_date, status, status_desc, remark, prod_line, work_ord, erp_work_ord, department, department_desc, applicant_name, item_num, item_name, item_spec, um, location_to, location_to_desc, lot_serial, qty_rec, qty_to, location_from, location_from_desc, ord_nbr, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(n.tenant_id, 0), 1, IFNULL(NULLIF(n.source_system,''), 'WMS_API'), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.RecID')) AS SIGNED), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) AS SIGNED), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Line')) AS SIGNED), STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Date')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'), UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')), UPPER(IFNULL(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Status')), '')), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Remark')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ProdLine')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.WorkOrd')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Address')), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Department')), JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Name')), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.ItemNum')), NULL, NULL, NULL, JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationTo')), NULL, CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LotSerial')) AS CHAR), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyRec')) AS DECIMAL(18,5)), CAST(JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.QtyTo')) AS DECIMAL(18,5)), JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.LocationFrom')), NULL, JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.OrdNbr')), CONCAT(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')), '#', JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID'))), %s, NOW() FROM mdp_stg_production_receipt n INNER JOIN mdp_stg_production_receipt d ON d.source_table='NbrDetail' AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Domain')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Domain')) AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Nbr')) = JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.Nbr')) AND d.source_system='WMS_API' AND d.sync_batch_id=%s WHERE n.source_table='NbrMaster' AND n.source_system='WMS_API' AND n.sync_batch_id=%s AND JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.Type'))='WOI' AND CAST(JSON_UNQUOTE(JSON_EXTRACT(n.raw_data,'$.IsActive')) AS SIGNED)=1 AND JSON_UNQUOTE(JSON_EXTRACT(d.raw_data,'$.RecID')) IS NOT NULL ON DUPLICATE KEY UPDATE qty_rec=VALUES(qty_rec), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP """, (BATCH, BATCH, BATCH), ) def main() -> int: masters = fetch("/api/production-receipt") details = fetch("/api/production-receipt-detail") conn = get_conn() try: with conn.cursor() as cur: for r in masters: upsert_stg(cur, "NbrMaster", r, ["Domain", "Nbr"]) for r in details: upsert_stg(cur, "NbrDetail", r, ["Domain", "Nbr", "Line"]) transform_std(cur) cur.execute( "SELECT COUNT(*) AS c, COUNT(DISTINCT source_biz_key) AS k " "FROM mdp_std_production_receipt WHERE sync_batch_id=%s", (BATCH,), ) row = cur.fetchone() print(f"[PASS] batch={BATCH} std rows={row['c']} keys={row['k']}") return 0 if row["c"] > 0 else 1 except Exception as e: print(f"[FAIL] {e}", file=sys.stderr) return 1 finally: conn.close() if __name__ == "__main__": sys.exit(main())