| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270 |
- """DB_SYNC 模拟源库(任务书 P1-C)。
- 职责边界:
- - 只在独立 MySQL 模拟 schema(默认 aidop_integration_sim)建表、装载、清理;
- - 绝不连接 Ai-DOP 库写贴源/标准/KPI 表;库名含 aidop 且非模拟库名时硬阻断;
- - 不使用 Ai-DOP 主库账号;连接凭据从环境变量读取(密码 AIDOP_SIM_MYSQL_PASSWORD);
- - 生成 Ai-DOP 配置清单(mdp_source / mdp_entity 建议行),仅供人工在配置页登记,不自动写入。
- """
- from __future__ import annotations
- import json
- import os
- from datetime import datetime
- # 允许创建的模拟库名;其余含 aidop 的库名一律阻断
- SIM_DB_DEFAULT = "aidop_integration_sim"
- _BLOCKED_DB_HINTS = ("aidopdev", "aidop_uat", "aidopprod", "aidop_prod", "admin", "mdp")
- class BlockedDatabaseError(RuntimeError):
- pass
- def assert_sim_database(database: str) -> None:
- db = (database or "").strip().lower()
- if not db:
- raise BlockedDatabaseError("database name required")
- if db != SIM_DB_DEFAULT and any(h in db for h in _BLOCKED_DB_HINTS):
- raise BlockedDatabaseError(
- f"database '{database}' looks like an Ai-DOP database; simulator may only use '{SIM_DB_DEFAULT}'")
- if db != SIM_DB_DEFAULT and not db.startswith("aidop_integration_sim"):
- # 允许 aidop_integration_sim_* 派生测试库,其余含 aidop 的阻断
- if "aidop" in db:
- raise BlockedDatabaseError(
- f"database '{database}' contains 'aidop' but is not the simulator schema")
- def mysql_conn_params(database: str | None = None) -> dict:
- db = database or os.environ.get("AIDOP_SIM_MYSQL_DATABASE", SIM_DB_DEFAULT)
- assert_sim_database(db)
- return {
- "host": os.environ.get("AIDOP_SIM_MYSQL_HOST", "127.0.0.1"),
- "port": int(os.environ.get("AIDOP_SIM_MYSQL_PORT", "3306")),
- "user": os.environ.get("AIDOP_SIM_MYSQL_USER", "root"),
- "password": os.environ.get("AIDOP_SIM_MYSQL_PASSWORD", ""),
- "database": db,
- "charset": "utf8mb4",
- }
- def connect(database: str | None = None):
- import pymysql # 延迟导入:单元测试不依赖真实驱动
- params = mysql_conn_params(database)
- return pymysql.connect(cursorclass=pymysql.cursors.DictCursor, autocommit=True, **params)
- # ---------------------------------------------------------------------------
- # DDL 生成(纯函数,可单测)
- # ---------------------------------------------------------------------------
- def infer_mysql_type(values: list) -> str:
- has_text = has_int = has_float = has_bool = False
- max_len = 0
- for v in values:
- if v is None:
- continue
- if isinstance(v, bool):
- has_bool = True
- elif isinstance(v, int):
- has_int = True
- elif isinstance(v, float):
- has_float = True
- else:
- has_text = True
- max_len = max(max_len, len(str(v)))
- if has_text:
- return f"VARCHAR({max(64, min(512, (max_len // 32 + 1) * 32))})"
- if has_bool:
- return "TINYINT(1)"
- if has_float:
- return "DECIMAL(18,6)"
- if has_int:
- return "BIGINT"
- return "VARCHAR(191)"
- def generate_ddl(table: str, rows: list[dict], business_key: list[str] | None = None,
- increment_column: str = "sourceUpdatedAt") -> str:
- """按样例行推断列类型;业务键唯一索引;增量列建普通索引。列名保持样例原样(反引号包裹)。"""
- if not rows:
- raise ValueError(f"cannot generate DDL for {table} from empty rows")
- business_key = business_key or ["bizKey"]
- columns: dict[str, list] = {}
- for row in rows:
- for k, v in row.items():
- columns.setdefault(k, []).append(v)
- lines = []
- for name, values in columns.items():
- lines.append(f" `{name}` {infer_mysql_type(values)} NULL")
- uk = ", ".join(f"`{k}`" for k in business_key if k in columns)
- if uk:
- lines.append(f" UNIQUE KEY `uk_{table}` ({uk})")
- if increment_column in columns:
- lines.append(f" KEY `ix_{table}_incr` (`{increment_column}`)")
- return (f"CREATE TABLE IF NOT EXISTS `{table}` (\n" + ",\n".join(lines) +
- "\n) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4")
- # ---------------------------------------------------------------------------
- # 库操作
- # ---------------------------------------------------------------------------
- def _col_names(rows: list[dict]) -> list[str]:
- cols: list[str] = []
- for row in rows:
- for k in row:
- if k not in cols:
- cols.append(k)
- return cols
- def _insert_sql(table: str, cols: list[str]) -> str:
- names = ", ".join(f"`{c}`" for c in cols)
- marks = ", ".join(["%s"] * len(cols))
- updates = ", ".join(f"`{c}`=VALUES(`{c}`)" for c in cols)
- return f"INSERT INTO `{table}` ({names}) VALUES ({marks}) ON DUPLICATE KEY UPDATE {updates}"
- def _convert(v):
- if isinstance(v, (dict, list)):
- return json.dumps(v, ensure_ascii=False)
- if isinstance(v, bool):
- return int(v)
- return v
- def init_schema(tables: dict[str, dict]) -> dict:
- """tables: {table: {"rows": [...], "businessKey": [...], "incrementColumn": str}}。建库建表。"""
- result = {"database": None, "tables": []}
- import pymysql
- params = mysql_conn_params()
- database = params.pop("database")
- conn = pymysql.connect(cursorclass=pymysql.cursors.DictCursor, autocommit=True, **params)
- try:
- with conn.cursor() as cur:
- cur.execute(f"CREATE DATABASE IF NOT EXISTS `{database}` DEFAULT CHARACTER SET utf8mb4")
- cur.execute(f"USE `{database}`")
- result["database"] = database
- for table, spec in tables.items():
- ddl = generate_ddl(table, spec["rows"], spec.get("businessKey"),
- spec.get("incrementColumn", "sourceUpdatedAt"))
- cur.execute(ddl)
- result["tables"].append(table)
- finally:
- conn.close()
- return result
- def _use_db():
- import pymysql
- params = mysql_conn_params()
- return pymysql.connect(cursorclass=pymysql.cursors.DictCursor, autocommit=True, **params)
- def reset_table(table: str) -> None:
- conn = _use_db()
- try:
- with conn.cursor() as cur:
- cur.execute(f"TRUNCATE TABLE `{table}`")
- finally:
- conn.close()
- def upsert_rows(table: str, rows: list[dict]) -> int:
- if not rows:
- return 0
- cols = _col_names(rows)
- sql = _insert_sql(table, cols)
- conn = _use_db()
- try:
- with conn.cursor() as cur:
- cur.executemany(sql, [tuple(_convert(row.get(c)) for c in cols) for row in rows])
- return len(rows)
- finally:
- conn.close()
- def update_rows(table: str, updates: dict, where: dict,
- increment_column: str = "sourceUpdatedAt",
- new_increment: str | None = None) -> int:
- """更新已有数据并推进增量列(DB_SYNC 增量更新用例)。"""
- sets = {**updates}
- if increment_column:
- sets[increment_column] = new_increment or datetime.now().strftime("%Y-%m-%dT%H:%M:%S")
- set_clause = ", ".join(f"`{k}`=%s" for k in sets)
- where_clause = " AND ".join(f"`{k}`=%s" for k in where)
- sql = f"UPDATE `{table}` SET {set_clause} WHERE {where_clause}"
- conn = _use_db()
- try:
- with conn.cursor() as cur:
- cur.execute(sql, tuple(sets.values()) + tuple(where.values()))
- return cur.rowcount
- finally:
- conn.close()
- def table_stats(table: str, increment_column: str = "sourceUpdatedAt") -> dict:
- conn = _use_db()
- try:
- with conn.cursor() as cur:
- cur.execute(f"SELECT COUNT(*) AS cnt, MAX(`{increment_column}`) AS max_incr FROM `{table}`")
- row = cur.fetchone() or {}
- return {"table": table, "rowCount": row.get("cnt", 0),
- "maxIncrement": str(row.get("max_incr") or "")}
- finally:
- conn.close()
- # ---------------------------------------------------------------------------
- # Ai-DOP 配置清单生成(只输出建议,不写库)
- # ---------------------------------------------------------------------------
- SIM_SYSTEM_CODE = "SIM_MODE3"
- SIM_SOURCES = {
- "DB_SYNC": "SIM_MODE3_DB",
- "API_PULL": "SIM_MODE3_API",
- "API_INBOUND": "SIM_MODE3_INBOUND",
- }
- def generate_aidop_config(obj: dict, mock_api_base_url: str = "http://127.0.0.1:8018") -> dict:
- """按对象生成 mdp_source / mdp_entity 建议配置,供人工在 Ai-DOP 配置页登记。"""
- db_sync = obj.get("dbSync") or {}
- api_pull = obj.get("apiPull") or {}
- inbound = obj.get("apiInbound") or {}
- config: dict = {"objectCode": obj["objectCode"], "systemCode": SIM_SYSTEM_CODE, "sources": [], "entities": []}
- if db_sync.get("supported"):
- config["sources"].append({
- "source_code": SIM_SOURCES["DB_SYNC"], "source_type": "DB", "conn_mode": "EXTERNAL",
- "hint": "连接指向模拟库 aidop_integration_sim(独立账号,勿用 Ai-DOP 主库账号)",
- })
- config["entities"].append({
- "entity_code": f"SIM_{obj['objectCode']}_DB", "source_code": SIM_SOURCES["DB_SYNC"],
- "source_table_name": db_sync.get("table"), "target_table_name": f"mdp_stg_sim_{db_sync.get('table')}",
- "biz_key_expr": ",".join(db_sync.get("businessKey") or ["bizKey"]),
- "incr_column": db_sync.get("incrementColumn"),
- })
- if api_pull.get("supported"):
- config["sources"].append({
- "source_code": SIM_SOURCES["API_PULL"], "source_type": "API",
- "api_base_url": mock_api_base_url, "api_auth_type": "TOKEN",
- "hint": "token 值即 Mock 控制台当前 TOKEN(环境变量/控制台设置,不落盘)",
- })
- config["entities"].append({
- "entity_code": f"SIM_{obj['objectCode']}_API", "source_code": SIM_SOURCES["API_PULL"],
- "source_api_path": api_pull.get("path"), "target_table_name": f"mdp_stg_sim_{db_sync.get('table')}",
- "response_data_path": api_pull.get("responsePath"), "dedup_key_path": api_pull.get("dedupKeyPath"),
- "biz_key_expr": ",".join(db_sync.get("businessKey") or ["bizKey"]),
- "incr_column": "cursor",
- })
- if inbound.get("supported"):
- config["sources"].append({
- "source_code": SIM_SOURCES["API_INBOUND"], "source_type": "API_INBOUND",
- "hint": "需在 mdp_inbound_grant 为测试 AccessKey 授权该 entityCode,并绑定测试租户开放身份",
- })
- config["entities"].append({
- "entity_code": inbound.get("entityCode"), "source_code": SIM_SOURCES["API_INBOUND"],
- "inbound_enabled": 1,
- })
- return config
|