#!/usr/bin/env python3 """Build and execute tenant-scoped UAT migration batches. The source SQL files contain extracted, anonymized S1-S4 chains. This runner assigns target-specific IDs and business keys, persists old/new mappings, and keeps credentials outside the repository by reading the local App.json. """ from __future__ import annotations import argparse import json import re from dataclasses import dataclass from datetime import datetime from pathlib import Path from typing import Iterable import pymysql from pymysql.constants import CLIENT ROOT = Path(__file__).resolve().parents[4] SQL_ROOT = ROOT / "doc" / "plan" / "sql" DATABASE_JSON = ROOT / "server" / "Admin.NET.Application" / "Configuration" / "Database.json" SOURCE_TENANT = 797403760988229 @dataclass(frozen=True) class Target: code: str tenant_id: int factory_id: int operator_id: int source_file: str source_batch: str target_batch: str id_base: int key_prefix: str TARGETS = { "A": Target( "A", 838257186181189, 838257186320453, 838257187360837, "S1S4_UAT_P0_execute_write.sql", "S1S4_UAT_20260604_RQ_V1", "UAT_WP1_A_20260817_V1", 9206081700000001, "UATA", ), "B": Target( "B", 838257212780613, 838257212858437, 838257213620293, "S1S4_UAT_P0_execute_write.sql", "S1S4_UAT_20260604_RQ_V1", "UAT_WP1_B_20260817_V1", 9206081800000001, "UATB", ), "DEMO": Target( "DEMO", 838257237606469, 838257237676101, 838257238302789, "S1S4_UAT_FULL_execute_write.sql", "S1S4_UAT_FULL_20260605_V1", "UAT_WP1_DEMO_20260817_V1", 9206081900000001, "DEMO", ), } def walk_strings(value: object) -> Iterable[str]: if isinstance(value, str): yield value elif isinstance(value, dict): for child in value.values(): yield from walk_strings(child) elif isinstance(value, list): for child in value: yield from walk_strings(child) def load_connection() -> pymysql.Connection: raw = DATABASE_JSON.read_text(encoding="utf-8-sig") candidates = [ s for s in re.findall(r'(?m)^\s*"ConnectionString"\s*:\s*"([^"]+)"', raw) if "Database=aidopdev" in s ] if not candidates: raise RuntimeError("Database.json 中未找到启用的 aidopdev 连接串") parts = {} for pair in candidates[0].split(";"): if "=" in pair: key, val = pair.split("=", 1) parts[key.strip().lower()] = val.strip() return pymysql.connect( host=parts.get("server", "127.0.0.1"), port=int(parts.get("port", "3306")), user=parts.get("uid") or parts.get("user id"), password=parts.get("pwd") or parts.get("password"), database=parts.get("database", "aidopdev"), charset="utf8mb4", autocommit=True, client_flag=CLIENT.MULTI_STATEMENTS, connect_timeout=20, read_timeout=600, write_timeout=600, ) def ordered_tokens(sql: str, pattern: str) -> list[str]: return list(dict.fromkeys(re.findall(pattern, sql))) def build_key_map(sql: str, target: Target) -> dict[str, tuple[str, str]]: mapping: dict[str, tuple[str, str]] = {} orders = ordered_tokens(sql, r"'((?:MPO|SO)\d{10,})'") works = ordered_tokens(sql, r"'(M\d{9,12})'") prs = ordered_tokens(sql, r"'(PR-UAT-[^']+)'") pos = ordered_tokens(sql, r"'(PO-UAT-[^']+)'") ships = ordered_tokens(sql, r"'(SH-UAT-[^']+)'") for idx, old in enumerate(orders, 1): scenario = f"A{idx:02d}" if target.code == "A" else f"CP{idx:02d}" new = ( f"UATA-{scenario}-SO" if target.code == "A" else f"UATB-{scenario}-SO" if target.code == "B" else f"DEMO-SO-{idx:03d}" ) mapping[old] = ("SALES_ORDER", new) for idx, old in enumerate(works, 1): new = f"{target.key_prefix}-WO-{idx:03d}" mapping[old] = ("WORK_ORDER", new) for kind, tokens, label in ( ("PURCHASE_REQUEST", prs, "PR"), ("PURCHASE_ORDER", pos, "PO"), ("SUPPLIER_SHIPMENT", ships, "SH"), ): for idx, old in enumerate(tokens, 1): mapping[old] = (kind, f"{target.key_prefix}-{label}-{idx:03d}") return mapping def patch_tenant_columns(sql: str) -> str: sql = sql.replace( "need_number, morder_production_number, IsDeleted,\n" " create_by_name, create_time, update_by_name, update_time,\n" " tenant_id, factory_id\n" ")\n" "SELECT target_mo_id, morder_no, @now, 'Released', product_code, product_code,\n" " need_number, need_number, 0,", "need_number, morder_production_number, IsDeleted, urgent,\n" " create_by_name, create_time, update_by_name, update_time,\n" " tenant_id, factory_id\n" ")\n" "SELECT target_mo_id, morder_no, @now, 'Released', product_code, product_code,\n" " need_number, need_number, 0, 0,", ) sql = sql.replace( "INSERT INTO scm_shd (id, po_billno, shddh, sh_purchase_num, " "estimated_delivery_date, tjrxm, tjrq, state)", "INSERT INTO scm_shd (id, po_billno, shddh, sh_purchase_num, " "estimated_delivery_date, tjrxm, tjrq, state, tenant_id)", ) sql = sql.replace( "DATE_FORMAT(@now,'%Y-%m-%d'), 0\nFROM tmp_s4_uat_ship", "DATE_FORMAT(@now,'%Y-%m-%d'), 0, @tenant_id\nFROM tmp_s4_uat_ship", ) sql = sql.replace( "INSERT INTO scm_shdzb (id, glid, sh_material_code, " "sh_delivery_quantity, po_bill, po_billline)", "INSERT INTO scm_shdzb (id, glid, sh_material_code, " "sh_delivery_quantity, po_bill, po_billline, tenant_id)", ) sql = sql.replace( "pl.item_num, pl.qty, pl.pur_ord, CAST(pl.line_no AS CHAR)\n" "FROM tmp_s3_uat_po_line", "pl.item_num, pl.qty, pl.pur_ord, CAST(pl.line_no AS CHAR), @tenant_id\n" "FROM tmp_s3_uat_po_line", ) sql += """ -- Keep both work-order representations aligned with the linked sales-order item. UPDATE mes_morder m JOIN mes_moentry me ON me.tenant_id=m.tenant_id AND me.moentry_mono=m.morder_no JOIN crm_seorderentry e ON e.tenant_id=me.tenant_id AND e.Id=me.soentry_id SET m.product_code=e.item_number, m.product_name=COALESCE(NULLIF(e.item_name,''),e.item_number) WHERE m.tenant_id=@tenant_id; UPDATE WorkOrdMaster w JOIN mes_moentry me ON me.tenant_id=w.tenant_id AND me.moentry_mono=w.WorkOrd JOIN crm_seorderentry e ON e.tenant_id=me.tenant_id AND e.Id=me.soentry_id SET w.ItemNum=e.item_number WHERE w.tenant_id=@tenant_id; UPDATE PurOrdDetail d JOIN PurOrdMaster h ON h.tenant_id=d.tenant_id AND h.PurOrd=d.PurOrd SET d.PurOrdRecID=h.RecID WHERE d.tenant_id=@tenant_id AND d.PurOrdRecID=0; """ return sql def split_statements(sql: str) -> list[str]: """Split these migration files without breaking on comment semicolons.""" body = "\n".join( line for line in sql.splitlines() if not line.lstrip().startswith("--") ) return [statement.strip() for statement in body.split(";") if statement.strip()] def adapt_sql(target: Target) -> tuple[str, dict[str, tuple[str, str]]]: source = (SQL_ROOT / target.source_file).read_text(encoding="utf-8-sig") mapping = build_key_map(source, target) sql = source sql = re.sub(r"SET @tenant_id\s*:=\s*\d+;", f"SET @tenant_id := {target.tenant_id};", sql) sql = re.sub( r"SET @factory_id\s*:=\s*\d+;", f"SET @factory_id := {target.factory_id};", sql, ) sql = re.sub(r"SET @operator_id\s*:=\s*\d+;", f"SET @operator_id := {target.operator_id};", sql) sql = re.sub( r"SET @operator_name\s*:=\s*'[^']*';", f"SET @operator_name := 'UAT-{target.code}-Migration';", sql, ) sql = re.sub(r"SET @batch_no\s*:=\s*'[^']+';", f"SET @batch_no := '{target.target_batch}';", sql) sql = re.sub(r"SET @id_base\s*:=\s*\d+;", f"SET @id_base := {target.id_base};", sql) # FULL uses literal IDs in its temporary mapping rows. if target.source_file.startswith("S1S4_UAT_FULL"): sql = sql.replace("91060005", str(target.id_base)[:8]) # The historical template was slimmed; use an extant anonymized source PO. sql = sql.replace( "WHERE PurOrd='PW202605100002' AND tenant_id=@tenant_id", f"WHERE PurOrd='DO20260711001' AND tenant_id={SOURCE_TENANT}", ) for old, (_, new) in sorted(mapping.items(), key=lambda item: len(item[0]), reverse=True): sql = sql.replace(f"'{old}'", f"'{new}'") sql = patch_tenant_columns(sql) return sql, mapping def ensure_audit_tables(cur: pymysql.cursors.Cursor) -> None: cur.execute( """ CREATE TABLE IF NOT EXISTS uat_data_migration_batch ( batch_code varchar(100) NOT NULL PRIMARY KEY, target_tenant_id bigint NOT NULL, target_set varchar(16) NOT NULL, source_file varchar(255) NOT NULL, status varchar(20) NOT NULL, started_at datetime(3) NOT NULL, finished_at datetime(3) NULL, summary_json json NULL ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 """ ) cur.execute( """ CREATE TABLE IF NOT EXISTS uat_data_migration_map ( batch_code varchar(100) NOT NULL, object_type varchar(40) NOT NULL, source_key varchar(200) NOT NULL, target_key varchar(200) NOT NULL, target_tenant_id bigint NOT NULL, created_at datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), PRIMARY KEY(batch_code, object_type, source_key), KEY ix_uat_map_target(target_tenant_id, object_type, target_key) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 """ ) def execute_target(target: Target, dry_run: bool) -> dict[str, object]: sql, mapping = adapt_sql(target) generated = Path(__file__).with_name(f"{target.target_batch}.sql") generated.write_text(sql, encoding="utf-8") if dry_run: return {"target": target.code, "generated": str(generated), "mapping_count": len(mapping)} conn = load_connection() try: with conn.cursor() as cur: ensure_audit_tables(cur) cur.execute( "SELECT status FROM uat_data_migration_batch WHERE batch_code=%s", (target.target_batch,), ) existing = cur.fetchone() if existing and existing[0] == "SUCCESS": raise RuntimeError(f"{target.target_batch} 已成功执行,拒绝重复导入") cur.execute( """ INSERT INTO uat_data_migration_batch (batch_code,target_tenant_id,target_set,source_file,status,started_at) VALUES(%s,%s,%s,%s,'RUNNING',NOW(3)) ON DUPLICATE KEY UPDATE status='RUNNING',started_at=NOW(3),finished_at=NULL """, (target.target_batch, target.tenant_id, target.code, target.source_file), ) try: for index, statement in enumerate(split_statements(sql), 1): try: cur.execute(statement) except Exception as exc: cur.execute("ROLLBACK") excerpt = " ".join(statement.split())[:240] raise RuntimeError( f"SQL #{index} 执行失败: {excerpt}: {exc}" ) from exc cur.executemany( """ INSERT INTO uat_data_migration_map (batch_code,object_type,source_key,target_key,target_tenant_id) VALUES(%s,%s,%s,%s,%s) ON DUPLICATE KEY UPDATE target_key=VALUES(target_key) """, [ (target.target_batch, kind, old, new, target.tenant_id) for old, (kind, new) in mapping.items() ], ) summary = verify(cur, target) cur.execute( """ UPDATE uat_data_migration_batch SET status='SUCCESS',finished_at=NOW(3),summary_json=%s WHERE batch_code=%s """, (json.dumps(summary, ensure_ascii=False), target.target_batch), ) return summary except Exception as exc: cur.execute( """ UPDATE uat_data_migration_batch SET status='FAILED',finished_at=NOW(3),summary_json=%s WHERE batch_code=%s """, (json.dumps({"error": str(exc)}, ensure_ascii=False), target.target_batch), ) raise finally: conn.close() def verify(cur: pymysql.cursors.Cursor, target: Target) -> dict[str, object]: queries = { "orders": ( "SELECT COUNT(*) FROM crm_seorder WHERE tenant_id=%s AND bill_from LIKE %s", (target.tenant_id, f"%{target.target_batch}%"), ), "order_lines": ( """ SELECT COUNT(*) FROM crm_seorderentry e JOIN crm_seorder h ON h.Id=e.seorder_id AND h.tenant_id=e.tenant_id WHERE h.tenant_id=%s AND h.bill_from LIKE %s """, (target.tenant_id, f"%{target.target_batch}%"), ), "work_orders": ( "SELECT COUNT(*) FROM mes_morder WHERE tenant_id=%s AND morder_no LIKE %s", (target.tenant_id, f"{target.key_prefix}-WO-%"), ), "purchase_requests": ( "SELECT COUNT(*) FROM srm_pr_main WHERE tenant_id=%s AND pr_billno LIKE %s", (target.tenant_id, f"{target.key_prefix}-PR-%"), ), "purchase_orders": ( "SELECT COUNT(*) FROM PurOrdMaster WHERE tenant_id=%s AND PurOrd LIKE %s", (target.tenant_id, f"{target.key_prefix}-PO-%"), ), "supplier_shipments": ( "SELECT COUNT(*) FROM scm_shd WHERE tenant_id=%s AND shddh LIKE %s", (target.tenant_id, f"{target.key_prefix}-SH-%"), ), } result: dict[str, object] = {"target": target.code, "tenant_id": target.tenant_id} for key, (query, params) in queries.items(): cur.execute(query, params) result[key] = int(cur.fetchone()[0]) return result def table_columns(cur: pymysql.cursors.Cursor, table: str) -> tuple[list[str], str]: cur.execute(f"SHOW COLUMNS FROM `{table}`") rows = cur.fetchall() writable = [row[0] for row in rows if "auto_increment" not in (row[5] or "")] primary = next(row[0] for row in rows if row[3] == "PRI") return writable, primary def ensure_master_rows( cur: pymysql.cursors.Cursor, target: Target, table: str, key_column: str, keys: list[str], fallback: bool = False, ) -> int: columns, primary = table_columns(cur, table) inserted = 0 for key in keys: cur.execute( f"SELECT `{primary}` FROM `{table}` WHERE tenant_id=%s AND `{key_column}`=%s LIMIT 1", (target.tenant_id, key), ) if cur.fetchone(): continue cur.execute( f""" SELECT {",".join(f"`{column}`" for column in columns)} FROM `{table}` WHERE tenant_id IN (%s,%s) AND `{key_column}`=%s ORDER BY FIELD(tenant_id,%s,%s) LIMIT 1 """, (SOURCE_TENANT, 824585161322565, key, SOURCE_TENANT, 824585161322565), ) row = cur.fetchone() if row is None and fallback: cur.execute( f""" SELECT {",".join(f"`{column}`" for column in columns)} FROM `{table}` WHERE tenant_id=%s ORDER BY `{primary}` LIMIT 1 """, (SOURCE_TENANT,), ) row = cur.fetchone() if row is None: raise RuntimeError(f"{table}.{key_column}={key} 无可用脱敏源记录") values = dict(zip(columns, row)) values[key_column] = key for column in ("tenant_id", "TenantId"): if column in values: values[column] = target.tenant_id for column in ( "company_ref_id", "factory_ref_id", "CompanyRefId", "FactoryRefId", "company_id", "factory_id", "org_id", ): if column in values: values[column] = target.factory_id cur.execute( f""" INSERT INTO `{table}` ({",".join(f"`{column}`" for column in columns)}) VALUES ({",".join(["%s"] * len(columns))}) """, tuple(values[column] for column in columns), ) new_id = cur.lastrowid cur.execute( """ INSERT INTO uat_data_migration_map (batch_code,object_type,source_key,target_key,target_tenant_id) VALUES(%s,%s,%s,%s,%s) ON DUPLICATE KEY UPDATE target_key=VALUES(target_key) """, ( target.target_batch, f"MASTER_{table.upper()}", key, f"{key}:{new_id}", target.tenant_id, ), ) inserted += 1 return inserted def scoped_values( columns: list[str], row: tuple[object, ...], target: Target ) -> dict[str, object]: values = dict(zip(columns, row)) for column in ("tenant_id", "TenantId"): if column in values: values[column] = target.tenant_id for column in ( "company_ref_id", "factory_ref_id", "CompanyRefId", "FactoryRefId", "company_id", "factory_id", "org_id", ): if column in values: values[column] = target.factory_id return values def insert_values( cur: pymysql.cursors.Cursor, table: str, columns: list[str], values: dict[str, object], ) -> int: cur.execute( f""" INSERT INTO `{table}` ({",".join(f"`{column}`" for column in columns)}) VALUES ({",".join(["%s"] * len(columns))}) """, tuple(values[column] for column in columns), ) return int(cur.lastrowid) def execute_structure_closure( cur: pymysql.cursors.Cursor, target: Target, initial_items: list[str] ) -> dict[str, int]: # Expand one BOM level so every copied BOM component also has an item master. placeholders = ",".join(["%s"] * len(initial_items)) component_items: set[str] = set() if initial_items: cur.execute( f""" SELECT DISTINCT ComponentItem FROM ProductStructureMaster WHERE tenant_id IN (%s,%s) AND ParentItem IN ({placeholders}) AND IFNULL(ComponentItem,'')<>'' """, (SOURCE_TENANT, 824585161322565, *initial_items), ) component_items = {str(row[0]) for row in cur.fetchall()} expanded_items = sorted(set(initial_items) | component_items) item_inserted = ensure_master_rows( cur, target, "ItemMaster", "ItemNum", expanded_items, fallback=True ) cur.execute( "SELECT ItemNum,RecID FROM ItemMaster WHERE tenant_id=%s", (target.tenant_id,), ) target_item_ids = {str(row[0]): int(row[1]) for row in cur.fetchall()} bom_columns, bom_primary = table_columns(cur, "ProductStructureMaster") bom_maps: dict[int, int] = {} seen_bom: set[tuple[object, ...]] = set() bom_inserted = 0 if initial_items: for source_tenant in (SOURCE_TENANT, 824585161322565): cur.execute( f""" SELECT `{bom_primary}`,{",".join(f"`{column}`" for column in bom_columns)} FROM ProductStructureMaster WHERE tenant_id=%s AND ParentItem IN ({placeholders}) ORDER BY `{bom_primary}` """, (source_tenant, *initial_items), ) for source_id, *row in cur.fetchall(): values = scoped_values(bom_columns, tuple(row), target) key = ( values.get("ParentItem"), values.get("ComponentItem"), target.factory_id, ) if key in seen_bom: continue seen_bom.add(key) if "ParentMaterialId" in values: values["ParentMaterialId"] = target_item_ids.get( str(values.get("ParentItem")) ) if "ComponentMaterialId" in values: values["ComponentMaterialId"] = target_item_ids.get( str(values.get("ComponentItem")) ) inserted = False try: target_id = insert_values( cur, "ProductStructureMaster", bom_columns, values ) inserted = True except pymysql.IntegrityError as exc: if exc.args[0] != 1062: raise cur.execute( """ SELECT RecID FROM ProductStructureMaster WHERE tenant_id=%s AND CompanyRefId=%s AND FactoryRefId=%s AND ParentItem=%s AND ComponentItem=%s LIMIT 1 """, ( target.tenant_id, target.factory_id, target.factory_id, values.get("ParentItem"), values.get("ComponentItem"), ), ) existing = cur.fetchone() if not existing: raise target_id = int(existing[0]) bom_maps[int(source_id)] = target_id if inserted: bom_inserted += 1 op_inserted = 0 if bom_maps: op_columns, op_primary = table_columns(cur, "ProductStructureOp") ids = list(bom_maps) for start in range(0, len(ids), 500): chunk = ids[start : start + 500] cur.execute( f""" SELECT `{op_primary}`,{",".join(f"`{column}`" for column in op_columns)} FROM ProductStructureOp WHERE ProductStructureMasterId IN ({",".join(["%s"] * len(chunk))}) """, tuple(chunk), ) for _, *row in cur.fetchall(): values = scoped_values(op_columns, tuple(row), target) old_master = int(values["ProductStructureMasterId"]) values["ProductStructureMasterId"] = bom_maps[old_master] try: insert_values(cur, "ProductStructureOp", op_columns, values) op_inserted += 1 except pymysql.IntegrityError as exc: if exc.args[0] != 1062: raise routing_columns, routing_primary = table_columns(cur, "RoutingOpDetail") routing_inserted = 0 seen_routing: set[tuple[object, ...]] = set() if expanded_items: routing_placeholders = ",".join(["%s"] * len(expanded_items)) for source_tenant in (SOURCE_TENANT, 824585161322565): cur.execute( f""" SELECT `{routing_primary}`,{",".join(f"`{column}`" for column in routing_columns)} FROM RoutingOpDetail WHERE tenant_id=%s AND MaterialCode IN ({routing_placeholders}) ORDER BY `{routing_primary}` """, (source_tenant, *expanded_items), ) for _, *row in cur.fetchall(): values = scoped_values(routing_columns, tuple(row), target) key = tuple( values.get(column) for column in ("RouteCode", "MaterialCode", "OperationCode", "SortNo") ) if key in seen_routing: continue seen_routing.add(key) try: insert_values(cur, "RoutingOpDetail", routing_columns, values) routing_inserted += 1 except pymysql.IntegrityError as exc: if exc.args[0] != 1062: raise return { "bom_component_items": len(component_items), "bom_component_items_inserted": item_inserted, "bom_rows_inserted": bom_inserted, "bom_operations_inserted": op_inserted, "routing_rows_inserted": routing_inserted, } def execute_resource_closure( cur: pymysql.cursors.Cursor, target: Target ) -> dict[str, int]: cur.execute( """ SELECT Department FROM DepartmentMaster WHERE tenant_id=%s AND IFNULL(Department,'')<>'' GROUP BY Department ORDER BY MIN(RecID) LIMIT 20 """, (SOURCE_TENANT,), ) departments = [str(row[0]) for row in cur.fetchall()] department_inserted = ensure_master_rows( cur, target, "DepartmentMaster", "Department", departments, fallback=True, ) cur.execute( """ SELECT DISTINCT WorkCenterCode FROM RoutingOpDetail WHERE tenant_id=%s AND IFNULL(WorkCenterCode,'')<>'' """, (target.tenant_id,), ) work_centers = [str(row[0]) for row in cur.fetchall()] cur.execute( """ SELECT DISTINCT COALESCE(NULLIF(StdOp,''),NULLIF(OperationCode,'')) FROM RoutingOpDetail WHERE tenant_id=%s AND COALESCE(NULLIF(StdOp,''),NULLIF(OperationCode,'')) IS NOT NULL """, (target.tenant_id,), ) operations = [str(row[0]) for row in cur.fetchall()] work_center_inserted = ensure_master_rows( cur, target, "WorkCtrMaster", "WorkCtr", work_centers, fallback=True ) operation_inserted = ensure_master_rows( cur, target, "StdOpMaster", "StdOp", operations, fallback=True ) cur.execute( """ SELECT Line FROM LineMaster WHERE tenant_id=%s AND IFNULL(Line,'')<>'' GROUP BY Line ORDER BY MIN(RecID) LIMIT 20 """, (SOURCE_TENANT,), ) lines = [str(row[0]) for row in cur.fetchall()] line_inserted = ensure_master_rows( cur, target, "LineMaster", "Line", lines, fallback=True ) calendar_columns, calendar_primary = table_columns(cur, "ShopCalendarWorkCtr") calendar_inserted = 0 if work_centers or lines: work_placeholders = ",".join(["%s"] * max(1, len(work_centers))) line_placeholders = ",".join(["%s"] * max(1, len(lines))) cur.execute( f""" SELECT {",".join(f"`{column}`" for column in calendar_columns)} FROM ShopCalendarWorkCtr WHERE tenant_id=%s AND (WorkCtr IN ({work_placeholders}) OR ProdLine IN ({line_placeholders})) ORDER BY `{calendar_primary}` """, ( SOURCE_TENANT, *(work_centers or [""]), *(lines or [""]), ), ) seen: set[tuple[object, ...]] = set() for row in cur.fetchall(): values = scoped_values(calendar_columns, tuple(row), target) key = ( values.get("Domain"), values.get("Site"), values.get("WorkCtr"), values.get("WeekDay"), values.get("ProdLine"), ) if key in seen: continue seen.add(key) cur.execute( """ SELECT 1 FROM ShopCalendarWorkCtr WHERE tenant_id=%s AND COALESCE(Domain,'')=COALESCE(%s,'') AND COALESCE(Site,'')=COALESCE(%s,'') AND COALESCE(WorkCtr,'')=COALESCE(%s,'') AND COALESCE(WeekDay,'')=COALESCE(%s,'') AND COALESCE(ProdLine,'')=COALESCE(%s,'') LIMIT 1 """, ( target.tenant_id, values.get("Domain"), values.get("Site"), values.get("WorkCtr"), values.get("WeekDay"), values.get("ProdLine"), ), ) if cur.fetchone(): continue insert_values(cur, "ShopCalendarWorkCtr", calendar_columns, values) calendar_inserted += 1 return { "departments_required": len(departments), "departments_inserted": department_inserted, "work_centers_required": len(work_centers), "work_centers_inserted": work_center_inserted, "operations_required": len(operations), "operations_inserted": operation_inserted, "lines_required": len(lines), "lines_inserted": line_inserted, "calendar_rows_inserted": calendar_inserted, } def execute_work_order_closure( cur: pymysql.cursors.Cursor, target: Target ) -> dict[str, int]: cur.execute( """ INSERT INTO WorkOrdDetail (Domain,Site,LineNum,WorkOrd,Op,ItemNum, QtyRequired,QtytoIssue,QtyIssued,QtyPicked, WorkOrdMasterRecID,tenant_id) SELECT w.Domain,w.Site, ROW_NUMBER() OVER(PARTITION BY w.RecID ORDER BY b.RecID), w.WorkOrd,10,b.ComponentItem, COALESCE(NULLIF(b.Qty,0),1)*COALESCE(NULLIF(w.QtyOrded,0),1), COALESCE(NULLIF(b.Qty,0),1)*COALESCE(NULLIF(w.QtyOrded,0),1), 0,0,w.RecID,w.tenant_id FROM WorkOrdMaster w JOIN ProductStructureMaster b ON b.tenant_id=w.tenant_id AND b.ParentItem=w.ItemNum WHERE w.tenant_id=%s AND NOT EXISTS ( SELECT 1 FROM WorkOrdDetail d WHERE d.tenant_id=w.tenant_id AND d.WorkOrdMasterRecID=w.RecID ) """, (target.tenant_id,), ) material_rows = int(cur.rowcount) cur.execute( """ INSERT INTO WorkOrdRouting (StdOp,Domain,RunCrew,WorkOrd,OP,ItemNum, QtyComplete,QtyOrded,WorkCtr,ParentOp, WorkOrdMasterRecID,tenant_id) SELECT COALESCE(NULLIF(r.StdOp,''),r.OperationCode), LEFT(COALESCE(r.Domain,''),8),1, w.WorkOrd, CASE WHEN r.OperationCode REGEXP '^[0-9]+$' THEN CAST(r.OperationCode AS UNSIGNED) ELSE ROW_NUMBER() OVER(PARTITION BY w.RecID ORDER BY r.SortNo,r.RecID)*10 END, w.ItemNum,0,COALESCE(NULLIF(w.QtyOrded,0),1), LEFT(COALESCE(r.WorkCenterCode,''),8),0, w.RecID,w.tenant_id FROM WorkOrdMaster w JOIN RoutingOpDetail r ON r.tenant_id=w.tenant_id AND r.MaterialCode=w.ItemNum WHERE w.tenant_id=%s AND NOT EXISTS ( SELECT 1 FROM WorkOrdRouting wr WHERE wr.tenant_id=w.tenant_id AND wr.WorkOrdMasterRecID=w.RecID ) """, (target.tenant_id,), ) routing_rows = int(cur.rowcount) return { "work_order_material_rows_inserted": material_rows, "work_order_routing_rows_inserted": routing_rows, } def execute_master_closure(target: Target) -> dict[str, object]: conn = load_connection() try: with conn.cursor() as cur: ensure_audit_tables(cur) cur.execute( """ SELECT DISTINCT item_number FROM crm_seorderentry WHERE tenant_id=%s AND IFNULL(item_number,'')<>'' UNION SELECT DISTINCT product_code FROM mes_morder WHERE tenant_id=%s AND IFNULL(product_code,'')<>'' UNION SELECT DISTINCT icitem_name FROM srm_pr_main WHERE tenant_id=%s AND IFNULL(icitem_name,'')<>'' """, (target.tenant_id, target.tenant_id, target.tenant_id), ) items = [str(row[0]) for row in cur.fetchall()] cur.execute( "SELECT DISTINCT custom_no FROM crm_seorder WHERE tenant_id=%s AND IFNULL(custom_no,'')<>''", (target.tenant_id,), ) customers = [str(row[0]) for row in cur.fetchall()] cur.execute( "SELECT DISTINCT Supp FROM PurOrdMaster WHERE tenant_id=%s AND IFNULL(Supp,'')<>''", (target.tenant_id,), ) suppliers = [str(row[0]) for row in cur.fetchall()] result = { "target": target.code, "items_required": len(items), "items_inserted": ensure_master_rows( cur, target, "ItemMaster", "ItemNum", items, fallback=True ), "customers_required": len(customers), "customers_inserted": ensure_master_rows( cur, target, "CustMaster", "Cust", customers, fallback=True ), "suppliers_required": len(suppliers), "suppliers_inserted": ensure_master_rows( cur, target, "SuppMaster", "Supp", suppliers, fallback=True ), } result.update(execute_structure_closure(cur, target, items)) result.update(execute_resource_closure(cur, target)) result.update(execute_work_order_closure(cur, target)) return result finally: conn.close() def main() -> None: parser = argparse.ArgumentParser() parser.add_argument("targets", nargs="+", choices=TARGETS.keys()) parser.add_argument("--dry-run", action="store_true") parser.add_argument("--master-closure", action="store_true") parser.add_argument("--evidence", type=Path) args = parser.parse_args() if args.master_closure: results = [execute_master_closure(TARGETS[code]) for code in args.targets] else: results = [execute_target(TARGETS[code], args.dry_run) for code in args.targets] payload = {"executed_at": datetime.now().isoformat(timespec="seconds"), "results": results} if args.evidence: args.evidence.parent.mkdir(parents=True, exist_ok=True) args.evidence.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") print(json.dumps(payload, ensure_ascii=False, indent=2)) if __name__ == "__main__": main()