MdpSourceScopeFactory.cs 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155
  1. using Admin.NET.Core;
  2. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  3. using SqlSugar;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  5. /// <summary>
  6. /// 按 mdp_source 动态获取只读 SqlSugar 连接域。
  7. /// ConfigId 约定:mdp-src-{sourceCode}。连接只看 conn_mode 与 db_*,不按来源编码分支。
  8. /// </summary>
  9. public sealed class MdpSourceScopeFactory : ITransient
  10. {
  11. private const string ConfigIdPrefix = "mdp-src-";
  12. private readonly ISqlSugarClient _db;
  13. public MdpSourceScopeFactory(ISqlSugarClient db)
  14. {
  15. _db = db;
  16. }
  17. /// <summary>按 sourceCode 获取连接域;优先复用已注册 ConfigId。</summary>
  18. public async Task<ISqlSugarClient> GetScopeAsync(string sourceCode, CancellationToken cancellationToken = default)
  19. {
  20. if (string.IsNullOrWhiteSpace(sourceCode))
  21. throw new ArgumentException("sourceCode 不能为空", nameof(sourceCode));
  22. var code = sourceCode.Trim();
  23. var source = await _db.Queryable<MdpSource>()
  24. .Where(x => x.SourceCode == code && x.Status == 1)
  25. .FirstAsync(cancellationToken)
  26. ?? throw new InvalidOperationException($"mdp_source 未找到启用源:{code}");
  27. return OpenScope(source);
  28. }
  29. /// <summary>
  30. /// 按 system_code 取唯一一条启用的 DB 源。0 条或多条都抛错,禁止取第一条。
  31. /// </summary>
  32. public async Task<ISqlSugarClient> GetScopeBySystemCodeAsync(string systemCode, CancellationToken cancellationToken = default)
  33. {
  34. if (string.IsNullOrWhiteSpace(systemCode))
  35. throw new ArgumentException("systemCode 不能为空", nameof(systemCode));
  36. var code = systemCode.Trim();
  37. var all = await _db.Queryable<MdpSource>()
  38. .Where(x => x.Status == 1 && x.SourceType == "DB")
  39. .ToListAsync(cancellationToken);
  40. var rows = all
  41. .Where(x => string.Equals(
  42. string.IsNullOrWhiteSpace(x.SystemCode) ? x.SourceCode : x.SystemCode,
  43. code,
  44. StringComparison.OrdinalIgnoreCase))
  45. .ToList();
  46. if (rows.Count == 0)
  47. throw new InvalidOperationException($"mdp_source 未找到启用的 DB 源:system_code={code}");
  48. if (rows.Count > 1)
  49. throw new InvalidOperationException(
  50. $"system_code={code} 命中多条 DB 源:{string.Join(",", rows.Select(x => x.SourceCode))}。请改用 GetScopeAsync(sourceCode)。");
  51. return OpenScope(rows[0]);
  52. }
  53. private ISqlSugarClient OpenScope(MdpSource source)
  54. {
  55. if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  56. throw new InvalidOperationException($"源 {source.SourceCode} 的 source_type={source.SourceType},不是 DB,无法建 SqlSugar 连接");
  57. var mode = (source.ConnMode ?? "").Trim().ToUpperInvariant();
  58. if (mode == "SELF")
  59. return _db;
  60. if (mode != "EXTERNAL")
  61. throw new InvalidOperationException($"源 {source.SourceCode} 的 conn_mode='{source.ConnMode}' 无法解析。只接受 SELF / EXTERNAL。");
  62. if (string.IsNullOrWhiteSpace(source.DbHost) || string.IsNullOrWhiteSpace(source.DbName) || string.IsNullOrWhiteSpace(source.DbUser))
  63. throw new InvalidOperationException($"源 {source.SourceCode} 是 EXTERNAL,但缺少 db_host / db_name / db_user。请在数据源页面补齐,禁止回退主库。");
  64. var configId = ConfigIdPrefix + source.SourceCode;
  65. if (_db.AsTenant().IsAnyConnection(configId))
  66. return _db.AsTenant().GetConnectionScope(configId);
  67. var connStr = BuildConnectionString(source);
  68. var dbType = MapDbType(source.DbType);
  69. _db.AsTenant().AddConnection(new ConnectionConfig
  70. {
  71. ConfigId = configId,
  72. DbType = dbType,
  73. ConnectionString = connStr,
  74. InitKeyType = InitKeyType.Attribute,
  75. IsAutoCloseConnection = true,
  76. MoreSettings = new ConnMoreSettings
  77. {
  78. IsAutoRemoveDataCache = true
  79. }
  80. });
  81. return _db.AsTenant().GetConnectionScope(configId);
  82. }
  83. public ISqlSugarClient GetScope(string sourceCode) =>
  84. GetScopeAsync(sourceCode).GetAwaiter().GetResult();
  85. private static DbType MapDbType(string? dbType)
  86. {
  87. if (string.IsNullOrWhiteSpace(dbType)) return DbType.MySql;
  88. return dbType.Trim().ToUpperInvariant() switch
  89. {
  90. "MYSQL" => DbType.MySql,
  91. "SQLSERVER" or "MSSQL" => DbType.SqlServer,
  92. "ORACLE" => DbType.Oracle,
  93. "POSTGRES" or "POSTGRESQL" => DbType.PostgreSQL,
  94. _ => throw new NotSupportedException($"不支持的 db_type:{dbType}")
  95. };
  96. }
  97. private static string BuildConnectionString(MdpSource source)
  98. {
  99. if (string.IsNullOrWhiteSpace(source.DbHost) || string.IsNullOrWhiteSpace(source.DbName))
  100. throw new InvalidOperationException($"源 {source.SourceCode} 缺少 db_host/db_name");
  101. var password = DecryptPassword(source.DbPasswordEnc);
  102. var port = source.DbPort;
  103. var user = source.DbUser ?? "";
  104. var extra = string.IsNullOrWhiteSpace(source.DbExtraParams) ? "" : ";" + source.DbExtraParams.Trim().TrimStart(';');
  105. return MapDbType(source.DbType) switch
  106. {
  107. DbType.SqlServer =>
  108. $"Server={source.DbHost}{(port is > 0 ? $",{port}" : "")};Database={source.DbName};User Id={user};Password={password};TrustServerCertificate=true;Encrypt=false{extra}",
  109. DbType.MySql =>
  110. $"Server={source.DbHost};Port={(port is > 0 ? port : 3306)};Database={source.DbName};Uid={user};Pwd={password};CharSet=utf8mb4;AllowLoadLocalInfile=true{extra}",
  111. DbType.PostgreSQL =>
  112. $"Host={source.DbHost};Port={(port is > 0 ? port : 5432)};Database={source.DbName};Username={user};Password={password}{extra}",
  113. DbType.Oracle =>
  114. $"Data Source={source.DbHost}:{(port is > 0 ? port : 1521)}/{source.DbName};User Id={user};Password={password}{extra}",
  115. _ => throw new NotSupportedException($"不支持的 db_type:{source.DbType}")
  116. };
  117. }
  118. private static string DecryptPassword(string? enc)
  119. {
  120. if (string.IsNullOrEmpty(enc)) return "";
  121. try
  122. {
  123. var plain = CryptogramUtil.Decrypt(enc);
  124. return string.IsNullOrEmpty(plain) ? enc : plain;
  125. }
  126. catch
  127. {
  128. // 明文或非本系统密文时原样使用(与 Database.json EnableConnEncrypt=false 一致)
  129. return enc;
  130. }
  131. }
  132. }