using System.Text;
using System.Text.RegularExpressions;
using Admin.NET.Plugin.AiDOP.DataPlatform;
using Microsoft.Extensions.Logging;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
///
/// 中台表(mdp_std_* / mdp_stg_* / dwd_*)的增量对齐。
/// 表不存在则按声明建表;已存在则只 ADD COLUMN / ADD KEY。不 DROP、不改类型、不改已有索引。
///
public static class MdpSchemaAligner
{
private static readonly Regex CreateHead = new(
@"^CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\s+`?(?[A-Za-z0-9_]+)`?(?:\s+LIKE\s+`?(?[A-Za-z0-9_]+)`?)?",
RegexOptions.IgnoreCase | RegexOptions.Compiled | RegexOptions.Singleline);
private static readonly Regex IndexHead = new(
@"^(?UNIQUE\s+)?(?:KEY|INDEX)\s+`?(?[A-Za-z0-9_]+)`?\s*\((?[^)]+)\)",
RegexOptions.IgnoreCase | RegexOptions.Compiled);
public static bool IsNeutralTable(string name) =>
name.StartsWith("mdp_std_", StringComparison.OrdinalIgnoreCase)
|| name.StartsWith("mdp_stg_", StringComparison.OrdinalIgnoreCase)
|| name.StartsWith("dwd_", StringComparison.OrdinalIgnoreCase);
///
/// 非中立建表语句原样交给数据库;中立建表改为增量对齐,避免 CREATE IF NOT EXISTS 在表已存在时不补列。
///
public static async Task ExecuteAsync(ISqlSugarClient db, string sql, params SugarParameter[] parameters)
{
if (!ContainsNeutralCreate(sql))
return await db.Ado.ExecuteCommandAsync(sql, parameters);
var applied = 0;
foreach (var statement in SplitStatements(sql))
{
if (IsNeutralCreate(statement))
applied += await AlignOneAsync(db, statement);
else
applied += await db.Ado.ExecuteCommandAsync(statement, parameters);
}
return applied;
}
public static Task ExecuteAsync(ISqlSugarClient db, string sql, List parameters) =>
ExecuteAsync(db, sql, parameters?.ToArray() ?? []);
public static async Task EnsureWrittenByColumnAsync(ISqlSugarClient db, string table)
{
if (!IsSafeName(table))
throw new InvalidOperationException($"非法表名:{table}");
var hasWritten = await db.Ado.GetIntAsync(
"""
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='written_by'
""",
new SugarParameter("@t", table));
if (hasWritten > 0) return;
var ddl =
$"ALTER TABLE `{table}` ADD COLUMN written_by VARCHAR(32) NOT NULL DEFAULT 'LEGACY_UNVERIFIED' COMMENT 'DB_SYNC/API_INBOUND/PLATFORM_FORM/SEED/LEGACY_UNVERIFIED'";
await db.Ado.ExecuteCommandAsync(ddl);
Log(ddl);
}
///
/// 给已有 source_system 的 mdp_std_* 补 written_by,并按贴源表名登记真实来源。
/// 确证不了的行不改 source_system,只保持 written_by=LEGACY_UNVERIFIED。
///
public static async Task EnsureSourceIdentityAsync(ISqlSugarClient db)
{
await db.Ado.ExecuteCommandAsync(
"""
CREATE TABLE IF NOT EXISTS mdp_source_table_registry (
id BIGINT NOT NULL AUTO_INCREMENT,
source_table VARCHAR(64) NOT NULL COMMENT '贴源记录的原始表名',
source_system VARCHAR(50) NOT NULL COMMENT '真实来源系统码',
remark VARCHAR(255) NULL,
PRIMARY KEY (id),
UNIQUE KEY uk_source_table (source_table)
) COMMENT='原始表名 → 真实来源系统 的唯一对照'
""");
var stgTables = await db.Ado.SqlQueryAsync(
"""
SELECT TABLE_NAME AS Name FROM information_schema.TABLES
WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME LIKE 'mdp_stg%'
""");
foreach (var table in stgTables)
{
if (!IsSafeName(table.Name)) continue;
var hasSourceTable = await db.Ado.GetIntAsync(
"""
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='source_table'
""",
new SugarParameter("@t", table.Name));
if (hasSourceTable == 0) continue;
var names = await db.Ado.SqlQueryAsync(
$"SELECT DISTINCT source_table AS Name FROM `{table.Name}` WHERE source_table IS NOT NULL AND TRIM(source_table)<>''");
foreach (var row in names)
{
var system = MdpSourceIdentity.ClassifySourceTable(row.Name);
if (system == null)
{
Log($"来源表未登记(命名规则无法确证,不猜测): {row.Name}");
continue;
}
await db.Ado.ExecuteCommandAsync(
"""
INSERT INTO mdp_source_table_registry (source_table, source_system, remark)
SELECT @t, @s, '1.0.565 按表名规则登记'
FROM DUAL
WHERE NOT EXISTS (SELECT 1 FROM mdp_source_table_registry WHERE source_table=@t)
""",
new SugarParameter("@t", row.Name.Trim()),
new SugarParameter("@s", system));
}
}
var stdTables = await db.Ado.SqlQueryAsync(
"""
SELECT TABLE_NAME AS Name FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND COLUMN_NAME='source_system' AND TABLE_NAME LIKE 'mdp_std%'
""");
foreach (var table in stdTables.Select(t => t.Name).Distinct(StringComparer.OrdinalIgnoreCase))
{
if (!IsSafeName(table)) continue;
var hasWritten = await db.Ado.GetIntAsync(
"""
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='written_by'
""",
new SugarParameter("@t", table));
if (hasWritten == 0)
{
var ddl =
$"ALTER TABLE `{table}` ADD COLUMN written_by VARCHAR(32) NOT NULL DEFAULT 'LEGACY_UNVERIFIED' COMMENT 'DB_SYNC/API_INBOUND/PLATFORM_FORM/SEED/LEGACY_UNVERIFIED'";
await db.Ado.ExecuteCommandAsync(ddl);
Log(ddl);
}
var defaults = await db.Ado.SqlQueryAsync(
"""
SELECT COLUMN_DEFAULT AS Name FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='source_system'
""",
new SugarParameter("@t", table));
var currentDefault = defaults.FirstOrDefault()?.Name?.Trim().Trim('\'') ?? "";
if (currentDefault == "AIDOP")
{
var ddl = $"ALTER TABLE `{table}` ALTER COLUMN source_system SET DEFAULT ''";
await db.Ado.ExecuteCommandAsync(ddl);
Log(ddl);
}
await db.Ado.ExecuteCommandAsync(
$"""
UPDATE `{table}` SET written_by='DB_SYNC'
WHERE written_by='LEGACY_UNVERIFIED'
AND source_system IN ('DOPDEMORQ_SQLSERVER','T8','AIDOPDEV_MYSQL')
""");
await db.Ado.ExecuteCommandAsync(
$"""
UPDATE `{table}` SET written_by='SEED'
WHERE written_by='LEGACY_UNVERIFIED'
AND source_system IN ('UAT_GENERATOR','UAT','DEMO')
""");
await db.Ado.ExecuteCommandAsync(
$"""
UPDATE `{table}` SET written_by='API_INBOUND'
WHERE written_by='LEGACY_UNVERIFIED' AND source_system='API'
""");
await db.Ado.ExecuteCommandAsync(
$"""
UPDATE `{table}` SET written_by='PLATFORM_FORM'
WHERE written_by='LEGACY_UNVERIFIED' AND source_system='AIDOP_NATIVE'
""");
}
await BackfillProvenSourceAsync(db, stdTables.Select(t => t.Name).Distinct(StringComparer.OrdinalIgnoreCase));
}
private static async Task BackfillProvenSourceAsync(ISqlSugarClient db, IEnumerable stdTables)
{
var stg = await db.Ado.SqlQueryAsync(
"""
SELECT c.TABLE_NAME AS Name
FROM information_schema.COLUMNS c
WHERE c.TABLE_SCHEMA=DATABASE() AND c.TABLE_NAME LIKE 'mdp_stg%' AND c.COLUMN_NAME='source_biz_key'
AND EXISTS (
SELECT 1 FROM information_schema.COLUMNS t
WHERE t.TABLE_SCHEMA=DATABASE() AND t.TABLE_NAME=c.TABLE_NAME AND t.COLUMN_NAME='source_table')
AND EXISTS (
SELECT 1 FROM information_schema.COLUMNS t
WHERE t.TABLE_SCHEMA=DATABASE() AND t.TABLE_NAME=c.TABLE_NAME AND t.COLUMN_NAME='tenant_id')
""");
var stgNames = stg.Select(x => x.Name).Where(IsSafeName).Distinct(StringComparer.OrdinalIgnoreCase).ToList();
if (stgNames.Count == 0) return;
var union = string.Join(" UNION ALL ", stgNames.Select(n =>
$"SELECT g.tenant_id, g.source_biz_key, r.source_system FROM `{n}` g JOIN mdp_source_table_registry r ON r.source_table=g.source_table"));
var proven =
$"""
SELECT tenant_id, source_biz_key, MIN(source_system) AS source_system
FROM ({union}) u
GROUP BY tenant_id, source_biz_key
HAVING COUNT(DISTINCT source_system)=1
""";
foreach (var table in stdTables)
{
if (!IsSafeName(table)) continue;
var ready = await db.Ado.GetIntAsync(
"""
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
AND COLUMN_NAME IN ('tenant_id','source_biz_key','source_system','written_by','id')
""",
new SugarParameter("@t", table));
if (ready < 5) continue;
try
{
await db.Ado.ExecuteCommandAsync(
$"""
UPDATE `{table}` s
JOIN ({proven}) x ON x.tenant_id=s.tenant_id AND x.source_biz_key=s.source_biz_key
LEFT JOIN `{table}` o
ON o.tenant_id=s.tenant_id AND o.source_system=x.source_system
AND o.source_biz_key=s.source_biz_key AND o.id<>s.id
SET s.source_system=x.source_system, s.written_by='DB_SYNC'
WHERE s.source_system IN ('AIDOP','','AIDOP_LEGACY')
AND o.id IS NULL
""");
}
catch (Exception ex)
{
Log($"来源回填跳过 {table}: {ex.Message}");
}
}
}
private static bool IsSafeName(string? name) =>
!string.IsNullOrWhiteSpace(name) && Regex.IsMatch(name, "^[A-Za-z0-9_]+$");
public static IReadOnlyList Plan(string statement, SchemaSnapshot snapshot)
{
var text = statement.Trim().TrimEnd(';').Trim();
var head = CreateHead.Match(text);
if (!head.Success || !IsNeutralTable(head.Groups["name"].Value))
return [text];
var table = head.Groups["name"].Value;
if (head.Groups["like"].Success)
return snapshot.Exists ? [] : [text];
if (!snapshot.Exists)
return [text];
var alters = new List();
foreach (var part in SplitTopLevel(Body(text)))
{
var line = part.Trim().TrimEnd(',').Trim();
if (line.Length == 0) continue;
if (line.StartsWith("PRIMARY", StringComparison.OrdinalIgnoreCase)
|| line.StartsWith("CONSTRAINT", StringComparison.OrdinalIgnoreCase)
|| line.StartsWith("FULLTEXT", StringComparison.OrdinalIgnoreCase)
|| line.StartsWith("SPATIAL", StringComparison.OrdinalIgnoreCase))
continue;
var index = IndexHead.Match(line);
if (index.Success)
{
var indexName = index.Groups["name"].Value;
if (snapshot.Indexes.Contains(indexName)) continue;
var unique = index.Groups["unique"].Success ? "UNIQUE KEY" : "KEY";
alters.Add(
$"ALTER TABLE `{table}` ADD {unique} `{indexName}` ({index.Groups["cols"].Value.Trim()})");
continue;
}
if (line.StartsWith("UNIQUE", StringComparison.OrdinalIgnoreCase)
|| line.StartsWith("KEY", StringComparison.OrdinalIgnoreCase)
|| line.StartsWith("INDEX", StringComparison.OrdinalIgnoreCase))
continue;
var column = ColumnName(line);
if (column.Length == 0 || snapshot.Columns.Contains(column)) continue;
alters.Add($"ALTER TABLE `{table}` ADD COLUMN {line}");
}
return alters;
}
private static async Task AlignOneAsync(ISqlSugarClient db, string statement)
{
var head = CreateHead.Match(statement.Trim());
var table = head.Groups["name"].Value;
var snapshot = await LoadAsync(db, table);
var planned = Plan(statement, snapshot);
var applied = 0;
foreach (var ddl in planned)
{
if (ddl.Contains("DROP ", StringComparison.OrdinalIgnoreCase))
throw new InvalidOperationException("中台结构对齐器禁止执行 DROP");
await db.Ado.ExecuteCommandAsync(ddl);
Log(ddl);
applied++;
}
return applied;
}
private static async Task LoadAsync(ISqlSugarClient db, string table)
{
var exists = await db.Ado.GetIntAsync(
"""
SELECT COUNT(*) FROM information_schema.TABLES
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
""",
new SugarParameter("@t", table));
if (exists == 0)
return SchemaSnapshot.Missing;
var columns = await db.Ado.SqlQueryAsync(
"""
SELECT COLUMN_NAME AS Name FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
""",
new SugarParameter("@t", table));
var indexes = await db.Ado.SqlQueryAsync(
"""
SELECT DISTINCT INDEX_NAME AS Name FROM information_schema.STATISTICS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
""",
new SugarParameter("@t", table));
return new SchemaSnapshot(
true,
columns.Select(x => x.Name),
indexes.Select(x => x.Name));
}
private static void Log(string ddl)
{
try
{
App.GetService>()?.LogInformation("中台结构对齐 {Ddl}", ddl);
}
catch
{
// 单测与未建容器时不记日志
}
}
internal static bool ContainsNeutralCreate(string sql) =>
SplitStatements(sql).Any(IsNeutralCreate);
private static bool IsNeutralCreate(string statement)
{
var head = CreateHead.Match(statement.Trim());
return head.Success && IsNeutralTable(head.Groups["name"].Value);
}
internal static List SplitStatements(string sql)
{
var list = new List();
var sb = new StringBuilder();
var depth = 0;
var quote = false;
for (var i = 0; i < sql.Length; i++)
{
var c = sql[i];
if (c == '\'')
{
if (quote && i + 1 < sql.Length && sql[i + 1] == '\'')
{
sb.Append("''");
i++;
continue;
}
quote = !quote;
sb.Append(c);
continue;
}
if (!quote && c == '(') depth++;
if (!quote && c == ')' && depth > 0) depth--;
if (!quote && depth == 0 && c == ';')
{
Add(list, sb);
continue;
}
sb.Append(c);
}
Add(list, sb);
return list;
}
private static void Add(List list, StringBuilder sb)
{
var text = sb.ToString().Trim();
if (text.Length > 0) list.Add(text);
sb.Clear();
}
private static string Body(string create)
{
var start = create.IndexOf('(');
var end = create.LastIndexOf(')');
if (start < 0 || end <= start) return string.Empty;
return create.Substring(start + 1, end - start - 1);
}
private static List SplitTopLevel(string body)
{
var list = new List();
var sb = new StringBuilder();
var depth = 0;
var quote = false;
foreach (var c in body)
{
if (c == '\'') quote = !quote;
if (!quote && c == '(') depth++;
if (!quote && c == ')' && depth > 0) depth--;
if (!quote && depth == 0 && c == ',')
{
list.Add(sb.ToString());
sb.Clear();
continue;
}
sb.Append(c);
}
if (sb.Length > 0) list.Add(sb.ToString());
return list;
}
private static string ColumnName(string line)
{
var token = line.Split([' ', '\t', '\r', '\n'], 2, StringSplitOptions.RemoveEmptyEntries)[0];
return token.Trim('`');
}
private sealed class NameRow
{
public string Name { get; set; } = string.Empty;
}
/// 仅用于取日志泛型,避免把静态类当作日志类别。
private sealed class MdpSchemaAlignerLog;
}
public sealed class SchemaSnapshot
{
public static readonly SchemaSnapshot Missing = new(false, [], []);
public SchemaSnapshot(bool exists, IEnumerable columns, IEnumerable indexes)
{
Exists = exists;
Columns = new HashSet(columns.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase);
Indexes = new HashSet(indexes.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase);
}
public bool Exists { get; }
public HashSet Columns { get; }
public HashSet Indexes { get; }
}