| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113 |
- """S6 报工契约:Mock /api/mes/report 以 DB/API 标签落 stg→std,比对 source_biz_key。"""
- 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"
- TS = datetime.now().strftime("%Y%m%d%H%M%S")
- PATH = "/api/mes/report"
- STG = "mdp_stg_s6_report"
- STD = "mdp_std_s6_report"
- FIELDS = ["noid", "kgdate"]
- def fetch() -> 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("invalid list")
- return rows
- def biz_key(row: dict) -> str:
- vals = []
- for f in FIELDS:
- alt = next((k for k in row if k.lower() == f.lower()), None)
- if alt is None or row[alt] in (None, ""):
- return str(row.get("bizKey") or row.get("Id") or row.get("id"))
- vals.append(str(row[alt]))
- return "#".join(vals)
- def upsert(cur, source_system: str, row: dict, batch: str):
- biz = biz_key(row)
- rid = str(row.get("Id") or row.get("id") or biz)
- cur.execute(
- f"""
- INSERT INTO {STG}
- (tenant_id, source_system, source_table, source_row_id, source_biz_key,
- raw_data, sync_batch_id, sync_time, process_status)
- VALUES (0, %s, 'Cj_Bg_Head_Rep', %s, %s, %s, %s, NOW(), 'PENDING')
- ON DUPLICATE KEY UPDATE
- raw_data=VALUES(raw_data), sync_batch_id=VALUES(sync_batch_id),
- sync_time=VALUES(sync_time), process_status='PENDING'
- """,
- (source_system, rid, biz, json.dumps(row, ensure_ascii=False), batch),
- )
- def transform(cur, source_system: str, batch: str):
- cur.execute(
- f"""
- INSERT INTO {STD}
- (tenant_id, factory_id, source_system, work_order_no, report_date, report_qty, ztid,
- source_row_id, source_biz_key, sync_batch_id, sync_time)
- SELECT 0, 1, %s,
- IFNULL(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.noid')),'null'), source_biz_key),
- STR_TO_DATE(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.kgdate')),'null'),''), '%%Y-%%m-%%d %%H:%%i:%%s'),
- CAST(NULLIF(NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.sl')),'null'),'') AS DECIMAL(18,6)),
- NULLIF(JSON_UNQUOTE(JSON_EXTRACT(raw_data,'$.ztid')),'null'),
- source_row_id, source_biz_key, %s, NOW()
- FROM {STG} WHERE sync_batch_id=%s
- ON DUPLICATE KEY UPDATE
- work_order_no=VALUES(work_order_no), report_date=VALUES(report_date),
- report_qty=VALUES(report_qty), sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
- """,
- (source_system, batch, batch),
- )
- def keys(cur, batch: str) -> set[str]:
- cur.execute(f"SELECT source_biz_key AS k FROM {STD} WHERE sync_batch_id=%s", (batch,))
- return {str(r["k"]) for r in cur.fetchall() if r["k"] is not None}
- def main() -> int:
- rows = fetch()
- db_batch = f"RPT_DB_{TS}"
- api_batch = f"RPT_API_{TS}"
- conn = get_conn()
- try:
- with conn.cursor() as cur:
- for r in rows:
- upsert(cur, "T8_V5_SQLSERVER", r, db_batch)
- transform(cur, "T8_V5_SQLSERVER", db_batch)
- kdb = keys(cur, db_batch)
- cur.execute(f"DELETE FROM {STD} WHERE sync_batch_id=%s", (db_batch,))
- cur.execute(f"DELETE FROM {STG} WHERE sync_batch_id=%s", (db_batch,))
- for r in rows:
- upsert(cur, "WMS_API", r, api_batch)
- transform(cur, "WMS_API", api_batch)
- kapi = keys(cur, api_batch)
- finally:
- conn.close()
- print(f"[S6_REPORT] db={len(kdb)} api={len(kapi)}")
- if not kdb or kdb != kapi:
- print(f"[FAIL] mismatch db={kdb} api={kapi}", file=sys.stderr)
- return 1
- print(f"[PASS] S6 report std contract (ts={TS})")
- return 0
- if __name__ == "__main__":
- raise SystemExit(main())
|