| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379 |
- 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;
- /// <summary>
- /// S8-RULE-GOVERNANCE-BATCH2:把**代码定义的规则**投射成**每租户一条运行策略**。
- ///
- /// <para><b>解决什么</b>:Batch 1 之后规则定义在代码里,但一个新租户并不会因此拥有
- /// <c>ado_s8_watch_rule</c> 上的运行策略行 —— 现存的 Rule 01 那一行是历史上人在页面里手工建的。
- /// 没有本服务,Batch 3 一旦封掉 Create API,新租户将永远拿不到任何规则。</para>
- ///
- /// <para><b>身份口径(本批最关键的一条)</b>:运行策略的逻辑唯一键是
- /// <c>(TenantId, RuleCode)</c>,<b>不含 FactoryId</b>。
- /// 一个租户有 0 个、1 个还是 3 个工厂,都只应有**一条** Rule 01 运行策略。
- /// 「Tenant → Factory → Rule」这个形状一旦写回来,规则实例数会随工厂数膨胀,
- /// 而 Tenant-only 改造(Batch 5/6)之后又要再清理一次。</para>
- ///
- /// <para><b>刻意不复用 <see cref="S8WatchRuleService.CreateAsync"/></b>:那是业务自助建规则的旧路径,
- /// Batch 3 就要退役。Provisioning 是系统内部写路径,依赖方向必须是
- /// <c>Catalog → Provisioning → DB</c>,而不是架在一条即将删除的 API 之上。</para>
- ///
- /// <para><b>幂等性边界(不得夸大)</b>:本批只保证**应用层顺序幂等** ——
- /// 同一进程内重复执行不会重复建行。DB 上 <c>UNIQUE(tenant_id, rule_code)</c> 尚未建立
- /// (属 Tenant-only migration 批次),因此两个实例同时 provisioning 理论上仍存在 race。
- /// 该缺口由后续唯一索引最终封死,本批不引入分布式锁去假装解决它。</para>
- /// </summary>
- public class S8RuleProvisioningService : ITransient
- {
- /// <summary>
- /// 新行的 <c>factory_id</c> 兼容值。
- ///
- /// <para>该列在库中是 <c>bigint NOT NULL</c>(实测 information_schema),因此不能写 NULL;
- /// 取 0 表示「无工厂作用域」。<b>刻意不去解析租户的唯一工厂</b>:那会让
- /// 零工厂 / 多工厂租户直接 provisioning 失败,也会把 Factory 重新拉回规则身份。
- /// 0 只是过渡期的兼容 metadata,不参与任何判定。</para>
- ///
- /// <para>已知过渡态:Batch 5 之前,<c>ListEnabledScopesAsync</c> 仍带
- /// <c>AND r.factory_id > 0</c>,因此 factory_id=0 的新行即使被启用也不会被调度拾取。
- /// 这是可接受的 —— 新 provision 的规则默认停用,正式部署要等 Tenant-only 完成。
- /// 不得为了让过渡期能启用而把工厂作用域写回 Provisioning。</para>
- /// </summary>
- internal const long CompatibilityFactoryId = 0L;
- private readonly SqlSugarRepository<AdoS8WatchRule> _rep;
- private readonly IS8RuleCatalog _catalog;
- private readonly ILogger<S8RuleProvisioningService> _logger;
- public S8RuleProvisioningService(
- SqlSugarRepository<AdoS8WatchRule> rep,
- IS8RuleCatalog catalog,
- ILogger<S8RuleProvisioningService> logger)
- {
- _rep = rep;
- _catalog = catalog;
- _logger = logger;
- }
- /// <summary>全量对账:所有启用租户 × 所有代码定义。启动期与运维触发使用。</summary>
- public async Task<S8RuleProvisioningResult> SyncAsync(CancellationToken cancellationToken = default)
- {
- // Status = 1(StatusEnum.Enable)是本仓判定"租户有效"的既有口径:
- // S8WatchSchedulerService.ListEnabledScopesAsync、InventoryReconService、
- // SourceDomainTenantResolver 三处都是 `SysTenant t ... AND t.Status = 1`。
- // 这里沿用同一口径,不另造一套。
- var tenantIds = await _rep.Context.Queryable<SysTenant>()
- .Where(t => t.Status == StatusEnum.Enable)
- .Select(t => t.Id)
- .ToListAsync(cancellationToken);
- return await SyncTenantsAsync(tenantIds, cancellationToken);
- }
- /// <summary>
- /// 单租户对账。供新租户初始化 / 管理端手动补齐 / 测试使用。
- /// <b>只处理该租户</b>,不会连带修改其他租户的任何一行。
- /// </summary>
- public async Task<S8RuleProvisioningResult> SyncTenantAsync(long tenantId, CancellationToken cancellationToken = default)
- {
- if (tenantId <= 0) throw new S8BizException("租户作用域非法");
- var active = await _rep.Context.Queryable<SysTenant>()
- .Where(t => t.Id == tenantId && t.Status == StatusEnum.Enable)
- .AnyAsync(cancellationToken);
- if (!active)
- throw new S8BizException("租户不存在或已停用,不予供给监控规则运行策略");
- return await SyncTenantsAsync(new List<long> { tenantId }, cancellationToken);
- }
- private async Task<S8RuleProvisioningResult> SyncTenantsAsync(
- List<long> 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<AdoS8WatchRule>()
- .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 是同一动机:
- // 藏在必须有仓储才能跑的方法里的规则,实际上没人验证过。
- // ================================================================================
- /// <summary>
- /// 计算供给方案。<paramref name="existingRows"/> 应为这些租户的**全部**规则行。
- /// </summary>
- internal static S8RuleProvisioningPlan Plan(
- IReadOnlyCollection<long> tenantIds,
- IReadOnlyCollection<S8RuleDefinition> definitions,
- IReadOnlyCollection<AdoS8WatchRule> 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;
- }
- /// <summary>
- /// 构造新租户运行策略行。
- /// <b>只写定义投影 + 运行参数默认值</b>;运行态列全部保持实体的"尚未运行"初值,
- /// 不伪造 last_run_at / last_status 之类的运行历史。
- /// </summary>
- 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;
- }
- /// <summary>
- /// 把定义投影刷到既有行上。<b>只碰投影列</b>;返回是否确有变化。
- ///
- /// <para>刻意<b>不碰</b>:<c>Enabled</c> / <c>Severity</c> / <c>PollIntervalSeconds</c> /
- /// <c>TriggerCountRequired</c> / <c>RecoverCountRequired</c> / <c>ParamsJson</c> /
- /// <c>FactoryId</c> / 全部运行态列。把租户调过的参数刷回默认值是本服务最大的风险,
- /// 这条边界必须在类型层面而不是靠自觉守住。</para>
- ///
- /// <para><c>ParamsJson</c> 尤其不能重写:既有行可能仍是旧的 A+B 混装 JSON,
- /// 而 Batch 1 已保证 A 类字段在运行期被忽略;重写整块反而可能误伤租户的 B 类参数。
- /// 同理,<c>ParamsJson = NULL</c>(G1 造成的历史状态)也不自动补写 ——
- /// 定义在代码里,NULL 时运行参数会正常回落默认值。</para>
- /// </summary>
- internal static bool ApplyProjection(AdoS8WatchRule row, S8RuleDefinition definition)
- {
- var changed = false;
- void SetString(Func<AdoS8WatchRule, string?> get, Action<string?> 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;
- }
- /// <summary>(TenantId, RuleCode) 身份比较器。RuleCode 按不敏感比较,与 Catalog 查找一致。</summary>
- 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));
- }
- }
- /// <summary>供给方案(纯数据,不含任何 DB 访问)。</summary>
- internal sealed class S8RuleProvisioningPlan
- {
- /// <summary>待新建的运行策略行。</summary>
- public List<AdoS8WatchRule> Inserts { get; } = new();
- /// <summary>待刷新投影的既有行(投影列已就地更新)。</summary>
- public List<AdoS8WatchRule> Updates { get; } = new();
- /// <summary>投影已一致、无需写入的数量。</summary>
- public int UnchangedCount { get; set; }
- /// <summary>库中存在但当前代码版本没有定义的规则行数。</summary>
- public int OrphanedCount { get; set; }
- /// <summary>同 (tenant, rule_code) 的多余行数(唯一索引建立前的历史遗留)。</summary>
- public int DuplicateCount { get; set; }
- }
- /// <summary>
- /// 供给结果。刻意给出分项计数而不是只回 200 ——
- /// 调用方必须能知道「这次到底做了什么」,否则运维无法判断新规则是否已经铺开。
- /// </summary>
- public sealed class S8RuleProvisioningResult
- {
- /// <summary>本次处理的租户数。</summary>
- public int TenantCount { get; set; }
- /// <summary>当前代码版本中的规则定义数。</summary>
- public int DefinitionCount { get; set; }
- /// <summary>新建的运行策略行数。</summary>
- public int CreatedCount { get; set; }
- /// <summary>投影被纠正的行数。</summary>
- public int RefreshedCount { get; set; }
- /// <summary>投影已一致、未写入的行数。</summary>
- public int UnchangedCount { get; set; }
- /// <summary>无代码定义的孤儿行数(保持原样)。</summary>
- public int OrphanedCount { get; set; }
- /// <summary>同 (tenant, rule_code) 的多余行数(未删除,仅上报)。</summary>
- public int DuplicateCount { get; set; }
- }
|