MdpSourceScopeFactory.cs 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137
  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};兼容已注册的 t8_v5(source_code=T8_ERP/T8_V5)。
  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. // 兼容既有硬编码 T8:未在 mdp_source 登记前也可工作
  24. if (IsT8Alias(code) && _db.AsTenant().IsAnyConnection("t8_v5"))
  25. return _db.AsTenant().GetConnectionScope("t8_v5");
  26. // 本库样板源:直接复用主库连接(mdp_source 可不落账号密文)
  27. if (IsLocalMysqlAlias(code))
  28. return _db;
  29. var configId = ConfigIdPrefix + code;
  30. if (_db.AsTenant().IsAnyConnection(configId))
  31. return _db.AsTenant().GetConnectionScope(configId);
  32. var source = await _db.Queryable<MdpSource>()
  33. .Where(x => x.SourceCode == code && x.Status == 1)
  34. .FirstAsync(cancellationToken)
  35. ?? throw new InvalidOperationException($"mdp_source 未找到启用源:{code}");
  36. if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  37. throw new InvalidOperationException($"源 {code} 的 source_type={source.SourceType},不是 DB,无法建 SqlSugar 连接");
  38. // 未配置账号时视为与主库同库,避免样板源建空 Uid 连接失败
  39. if (string.IsNullOrWhiteSpace(source.DbUser))
  40. return _db;
  41. var connStr = BuildConnectionString(source);
  42. var dbType = MapDbType(source.DbType);
  43. _db.AsTenant().AddConnection(new ConnectionConfig
  44. {
  45. ConfigId = configId,
  46. DbType = dbType,
  47. ConnectionString = connStr,
  48. InitKeyType = InitKeyType.Attribute,
  49. IsAutoCloseConnection = true,
  50. MoreSettings = new ConnMoreSettings
  51. {
  52. IsAutoRemoveDataCache = true
  53. }
  54. });
  55. return _db.AsTenant().GetConnectionScope(configId);
  56. }
  57. public ISqlSugarClient GetScope(string sourceCode) =>
  58. GetScopeAsync(sourceCode).GetAwaiter().GetResult();
  59. private static bool IsT8Alias(string code) =>
  60. code.Equals("T8_ERP", StringComparison.OrdinalIgnoreCase)
  61. || code.Equals("T8_V5", StringComparison.OrdinalIgnoreCase)
  62. || code.Equals("T8_V5_SQLSERVER", StringComparison.OrdinalIgnoreCase)
  63. || code.Equals("t8_v5", StringComparison.OrdinalIgnoreCase);
  64. private static bool IsLocalMysqlAlias(string code) =>
  65. code.Equals("AIDOPDEV_MYSQL", StringComparison.OrdinalIgnoreCase)
  66. || code.Equals("LOCAL_MYSQL", StringComparison.OrdinalIgnoreCase)
  67. || code.Equals("AIDOP_MYSQL", StringComparison.OrdinalIgnoreCase);
  68. private static DbType MapDbType(string? dbType)
  69. {
  70. if (string.IsNullOrWhiteSpace(dbType)) return DbType.MySql;
  71. return dbType.Trim().ToUpperInvariant() switch
  72. {
  73. "MYSQL" => DbType.MySql,
  74. "SQLSERVER" or "MSSQL" => DbType.SqlServer,
  75. "ORACLE" => DbType.Oracle,
  76. "POSTGRES" or "POSTGRESQL" => DbType.PostgreSQL,
  77. _ => throw new NotSupportedException($"不支持的 db_type:{dbType}")
  78. };
  79. }
  80. private static string BuildConnectionString(MdpSource source)
  81. {
  82. if (string.IsNullOrWhiteSpace(source.DbHost) || string.IsNullOrWhiteSpace(source.DbName))
  83. throw new InvalidOperationException($"源 {source.SourceCode} 缺少 db_host/db_name");
  84. var password = DecryptPassword(source.DbPasswordEnc);
  85. var port = source.DbPort;
  86. var user = source.DbUser ?? "";
  87. var extra = string.IsNullOrWhiteSpace(source.DbExtraParams) ? "" : ";" + source.DbExtraParams.Trim().TrimStart(';');
  88. return MapDbType(source.DbType) switch
  89. {
  90. DbType.SqlServer =>
  91. $"Server={source.DbHost}{(port is > 0 ? $",{port}" : "")};Database={source.DbName};User Id={user};Password={password};TrustServerCertificate=true;Encrypt=false{extra}",
  92. DbType.MySql =>
  93. $"Server={source.DbHost};Port={(port is > 0 ? port : 3306)};Database={source.DbName};Uid={user};Pwd={password};CharSet=utf8mb4;AllowLoadLocalInfile=true{extra}",
  94. DbType.PostgreSQL =>
  95. $"Host={source.DbHost};Port={(port is > 0 ? port : 5432)};Database={source.DbName};Username={user};Password={password}{extra}",
  96. DbType.Oracle =>
  97. $"Data Source={source.DbHost}:{(port is > 0 ? port : 1521)}/{source.DbName};User Id={user};Password={password}{extra}",
  98. _ => throw new NotSupportedException($"不支持的 db_type:{source.DbType}")
  99. };
  100. }
  101. private static string DecryptPassword(string? enc)
  102. {
  103. if (string.IsNullOrEmpty(enc)) return "";
  104. try
  105. {
  106. var plain = CryptogramUtil.Decrypt(enc);
  107. return string.IsNullOrEmpty(plain) ? enc : plain;
  108. }
  109. catch
  110. {
  111. // 明文或非本系统密文时原样使用(与 Database.json EnableConnEncrypt=false 一致)
  112. return enc;
  113. }
  114. }
  115. }