AdoSmartOpsKpiCalcConfigService.cs 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372
  1. using Admin.NET.Core;
  2. using Admin.NET.Plugin.AiDOP.Dto.SmartOps;
  3. using Admin.NET.Plugin.AiDOP.Entity;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure;
  5. using SqlSugar;
  6. namespace Admin.NET.Plugin.AiDOP.SmartOps;
  7. /// <summary>
  8. /// KPI 计算配置服务:CRUD + 校验 + 试算 + 发布/停用/激活(发布状态机,单生效版本事务保证)。
  9. /// 租户显式控制(ClearFilter):配置落在"分发器运行时查找的租户"= KPI 数据租户(S5=AidopSourceTenantMap 解析),
  10. /// 而非 API 调用者 JWT 租户;否则运行时找不到配置。SqlScript 属高危配置,接口须登录鉴权(Controller 层)。
  11. /// </summary>
  12. public sealed class AdoSmartOpsKpiCalcConfigService : ITransient
  13. {
  14. private const string EngineLegacyCode = "LEGACY_CODE";
  15. private const string EngineConfigSql = "CONFIG_SQL";
  16. private const string EngineLegacyTvf = "LEGACY_TVF";
  17. private const string StatusDraft = "DRAFT";
  18. private const string StatusPublished = "PUBLISHED";
  19. private const string StatusRetired = "RETIRED";
  20. private static readonly string[] TenantScopedModules = { "S5", "S6", "S7" };
  21. private readonly ISqlSugarClient _db;
  22. private readonly KpiSqlReadOnlyExecutor _executor;
  23. private readonly AdoSmartOpsKpiBusinessInputService _businessInput;
  24. public AdoSmartOpsKpiCalcConfigService(ISqlSugarClient db, KpiSqlReadOnlyExecutor executor, AdoSmartOpsKpiBusinessInputService businessInput)
  25. {
  26. _db = db;
  27. _executor = executor;
  28. _businessInput = businessInput;
  29. }
  30. /// <summary>
  31. /// CONFIG_SQL 业务门禁:仅 BusinessInputStatus=READY_FOR_CONFIG 放行 试算/发布/激活。
  32. /// 未就绪则明确拒绝——不执行 SQL、不写运行日志、不改版本状态/IsCurrent。
  33. /// </summary>
  34. private async Task EnsureBusinessReadyAsync(long tenantId, string metricCode, string action)
  35. {
  36. var status = await _businessInput.GetStatusAsync(tenantId, metricCode);
  37. if (status != AdoSmartOpsKpiBusinessInputService.StatusReady)
  38. throw Oops.Bah($"业务来源未就绪(当前 {status}):{action} 需先在「业务来源登记」补齐来源系统/表/字段/SQL来源并置为 READY_FOR_CONFIG");
  39. }
  40. /// <summary>解析配置应归属的租户(= 运行时 KPI 落库租户)。S5/S6/S7 走 T8 账套映射(pbxfxp→AIDOP)。</summary>
  41. public static long ResolveKpiTenantId(string moduleCode)
  42. {
  43. if (TenantScopedModules.Contains((moduleCode ?? "").Trim().ToUpperInvariant()))
  44. return AidopSourceTenantMap.ResolveTenantId("pbxfxp");
  45. return SqlSugarConst.DefaultTenantId;
  46. }
  47. private ISugarQueryable<AdoSmartOpsKpiCalcConfig> Query() =>
  48. _db.Queryable<AdoSmartOpsKpiCalcConfig>().ClearFilter<ITenantIdFilter>();
  49. /// <summary>
  50. /// 写路径一律走原生参数化 SQL(表名与实体 SugarTable 一致,列名 PascalCase:EnableUnderLine=false)。
  51. /// 原因:全局 MoreSettings.IsAutoUpdateQueryFilter / IsAutoDeleteQueryFilter=true 会给 ORM 的 UPDATE/DELETE
  52. /// 自动追加「登录 JWT 租户」条件,而本表按 <see cref="ResolveKpiTenantId"/> 落在 KPI 数据租户,
  53. /// 两者不一致时静默 0 行(UAT-S0-07 / UAT-S9-01);SqlSugar 5.1.4 的 IUpdateable/IDeleteable
  54. /// 又无法关闭该过滤器(EnableQueryFilter 只能开不能关,且 ClearFilter 仅存在于 ISugarQueryable)。
  55. /// 每条写语句都必须显式带 TenantId + 业务状态条件,**禁止按 Id 裸写**。
  56. /// </summary>
  57. private const string TableName = "ado_smart_ops_kpi_calc_config";
  58. /// <summary>取配置归属租户;缺失即数据异常,直接拒绝,避免降级成按 Id 跨租户裸写。</summary>
  59. private static long RequireTenantId(AdoSmartOpsKpiCalcConfig e) =>
  60. e.TenantId ?? throw Oops.Bah($"配置版本 Id={e.Id} 缺少 TenantId,数据异常,拒绝写入");
  61. /// <summary>
  62. /// 0 行写入归因:区分「目标不存在 / 租户不匹配 / 状态不满足 / 并发状态变化」,绝不静默返回成功。
  63. /// </summary>
  64. private async Task<Exception> ExplainZeroRowAsync(
  65. long id, long? expectedTenantId, string action, string statusRequirement, params string[] acceptableStatuses)
  66. {
  67. var now = await Query().Where(x => x.Id == id).FirstAsync();
  68. if (now == null)
  69. return Oops.Bah($"{action}失败:配置版本 Id={id} 不存在(可能已被并发删除)");
  70. if (now.TenantId != expectedTenantId)
  71. return Oops.Bah($"{action}失败:配置版本 Id={id} 归属租户 {now.TenantId},与本次作用域租户 {expectedTenantId} 不一致");
  72. if (acceptableStatuses.Length > 0 && !acceptableStatuses.Contains(now.PublishStatus))
  73. return Oops.Bah($"{action}失败:配置版本 Id={id} 当前状态为 {now.PublishStatus},不满足要求的 {statusRequirement}(并发修改,请刷新后重试)");
  74. return Oops.Bah($"{action}失败:配置版本 Id={id}(租户 {now.TenantId} / 状态 {now.PublishStatus})未命中任何行,疑似并发状态变化,请刷新后重试");
  75. }
  76. /// <summary>某 KPI 的全部版本(按版本号倒序)。</summary>
  77. public async Task<List<KpiCalcConfigDto>> GetByMetricAsync(string metricCode, string moduleCode)
  78. {
  79. var tenantId = ResolveKpiTenantId(moduleCode);
  80. var list = await Query()
  81. .Where(x => x.MetricCode == metricCode && x.TenantId == tenantId)
  82. .OrderBy(x => x.VersionNo, OrderByType.Desc)
  83. .ToListAsync();
  84. return list.Select(ToDto).ToList();
  85. }
  86. /// <summary>当前生效配置(PUBLISHED + IsCurrent),无则 null(= 走 legacy)。</summary>
  87. public async Task<AdoSmartOpsKpiCalcConfig?> GetActiveAsync(long tenantId, string metricCode)
  88. {
  89. return await Query()
  90. .Where(x => x.MetricCode == metricCode && x.TenantId == tenantId
  91. && x.IsCurrent && x.PublishStatus == StatusPublished)
  92. .FirstAsync();
  93. }
  94. /// <summary>新增草稿(VersionNo = 现有最大 + 1)。</summary>
  95. public async Task<KpiCalcConfigDto> CreateDraftAsync(KpiCalcConfigUpsertDto dto, string operatorName)
  96. {
  97. ValidateEngineType(dto.CalcEngineType);
  98. var tenantId = ResolveKpiTenantId(dto.ModuleCode);
  99. var maxVer = await Query()
  100. .Where(x => x.MetricCode == dto.MetricCode && x.TenantId == tenantId)
  101. .MaxAsync(x => (int?)x.VersionNo) ?? 0;
  102. var entity = new AdoSmartOpsKpiCalcConfig
  103. {
  104. TenantId = tenantId,
  105. MetricCode = dto.MetricCode,
  106. ModuleCode = dto.ModuleCode,
  107. VersionNo = maxVer + 1,
  108. CalcEngineType = dto.CalcEngineType,
  109. DataSourceCode = dto.DataSourceCode,
  110. SqlScript = dto.SqlScript,
  111. SqlParametersJson = dto.SqlParametersJson,
  112. TimeoutSeconds = NormalizeTimeout(dto.TimeoutSeconds),
  113. PublishStatus = StatusDraft,
  114. IsCurrent = false,
  115. Remark = dto.Remark,
  116. CreatedBy = operatorName,
  117. CreatedAt = DateTime.Now,
  118. };
  119. entity.Id = await _db.Insertable(entity).ExecuteReturnBigIdentityAsync();
  120. return ToDto(entity);
  121. }
  122. /// <summary>编辑草稿(仅 DRAFT 可改;已发布/停用版本不可改,须新建版本)。</summary>
  123. public async Task UpdateDraftAsync(long id, KpiCalcConfigUpsertDto dto, string operatorName)
  124. {
  125. var entity = await Query().Where(x => x.Id == id).FirstAsync()
  126. ?? throw Oops.Bah("配置版本不存在");
  127. if (entity.PublishStatus != StatusDraft)
  128. throw Oops.Bah("仅草稿(DRAFT)可编辑;已发布/停用版本请新建版本");
  129. ValidateEngineType(dto.CalcEngineType);
  130. var tenantId = RequireTenantId(entity);
  131. var affected = await _db.Ado.ExecuteCommandAsync(
  132. $@"UPDATE {TableName}
  133. SET CalcEngineType=@engine, DataSourceCode=@ds, SqlScript=@sql, SqlParametersJson=@pars,
  134. TimeoutSeconds=@timeout, Remark=@remark, UpdatedBy=@op, UpdatedAt=@now
  135. WHERE Id=@id AND TenantId=@tenant AND PublishStatus=@draft",
  136. new SugarParameter("@engine", dto.CalcEngineType),
  137. new SugarParameter("@ds", dto.DataSourceCode),
  138. new SugarParameter("@sql", dto.SqlScript),
  139. new SugarParameter("@pars", dto.SqlParametersJson),
  140. new SugarParameter("@timeout", NormalizeTimeout(dto.TimeoutSeconds)),
  141. new SugarParameter("@remark", dto.Remark),
  142. new SugarParameter("@op", operatorName),
  143. new SugarParameter("@now", DateTime.Now),
  144. new SugarParameter("@id", id),
  145. new SugarParameter("@tenant", tenantId),
  146. new SugarParameter("@draft", StatusDraft));
  147. if (affected <= 0) throw await ExplainZeroRowAsync(id, tenantId, "编辑草稿", "DRAFT", StatusDraft);
  148. }
  149. /// <summary>删除草稿(仅 DRAFT)。</summary>
  150. public async Task DeleteAsync(long id)
  151. {
  152. var entity = await Query().Where(x => x.Id == id).FirstAsync()
  153. ?? throw Oops.Bah("配置版本不存在");
  154. if (entity.PublishStatus != StatusDraft)
  155. throw Oops.Bah("仅草稿(DRAFT)可删除;已发布版本请用停用");
  156. var tenantId = RequireTenantId(entity);
  157. var affected = await _db.Ado.ExecuteCommandAsync(
  158. $"DELETE FROM {TableName} WHERE Id=@id AND TenantId=@tenant AND PublishStatus=@draft",
  159. new SugarParameter("@id", id),
  160. new SugarParameter("@tenant", tenantId),
  161. new SugarParameter("@draft", StatusDraft));
  162. if (affected <= 0) throw await ExplainZeroRowAsync(id, tenantId, "删除草稿", "DRAFT", StatusDraft);
  163. }
  164. /// <summary>SQL 安全校验(多层非正则)。</summary>
  165. public KpiSqlValidateResultDto Validate(string? sql)
  166. {
  167. var r = KpiSqlSecurityValidator.Validate(sql);
  168. return new KpiSqlValidateResultDto
  169. {
  170. Ok = r.Ok,
  171. ValidationStatus = r.Ok ? "VALID" : "INVALID",
  172. ErrorCode = r.ErrorCode,
  173. ErrorMessage = r.ErrorMessage,
  174. ReferencedTables = r.ReferencedTables,
  175. };
  176. }
  177. /// <summary>试算(服务端重校验 + 只读执行;不写正式 KPI 值)。</summary>
  178. public async Task<KpiSqlPreviewResultDto> PreviewAsync(KpiSqlPreviewDto dto, CancellationToken ct)
  179. {
  180. var bizDate = (dto.BizDate ?? DateTime.Today.AddDays(-1)).Date;
  181. var res = new KpiSqlPreviewResultDto
  182. {
  183. MetricCode = dto.MetricCode,
  184. EngineType = EngineConfigSql,
  185. DataSourceCode = dto.DataSourceCode,
  186. BizDate = bizDate.ToString("yyyy-MM-dd"),
  187. };
  188. var val = KpiSqlSecurityValidator.Validate(dto.SqlScript);
  189. res.ValidationStatus = val.Ok ? "VALID" : "INVALID";
  190. if (!val.Ok)
  191. {
  192. res.Ok = false;
  193. res.ResultStatus = "FAILED";
  194. res.ErrorCode = val.ErrorCode;
  195. res.ErrorMessage = val.ErrorMessage;
  196. return res;
  197. }
  198. var tenantId = ResolveKpiTenantId(dto.ModuleCode);
  199. // 业务门禁:试算是真实数据只读执行,未 READY 直接拒绝(不落只读事务、不占连接)。
  200. await EnsureBusinessReadyAsync(tenantId, dto.MetricCode, "试算");
  201. var pars = new KpiSqlRunParams
  202. {
  203. TenantId = tenantId,
  204. FactoryId = 1,
  205. ModuleCode = dto.ModuleCode,
  206. MetricCode = dto.MetricCode,
  207. BizDate = bizDate,
  208. PeriodStart = bizDate,
  209. PeriodEnd = bizDate.AddDays(1).AddSeconds(-1),
  210. SourceZtid = "pbxfxp",
  211. };
  212. var exec = await _executor.ExecuteAsync(dto.DataSourceCode, dto.SqlScript!, dto.TimeoutSeconds, true, pars, ct);
  213. res.Ok = exec.Status != "FAILED";
  214. res.ResultStatus = exec.Status;
  215. res.DurationMs = exec.DurationMs;
  216. res.RowCount = exec.RowCount;
  217. res.MetricValue = exec.MetricValue;
  218. res.NumeratorValue = exec.NumeratorValue;
  219. res.DenominatorValue = exec.DenominatorValue;
  220. res.ResultMessage = exec.ResultMessage;
  221. res.ErrorCode = exec.ErrorCode;
  222. res.ErrorMessage = exec.ErrorMessage;
  223. res.ExecutedParameters =
  224. $"tenant_id={pars.TenantId}, factory_id={pars.FactoryId}, module_code={pars.ModuleCode}, " +
  225. $"metric_code={pars.MetricCode}, biz_date={pars.BizDate:yyyy-MM-dd}, ztid={pars.SourceZtid}";
  226. return res;
  227. }
  228. /// <summary>发布:事务内校验→旧 current 置 0→本版 PUBLISHED+IsCurrent=1(保证单生效版本)。</summary>
  229. public async Task PublishAsync(long id, string operatorName)
  230. {
  231. var entity = await Query().Where(x => x.Id == id).FirstAsync()
  232. ?? throw Oops.Bah("配置版本不存在");
  233. if (entity.PublishStatus == StatusRetired)
  234. throw Oops.Bah("已停用版本不可发布,请新建版本");
  235. // CONFIG_SQL 发布前必须再过一次安全校验
  236. if (entity.CalcEngineType == EngineConfigSql)
  237. {
  238. var val = KpiSqlSecurityValidator.Validate(entity.SqlScript);
  239. if (!val.Ok) throw Oops.Bah($"SQL 安全校验未通过:{val.ErrorCode} {val.ErrorMessage}");
  240. if (string.IsNullOrWhiteSpace(entity.DataSourceCode))
  241. throw Oops.Bah("CONFIG_SQL 必须指定数据源");
  242. await EnsureBusinessReadyAsync(entity.TenantId ?? 0, entity.MetricCode, "发布");
  243. }
  244. var tenantId = RequireTenantId(entity);
  245. var tran = await _db.AsTenant().UseTranAsync(async () =>
  246. {
  247. // 旧 current 置 0:0 行是合法结果(本 KPI 此前无生效版本),故不做 affected 断言
  248. await _db.Ado.ExecuteCommandAsync(
  249. $@"UPDATE {TableName} SET IsCurrent=0
  250. WHERE MetricCode=@metric AND TenantId=@tenant AND Id<>@id AND IsCurrent=1",
  251. new SugarParameter("@metric", entity.MetricCode),
  252. new SugarParameter("@tenant", tenantId),
  253. new SugarParameter("@id", id));
  254. var affected = await _db.Ado.ExecuteCommandAsync(
  255. $@"UPDATE {TableName}
  256. SET PublishStatus=@published, IsCurrent=1, PublishedBy=@op, PublishedAt=@now
  257. WHERE Id=@id AND TenantId=@tenant AND PublishStatus<>@retired",
  258. new SugarParameter("@published", StatusPublished),
  259. new SugarParameter("@op", operatorName),
  260. new SugarParameter("@now", DateTime.Now),
  261. new SugarParameter("@id", id),
  262. new SugarParameter("@tenant", tenantId),
  263. new SugarParameter("@retired", StatusRetired));
  264. if (affected <= 0)
  265. throw await ExplainZeroRowAsync(id, tenantId, "发布", "非 RETIRED", StatusDraft, StatusPublished);
  266. });
  267. if (!tran.IsSuccess) throw tran.ErrorException;
  268. }
  269. /// <summary>激活/回滚:把某历史版本设为当前生效(其它同 KPI 版本 IsCurrent 置 0),事务内。仅对已发布版本。</summary>
  270. public async Task ActivateAsync(long id, string operatorName)
  271. {
  272. var entity = await Query().Where(x => x.Id == id).FirstAsync()
  273. ?? throw Oops.Bah("配置版本不存在");
  274. if (entity.PublishStatus != StatusPublished)
  275. throw Oops.Bah("只能激活已发布(PUBLISHED)版本;回滚是激活历史已发布版本,不是复制重发");
  276. if (entity.CalcEngineType == EngineConfigSql)
  277. await EnsureBusinessReadyAsync(entity.TenantId ?? 0, entity.MetricCode, "激活");
  278. var tenantId = RequireTenantId(entity);
  279. var tran = await _db.AsTenant().UseTranAsync(async () =>
  280. {
  281. // 旧 current 置 0:0 行是合法结果(本 KPI 此前无生效版本),故不做 affected 断言
  282. await _db.Ado.ExecuteCommandAsync(
  283. $@"UPDATE {TableName} SET IsCurrent=0
  284. WHERE MetricCode=@metric AND TenantId=@tenant AND Id<>@id AND IsCurrent=1",
  285. new SugarParameter("@metric", entity.MetricCode),
  286. new SugarParameter("@tenant", tenantId),
  287. new SugarParameter("@id", id));
  288. var affected = await _db.Ado.ExecuteCommandAsync(
  289. $@"UPDATE {TableName} SET IsCurrent=1, UpdatedBy=@op, UpdatedAt=@now
  290. WHERE Id=@id AND TenantId=@tenant AND PublishStatus=@published",
  291. new SugarParameter("@op", operatorName),
  292. new SugarParameter("@now", DateTime.Now),
  293. new SugarParameter("@id", id),
  294. new SugarParameter("@tenant", tenantId),
  295. new SugarParameter("@published", StatusPublished));
  296. if (affected <= 0)
  297. throw await ExplainZeroRowAsync(id, tenantId, "激活", "PUBLISHED", StatusPublished);
  298. });
  299. if (!tran.IsSuccess) throw tran.ErrorException;
  300. }
  301. /// <summary>停用(RETIRED,IsCurrent=0)。停用当前生效版本后无 current → 运行时回落 legacy。</summary>
  302. public async Task RetireAsync(long id, string operatorName)
  303. {
  304. var entity = await Query().Where(x => x.Id == id).FirstAsync()
  305. ?? throw Oops.Bah("配置版本不存在");
  306. var tenantId = RequireTenantId(entity);
  307. var affected = await _db.Ado.ExecuteCommandAsync(
  308. $@"UPDATE {TableName}
  309. SET PublishStatus=@retired, IsCurrent=0, RetiredBy=@op, RetiredAt=@now
  310. WHERE Id=@id AND TenantId=@tenant",
  311. new SugarParameter("@retired", StatusRetired),
  312. new SugarParameter("@op", operatorName),
  313. new SugarParameter("@now", DateTime.Now),
  314. new SugarParameter("@id", id),
  315. new SugarParameter("@tenant", tenantId));
  316. if (affected <= 0) throw await ExplainZeroRowAsync(id, tenantId, "停用", "任意状态");
  317. }
  318. /// <summary>已登记数据源列表(第一版仅本地中台库)。</summary>
  319. public List<object> ListDataSources() => new()
  320. {
  321. new { code = KpiSqlReadOnlyExecutor.LocalDataSourceCode, name = "本地中台库(aidopdev / mdp_std_* / dwd_*)", readOnly = true },
  322. };
  323. private static void ValidateEngineType(string engine)
  324. {
  325. if (engine != EngineLegacyCode && engine != EngineConfigSql && engine != EngineLegacyTvf)
  326. throw Oops.Bah($"非法引擎类型:{engine}(LEGACY_CODE/CONFIG_SQL/LEGACY_TVF)");
  327. }
  328. private static int NormalizeTimeout(int t) =>
  329. t <= 0 ? 30 : Math.Min(t, KpiSqlReadOnlyExecutor.SystemMaxTimeoutSeconds);
  330. private static KpiCalcConfigDto ToDto(AdoSmartOpsKpiCalcConfig e) => new()
  331. {
  332. Id = e.Id, TenantId = e.TenantId, MetricCode = e.MetricCode, ModuleCode = e.ModuleCode,
  333. VersionNo = e.VersionNo, CalcEngineType = e.CalcEngineType, DataSourceCode = e.DataSourceCode,
  334. SqlScript = e.SqlScript, SqlParametersJson = e.SqlParametersJson, TimeoutSeconds = e.TimeoutSeconds,
  335. PublishStatus = e.PublishStatus, IsCurrent = e.IsCurrent, Remark = e.Remark,
  336. CreatedBy = e.CreatedBy, CreatedAt = e.CreatedAt, PublishedBy = e.PublishedBy, PublishedAt = e.PublishedAt,
  337. };
  338. }