using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; namespace Admin.NET.Plugin.AiDOP.DataPlatform; /// 数据源配置 API(供配置中心 sources 页读取与保存)。 [ApiDescriptionSettings(Order = 324, Description = "数据中台数据源")] [Route("api/DataPlatform")] [NonUnify] public class MdpSourceConfigService : IDynamicApiController, ITransient { private readonly ISqlSugarClient _db; private readonly UserManager _user; private readonly MdpSourceScopeFactory _scopes; public MdpSourceConfigService(ISqlSugarClient db, UserManager user, MdpSourceScopeFactory scopes) { _db = db; _user = user; _scopes = scopes; } [DisplayName("数据源列表")] [HttpGet("sources")] public async Task GetList([FromQuery] string? keyword, [FromQuery] string? type) { var q = _db.Queryable(); if (!_user.SuperAdmin && !_user.SysAdmin) q = q.Where(x => x.TenantId == _user.TenantId); if (!string.IsNullOrWhiteSpace(keyword)) { var kw = keyword.Trim(); q = q.Where(x => x.SourceCode.Contains(kw) || x.SourceName.Contains(kw)); } if (!string.IsNullOrWhiteSpace(type)) q = q.Where(x => x.SourceType == type.Trim()); var list = await q.OrderBy(x => x.SourceCode).ToListAsync(); return new { list = list.Select(x => new { id = x.Id, sourceCode = x.SourceCode, systemCode = string.IsNullOrWhiteSpace(x.SystemCode) ? x.SourceCode : x.SystemCode, connMode = string.IsNullOrWhiteSpace(x.ConnMode) ? "EXTERNAL" : x.ConnMode, sourceName = x.SourceName, sourceType = x.SourceType, status = x.Status, direction = string.Equals(x.SourceType, "API", StringComparison.OrdinalIgnoreCase) ? "拉取" : string.Equals(x.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase) ? "文件入站" : "入站", dbType = x.DbType ?? "-", dbHost = x.DbHost, dbPort = x.DbPort, dbName = x.DbName, dbUser = x.DbUser, apiBaseUrl = x.ApiBaseUrl, host = string.Equals(x.SourceType, "API", StringComparison.OrdinalIgnoreCase) ? (x.ApiBaseUrl ?? "-") : string.Equals(x.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase) ? "-" : $"{x.DbHost}/{x.DbName}", health = string.Equals(x.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase) ? (string.IsNullOrWhiteSpace(x.SourceCode) ? "配置异常" : "可用") : x.HealthStatus switch { 1 => "正常", 0 => "异常", _ => "未知" }, healthType = x.HealthStatus switch { 1 => "success", 0 => "danger", _ => "info" }, lastCheck = x.LastHealthCheck?.ToString("yyyy-MM-dd HH:mm:ss") ?? "-", healthMsg = x.HealthMsg }) }; } [DisplayName("新建或更新数据源")] [HttpPost("sources")] public async Task Save([FromBody] MdpSourceSaveInput input) { if (string.IsNullOrWhiteSpace(input.SourceCode) || string.IsNullOrWhiteSpace(input.SourceName) || string.IsNullOrWhiteSpace(input.SourceType)) throw Oops.Oh("sourceCode、sourceName、sourceType 必填"); var code = input.SourceCode.Trim(); var mode = string.IsNullOrWhiteSpace(input.ConnMode) ? "EXTERNAL" : input.ConnMode.Trim().ToUpperInvariant(); if (mode is not ("SELF" or "EXTERNAL")) throw Oops.Oh("connMode 只能是 SELF 或 EXTERNAL"); var tenant = _user.TenantId; var existing = await _db.Queryable().Where(x => x.SourceCode == code).FirstAsync(); if (existing == null) { await _db.Insertable(new MdpSource { TenantId = tenant, SourceCode = code, SourceName = input.SourceName.Trim(), SourceType = input.SourceType.Trim(), SystemCode = input.SystemCode?.Trim(), ConnMode = mode, Status = 1, DbType = input.DbType, DbHost = input.DbHost, DbPort = input.DbPort, DbName = input.DbName, DbUser = input.DbUser, DbPasswordEnc = input.DbPassword, ApiBaseUrl = input.ApiBaseUrl, CreateTime = DateTime.Now, UpdateTime = DateTime.Now }).ExecuteCommandAsync(); } else { existing.SourceName = input.SourceName.Trim(); existing.SourceType = input.SourceType.Trim(); existing.SystemCode = input.SystemCode?.Trim(); existing.ConnMode = mode; existing.DbType = input.DbType; existing.DbHost = input.DbHost; existing.DbPort = input.DbPort; existing.DbName = input.DbName; existing.DbUser = input.DbUser; if (!string.IsNullOrEmpty(input.DbPassword)) existing.DbPasswordEnc = input.DbPassword; existing.ApiBaseUrl = input.ApiBaseUrl; existing.UpdateTime = DateTime.Now; await _db.Updateable(existing).ExecuteCommandAsync(); } return new { sourceCode = code }; } [DisplayName("测试数据源连接")] [HttpPost("sources/{sourceCode}/test")] public async Task Test([FromRoute] string sourceCode) { var source = await FindOwnSourceAsync(sourceCode); if (string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase)) { try { var scope = await _scopes.GetScopeAsync(source.SourceCode); await scope.Ado.GetScalarAsync("SELECT 1"); return new { ok = true, message = "数据库连接成功" }; } catch (Exception ex) { return new { ok = false, message = ex.Message }; } } if (string.Equals(source.SourceType, "API", StringComparison.OrdinalIgnoreCase)) { if (string.IsNullOrWhiteSpace(source.ApiBaseUrl)) return new { ok = false, message = "未填写接口地址" }; try { using var http = new HttpClient { Timeout = TimeSpan.FromSeconds(5) }; using var resp = await http.GetAsync(source.ApiBaseUrl); return new { ok = resp.IsSuccessStatusCode, message = $"HTTP {(int)resp.StatusCode}" }; } catch (Exception ex) { return new { ok = false, message = ex.Message }; } } return new { ok = true, message = "对方推送不需要探测对端连接" }; } [DisplayName("停用数据源前的登记检查")] [HttpGet("sources/{sourceCode}/bindings")] public async Task Bindings([FromRoute] string sourceCode) { var source = await FindOwnSourceAsync(sourceCode); var system = string.IsNullOrWhiteSpace(source.SystemCode) ? source.SourceCode : source.SystemCode; var rows = await _db.Ado.SqlQueryAsync( """ SELECT tenant_id AS tenantId, std_object AS stdObject FROM mdp_tenant_std_source WHERE source_system=@sys ORDER BY tenant_id, std_object """, new { sys = system }); return new { systemCode = system, bindings = rows }; } [DisplayName("停用数据源")] [HttpPost("sources/{sourceCode}/disable")] public async Task Disable([FromRoute] string sourceCode, [FromBody] MdpSourceDisableInput input) { var source = await FindOwnSourceAsync(sourceCode); var system = string.IsNullOrWhiteSpace(source.SystemCode) ? source.SourceCode : source.SystemCode; var rows = await _db.Ado.SqlQueryAsync( "SELECT std_object AS stdObject FROM mdp_tenant_std_source WHERE source_system=@sys", new { sys = system }); if (rows.Count > 0 && input.Confirm != true) return new { disabled = false, needConfirm = true, bindings = rows }; source.Status = 0; source.UpdateTime = DateTime.Now; await _db.Updateable(source).UpdateColumns(x => new { x.Status, x.UpdateTime }).ExecuteCommandAsync(); return new { disabled = true, sourceCode = source.SourceCode }; } private async Task FindOwnSourceAsync(string sourceCode) { var code = sourceCode.Trim(); var source = await _db.Queryable().Where(x => x.SourceCode == code).FirstAsync() ?? throw Oops.Oh($"未找到来源 {code}"); if (!_user.SuperAdmin && !_user.SysAdmin && source.TenantId != _user.TenantId) throw Oops.Oh("不能操作其他租户的数据源"); return source; } } public sealed class MdpSourceSaveInput { public string SourceCode { get; set; } = ""; public string SourceName { get; set; } = ""; public string SourceType { get; set; } = "DB"; public string? SystemCode { get; set; } public string? ConnMode { get; set; } public string? DbType { get; set; } public string? DbHost { get; set; } public int? DbPort { get; set; } public string? DbName { get; set; } public string? DbUser { get; set; } public string? DbPassword { get; set; } public string? ApiBaseUrl { get; set; } } public sealed class MdpSourceDisableInput { public bool? Confirm { get; set; } }