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; }
}