run_uat_migration.py 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951
  1. #!/usr/bin/env python3
  2. """Build and execute tenant-scoped UAT migration batches.
  3. The source SQL files contain extracted, anonymized S1-S4 chains. This runner
  4. assigns target-specific IDs and business keys, persists old/new mappings, and
  5. keeps credentials outside the repository by reading the local App.json.
  6. """
  7. from __future__ import annotations
  8. import argparse
  9. import json
  10. import re
  11. from dataclasses import dataclass
  12. from datetime import datetime
  13. from pathlib import Path
  14. from typing import Iterable
  15. import pymysql
  16. from pymysql.constants import CLIENT
  17. ROOT = Path(__file__).resolve().parents[4]
  18. SQL_ROOT = ROOT / "doc" / "plan" / "sql"
  19. DATABASE_JSON = ROOT / "server" / "Admin.NET.Application" / "Configuration" / "Database.json"
  20. SOURCE_TENANT = 797403760988229
  21. @dataclass(frozen=True)
  22. class Target:
  23. code: str
  24. tenant_id: int
  25. factory_id: int
  26. operator_id: int
  27. source_file: str
  28. source_batch: str
  29. target_batch: str
  30. id_base: int
  31. key_prefix: str
  32. TARGETS = {
  33. "A": Target(
  34. "A",
  35. 838257186181189,
  36. 838257186320453,
  37. 838257187360837,
  38. "S1S4_UAT_P0_execute_write.sql",
  39. "S1S4_UAT_20260604_RQ_V1",
  40. "UAT_WP1_A_20260817_V1",
  41. 9206081700000001,
  42. "UATA",
  43. ),
  44. "B": Target(
  45. "B",
  46. 838257212780613,
  47. 838257212858437,
  48. 838257213620293,
  49. "S1S4_UAT_P0_execute_write.sql",
  50. "S1S4_UAT_20260604_RQ_V1",
  51. "UAT_WP1_B_20260817_V1",
  52. 9206081800000001,
  53. "UATB",
  54. ),
  55. "DEMO": Target(
  56. "DEMO",
  57. 838257237606469,
  58. 838257237676101,
  59. 838257238302789,
  60. "S1S4_UAT_FULL_execute_write.sql",
  61. "S1S4_UAT_FULL_20260605_V1",
  62. "UAT_WP1_DEMO_20260817_V1",
  63. 9206081900000001,
  64. "DEMO",
  65. ),
  66. }
  67. def walk_strings(value: object) -> Iterable[str]:
  68. if isinstance(value, str):
  69. yield value
  70. elif isinstance(value, dict):
  71. for child in value.values():
  72. yield from walk_strings(child)
  73. elif isinstance(value, list):
  74. for child in value:
  75. yield from walk_strings(child)
  76. def load_connection() -> pymysql.Connection:
  77. raw = DATABASE_JSON.read_text(encoding="utf-8-sig")
  78. candidates = [
  79. s
  80. for s in re.findall(r'(?m)^\s*"ConnectionString"\s*:\s*"([^"]+)"', raw)
  81. if "Database=aidopdev" in s
  82. ]
  83. if not candidates:
  84. raise RuntimeError("Database.json 中未找到启用的 aidopdev 连接串")
  85. parts = {}
  86. for pair in candidates[0].split(";"):
  87. if "=" in pair:
  88. key, val = pair.split("=", 1)
  89. parts[key.strip().lower()] = val.strip()
  90. return pymysql.connect(
  91. host=parts.get("server", "127.0.0.1"),
  92. port=int(parts.get("port", "3306")),
  93. user=parts.get("uid") or parts.get("user id"),
  94. password=parts.get("pwd") or parts.get("password"),
  95. database=parts.get("database", "aidopdev"),
  96. charset="utf8mb4",
  97. autocommit=True,
  98. client_flag=CLIENT.MULTI_STATEMENTS,
  99. connect_timeout=20,
  100. read_timeout=600,
  101. write_timeout=600,
  102. )
  103. def ordered_tokens(sql: str, pattern: str) -> list[str]:
  104. return list(dict.fromkeys(re.findall(pattern, sql)))
  105. def build_key_map(sql: str, target: Target) -> dict[str, tuple[str, str]]:
  106. mapping: dict[str, tuple[str, str]] = {}
  107. orders = ordered_tokens(sql, r"'((?:MPO|SO)\d{10,})'")
  108. works = ordered_tokens(sql, r"'(M\d{9,12})'")
  109. prs = ordered_tokens(sql, r"'(PR-UAT-[^']+)'")
  110. pos = ordered_tokens(sql, r"'(PO-UAT-[^']+)'")
  111. ships = ordered_tokens(sql, r"'(SH-UAT-[^']+)'")
  112. for idx, old in enumerate(orders, 1):
  113. scenario = f"A{idx:02d}" if target.code == "A" else f"CP{idx:02d}"
  114. new = (
  115. f"UATA-{scenario}-SO"
  116. if target.code == "A"
  117. else f"UATB-{scenario}-SO"
  118. if target.code == "B"
  119. else f"DEMO-SO-{idx:03d}"
  120. )
  121. mapping[old] = ("SALES_ORDER", new)
  122. for idx, old in enumerate(works, 1):
  123. new = f"{target.key_prefix}-WO-{idx:03d}"
  124. mapping[old] = ("WORK_ORDER", new)
  125. for kind, tokens, label in (
  126. ("PURCHASE_REQUEST", prs, "PR"),
  127. ("PURCHASE_ORDER", pos, "PO"),
  128. ("SUPPLIER_SHIPMENT", ships, "SH"),
  129. ):
  130. for idx, old in enumerate(tokens, 1):
  131. mapping[old] = (kind, f"{target.key_prefix}-{label}-{idx:03d}")
  132. return mapping
  133. def patch_tenant_columns(sql: str) -> str:
  134. sql = sql.replace(
  135. "need_number, morder_production_number, IsDeleted,\n"
  136. " create_by_name, create_time, update_by_name, update_time,\n"
  137. " tenant_id, factory_id\n"
  138. ")\n"
  139. "SELECT target_mo_id, morder_no, @now, 'Released', product_code, product_code,\n"
  140. " need_number, need_number, 0,",
  141. "need_number, morder_production_number, IsDeleted, urgent,\n"
  142. " create_by_name, create_time, update_by_name, update_time,\n"
  143. " tenant_id, factory_id\n"
  144. ")\n"
  145. "SELECT target_mo_id, morder_no, @now, 'Released', product_code, product_code,\n"
  146. " need_number, need_number, 0, 0,",
  147. )
  148. sql = sql.replace(
  149. "INSERT INTO scm_shd (id, po_billno, shddh, sh_purchase_num, "
  150. "estimated_delivery_date, tjrxm, tjrq, state)",
  151. "INSERT INTO scm_shd (id, po_billno, shddh, sh_purchase_num, "
  152. "estimated_delivery_date, tjrxm, tjrq, state, tenant_id)",
  153. )
  154. sql = sql.replace(
  155. "DATE_FORMAT(@now,'%Y-%m-%d'), 0\nFROM tmp_s4_uat_ship",
  156. "DATE_FORMAT(@now,'%Y-%m-%d'), 0, @tenant_id\nFROM tmp_s4_uat_ship",
  157. )
  158. sql = sql.replace(
  159. "INSERT INTO scm_shdzb (id, glid, sh_material_code, "
  160. "sh_delivery_quantity, po_bill, po_billline)",
  161. "INSERT INTO scm_shdzb (id, glid, sh_material_code, "
  162. "sh_delivery_quantity, po_bill, po_billline, tenant_id)",
  163. )
  164. sql = sql.replace(
  165. "pl.item_num, pl.qty, pl.pur_ord, CAST(pl.line_no AS CHAR)\n"
  166. "FROM tmp_s3_uat_po_line",
  167. "pl.item_num, pl.qty, pl.pur_ord, CAST(pl.line_no AS CHAR), @tenant_id\n"
  168. "FROM tmp_s3_uat_po_line",
  169. )
  170. sql += """
  171. -- Keep both work-order representations aligned with the linked sales-order item.
  172. UPDATE mes_morder m
  173. JOIN mes_moentry me
  174. ON me.tenant_id=m.tenant_id AND me.moentry_mono=m.morder_no
  175. JOIN crm_seorderentry e
  176. ON e.tenant_id=me.tenant_id AND e.Id=me.soentry_id
  177. SET m.product_code=e.item_number,
  178. m.product_name=COALESCE(NULLIF(e.item_name,''),e.item_number)
  179. WHERE m.tenant_id=@tenant_id;
  180. UPDATE WorkOrdMaster w
  181. JOIN mes_moentry me
  182. ON me.tenant_id=w.tenant_id AND me.moentry_mono=w.WorkOrd
  183. JOIN crm_seorderentry e
  184. ON e.tenant_id=me.tenant_id AND e.Id=me.soentry_id
  185. SET w.ItemNum=e.item_number
  186. WHERE w.tenant_id=@tenant_id;
  187. UPDATE PurOrdDetail d
  188. JOIN PurOrdMaster h
  189. ON h.tenant_id=d.tenant_id AND h.PurOrd=d.PurOrd
  190. SET d.PurOrdRecID=h.RecID
  191. WHERE d.tenant_id=@tenant_id AND d.PurOrdRecID=0;
  192. """
  193. return sql
  194. def split_statements(sql: str) -> list[str]:
  195. """Split these migration files without breaking on comment semicolons."""
  196. body = "\n".join(
  197. line for line in sql.splitlines() if not line.lstrip().startswith("--")
  198. )
  199. return [statement.strip() for statement in body.split(";") if statement.strip()]
  200. def adapt_sql(target: Target) -> tuple[str, dict[str, tuple[str, str]]]:
  201. source = (SQL_ROOT / target.source_file).read_text(encoding="utf-8-sig")
  202. mapping = build_key_map(source, target)
  203. sql = source
  204. sql = re.sub(r"SET @tenant_id\s*:=\s*\d+;", f"SET @tenant_id := {target.tenant_id};", sql)
  205. sql = re.sub(
  206. r"SET @factory_id\s*:=\s*\d+;",
  207. f"SET @factory_id := {target.factory_id};",
  208. sql,
  209. )
  210. sql = re.sub(r"SET @operator_id\s*:=\s*\d+;", f"SET @operator_id := {target.operator_id};", sql)
  211. sql = re.sub(
  212. r"SET @operator_name\s*:=\s*'[^']*';",
  213. f"SET @operator_name := 'UAT-{target.code}-Migration';",
  214. sql,
  215. )
  216. sql = re.sub(r"SET @batch_no\s*:=\s*'[^']+';", f"SET @batch_no := '{target.target_batch}';", sql)
  217. sql = re.sub(r"SET @id_base\s*:=\s*\d+;", f"SET @id_base := {target.id_base};", sql)
  218. # FULL uses literal IDs in its temporary mapping rows.
  219. if target.source_file.startswith("S1S4_UAT_FULL"):
  220. sql = sql.replace("91060005", str(target.id_base)[:8])
  221. # The historical template was slimmed; use an extant anonymized source PO.
  222. sql = sql.replace(
  223. "WHERE PurOrd='PW202605100002' AND tenant_id=@tenant_id",
  224. f"WHERE PurOrd='DO20260711001' AND tenant_id={SOURCE_TENANT}",
  225. )
  226. for old, (_, new) in sorted(mapping.items(), key=lambda item: len(item[0]), reverse=True):
  227. sql = sql.replace(f"'{old}'", f"'{new}'")
  228. sql = patch_tenant_columns(sql)
  229. return sql, mapping
  230. def ensure_audit_tables(cur: pymysql.cursors.Cursor) -> None:
  231. cur.execute(
  232. """
  233. CREATE TABLE IF NOT EXISTS uat_data_migration_batch (
  234. batch_code varchar(100) NOT NULL PRIMARY KEY,
  235. target_tenant_id bigint NOT NULL,
  236. target_set varchar(16) NOT NULL,
  237. source_file varchar(255) NOT NULL,
  238. status varchar(20) NOT NULL,
  239. started_at datetime(3) NOT NULL,
  240. finished_at datetime(3) NULL,
  241. summary_json json NULL
  242. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
  243. """
  244. )
  245. cur.execute(
  246. """
  247. CREATE TABLE IF NOT EXISTS uat_data_migration_map (
  248. batch_code varchar(100) NOT NULL,
  249. object_type varchar(40) NOT NULL,
  250. source_key varchar(200) NOT NULL,
  251. target_key varchar(200) NOT NULL,
  252. target_tenant_id bigint NOT NULL,
  253. created_at datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
  254. PRIMARY KEY(batch_code, object_type, source_key),
  255. KEY ix_uat_map_target(target_tenant_id, object_type, target_key)
  256. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4
  257. """
  258. )
  259. def execute_target(target: Target, dry_run: bool) -> dict[str, object]:
  260. sql, mapping = adapt_sql(target)
  261. generated = Path(__file__).with_name(f"{target.target_batch}.sql")
  262. generated.write_text(sql, encoding="utf-8")
  263. if dry_run:
  264. return {"target": target.code, "generated": str(generated), "mapping_count": len(mapping)}
  265. conn = load_connection()
  266. try:
  267. with conn.cursor() as cur:
  268. ensure_audit_tables(cur)
  269. cur.execute(
  270. "SELECT status FROM uat_data_migration_batch WHERE batch_code=%s",
  271. (target.target_batch,),
  272. )
  273. existing = cur.fetchone()
  274. if existing and existing[0] == "SUCCESS":
  275. raise RuntimeError(f"{target.target_batch} 已成功执行,拒绝重复导入")
  276. cur.execute(
  277. """
  278. INSERT INTO uat_data_migration_batch
  279. (batch_code,target_tenant_id,target_set,source_file,status,started_at)
  280. VALUES(%s,%s,%s,%s,'RUNNING',NOW(3))
  281. ON DUPLICATE KEY UPDATE status='RUNNING',started_at=NOW(3),finished_at=NULL
  282. """,
  283. (target.target_batch, target.tenant_id, target.code, target.source_file),
  284. )
  285. try:
  286. for index, statement in enumerate(split_statements(sql), 1):
  287. try:
  288. cur.execute(statement)
  289. except Exception as exc:
  290. cur.execute("ROLLBACK")
  291. excerpt = " ".join(statement.split())[:240]
  292. raise RuntimeError(
  293. f"SQL #{index} 执行失败: {excerpt}: {exc}"
  294. ) from exc
  295. cur.executemany(
  296. """
  297. INSERT INTO uat_data_migration_map
  298. (batch_code,object_type,source_key,target_key,target_tenant_id)
  299. VALUES(%s,%s,%s,%s,%s)
  300. ON DUPLICATE KEY UPDATE target_key=VALUES(target_key)
  301. """,
  302. [
  303. (target.target_batch, kind, old, new, target.tenant_id)
  304. for old, (kind, new) in mapping.items()
  305. ],
  306. )
  307. summary = verify(cur, target)
  308. cur.execute(
  309. """
  310. UPDATE uat_data_migration_batch
  311. SET status='SUCCESS',finished_at=NOW(3),summary_json=%s
  312. WHERE batch_code=%s
  313. """,
  314. (json.dumps(summary, ensure_ascii=False), target.target_batch),
  315. )
  316. return summary
  317. except Exception as exc:
  318. cur.execute(
  319. """
  320. UPDATE uat_data_migration_batch
  321. SET status='FAILED',finished_at=NOW(3),summary_json=%s
  322. WHERE batch_code=%s
  323. """,
  324. (json.dumps({"error": str(exc)}, ensure_ascii=False), target.target_batch),
  325. )
  326. raise
  327. finally:
  328. conn.close()
  329. def verify(cur: pymysql.cursors.Cursor, target: Target) -> dict[str, object]:
  330. queries = {
  331. "orders": (
  332. "SELECT COUNT(*) FROM crm_seorder WHERE tenant_id=%s AND bill_from LIKE %s",
  333. (target.tenant_id, f"%{target.target_batch}%"),
  334. ),
  335. "order_lines": (
  336. """
  337. SELECT COUNT(*) FROM crm_seorderentry e
  338. JOIN crm_seorder h ON h.Id=e.seorder_id AND h.tenant_id=e.tenant_id
  339. WHERE h.tenant_id=%s AND h.bill_from LIKE %s
  340. """,
  341. (target.tenant_id, f"%{target.target_batch}%"),
  342. ),
  343. "work_orders": (
  344. "SELECT COUNT(*) FROM mes_morder WHERE tenant_id=%s AND morder_no LIKE %s",
  345. (target.tenant_id, f"{target.key_prefix}-WO-%"),
  346. ),
  347. "purchase_requests": (
  348. "SELECT COUNT(*) FROM srm_pr_main WHERE tenant_id=%s AND pr_billno LIKE %s",
  349. (target.tenant_id, f"{target.key_prefix}-PR-%"),
  350. ),
  351. "purchase_orders": (
  352. "SELECT COUNT(*) FROM PurOrdMaster WHERE tenant_id=%s AND PurOrd LIKE %s",
  353. (target.tenant_id, f"{target.key_prefix}-PO-%"),
  354. ),
  355. "supplier_shipments": (
  356. "SELECT COUNT(*) FROM scm_shd WHERE tenant_id=%s AND shddh LIKE %s",
  357. (target.tenant_id, f"{target.key_prefix}-SH-%"),
  358. ),
  359. }
  360. result: dict[str, object] = {"target": target.code, "tenant_id": target.tenant_id}
  361. for key, (query, params) in queries.items():
  362. cur.execute(query, params)
  363. result[key] = int(cur.fetchone()[0])
  364. return result
  365. def table_columns(cur: pymysql.cursors.Cursor, table: str) -> tuple[list[str], str]:
  366. cur.execute(f"SHOW COLUMNS FROM `{table}`")
  367. rows = cur.fetchall()
  368. writable = [row[0] for row in rows if "auto_increment" not in (row[5] or "")]
  369. primary = next(row[0] for row in rows if row[3] == "PRI")
  370. return writable, primary
  371. def ensure_master_rows(
  372. cur: pymysql.cursors.Cursor,
  373. target: Target,
  374. table: str,
  375. key_column: str,
  376. keys: list[str],
  377. fallback: bool = False,
  378. ) -> int:
  379. columns, primary = table_columns(cur, table)
  380. inserted = 0
  381. for key in keys:
  382. cur.execute(
  383. f"SELECT `{primary}` FROM `{table}` WHERE tenant_id=%s AND `{key_column}`=%s LIMIT 1",
  384. (target.tenant_id, key),
  385. )
  386. if cur.fetchone():
  387. continue
  388. cur.execute(
  389. f"""
  390. SELECT {",".join(f"`{column}`" for column in columns)}
  391. FROM `{table}`
  392. WHERE tenant_id IN (%s,%s) AND `{key_column}`=%s
  393. ORDER BY FIELD(tenant_id,%s,%s)
  394. LIMIT 1
  395. """,
  396. (SOURCE_TENANT, 824585161322565, key, SOURCE_TENANT, 824585161322565),
  397. )
  398. row = cur.fetchone()
  399. if row is None and fallback:
  400. cur.execute(
  401. f"""
  402. SELECT {",".join(f"`{column}`" for column in columns)}
  403. FROM `{table}` WHERE tenant_id=%s ORDER BY `{primary}` LIMIT 1
  404. """,
  405. (SOURCE_TENANT,),
  406. )
  407. row = cur.fetchone()
  408. if row is None:
  409. raise RuntimeError(f"{table}.{key_column}={key} 无可用脱敏源记录")
  410. values = dict(zip(columns, row))
  411. values[key_column] = key
  412. for column in ("tenant_id", "TenantId"):
  413. if column in values:
  414. values[column] = target.tenant_id
  415. for column in (
  416. "company_ref_id",
  417. "factory_ref_id",
  418. "CompanyRefId",
  419. "FactoryRefId",
  420. "company_id",
  421. "factory_id",
  422. "org_id",
  423. ):
  424. if column in values:
  425. values[column] = target.factory_id
  426. cur.execute(
  427. f"""
  428. INSERT INTO `{table}` ({",".join(f"`{column}`" for column in columns)})
  429. VALUES ({",".join(["%s"] * len(columns))})
  430. """,
  431. tuple(values[column] for column in columns),
  432. )
  433. new_id = cur.lastrowid
  434. cur.execute(
  435. """
  436. INSERT INTO uat_data_migration_map
  437. (batch_code,object_type,source_key,target_key,target_tenant_id)
  438. VALUES(%s,%s,%s,%s,%s)
  439. ON DUPLICATE KEY UPDATE target_key=VALUES(target_key)
  440. """,
  441. (
  442. target.target_batch,
  443. f"MASTER_{table.upper()}",
  444. key,
  445. f"{key}:{new_id}",
  446. target.tenant_id,
  447. ),
  448. )
  449. inserted += 1
  450. return inserted
  451. def scoped_values(
  452. columns: list[str], row: tuple[object, ...], target: Target
  453. ) -> dict[str, object]:
  454. values = dict(zip(columns, row))
  455. for column in ("tenant_id", "TenantId"):
  456. if column in values:
  457. values[column] = target.tenant_id
  458. for column in (
  459. "company_ref_id",
  460. "factory_ref_id",
  461. "CompanyRefId",
  462. "FactoryRefId",
  463. "company_id",
  464. "factory_id",
  465. "org_id",
  466. ):
  467. if column in values:
  468. values[column] = target.factory_id
  469. return values
  470. def insert_values(
  471. cur: pymysql.cursors.Cursor,
  472. table: str,
  473. columns: list[str],
  474. values: dict[str, object],
  475. ) -> int:
  476. cur.execute(
  477. f"""
  478. INSERT INTO `{table}` ({",".join(f"`{column}`" for column in columns)})
  479. VALUES ({",".join(["%s"] * len(columns))})
  480. """,
  481. tuple(values[column] for column in columns),
  482. )
  483. return int(cur.lastrowid)
  484. def execute_structure_closure(
  485. cur: pymysql.cursors.Cursor, target: Target, initial_items: list[str]
  486. ) -> dict[str, int]:
  487. # Expand one BOM level so every copied BOM component also has an item master.
  488. placeholders = ",".join(["%s"] * len(initial_items))
  489. component_items: set[str] = set()
  490. if initial_items:
  491. cur.execute(
  492. f"""
  493. SELECT DISTINCT ComponentItem
  494. FROM ProductStructureMaster
  495. WHERE tenant_id IN (%s,%s)
  496. AND ParentItem IN ({placeholders})
  497. AND IFNULL(ComponentItem,'')<>''
  498. """,
  499. (SOURCE_TENANT, 824585161322565, *initial_items),
  500. )
  501. component_items = {str(row[0]) for row in cur.fetchall()}
  502. expanded_items = sorted(set(initial_items) | component_items)
  503. item_inserted = ensure_master_rows(
  504. cur, target, "ItemMaster", "ItemNum", expanded_items, fallback=True
  505. )
  506. cur.execute(
  507. "SELECT ItemNum,RecID FROM ItemMaster WHERE tenant_id=%s",
  508. (target.tenant_id,),
  509. )
  510. target_item_ids = {str(row[0]): int(row[1]) for row in cur.fetchall()}
  511. bom_columns, bom_primary = table_columns(cur, "ProductStructureMaster")
  512. bom_maps: dict[int, int] = {}
  513. seen_bom: set[tuple[object, ...]] = set()
  514. bom_inserted = 0
  515. if initial_items:
  516. for source_tenant in (SOURCE_TENANT, 824585161322565):
  517. cur.execute(
  518. f"""
  519. SELECT `{bom_primary}`,{",".join(f"`{column}`" for column in bom_columns)}
  520. FROM ProductStructureMaster
  521. WHERE tenant_id=%s AND ParentItem IN ({placeholders})
  522. ORDER BY `{bom_primary}`
  523. """,
  524. (source_tenant, *initial_items),
  525. )
  526. for source_id, *row in cur.fetchall():
  527. values = scoped_values(bom_columns, tuple(row), target)
  528. key = (
  529. values.get("ParentItem"),
  530. values.get("ComponentItem"),
  531. target.factory_id,
  532. )
  533. if key in seen_bom:
  534. continue
  535. seen_bom.add(key)
  536. if "ParentMaterialId" in values:
  537. values["ParentMaterialId"] = target_item_ids.get(
  538. str(values.get("ParentItem"))
  539. )
  540. if "ComponentMaterialId" in values:
  541. values["ComponentMaterialId"] = target_item_ids.get(
  542. str(values.get("ComponentItem"))
  543. )
  544. inserted = False
  545. try:
  546. target_id = insert_values(
  547. cur, "ProductStructureMaster", bom_columns, values
  548. )
  549. inserted = True
  550. except pymysql.IntegrityError as exc:
  551. if exc.args[0] != 1062:
  552. raise
  553. cur.execute(
  554. """
  555. SELECT RecID FROM ProductStructureMaster
  556. WHERE tenant_id=%s AND CompanyRefId=%s AND FactoryRefId=%s
  557. AND ParentItem=%s AND ComponentItem=%s
  558. LIMIT 1
  559. """,
  560. (
  561. target.tenant_id,
  562. target.factory_id,
  563. target.factory_id,
  564. values.get("ParentItem"),
  565. values.get("ComponentItem"),
  566. ),
  567. )
  568. existing = cur.fetchone()
  569. if not existing:
  570. raise
  571. target_id = int(existing[0])
  572. bom_maps[int(source_id)] = target_id
  573. if inserted:
  574. bom_inserted += 1
  575. op_inserted = 0
  576. if bom_maps:
  577. op_columns, op_primary = table_columns(cur, "ProductStructureOp")
  578. ids = list(bom_maps)
  579. for start in range(0, len(ids), 500):
  580. chunk = ids[start : start + 500]
  581. cur.execute(
  582. f"""
  583. SELECT `{op_primary}`,{",".join(f"`{column}`" for column in op_columns)}
  584. FROM ProductStructureOp
  585. WHERE ProductStructureMasterId IN ({",".join(["%s"] * len(chunk))})
  586. """,
  587. tuple(chunk),
  588. )
  589. for _, *row in cur.fetchall():
  590. values = scoped_values(op_columns, tuple(row), target)
  591. old_master = int(values["ProductStructureMasterId"])
  592. values["ProductStructureMasterId"] = bom_maps[old_master]
  593. try:
  594. insert_values(cur, "ProductStructureOp", op_columns, values)
  595. op_inserted += 1
  596. except pymysql.IntegrityError as exc:
  597. if exc.args[0] != 1062:
  598. raise
  599. routing_columns, routing_primary = table_columns(cur, "RoutingOpDetail")
  600. routing_inserted = 0
  601. seen_routing: set[tuple[object, ...]] = set()
  602. if expanded_items:
  603. routing_placeholders = ",".join(["%s"] * len(expanded_items))
  604. for source_tenant in (SOURCE_TENANT, 824585161322565):
  605. cur.execute(
  606. f"""
  607. SELECT `{routing_primary}`,{",".join(f"`{column}`" for column in routing_columns)}
  608. FROM RoutingOpDetail
  609. WHERE tenant_id=%s AND MaterialCode IN ({routing_placeholders})
  610. ORDER BY `{routing_primary}`
  611. """,
  612. (source_tenant, *expanded_items),
  613. )
  614. for _, *row in cur.fetchall():
  615. values = scoped_values(routing_columns, tuple(row), target)
  616. key = tuple(
  617. values.get(column)
  618. for column in ("RouteCode", "MaterialCode", "OperationCode", "SortNo")
  619. )
  620. if key in seen_routing:
  621. continue
  622. seen_routing.add(key)
  623. try:
  624. insert_values(cur, "RoutingOpDetail", routing_columns, values)
  625. routing_inserted += 1
  626. except pymysql.IntegrityError as exc:
  627. if exc.args[0] != 1062:
  628. raise
  629. return {
  630. "bom_component_items": len(component_items),
  631. "bom_component_items_inserted": item_inserted,
  632. "bom_rows_inserted": bom_inserted,
  633. "bom_operations_inserted": op_inserted,
  634. "routing_rows_inserted": routing_inserted,
  635. }
  636. def execute_resource_closure(
  637. cur: pymysql.cursors.Cursor, target: Target
  638. ) -> dict[str, int]:
  639. cur.execute(
  640. """
  641. SELECT Department FROM DepartmentMaster
  642. WHERE tenant_id=%s AND IFNULL(Department,'')<>''
  643. GROUP BY Department ORDER BY MIN(RecID) LIMIT 20
  644. """,
  645. (SOURCE_TENANT,),
  646. )
  647. departments = [str(row[0]) for row in cur.fetchall()]
  648. department_inserted = ensure_master_rows(
  649. cur,
  650. target,
  651. "DepartmentMaster",
  652. "Department",
  653. departments,
  654. fallback=True,
  655. )
  656. cur.execute(
  657. """
  658. SELECT DISTINCT WorkCenterCode FROM RoutingOpDetail
  659. WHERE tenant_id=%s AND IFNULL(WorkCenterCode,'')<>''
  660. """,
  661. (target.tenant_id,),
  662. )
  663. work_centers = [str(row[0]) for row in cur.fetchall()]
  664. cur.execute(
  665. """
  666. SELECT DISTINCT COALESCE(NULLIF(StdOp,''),NULLIF(OperationCode,''))
  667. FROM RoutingOpDetail
  668. WHERE tenant_id=%s
  669. AND COALESCE(NULLIF(StdOp,''),NULLIF(OperationCode,'')) IS NOT NULL
  670. """,
  671. (target.tenant_id,),
  672. )
  673. operations = [str(row[0]) for row in cur.fetchall()]
  674. work_center_inserted = ensure_master_rows(
  675. cur, target, "WorkCtrMaster", "WorkCtr", work_centers, fallback=True
  676. )
  677. operation_inserted = ensure_master_rows(
  678. cur, target, "StdOpMaster", "StdOp", operations, fallback=True
  679. )
  680. cur.execute(
  681. """
  682. SELECT Line FROM LineMaster
  683. WHERE tenant_id=%s AND IFNULL(Line,'')<>''
  684. GROUP BY Line ORDER BY MIN(RecID) LIMIT 20
  685. """,
  686. (SOURCE_TENANT,),
  687. )
  688. lines = [str(row[0]) for row in cur.fetchall()]
  689. line_inserted = ensure_master_rows(
  690. cur, target, "LineMaster", "Line", lines, fallback=True
  691. )
  692. calendar_columns, calendar_primary = table_columns(cur, "ShopCalendarWorkCtr")
  693. calendar_inserted = 0
  694. if work_centers or lines:
  695. work_placeholders = ",".join(["%s"] * max(1, len(work_centers)))
  696. line_placeholders = ",".join(["%s"] * max(1, len(lines)))
  697. cur.execute(
  698. f"""
  699. SELECT {",".join(f"`{column}`" for column in calendar_columns)}
  700. FROM ShopCalendarWorkCtr
  701. WHERE tenant_id=%s
  702. AND (WorkCtr IN ({work_placeholders}) OR ProdLine IN ({line_placeholders}))
  703. ORDER BY `{calendar_primary}`
  704. """,
  705. (
  706. SOURCE_TENANT,
  707. *(work_centers or [""]),
  708. *(lines or [""]),
  709. ),
  710. )
  711. seen: set[tuple[object, ...]] = set()
  712. for row in cur.fetchall():
  713. values = scoped_values(calendar_columns, tuple(row), target)
  714. key = (
  715. values.get("Domain"),
  716. values.get("Site"),
  717. values.get("WorkCtr"),
  718. values.get("WeekDay"),
  719. values.get("ProdLine"),
  720. )
  721. if key in seen:
  722. continue
  723. seen.add(key)
  724. cur.execute(
  725. """
  726. SELECT 1 FROM ShopCalendarWorkCtr
  727. WHERE tenant_id=%s
  728. AND COALESCE(Domain,'')=COALESCE(%s,'')
  729. AND COALESCE(Site,'')=COALESCE(%s,'')
  730. AND COALESCE(WorkCtr,'')=COALESCE(%s,'')
  731. AND COALESCE(WeekDay,'')=COALESCE(%s,'')
  732. AND COALESCE(ProdLine,'')=COALESCE(%s,'')
  733. LIMIT 1
  734. """,
  735. (
  736. target.tenant_id,
  737. values.get("Domain"),
  738. values.get("Site"),
  739. values.get("WorkCtr"),
  740. values.get("WeekDay"),
  741. values.get("ProdLine"),
  742. ),
  743. )
  744. if cur.fetchone():
  745. continue
  746. insert_values(cur, "ShopCalendarWorkCtr", calendar_columns, values)
  747. calendar_inserted += 1
  748. return {
  749. "departments_required": len(departments),
  750. "departments_inserted": department_inserted,
  751. "work_centers_required": len(work_centers),
  752. "work_centers_inserted": work_center_inserted,
  753. "operations_required": len(operations),
  754. "operations_inserted": operation_inserted,
  755. "lines_required": len(lines),
  756. "lines_inserted": line_inserted,
  757. "calendar_rows_inserted": calendar_inserted,
  758. }
  759. def execute_work_order_closure(
  760. cur: pymysql.cursors.Cursor, target: Target
  761. ) -> dict[str, int]:
  762. cur.execute(
  763. """
  764. INSERT INTO WorkOrdDetail
  765. (Domain,Site,LineNum,WorkOrd,Op,ItemNum,
  766. QtyRequired,QtytoIssue,QtyIssued,QtyPicked,
  767. WorkOrdMasterRecID,tenant_id)
  768. SELECT w.Domain,w.Site,
  769. ROW_NUMBER() OVER(PARTITION BY w.RecID ORDER BY b.RecID),
  770. w.WorkOrd,10,b.ComponentItem,
  771. COALESCE(NULLIF(b.Qty,0),1)*COALESCE(NULLIF(w.QtyOrded,0),1),
  772. COALESCE(NULLIF(b.Qty,0),1)*COALESCE(NULLIF(w.QtyOrded,0),1),
  773. 0,0,w.RecID,w.tenant_id
  774. FROM WorkOrdMaster w
  775. JOIN ProductStructureMaster b
  776. ON b.tenant_id=w.tenant_id AND b.ParentItem=w.ItemNum
  777. WHERE w.tenant_id=%s
  778. AND NOT EXISTS (
  779. SELECT 1 FROM WorkOrdDetail d
  780. WHERE d.tenant_id=w.tenant_id AND d.WorkOrdMasterRecID=w.RecID
  781. )
  782. """,
  783. (target.tenant_id,),
  784. )
  785. material_rows = int(cur.rowcount)
  786. cur.execute(
  787. """
  788. INSERT INTO WorkOrdRouting
  789. (StdOp,Domain,RunCrew,WorkOrd,OP,ItemNum,
  790. QtyComplete,QtyOrded,WorkCtr,ParentOp,
  791. WorkOrdMasterRecID,tenant_id)
  792. SELECT COALESCE(NULLIF(r.StdOp,''),r.OperationCode),
  793. LEFT(COALESCE(r.Domain,''),8),1,
  794. w.WorkOrd,
  795. CASE
  796. WHEN r.OperationCode REGEXP '^[0-9]+$'
  797. THEN CAST(r.OperationCode AS UNSIGNED)
  798. ELSE ROW_NUMBER() OVER(PARTITION BY w.RecID ORDER BY r.SortNo,r.RecID)*10
  799. END,
  800. w.ItemNum,0,COALESCE(NULLIF(w.QtyOrded,0),1),
  801. LEFT(COALESCE(r.WorkCenterCode,''),8),0,
  802. w.RecID,w.tenant_id
  803. FROM WorkOrdMaster w
  804. JOIN RoutingOpDetail r
  805. ON r.tenant_id=w.tenant_id AND r.MaterialCode=w.ItemNum
  806. WHERE w.tenant_id=%s
  807. AND NOT EXISTS (
  808. SELECT 1 FROM WorkOrdRouting wr
  809. WHERE wr.tenant_id=w.tenant_id AND wr.WorkOrdMasterRecID=w.RecID
  810. )
  811. """,
  812. (target.tenant_id,),
  813. )
  814. routing_rows = int(cur.rowcount)
  815. return {
  816. "work_order_material_rows_inserted": material_rows,
  817. "work_order_routing_rows_inserted": routing_rows,
  818. }
  819. def execute_master_closure(target: Target) -> dict[str, object]:
  820. conn = load_connection()
  821. try:
  822. with conn.cursor() as cur:
  823. ensure_audit_tables(cur)
  824. cur.execute(
  825. """
  826. SELECT DISTINCT item_number FROM crm_seorderentry
  827. WHERE tenant_id=%s AND IFNULL(item_number,'')<>''
  828. UNION
  829. SELECT DISTINCT product_code FROM mes_morder
  830. WHERE tenant_id=%s AND IFNULL(product_code,'')<>''
  831. UNION
  832. SELECT DISTINCT icitem_name FROM srm_pr_main
  833. WHERE tenant_id=%s AND IFNULL(icitem_name,'')<>''
  834. """,
  835. (target.tenant_id, target.tenant_id, target.tenant_id),
  836. )
  837. items = [str(row[0]) for row in cur.fetchall()]
  838. cur.execute(
  839. "SELECT DISTINCT custom_no FROM crm_seorder WHERE tenant_id=%s AND IFNULL(custom_no,'')<>''",
  840. (target.tenant_id,),
  841. )
  842. customers = [str(row[0]) for row in cur.fetchall()]
  843. cur.execute(
  844. "SELECT DISTINCT Supp FROM PurOrdMaster WHERE tenant_id=%s AND IFNULL(Supp,'')<>''",
  845. (target.tenant_id,),
  846. )
  847. suppliers = [str(row[0]) for row in cur.fetchall()]
  848. result = {
  849. "target": target.code,
  850. "items_required": len(items),
  851. "items_inserted": ensure_master_rows(
  852. cur, target, "ItemMaster", "ItemNum", items, fallback=True
  853. ),
  854. "customers_required": len(customers),
  855. "customers_inserted": ensure_master_rows(
  856. cur, target, "CustMaster", "Cust", customers, fallback=True
  857. ),
  858. "suppliers_required": len(suppliers),
  859. "suppliers_inserted": ensure_master_rows(
  860. cur, target, "SuppMaster", "Supp", suppliers, fallback=True
  861. ),
  862. }
  863. result.update(execute_structure_closure(cur, target, items))
  864. result.update(execute_resource_closure(cur, target))
  865. result.update(execute_work_order_closure(cur, target))
  866. return result
  867. finally:
  868. conn.close()
  869. def main() -> None:
  870. parser = argparse.ArgumentParser()
  871. parser.add_argument("targets", nargs="+", choices=TARGETS.keys())
  872. parser.add_argument("--dry-run", action="store_true")
  873. parser.add_argument("--master-closure", action="store_true")
  874. parser.add_argument("--evidence", type=Path)
  875. args = parser.parse_args()
  876. if args.master_closure:
  877. results = [execute_master_closure(TARGETS[code]) for code in args.targets]
  878. else:
  879. results = [execute_target(TARGETS[code], args.dry_run) for code in args.targets]
  880. payload = {"executed_at": datetime.now().isoformat(timespec="seconds"), "results": results}
  881. if args.evidence:
  882. args.evidence.parent.mkdir(parents=True, exist_ok=True)
  883. args.evidence.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
  884. print(json.dumps(payload, ensure_ascii=False, indent=2))
  885. if __name__ == "__main__":
  886. main()