"""S6 IPQC 端到端:Mock API → mdp_stg_ipqc_pull → mdp_std_ipqc_inspection。 用法: set MOCK_TOKEN=uat-mock-token python e2e_api_to_std_s6.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_S6_{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, row: dict): biz = str(row.get("djbh") or row.get("bizKey") or row.get("id")) rid = str(row.get("id") or biz) raw = json.dumps(row, ensure_ascii=False) cur.execute( """ INSERT INTO mdp_stg_ipqc_pull (tenant_id, source_system, source_table, source_row_id, source_biz_key, raw_data, sync_batch_id, sync_time, process_status, create_time) VALUES (1300000000001, 'WMS_API', 'qms_gcjyd', %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() """, (rid, biz, raw, BATCH), ) def transform_std(cur): cur.execute( """ INSERT INTO mdp_std_ipqc_inspection (tenant_id, factory_id, source_system, bill_no, product_model, production_batch_no, production_work_order, result_judgement, attachment, remark, inspector, process_code, process_name, production_person, sample_qty, form_no, version_no, effective_date, material_code, material_name, inspec_standard_version, inspec_standard_code, inspection_status, source_row_id, source_biz_key, sync_batch_id, sync_time) SELECT IFNULL(m.tenant_id, 1300000000001), 1, IFNULL(NULLIF(m.source_system,''), 'WMS_API'), IFNULL(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.cplx')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.scph')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.lydjbh')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgpd')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.fj')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bz')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jyr')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.gxbm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.gxmc')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.sczyry')), CAST(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.ybl')) AS DECIMAL(18,6)), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bdbh')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.bbh')), STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.sxrq')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.wlbm')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.wlmc')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgbb')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.jgbh')), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.status')), IFNULL(m.source_row_id, JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.id'))), IFNULL(NULLIF(m.source_biz_key,''), JSON_UNQUOTE(JSON_EXTRACT(m.raw_data,'$.djbh'))), %s, NOW() FROM mdp_stg_ipqc_pull m WHERE m.source_table='qms_gcjyd' AND m.source_system='WMS_API' AND m.sync_batch_id=%s ON DUPLICATE KEY UPDATE bill_no=VALUES(bill_no), result_judgement=VALUES(result_judgement), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time), update_time=CURRENT_TIMESTAMP """, (BATCH, BATCH), ) def main() -> int: rows = fetch("/api/ipqc") conn = get_conn() try: with conn.cursor() as cur: for r in rows: upsert_stg(cur, r) transform_std(cur) cur.execute( "SELECT COUNT(*) AS c, COUNT(DISTINCT source_biz_key) AS k " "FROM mdp_std_ipqc_inspection 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())