| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951 |
- #!/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()
|