MdpSourceConfigService.cs 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236
  1. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform;
  3. /// <summary>数据源配置 API(供配置中心 sources 页读取与保存)。</summary>
  4. [ApiDescriptionSettings(Order = 324, Description = "数据中台数据源")]
  5. [Route("api/DataPlatform")]
  6. [NonUnify]
  7. public class MdpSourceConfigService : IDynamicApiController, ITransient
  8. {
  9. private readonly ISqlSugarClient _db;
  10. private readonly UserManager _user;
  11. private readonly MdpSourceScopeFactory _scopes;
  12. public MdpSourceConfigService(ISqlSugarClient db, UserManager user, MdpSourceScopeFactory scopes)
  13. {
  14. _db = db;
  15. _user = user;
  16. _scopes = scopes;
  17. }
  18. [DisplayName("数据源列表")]
  19. [HttpGet("sources")]
  20. public async Task<object> GetList([FromQuery] string? keyword, [FromQuery] string? type)
  21. {
  22. var q = _db.Queryable<MdpSource>();
  23. if (!_user.SuperAdmin && !_user.SysAdmin)
  24. q = q.Where(x => x.TenantId == _user.TenantId);
  25. if (!string.IsNullOrWhiteSpace(keyword))
  26. {
  27. var kw = keyword.Trim();
  28. q = q.Where(x => x.SourceCode.Contains(kw) || x.SourceName.Contains(kw));
  29. }
  30. if (!string.IsNullOrWhiteSpace(type))
  31. q = q.Where(x => x.SourceType == type.Trim());
  32. var list = await q.OrderBy(x => x.SourceCode).ToListAsync();
  33. return new
  34. {
  35. list = list.Select(x => new
  36. {
  37. id = x.Id,
  38. sourceCode = x.SourceCode,
  39. systemCode = string.IsNullOrWhiteSpace(x.SystemCode) ? x.SourceCode : x.SystemCode,
  40. connMode = string.IsNullOrWhiteSpace(x.ConnMode) ? "EXTERNAL" : x.ConnMode,
  41. sourceName = x.SourceName,
  42. sourceType = x.SourceType,
  43. status = x.Status,
  44. direction = string.Equals(x.SourceType, "API", StringComparison.OrdinalIgnoreCase) ? "拉取"
  45. : string.Equals(x.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase) ? "文件入站" : "入站",
  46. dbType = x.DbType ?? "-",
  47. dbHost = x.DbHost,
  48. dbPort = x.DbPort,
  49. dbName = x.DbName,
  50. dbUser = x.DbUser,
  51. apiBaseUrl = x.ApiBaseUrl,
  52. host = string.Equals(x.SourceType, "API", StringComparison.OrdinalIgnoreCase)
  53. ? (x.ApiBaseUrl ?? "-")
  54. : string.Equals(x.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase)
  55. ? "-"
  56. : $"{x.DbHost}/{x.DbName}",
  57. health = string.Equals(x.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase)
  58. ? (string.IsNullOrWhiteSpace(x.SourceCode) ? "配置异常" : "可用")
  59. : x.HealthStatus switch
  60. {
  61. 1 => "正常",
  62. 0 => "异常",
  63. _ => "未知"
  64. },
  65. healthType = x.HealthStatus switch
  66. {
  67. 1 => "success",
  68. 0 => "danger",
  69. _ => "info"
  70. },
  71. lastCheck = x.LastHealthCheck?.ToString("yyyy-MM-dd HH:mm:ss") ?? "-",
  72. healthMsg = x.HealthMsg
  73. })
  74. };
  75. }
  76. [DisplayName("新建或更新数据源")]
  77. [HttpPost("sources")]
  78. public async Task<object> Save([FromBody] MdpSourceSaveInput input)
  79. {
  80. if (string.IsNullOrWhiteSpace(input.SourceCode) || string.IsNullOrWhiteSpace(input.SourceName) || string.IsNullOrWhiteSpace(input.SourceType))
  81. throw Oops.Oh("sourceCode、sourceName、sourceType 必填");
  82. var code = input.SourceCode.Trim();
  83. var mode = string.IsNullOrWhiteSpace(input.ConnMode) ? "EXTERNAL" : input.ConnMode.Trim().ToUpperInvariant();
  84. if (mode is not ("SELF" or "EXTERNAL"))
  85. throw Oops.Oh("connMode 只能是 SELF 或 EXTERNAL");
  86. var tenant = _user.TenantId;
  87. var existing = await _db.Queryable<MdpSource>().Where(x => x.SourceCode == code).FirstAsync();
  88. if (existing == null)
  89. {
  90. await _db.Insertable(new MdpSource
  91. {
  92. TenantId = tenant,
  93. SourceCode = code,
  94. SourceName = input.SourceName.Trim(),
  95. SourceType = input.SourceType.Trim(),
  96. SystemCode = input.SystemCode?.Trim(),
  97. ConnMode = mode,
  98. Status = 1,
  99. DbType = input.DbType,
  100. DbHost = input.DbHost,
  101. DbPort = input.DbPort,
  102. DbName = input.DbName,
  103. DbUser = input.DbUser,
  104. DbPasswordEnc = input.DbPassword,
  105. ApiBaseUrl = input.ApiBaseUrl,
  106. CreateTime = DateTime.Now,
  107. UpdateTime = DateTime.Now
  108. }).ExecuteCommandAsync();
  109. }
  110. else
  111. {
  112. existing.SourceName = input.SourceName.Trim();
  113. existing.SourceType = input.SourceType.Trim();
  114. existing.SystemCode = input.SystemCode?.Trim();
  115. existing.ConnMode = mode;
  116. existing.DbType = input.DbType;
  117. existing.DbHost = input.DbHost;
  118. existing.DbPort = input.DbPort;
  119. existing.DbName = input.DbName;
  120. existing.DbUser = input.DbUser;
  121. if (!string.IsNullOrEmpty(input.DbPassword))
  122. existing.DbPasswordEnc = input.DbPassword;
  123. existing.ApiBaseUrl = input.ApiBaseUrl;
  124. existing.UpdateTime = DateTime.Now;
  125. await _db.Updateable(existing).ExecuteCommandAsync();
  126. }
  127. return new { sourceCode = code };
  128. }
  129. [DisplayName("测试数据源连接")]
  130. [HttpPost("sources/{sourceCode}/test")]
  131. public async Task<object> Test([FromRoute] string sourceCode)
  132. {
  133. var source = await FindOwnSourceAsync(sourceCode);
  134. if (string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
  135. {
  136. try
  137. {
  138. var scope = await _scopes.GetScopeAsync(source.SourceCode);
  139. await scope.Ado.GetScalarAsync("SELECT 1");
  140. return new { ok = true, message = "数据库连接成功" };
  141. }
  142. catch (Exception ex)
  143. {
  144. return new { ok = false, message = ex.Message };
  145. }
  146. }
  147. if (string.Equals(source.SourceType, "API", StringComparison.OrdinalIgnoreCase))
  148. {
  149. if (string.IsNullOrWhiteSpace(source.ApiBaseUrl))
  150. return new { ok = false, message = "未填写接口地址" };
  151. try
  152. {
  153. using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(5) };
  154. using var resp = await http.GetAsync(source.ApiBaseUrl);
  155. return new { ok = resp.IsSuccessStatusCode, message = $"HTTP {(int)resp.StatusCode}" };
  156. }
  157. catch (Exception ex)
  158. {
  159. return new { ok = false, message = ex.Message };
  160. }
  161. }
  162. return new { ok = true, message = "对方推送不需要探测对端连接" };
  163. }
  164. [DisplayName("停用数据源前的登记检查")]
  165. [HttpGet("sources/{sourceCode}/bindings")]
  166. public async Task<object> Bindings([FromRoute] string sourceCode)
  167. {
  168. var source = await FindOwnSourceAsync(sourceCode);
  169. var system = string.IsNullOrWhiteSpace(source.SystemCode) ? source.SourceCode : source.SystemCode;
  170. var rows = await _db.Ado.SqlQueryAsync<dynamic>(
  171. """
  172. SELECT tenant_id AS tenantId, std_object AS stdObject
  173. FROM mdp_tenant_std_source
  174. WHERE source_system=@sys
  175. ORDER BY tenant_id, std_object
  176. """,
  177. new { sys = system });
  178. return new { systemCode = system, bindings = rows };
  179. }
  180. [DisplayName("停用数据源")]
  181. [HttpPost("sources/{sourceCode}/disable")]
  182. public async Task<object> Disable([FromRoute] string sourceCode, [FromBody] MdpSourceDisableInput input)
  183. {
  184. var source = await FindOwnSourceAsync(sourceCode);
  185. var system = string.IsNullOrWhiteSpace(source.SystemCode) ? source.SourceCode : source.SystemCode;
  186. var rows = await _db.Ado.SqlQueryAsync<dynamic>(
  187. "SELECT std_object AS stdObject FROM mdp_tenant_std_source WHERE source_system=@sys",
  188. new { sys = system });
  189. if (rows.Count > 0 && input.Confirm != true)
  190. return new { disabled = false, needConfirm = true, bindings = rows };
  191. source.Status = 0;
  192. source.UpdateTime = DateTime.Now;
  193. await _db.Updateable(source).UpdateColumns(x => new { x.Status, x.UpdateTime }).ExecuteCommandAsync();
  194. return new { disabled = true, sourceCode = source.SourceCode };
  195. }
  196. private async Task<MdpSource> FindOwnSourceAsync(string sourceCode)
  197. {
  198. var code = sourceCode.Trim();
  199. var source = await _db.Queryable<MdpSource>().Where(x => x.SourceCode == code).FirstAsync()
  200. ?? throw Oops.Oh($"未找到来源 {code}");
  201. if (!_user.SuperAdmin && !_user.SysAdmin && source.TenantId != _user.TenantId)
  202. throw Oops.Oh("不能操作其他租户的数据源");
  203. return source;
  204. }
  205. }
  206. public sealed class MdpSourceSaveInput
  207. {
  208. public string SourceCode { get; set; } = "";
  209. public string SourceName { get; set; } = "";
  210. public string SourceType { get; set; } = "DB";
  211. public string? SystemCode { get; set; }
  212. public string? ConnMode { get; set; }
  213. public string? DbType { get; set; }
  214. public string? DbHost { get; set; }
  215. public int? DbPort { get; set; }
  216. public string? DbName { get; set; }
  217. public string? DbUser { get; set; }
  218. public string? DbPassword { get; set; }
  219. public string? ApiBaseUrl { get; set; }
  220. }
  221. public sealed class MdpSourceDisableInput
  222. {
  223. public bool? Confirm { get; set; }
  224. }