using Admin.NET.Core; using Admin.NET.Plugin.AiDOP.Entity.S8; using Admin.NET.Plugin.AiDOP.Service.S8.Rules; using Admin.NET.Plugin.AiDOP.Service.S8.Rules.Definitions; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.Service.S8; /// /// S8-RULE-GOVERNANCE-BATCH2:把**代码定义的规则**投射成**每租户一条运行策略**。 /// /// 解决什么:Batch 1 之后规则定义在代码里,但一个新租户并不会因此拥有 /// ado_s8_watch_rule 上的运行策略行 —— 现存的 Rule 01 那一行是历史上人在页面里手工建的。 /// 没有本服务,Batch 3 一旦封掉 Create API,新租户将永远拿不到任何规则。 /// /// 身份口径(本批最关键的一条):运行策略的逻辑唯一键是 /// (TenantId, RuleCode)不含 FactoryId。 /// 一个租户有 0 个、1 个还是 3 个工厂,都只应有**一条** Rule 01 运行策略。 /// 「Tenant → Factory → Rule」这个形状一旦写回来,规则实例数会随工厂数膨胀, /// 而 Tenant-only 改造(Batch 5/6)之后又要再清理一次。 /// /// 刻意不复用 :那是业务自助建规则的旧路径, /// Batch 3 就要退役。Provisioning 是系统内部写路径,依赖方向必须是 /// Catalog → Provisioning → DB,而不是架在一条即将删除的 API 之上。 /// /// 幂等性边界(不得夸大):本批只保证**应用层顺序幂等** —— /// 同一进程内重复执行不会重复建行。DB 上 UNIQUE(tenant_id, rule_code) 尚未建立 /// (属 Tenant-only migration 批次),因此两个实例同时 provisioning 理论上仍存在 race。 /// 该缺口由后续唯一索引最终封死,本批不引入分布式锁去假装解决它。 /// public class S8RuleProvisioningService : ITransient { /// /// 新行的 factory_id 兼容值。 /// /// 该列在库中是 bigint NOT NULL(实测 information_schema),因此不能写 NULL; /// 取 0 表示「无工厂作用域」。刻意不去解析租户的唯一工厂:那会让 /// 零工厂 / 多工厂租户直接 provisioning 失败,也会把 Factory 重新拉回规则身份。 /// 0 只是过渡期的兼容 metadata,不参与任何判定。 /// /// 已知过渡态:Batch 5 之前,ListEnabledScopesAsync 仍带 /// AND r.factory_id > 0,因此 factory_id=0 的新行即使被启用也不会被调度拾取。 /// 这是可接受的 —— 新 provision 的规则默认停用,正式部署要等 Tenant-only 完成。 /// 不得为了让过渡期能启用而把工厂作用域写回 Provisioning。 /// internal const long CompatibilityFactoryId = 0L; private readonly SqlSugarRepository _rep; private readonly IS8RuleCatalog _catalog; private readonly ILogger _logger; public S8RuleProvisioningService( SqlSugarRepository rep, IS8RuleCatalog catalog, ILogger logger) { _rep = rep; _catalog = catalog; _logger = logger; } /// 全量对账:所有启用租户 × 所有代码定义。启动期与运维触发使用。 public async Task SyncAsync(CancellationToken cancellationToken = default) { // Status = 1(StatusEnum.Enable)是本仓判定"租户有效"的既有口径: // S8WatchSchedulerService.ListEnabledScopesAsync、InventoryReconService、 // SourceDomainTenantResolver 三处都是 `SysTenant t ... AND t.Status = 1`。 // 这里沿用同一口径,不另造一套。 var tenantIds = await _rep.Context.Queryable() .Where(t => t.Status == StatusEnum.Enable) .Select(t => t.Id) .ToListAsync(cancellationToken); return await SyncTenantsAsync(tenantIds, cancellationToken); } /// /// 单租户对账。供新租户初始化 / 管理端手动补齐 / 测试使用。 /// 只处理该租户,不会连带修改其他租户的任何一行。 /// public async Task SyncTenantAsync(long tenantId, CancellationToken cancellationToken = default) { if (tenantId <= 0) throw new S8BizException("租户作用域非法"); var active = await _rep.Context.Queryable() .Where(t => t.Id == tenantId && t.Status == StatusEnum.Enable) .AnyAsync(cancellationToken); if (!active) throw new S8BizException("租户不存在或已停用,不予供给监控规则运行策略"); return await SyncTenantsAsync(new List { tenantId }, cancellationToken); } private async Task SyncTenantsAsync( List tenantIds, CancellationToken cancellationToken) { var definitions = _catalog.Definitions.ToList(); var result = new S8RuleProvisioningResult { TenantCount = tenantIds.Count, DefinitionCount = definitions.Count }; if (tenantIds.Count == 0 || definitions.Count == 0) { _logger.LogInformation( "s8_rule_provisioning_skipped tenants={Tenants} definitions={Definitions}", tenantIds.Count, definitions.Count); return result; } // 一次性取回这些租户的全部规则行:既用于存在性判断,也用于孤儿与重复统计。 // **查询谓词只有 tenant_id**,绝不含 factory_id —— 那正是本批要杜绝的身份形状。 var existingRows = await _rep.AsQueryable() .Where(x => tenantIds.Contains(x.TenantId)) .ToListAsync(cancellationToken); var plan = Plan(tenantIds, definitions, existingRows, DateTime.Now); result.UnchangedCount = plan.UnchangedCount; result.OrphanedCount = plan.OrphanedCount; result.DuplicateCount = plan.DuplicateCount; foreach (var row in plan.Inserts) { cancellationToken.ThrowIfCancellationRequested(); row.Id = await _rep.AsInsertable(row).ExecuteReturnBigIdentityAsync(); result.CreatedCount++; _logger.LogInformation( "s8_rule_provisioned tenant={Tenant} rule={Rule} id={Id}", row.TenantId, row.RuleCode, row.Id); } foreach (var row in plan.Updates) { cancellationToken.ThrowIfCancellationRequested(); // 白名单更新:物理上只可能写到 Definition 投影列。 // 租户参数(enabled / severity / poll / trigger / recover / params_json) // 与全部 13 个调度运行态列都不在 SetColumns 里,因此不可能被覆盖。 await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { DatasetCode = row.DatasetCode, RuleType = row.RuleType, SourceObjectType = row.SourceObjectType, WatchObjectType = row.WatchObjectType, SceneCode = row.SceneCode, StageCode = row.StageCode, OrderFlowCode = row.OrderFlowCode, RuleMechanism = row.RuleMechanism, UpdatedAt = DateTime.Now }) .Where(x => x.Id == row.Id) .ExecuteCommandAsync(cancellationToken); result.RefreshedCount++; _logger.LogInformation( "s8_rule_projection_refreshed tenant={Tenant} rule={Rule} id={Id}", row.TenantId, row.RuleCode, row.Id); } if (result.OrphanedCount > 0) _logger.LogWarning( "s8_rule_orphaned count={Count};这些规则在当前代码版本中没有定义,保持原样不删除、不启用", result.OrphanedCount); if (result.DuplicateCount > 0) _logger.LogWarning( "s8_rule_duplicate_identity count={Count};同 (tenant, rule_code) 存在多行," + "在 UNIQUE(tenant_id, rule_code) 建立前只刷新 Id 最小的一行", result.DuplicateCount); _logger.LogInformation( "s8_rule_provisioning_done tenants={Tenants} definitions={Definitions} created={Created} refreshed={Refreshed} unchanged={Unchanged} orphaned={Orphaned} duplicate={Duplicate}", result.TenantCount, result.DefinitionCount, result.CreatedCount, result.RefreshedCount, result.UnchangedCount, result.OrphanedCount, result.DuplicateCount); return result; } // ================================================================================ // 纯决策核心 // // 抽成 static 是为了让「租户隔离」「工厂无关」「不覆盖租户参数」这三条不变量 // 能在**不接数据库**的情况下被逐字段断言。Batch 1 的 ApplyParameters 是同一动机: // 藏在必须有仓储才能跑的方法里的规则,实际上没人验证过。 // ================================================================================ /// /// 计算供给方案。 应为这些租户的**全部**规则行。 /// internal static S8RuleProvisioningPlan Plan( IReadOnlyCollection tenantIds, IReadOnlyCollection definitions, IReadOnlyCollection existingRows, DateTime now) { var plan = new S8RuleProvisioningPlan(); var definedCodes = definitions .Select(d => d.RuleCode) .ToHashSet(StringComparer.OrdinalIgnoreCase); // 身份 = (TenantId, RuleCode)。**没有 FactoryId。** var byIdentity = existingRows .GroupBy(r => (r.TenantId, Code: r.RuleCode ?? string.Empty), TenantRuleComparer.Instance) .ToDictionary(g => g.Key, g => g.OrderBy(r => r.Id).ToList(), TenantRuleComparer.Instance); foreach (var tenantId in tenantIds) { foreach (var definition in definitions) { var key = (tenantId, Code: definition.RuleCode); if (!byIdentity.TryGetValue(key, out var rows) || rows.Count == 0) { plan.Inserts.Add(BuildNewRow(tenantId, definition, now)); continue; } // 同身份多行是历史遗留(曾按工厂各建一条)。在唯一索引建立前, // 只把 Id 最小的一行当作 canonical 刷新,其余计数上报、不动、不删。 if (rows.Count > 1) plan.DuplicateCount += rows.Count - 1; var canonical = rows[0]; if (ApplyProjection(canonical, definition)) plan.Updates.Add(canonical); else plan.UnchangedCount++; } } // 孤儿:库里有、目录里没有。不删、不改、不启用。 plan.OrphanedCount = existingRows.Count(r => !definedCodes.Contains(r.RuleCode ?? string.Empty)); return plan; } /// /// 构造新租户运行策略行。 /// 只写定义投影 + 运行参数默认值;运行态列全部保持实体的"尚未运行"初值, /// 不伪造 last_run_at / last_status 之类的运行历史。 /// internal static AdoS8WatchRule BuildNewRow(long tenantId, S8RuleDefinition definition, DateTime now) { var policy = definition.Parameters ?? new S8RuleParameterPolicy(); var row = new AdoS8WatchRule { Id = 0, TenantId = tenantId, FactoryId = CompatibilityFactoryId, // ── Definition 投影 ── RuleCode = definition.RuleCode, DatasetCode = definition.DatasetCode, RuleType = definition.RuleType, SourceObjectType = definition.SourceObjectType, WatchObjectType = definition.SourceObjectType, SceneCode = definition.SceneCode, StageCode = definition.StageCode, OrderFlowCode = definition.OrderFlowCode, RuleMechanism = definition.RuleMechanism, // ── 运行参数默认值(全部来自定义声明的策略,不写死在这里)── Enabled = false, Severity = policy.SeverityDefault, PollIntervalSeconds = policy.PollIntervalSecondsDefault, TriggerCountRequired = policy.TriggerCountRequiredDefault, RecoverCountRequired = policy.RecoverCountRequiredDefault, CreatedAt = now, UpdatedAt = null }; // params_json 只承载 B 类运行参数;新行不写任何判定语义。 row.ParamsJson = new S8RuleRuntimeParameters { GraceMinutes = policy.GraceMinutesDefault }.ToParamsJson(); return row; } /// /// 把定义投影刷到既有行上。只碰投影列;返回是否确有变化。 /// /// 刻意不碰Enabled / Severity / PollIntervalSeconds / /// TriggerCountRequired / RecoverCountRequired / ParamsJson / /// FactoryId / 全部运行态列。把租户调过的参数刷回默认值是本服务最大的风险, /// 这条边界必须在类型层面而不是靠自觉守住。 /// /// ParamsJson 尤其不能重写:既有行可能仍是旧的 A+B 混装 JSON, /// 而 Batch 1 已保证 A 类字段在运行期被忽略;重写整块反而可能误伤租户的 B 类参数。 /// 同理,ParamsJson = NULL(G1 造成的历史状态)也不自动补写 —— /// 定义在代码里,NULL 时运行参数会正常回落默认值。 /// internal static bool ApplyProjection(AdoS8WatchRule row, S8RuleDefinition definition) { var changed = false; void SetString(Func get, Action set, string? next) { if (string.Equals(get(row), next, StringComparison.Ordinal)) return; set(next); changed = true; } SetString(r => r.DatasetCode, v => row.DatasetCode = v, definition.DatasetCode); SetString(r => r.RuleType, v => row.RuleType = v, definition.RuleType); SetString(r => r.SourceObjectType, v => row.SourceObjectType = v, definition.SourceObjectType); // WatchObjectType 是 SourceObjectType 的 legacy 同义列。它同样是定义派生的, // 留两份互相矛盾的"这条规则看什么对象"正是本批要消除的漂移。 SetString(r => r.WatchObjectType, v => row.WatchObjectType = v ?? string.Empty, definition.SourceObjectType); SetString(r => r.SceneCode, v => row.SceneCode = v ?? string.Empty, definition.SceneCode); SetString(r => r.StageCode, v => row.StageCode = v, definition.StageCode); SetString(r => r.OrderFlowCode, v => row.OrderFlowCode = v, definition.OrderFlowCode); SetString(r => r.RuleMechanism, v => row.RuleMechanism = v, definition.RuleMechanism); return changed; } /// (TenantId, RuleCode) 身份比较器。RuleCode 按不敏感比较,与 Catalog 查找一致。 private sealed class TenantRuleComparer : IEqualityComparer<(long TenantId, string Code)> { public static readonly TenantRuleComparer Instance = new(); public bool Equals((long TenantId, string Code) x, (long TenantId, string Code) y) => x.TenantId == y.TenantId && string.Equals(x.Code, y.Code, StringComparison.OrdinalIgnoreCase); public int GetHashCode((long TenantId, string Code) obj) => HashCode.Combine(obj.TenantId, StringComparer.OrdinalIgnoreCase.GetHashCode(obj.Code ?? string.Empty)); } } /// 供给方案(纯数据,不含任何 DB 访问)。 internal sealed class S8RuleProvisioningPlan { /// 待新建的运行策略行。 public List Inserts { get; } = new(); /// 待刷新投影的既有行(投影列已就地更新)。 public List Updates { get; } = new(); /// 投影已一致、无需写入的数量。 public int UnchangedCount { get; set; } /// 库中存在但当前代码版本没有定义的规则行数。 public int OrphanedCount { get; set; } /// 同 (tenant, rule_code) 的多余行数(唯一索引建立前的历史遗留)。 public int DuplicateCount { get; set; } } /// /// 供给结果。刻意给出分项计数而不是只回 200 —— /// 调用方必须能知道「这次到底做了什么」,否则运维无法判断新规则是否已经铺开。 /// public sealed class S8RuleProvisioningResult { /// 本次处理的租户数。 public int TenantCount { get; set; } /// 当前代码版本中的规则定义数。 public int DefinitionCount { get; set; } /// 新建的运行策略行数。 public int CreatedCount { get; set; } /// 投影被纠正的行数。 public int RefreshedCount { get; set; } /// 投影已一致、未写入的行数。 public int UnchangedCount { get; set; } /// 无代码定义的孤儿行数(保持原样)。 public int OrphanedCount { get; set; } /// 同 (tenant, rule_code) 的多余行数(未删除,仅上报)。 public int DuplicateCount { get; set; } }