S8RuleProvisioningService.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379
  1. using Admin.NET.Core;
  2. using Admin.NET.Plugin.AiDOP.Entity.S8;
  3. using Admin.NET.Plugin.AiDOP.Service.S8.Rules;
  4. using Admin.NET.Plugin.AiDOP.Service.S8.Rules.Definitions;
  5. using Microsoft.Extensions.Logging;
  6. namespace Admin.NET.Plugin.AiDOP.Service.S8;
  7. /// <summary>
  8. /// S8-RULE-GOVERNANCE-BATCH2:把**代码定义的规则**投射成**每租户一条运行策略**。
  9. ///
  10. /// <para><b>解决什么</b>:Batch 1 之后规则定义在代码里,但一个新租户并不会因此拥有
  11. /// <c>ado_s8_watch_rule</c> 上的运行策略行 —— 现存的 Rule 01 那一行是历史上人在页面里手工建的。
  12. /// 没有本服务,Batch 3 一旦封掉 Create API,新租户将永远拿不到任何规则。</para>
  13. ///
  14. /// <para><b>身份口径(本批最关键的一条)</b>:运行策略的逻辑唯一键是
  15. /// <c>(TenantId, RuleCode)</c>,<b>不含 FactoryId</b>。
  16. /// 一个租户有 0 个、1 个还是 3 个工厂,都只应有**一条** Rule 01 运行策略。
  17. /// 「Tenant → Factory → Rule」这个形状一旦写回来,规则实例数会随工厂数膨胀,
  18. /// 而 Tenant-only 改造(Batch 5/6)之后又要再清理一次。</para>
  19. ///
  20. /// <para><b>刻意不复用 <see cref="S8WatchRuleService.CreateAsync"/></b>:那是业务自助建规则的旧路径,
  21. /// Batch 3 就要退役。Provisioning 是系统内部写路径,依赖方向必须是
  22. /// <c>Catalog → Provisioning → DB</c>,而不是架在一条即将删除的 API 之上。</para>
  23. ///
  24. /// <para><b>幂等性边界(不得夸大)</b>:本批只保证**应用层顺序幂等** ——
  25. /// 同一进程内重复执行不会重复建行。DB 上 <c>UNIQUE(tenant_id, rule_code)</c> 尚未建立
  26. /// (属 Tenant-only migration 批次),因此两个实例同时 provisioning 理论上仍存在 race。
  27. /// 该缺口由后续唯一索引最终封死,本批不引入分布式锁去假装解决它。</para>
  28. /// </summary>
  29. public class S8RuleProvisioningService : ITransient
  30. {
  31. /// <summary>
  32. /// 新行的 <c>factory_id</c> 兼容值。
  33. ///
  34. /// <para>该列在库中是 <c>bigint NOT NULL</c>(实测 information_schema),因此不能写 NULL;
  35. /// 取 0 表示「无工厂作用域」。<b>刻意不去解析租户的唯一工厂</b>:那会让
  36. /// 零工厂 / 多工厂租户直接 provisioning 失败,也会把 Factory 重新拉回规则身份。
  37. /// 0 只是过渡期的兼容 metadata,不参与任何判定。</para>
  38. ///
  39. /// <para>已知过渡态:Batch 5 之前,<c>ListEnabledScopesAsync</c> 仍带
  40. /// <c>AND r.factory_id &gt; 0</c>,因此 factory_id=0 的新行即使被启用也不会被调度拾取。
  41. /// 这是可接受的 —— 新 provision 的规则默认停用,正式部署要等 Tenant-only 完成。
  42. /// 不得为了让过渡期能启用而把工厂作用域写回 Provisioning。</para>
  43. /// </summary>
  44. internal const long CompatibilityFactoryId = 0L;
  45. private readonly SqlSugarRepository<AdoS8WatchRule> _rep;
  46. private readonly IS8RuleCatalog _catalog;
  47. private readonly ILogger<S8RuleProvisioningService> _logger;
  48. public S8RuleProvisioningService(
  49. SqlSugarRepository<AdoS8WatchRule> rep,
  50. IS8RuleCatalog catalog,
  51. ILogger<S8RuleProvisioningService> logger)
  52. {
  53. _rep = rep;
  54. _catalog = catalog;
  55. _logger = logger;
  56. }
  57. /// <summary>全量对账:所有启用租户 × 所有代码定义。启动期与运维触发使用。</summary>
  58. public async Task<S8RuleProvisioningResult> SyncAsync(CancellationToken cancellationToken = default)
  59. {
  60. // Status = 1(StatusEnum.Enable)是本仓判定"租户有效"的既有口径:
  61. // S8WatchSchedulerService.ListEnabledScopesAsync、InventoryReconService、
  62. // SourceDomainTenantResolver 三处都是 `SysTenant t ... AND t.Status = 1`。
  63. // 这里沿用同一口径,不另造一套。
  64. var tenantIds = await _rep.Context.Queryable<SysTenant>()
  65. .Where(t => t.Status == StatusEnum.Enable)
  66. .Select(t => t.Id)
  67. .ToListAsync(cancellationToken);
  68. return await SyncTenantsAsync(tenantIds, cancellationToken);
  69. }
  70. /// <summary>
  71. /// 单租户对账。供新租户初始化 / 管理端手动补齐 / 测试使用。
  72. /// <b>只处理该租户</b>,不会连带修改其他租户的任何一行。
  73. /// </summary>
  74. public async Task<S8RuleProvisioningResult> SyncTenantAsync(long tenantId, CancellationToken cancellationToken = default)
  75. {
  76. if (tenantId <= 0) throw new S8BizException("租户作用域非法");
  77. var active = await _rep.Context.Queryable<SysTenant>()
  78. .Where(t => t.Id == tenantId && t.Status == StatusEnum.Enable)
  79. .AnyAsync(cancellationToken);
  80. if (!active)
  81. throw new S8BizException("租户不存在或已停用,不予供给监控规则运行策略");
  82. return await SyncTenantsAsync(new List<long> { tenantId }, cancellationToken);
  83. }
  84. private async Task<S8RuleProvisioningResult> SyncTenantsAsync(
  85. List<long> tenantIds, CancellationToken cancellationToken)
  86. {
  87. var definitions = _catalog.Definitions.ToList();
  88. var result = new S8RuleProvisioningResult
  89. {
  90. TenantCount = tenantIds.Count,
  91. DefinitionCount = definitions.Count
  92. };
  93. if (tenantIds.Count == 0 || definitions.Count == 0)
  94. {
  95. _logger.LogInformation(
  96. "s8_rule_provisioning_skipped tenants={Tenants} definitions={Definitions}",
  97. tenantIds.Count, definitions.Count);
  98. return result;
  99. }
  100. // 一次性取回这些租户的全部规则行:既用于存在性判断,也用于孤儿与重复统计。
  101. // **查询谓词只有 tenant_id**,绝不含 factory_id —— 那正是本批要杜绝的身份形状。
  102. var existingRows = await _rep.AsQueryable()
  103. .Where(x => tenantIds.Contains(x.TenantId))
  104. .ToListAsync(cancellationToken);
  105. var plan = Plan(tenantIds, definitions, existingRows, DateTime.Now);
  106. result.UnchangedCount = plan.UnchangedCount;
  107. result.OrphanedCount = plan.OrphanedCount;
  108. result.DuplicateCount = plan.DuplicateCount;
  109. foreach (var row in plan.Inserts)
  110. {
  111. cancellationToken.ThrowIfCancellationRequested();
  112. row.Id = await _rep.AsInsertable(row).ExecuteReturnBigIdentityAsync();
  113. result.CreatedCount++;
  114. _logger.LogInformation(
  115. "s8_rule_provisioned tenant={Tenant} rule={Rule} id={Id}", row.TenantId, row.RuleCode, row.Id);
  116. }
  117. foreach (var row in plan.Updates)
  118. {
  119. cancellationToken.ThrowIfCancellationRequested();
  120. // 白名单更新:物理上只可能写到 Definition 投影列。
  121. // 租户参数(enabled / severity / poll / trigger / recover / params_json)
  122. // 与全部 13 个调度运行态列都不在 SetColumns 里,因此不可能被覆盖。
  123. await _rep.Context.Updateable<AdoS8WatchRule>()
  124. .SetColumns(x => new AdoS8WatchRule
  125. {
  126. DatasetCode = row.DatasetCode,
  127. RuleType = row.RuleType,
  128. SourceObjectType = row.SourceObjectType,
  129. WatchObjectType = row.WatchObjectType,
  130. SceneCode = row.SceneCode,
  131. StageCode = row.StageCode,
  132. OrderFlowCode = row.OrderFlowCode,
  133. RuleMechanism = row.RuleMechanism,
  134. UpdatedAt = DateTime.Now
  135. })
  136. .Where(x => x.Id == row.Id)
  137. .ExecuteCommandAsync(cancellationToken);
  138. result.RefreshedCount++;
  139. _logger.LogInformation(
  140. "s8_rule_projection_refreshed tenant={Tenant} rule={Rule} id={Id}",
  141. row.TenantId, row.RuleCode, row.Id);
  142. }
  143. if (result.OrphanedCount > 0)
  144. _logger.LogWarning(
  145. "s8_rule_orphaned count={Count};这些规则在当前代码版本中没有定义,保持原样不删除、不启用",
  146. result.OrphanedCount);
  147. if (result.DuplicateCount > 0)
  148. _logger.LogWarning(
  149. "s8_rule_duplicate_identity count={Count};同 (tenant, rule_code) 存在多行,"
  150. + "在 UNIQUE(tenant_id, rule_code) 建立前只刷新 Id 最小的一行",
  151. result.DuplicateCount);
  152. _logger.LogInformation(
  153. "s8_rule_provisioning_done tenants={Tenants} definitions={Definitions} created={Created} refreshed={Refreshed} unchanged={Unchanged} orphaned={Orphaned} duplicate={Duplicate}",
  154. result.TenantCount, result.DefinitionCount, result.CreatedCount,
  155. result.RefreshedCount, result.UnchangedCount, result.OrphanedCount, result.DuplicateCount);
  156. return result;
  157. }
  158. // ================================================================================
  159. // 纯决策核心
  160. //
  161. // 抽成 static 是为了让「租户隔离」「工厂无关」「不覆盖租户参数」这三条不变量
  162. // 能在**不接数据库**的情况下被逐字段断言。Batch 1 的 ApplyParameters 是同一动机:
  163. // 藏在必须有仓储才能跑的方法里的规则,实际上没人验证过。
  164. // ================================================================================
  165. /// <summary>
  166. /// 计算供给方案。<paramref name="existingRows"/> 应为这些租户的**全部**规则行。
  167. /// </summary>
  168. internal static S8RuleProvisioningPlan Plan(
  169. IReadOnlyCollection<long> tenantIds,
  170. IReadOnlyCollection<S8RuleDefinition> definitions,
  171. IReadOnlyCollection<AdoS8WatchRule> existingRows,
  172. DateTime now)
  173. {
  174. var plan = new S8RuleProvisioningPlan();
  175. var definedCodes = definitions
  176. .Select(d => d.RuleCode)
  177. .ToHashSet(StringComparer.OrdinalIgnoreCase);
  178. // 身份 = (TenantId, RuleCode)。**没有 FactoryId。**
  179. var byIdentity = existingRows
  180. .GroupBy(r => (r.TenantId, Code: r.RuleCode ?? string.Empty), TenantRuleComparer.Instance)
  181. .ToDictionary(g => g.Key, g => g.OrderBy(r => r.Id).ToList(), TenantRuleComparer.Instance);
  182. foreach (var tenantId in tenantIds)
  183. {
  184. foreach (var definition in definitions)
  185. {
  186. var key = (tenantId, Code: definition.RuleCode);
  187. if (!byIdentity.TryGetValue(key, out var rows) || rows.Count == 0)
  188. {
  189. plan.Inserts.Add(BuildNewRow(tenantId, definition, now));
  190. continue;
  191. }
  192. // 同身份多行是历史遗留(曾按工厂各建一条)。在唯一索引建立前,
  193. // 只把 Id 最小的一行当作 canonical 刷新,其余计数上报、不动、不删。
  194. if (rows.Count > 1) plan.DuplicateCount += rows.Count - 1;
  195. var canonical = rows[0];
  196. if (ApplyProjection(canonical, definition))
  197. plan.Updates.Add(canonical);
  198. else
  199. plan.UnchangedCount++;
  200. }
  201. }
  202. // 孤儿:库里有、目录里没有。不删、不改、不启用。
  203. plan.OrphanedCount = existingRows.Count(r => !definedCodes.Contains(r.RuleCode ?? string.Empty));
  204. return plan;
  205. }
  206. /// <summary>
  207. /// 构造新租户运行策略行。
  208. /// <b>只写定义投影 + 运行参数默认值</b>;运行态列全部保持实体的"尚未运行"初值,
  209. /// 不伪造 last_run_at / last_status 之类的运行历史。
  210. /// </summary>
  211. internal static AdoS8WatchRule BuildNewRow(long tenantId, S8RuleDefinition definition, DateTime now)
  212. {
  213. var policy = definition.Parameters ?? new S8RuleParameterPolicy();
  214. var row = new AdoS8WatchRule
  215. {
  216. Id = 0,
  217. TenantId = tenantId,
  218. FactoryId = CompatibilityFactoryId,
  219. // ── Definition 投影 ──
  220. RuleCode = definition.RuleCode,
  221. DatasetCode = definition.DatasetCode,
  222. RuleType = definition.RuleType,
  223. SourceObjectType = definition.SourceObjectType,
  224. WatchObjectType = definition.SourceObjectType,
  225. SceneCode = definition.SceneCode,
  226. StageCode = definition.StageCode,
  227. OrderFlowCode = definition.OrderFlowCode,
  228. RuleMechanism = definition.RuleMechanism,
  229. // ── 运行参数默认值(全部来自定义声明的策略,不写死在这里)──
  230. Enabled = false,
  231. Severity = policy.SeverityDefault,
  232. PollIntervalSeconds = policy.PollIntervalSecondsDefault,
  233. TriggerCountRequired = policy.TriggerCountRequiredDefault,
  234. RecoverCountRequired = policy.RecoverCountRequiredDefault,
  235. CreatedAt = now,
  236. UpdatedAt = null
  237. };
  238. // params_json 只承载 B 类运行参数;新行不写任何判定语义。
  239. row.ParamsJson = new S8RuleRuntimeParameters
  240. {
  241. GraceMinutes = policy.GraceMinutesDefault
  242. }.ToParamsJson();
  243. return row;
  244. }
  245. /// <summary>
  246. /// 把定义投影刷到既有行上。<b>只碰投影列</b>;返回是否确有变化。
  247. ///
  248. /// <para>刻意<b>不碰</b>:<c>Enabled</c> / <c>Severity</c> / <c>PollIntervalSeconds</c> /
  249. /// <c>TriggerCountRequired</c> / <c>RecoverCountRequired</c> / <c>ParamsJson</c> /
  250. /// <c>FactoryId</c> / 全部运行态列。把租户调过的参数刷回默认值是本服务最大的风险,
  251. /// 这条边界必须在类型层面而不是靠自觉守住。</para>
  252. ///
  253. /// <para><c>ParamsJson</c> 尤其不能重写:既有行可能仍是旧的 A+B 混装 JSON,
  254. /// 而 Batch 1 已保证 A 类字段在运行期被忽略;重写整块反而可能误伤租户的 B 类参数。
  255. /// 同理,<c>ParamsJson = NULL</c>(G1 造成的历史状态)也不自动补写 ——
  256. /// 定义在代码里,NULL 时运行参数会正常回落默认值。</para>
  257. /// </summary>
  258. internal static bool ApplyProjection(AdoS8WatchRule row, S8RuleDefinition definition)
  259. {
  260. var changed = false;
  261. void SetString(Func<AdoS8WatchRule, string?> get, Action<string?> set, string? next)
  262. {
  263. if (string.Equals(get(row), next, StringComparison.Ordinal)) return;
  264. set(next);
  265. changed = true;
  266. }
  267. SetString(r => r.DatasetCode, v => row.DatasetCode = v, definition.DatasetCode);
  268. SetString(r => r.RuleType, v => row.RuleType = v, definition.RuleType);
  269. SetString(r => r.SourceObjectType, v => row.SourceObjectType = v, definition.SourceObjectType);
  270. // WatchObjectType 是 SourceObjectType 的 legacy 同义列。它同样是定义派生的,
  271. // 留两份互相矛盾的"这条规则看什么对象"正是本批要消除的漂移。
  272. SetString(r => r.WatchObjectType, v => row.WatchObjectType = v ?? string.Empty, definition.SourceObjectType);
  273. SetString(r => r.SceneCode, v => row.SceneCode = v ?? string.Empty, definition.SceneCode);
  274. SetString(r => r.StageCode, v => row.StageCode = v, definition.StageCode);
  275. SetString(r => r.OrderFlowCode, v => row.OrderFlowCode = v, definition.OrderFlowCode);
  276. SetString(r => r.RuleMechanism, v => row.RuleMechanism = v, definition.RuleMechanism);
  277. return changed;
  278. }
  279. /// <summary>(TenantId, RuleCode) 身份比较器。RuleCode 按不敏感比较,与 Catalog 查找一致。</summary>
  280. private sealed class TenantRuleComparer : IEqualityComparer<(long TenantId, string Code)>
  281. {
  282. public static readonly TenantRuleComparer Instance = new();
  283. public bool Equals((long TenantId, string Code) x, (long TenantId, string Code) y) =>
  284. x.TenantId == y.TenantId && string.Equals(x.Code, y.Code, StringComparison.OrdinalIgnoreCase);
  285. public int GetHashCode((long TenantId, string Code) obj) =>
  286. HashCode.Combine(obj.TenantId, StringComparer.OrdinalIgnoreCase.GetHashCode(obj.Code ?? string.Empty));
  287. }
  288. }
  289. /// <summary>供给方案(纯数据,不含任何 DB 访问)。</summary>
  290. internal sealed class S8RuleProvisioningPlan
  291. {
  292. /// <summary>待新建的运行策略行。</summary>
  293. public List<AdoS8WatchRule> Inserts { get; } = new();
  294. /// <summary>待刷新投影的既有行(投影列已就地更新)。</summary>
  295. public List<AdoS8WatchRule> Updates { get; } = new();
  296. /// <summary>投影已一致、无需写入的数量。</summary>
  297. public int UnchangedCount { get; set; }
  298. /// <summary>库中存在但当前代码版本没有定义的规则行数。</summary>
  299. public int OrphanedCount { get; set; }
  300. /// <summary>同 (tenant, rule_code) 的多余行数(唯一索引建立前的历史遗留)。</summary>
  301. public int DuplicateCount { get; set; }
  302. }
  303. /// <summary>
  304. /// 供给结果。刻意给出分项计数而不是只回 200 ——
  305. /// 调用方必须能知道「这次到底做了什么」,否则运维无法判断新规则是否已经铺开。
  306. /// </summary>
  307. public sealed class S8RuleProvisioningResult
  308. {
  309. /// <summary>本次处理的租户数。</summary>
  310. public int TenantCount { get; set; }
  311. /// <summary>当前代码版本中的规则定义数。</summary>
  312. public int DefinitionCount { get; set; }
  313. /// <summary>新建的运行策略行数。</summary>
  314. public int CreatedCount { get; set; }
  315. /// <summary>投影被纠正的行数。</summary>
  316. public int RefreshedCount { get; set; }
  317. /// <summary>投影已一致、未写入的行数。</summary>
  318. public int UnchangedCount { get; set; }
  319. /// <summary>无代码定义的孤儿行数(保持原样)。</summary>
  320. public int OrphanedCount { get; set; }
  321. /// <summary>同 (tenant, rule_code) 的多余行数(未删除,仅上报)。</summary>
  322. public int DuplicateCount { get; set; }
  323. }