using System.Text.Json; using Admin.NET.Plugin.AiDOP.Entity.S8; using Admin.NET.Plugin.AiDOP.Infrastructure; using Admin.NET.Plugin.AiDOP.Infrastructure.S8; using Admin.NET.Plugin.AiDOP.Service.S8.Rules; using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess; using Admin.NET.Plugin.AiDOP.Service.S8.Rules.Definitions; namespace Admin.NET.Plugin.AiDOP.Service.S8; public class S8WatchRuleService : ITransient { private readonly SqlSugarRepository _rep; private readonly S8DatasetEnableGate _datasetEnableGate; private readonly IS8RuleCatalog _ruleCatalog; public S8WatchRuleService( SqlSugarRepository rep, S8DatasetEnableGate datasetEnableGate, IS8RuleCatalog ruleCatalog) { _rep = rep; _datasetEnableGate = datasetEnableGate; _ruleCatalog = ruleCatalog; } public async Task> ListAsync(long tenantId, long factoryId) => await _rep.AsQueryable() .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId) .ToListAsync(); // ================================================================================ // S8-RULE-GOVERNANCE-BATCH3:业务侧建规则能力已退役。 // // 原实现是「客户端传一个 AdoS8WatchRule → 服务端盖章归属 → 校验词表 → 落库」, // 即调用方可以凭空定义一条规则的 rule_code / dataset_code / rule_type / // source_object_type / scene_code,再配一段 params_json 决定它判什么。 // 那正是 Batch 1 要消灭的形态:规则的业务语义必须来自代码定义,不能来自请求体。 // // 替代路径(Batch 2 已建立,且**不经过本方法**): // S8RuleCatalog(代码定义) // → S8RuleProvisioningService(系统内部写路径) // → 每租户一条 runtime policy,默认停用 // // ⚠️ 抛出发生在**任何 DB 访问之前**:不查重、不读场景、不碰仓储。 // 这既保证零写入,也让「能力没了」这件事与入参是否合法无关。 // // 保留方法名而非整个删掉:S8RuleCreationRetiredTests 用反射断言 // 「该名字存在但恒抛退役异常」,比断言「名字不存在」更能防住有人换个名字复活它。 // ================================================================================ public Task CreateAsync(AdoS8WatchRule body, S8TrustedScope scope) => throw new S8WriteRetiredException(S8WriteRetiredException.WatchRuleCreateMessage); // ================================================================================ // S8-STEP6E-CFG-WATCH-FIX-AND-SAFE-CERT-1:通用整实体 PUT 已退役。 // // 原实现 `_rep.UpdateAsync(body)` 是**整列更新**,而 body 直接由客户端 JSON 绑定 // (本实体即 DTO),服务端只重新盖章 5 个字段(Tenant/Factory/Id/CreatedAt/UpdatedAt)。 // 实体与仓储层均无 UpdateIgnoreColumns / IsOnlyIgnoreUpdate 保护,因此调用方可写入 // 全部 13 个 scheduler-owned 运行时列: // lock_token / locked_by / lock_until / running_started_at / // next_run_at / last_run_at / last_status / last_error / last_duration_ms / last_run_id / // consecutive_failure_count / paused_until / pause_reason // // 后果(按 S8WatchSchedulerService 的租约语义): // · 写 lock_token/lock_until → 窃取或作废活跃租约,令运行中实例的回写静默失败 // · 清 lock → 第二实例重复拾取同一规则 → 重复建单 // · lock_until 设远未来 → 该规则永不再被 PickReadyRulesAsync 选中(静默 DoS) // · paused_until 设远未来 → UI 仍显示「启用」但监控实际已停 // · 写 last_status/last_run_id/last_error → 伪造调度审计轨迹 // 且**无需恶意**:部分字段的 PUT body 会让这 13 列静默变 NULL(HTTP 200、无报错)。 // // 退役而非改白名单,是因为该入口没有正式消费方(已穷举:前端 s8ConfigApi.watchRules // 无 update;e2e 只用 GET/POST/DELETE;服务端唯一引用是本 controller), // 而正式配置修改已有 UpdateParamsAsync / UpdateScheduleAsync 等窄入口。 // // ⚠️ 抛出发生在**任何 DB 访问之前**:不做 LoadScopedAsync、不做重复性查询。 // 这既保证零 DB 触碰,也保证任意 id(含越权 id)一律 410 而非 404 // ——「这个能力没了」优先于「这条记录不属于你」,避免越权探测反推他租户数据是否存在。 // // 保留方法签名是硬约束:S8TenantIsolationContractTests 用反射断言带 S8TrustedScope // 的写入口存在、且无 scope 的旧重载不存在。 // ================================================================================ public Task UpdateAsync(long id, AdoS8WatchRule body, S8TrustedScope scope) => throw new S8WriteRetiredException(S8WriteRetiredException.WatchRuleUpdateMessage); // ================================================================================ // S8-RULE-GOVERNANCE-BATCH3:业务侧删规则能力已退役。 // // 规则定义的生命周期属于代码版本;运行策略行的删除同样不该由页面触发: // · 删掉会连带丢失该租户已调好的参数与运行历史(含 recovered/抗抖计数); // · 下一次供给对账又会把它按默认值建回来 —— 用户看到的是"删了又回来了"; // · 真正的诉求「不想让这条规则再跑」已经有正确表达:**停用**。 // // 规则撤销的唯一正当路径是代码版本移除 Definition,随后由治理流程处理 orphan 行 //(Provisioning 的 OrphanedCount 会持续上报,Scheduler 以 rule_definition_not_found 拒绝执行)。 // // ⚠️ 与 CreateAsync 同理:抛出在任何 DB 访问之前,任意 id(含越权 id)一律 410 而非 404。 // ================================================================================ public Task DeleteAsync(long id, S8TrustedScope scope) => throw new S8WriteRetiredException(S8WriteRetiredException.WatchRuleDeleteMessage); /// 按 Id + 可信作用域取行;不在作用域内一律按「不存在」处理,不泄露他租户资源是否存在。 private async Task LoadScopedAsync(long id, S8TrustedScope scope) => await _rep.AsQueryable() .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .FirstAsync() ?? throw new S8NotFoundException(); /// /// S8-RULE-GOVERNANCE-BATCH1:运行参数的**部分更新**。 /// /// 与被它取代的 UpdateParamsAsync 的三处根本差异: /// /// 不再接受 params_json 原文。判定语义(dueAtField / statusField / /// completedStates / objectIdField / exceptionTypeCode)现在只存在于代码定义里, /// 没有任何 API 能改到它们; /// 不再承担启停。enabled 走 / /// PATCH 而非整块覆盖。未提供的字段保持原值 —— 这是 G1 的直接修复: /// 旧实现里一次 {"enabled":false} 就会把 params_json 抹成 NULL。 /// /// /// 写入用 UpdateColumns 白名单,物理上无法触碰 Definition 投影列与 13 个调度运行态列。 /// public async Task UpdateParametersAsync(long id, S8RuleParametersPayload payload, S8TrustedScope scope) { // 前置拒绝:在任何 DB 访问之前判掉非法载荷,越权 id 也不会被用来探测记录是否存在。 // 同一道守卫在 ApplyParameters 内再做一次 —— 那里才是所有调用方的必经之处。 EnsureNoDefinitionFields(payload); var entity = await LoadScopedAsync(id, scope); var definition = _ruleCatalog.GetRequired(entity.RuleCode); var next = ApplyParameters(entity, payload, definition); await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { PollIntervalSeconds = next.PollIntervalSeconds, TriggerCountRequired = next.TriggerCountRequired, RecoverCountRequired = next.RecoverCountRequired, Severity = next.Severity, ParamsJson = next.ToParamsJson(), UpdatedAt = DateTime.Now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return await LoadScopedAsync(id, scope); } /// /// 纯函数:把 PATCH 载荷叠加到「当前生效参数」上,并按 Definition 声明的取值域校验。 /// /// 抽成 static 是为了让 G1 的核心不变量(未提供 = 不变)能在**不接数据库**的情况下被测试 /// 逐字段断言。之前那条缺陷之所以能活到生产,正是因为它藏在一个必须有仓储才能跑的方法里。 /// internal static S8RuleRuntimeParameters ApplyParameters( AdoS8WatchRule entity, S8RuleParametersPayload payload, S8RuleDefinition definition) { EnsureNoDefinitionFields(payload); var current = S8RuleRuntimeParameters.Resolve(entity, definition); var policy = definition.Parameters ?? new S8RuleParameterPolicy(); var poll = payload.PollIntervalSeconds ?? current.PollIntervalSeconds; var trigger = payload.TriggerCountRequired ?? current.TriggerCountRequired; var recover = payload.RecoverCountRequired ?? current.RecoverCountRequired; var grace = payload.GraceMinutes ?? current.GraceMinutes; var severity = payload.Severity?.Trim() ?? current.Severity; EnsureInRange("轮询间隔(秒)", poll, policy.PollIntervalSecondsMin, policy.PollIntervalSecondsMax); EnsureInRange("连续命中建单次数", trigger, policy.TriggerCountRequiredMin, policy.TriggerCountRequiredMax); EnsureInRange("连续未命中恢复次数", recover, policy.RecoverCountRequiredMin, policy.RecoverCountRequiredMax); EnsureInRange("宽限分钟", grace, policy.GraceMinutesMin, policy.GraceMinutesMax); if (!policy.AllowedSeverities.Contains(severity, StringComparer.Ordinal)) throw new S8BizException( $"不支持的严重度:{severity};当前规则仅支持 {string.Join(" / ", policy.AllowedSeverities)}"); var occurrenceDept = payload.DefaultOccurrenceDeptId ?? current.DefaultOccurrenceDeptId; var responsibleDept = payload.DefaultResponsibleDeptId ?? current.DefaultResponsibleDeptId; if (!policy.AllowsDepartmentDefaults && (occurrenceDept.HasValue || responsibleDept.HasValue)) throw new S8BizException("当前规则不支持配置部门兜底"); return new S8RuleRuntimeParameters { PollIntervalSeconds = poll, TriggerCountRequired = trigger, RecoverCountRequired = recover, Severity = severity, GraceMinutes = grace, DefaultOccurrenceDeptId = occurrenceDept, DefaultResponsibleDeptId = responsibleDept }; } /// /// 载荷里出现 Definition 字段(或已迁走的 enabled)一律**显式拒绝**。 /// /// 静默忽略等于告诉调用方"改成功了",而实际什么都没发生 —— 那比报错更危险, /// 因为调用方会据此认为规则已经按新口径运行。 /// /// 放在 内而不是只放在 API 层: /// 纯函数是所有写入路径的必经之处,守卫挂在这里才不会被下一个调用方绕过。 /// private static void EnsureNoDefinitionFields(S8RuleParametersPayload payload) { if (payload == null) throw new S8BizException("请求体不能为空"); var rejected = payload.RejectedDefinitionFields; if (rejected.Count > 0) throw new S8BizException( "以下字段由代码定义,不能通过参数接口修改:" + string.Join(" / ", rejected) + ";启停请使用 /enable 与 /disable"); } private static void EnsureInRange(string label, int value, int min, int max) { if (value < min || value > max) throw new S8BizException($"{label} 必须在 {min}–{max} 之间,当前值 {value}"); } /// /// 启用规则。只写 enabled 与调度触发时间,绝不触碰任何参数列或 Definition 投影列。 /// /// 幂等:已启用时直接返回,不产生写入 —— 重复调用不会重排下次执行时间, /// 也就不会被用来变相"插队"调度。 /// public async Task EnableAsync(long id, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); if (entity.Enabled) return entity; // ① 没有代码定义的规则不得启用。这是 Create API 仍然存在期间的安全过渡: // 业务即使造出一条任意 rule_code 的规则,也无法让调度器替它跑。 var definition = _ruleCatalog.GetRequired(entity.RuleCode); // ② 数据集侧完整运行条件。按 Definition 的 dataset_code / rule_type 判定, // 而不是 DB 上那两列 —— 后者是投影,可能被人为改过。 _datasetEnableGate.EnsureCanEnable( definition.DatasetCode, definition.RuleType, definition.RuleCode, scope.TenantId, scope.FactoryId); var now = DateTime.Now; await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { Enabled = true, NextRunAt = now, UpdatedAt = now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return await LoadScopedAsync(id, scope); } /// /// 停用规则。只写 enabled。 /// /// 幂等:已停用时直接返回。 /// 不动 lease / next_run_at:PickReadyRulesAsync 的候选谓词第一条就是 /// x.Enabled,停用后自然不会再被拾取;正在执行中的那一轮由 /// ResetExpiredLeasesAsync 按既有租约语义收尾。强行清租约反而会与正在跑的实例撕扯。 /// public async Task DisableAsync(long id, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); if (!entity.Enabled) return entity; // 停用**不过** Enable Gate:数据集出问题之后仍然必须能把规则关掉。 await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { Enabled = false, UpdatedAt = DateTime.Now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return await LoadScopedAsync(id, scope); } // S8-RULE-GOVERNANCE-BATCH3:以下三段随建规则能力一并删除 —— // · CanonicalRuleTypes / ValidateVocabularyForCreate:校验的是"调用方传来的 rule_type、 // scene_code、severity 是否在词表内",而这三项现在只可能来自代码定义, // 由 S8RuleCatalog 在构造期做更严格的校验(含 ordinal 精确匹配); // · ValidateParamsJsonShape:校验的是调用方提交的 params_json, // 而 params_json 已不再由任何接口整块写入; // · ValidateAsync:唯一调用方是已退役的 CreateAsync / TestAsync,且它要求 // ado_s8_scene_config 存在对应场景行 —— 对代码定义的规则那是一条无关的额外前提。 // 留着它们只会让人以为还存在"提交规则定义"这条路。 // S8-SCHED-FRONTEND-1:远未来手工暂停哨兵值(与 SqlSugar DateTime 兼容;前端按 paused_until > now 判定)。 private static readonly DateTime ManualPausedSentinel = new(9999, 12, 31, 23, 59, 59); /// /// S8-SCHED-FRONTEND-1:调度参数安全更新。仅修改 poll_interval_seconds / trigger_count_required / /// recover_count_required;不动 params_json / expression / rule_type / scene_code / data_source_id。 /// public Task UpdateScheduleAsync(long id, S8WatchRuleSchedulePayload payload, S8TrustedScope scope) => // S8-RULE-GOVERNANCE-BATCH1:三个调度字段已并入统一参数白名单,本入口只做形状转换后委托, // 不再各自维护一份取值域 —— 两处取值域一旦漂移,就会出现「A 接口存得下、B 接口存不下」。 UpdateParametersAsync(id, new S8RuleParametersPayload { PollIntervalSeconds = payload?.PollIntervalSeconds, TriggerCountRequired = payload?.TriggerCountRequired, RecoverCountRequired = payload?.RecoverCountRequired }, scope); /// /// S8-SCHED-FRONTEND-1:立即执行一次。把 next_run_at 置为 NOW,让下个 tick 拾取。 /// 不直接同步执行 evaluator;不阻塞请求;返回 200 + 提示。 /// public async Task RunNowAsync(long id, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); if (!entity.Enabled) throw new S8BizException("规则未启用,不能立即执行"); var now = DateTime.Now; if (entity.PausedUntil.HasValue && entity.PausedUntil.Value > now) throw new S8BizException("规则已暂停,请先恢复"); if (entity.LockUntil.HasValue && entity.LockUntil.Value > now) throw new S8BizException("规则正在执行中,请稍后再试"); await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { NextRunAt = now, UpdatedAt = now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return new { id, queued = true, message = "已排队,最长 1 分钟内执行" }; } /// /// S8-SCHED-FRONTEND-1:手工暂停。paused_until = 9999-12-31 哨兵 + pause_reason=MANUAL_PAUSED。 /// 不强杀正在执行的 lease;当前运行完成后下一轮自然不被拾取。 /// 不清 last_status / last_error。 /// public async Task PauseAsync(long id, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { PausedUntil = ManualPausedSentinel, PauseReason = "MANUAL_PAUSED", UpdatedAt = DateTime.Now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return new { id, paused = true, message = "已暂停" }; } /// /// S8-SCHED-FRONTEND-1:恢复。清 paused_until / pause_reason / last_error;归零 consecutive_failure_count; /// next_run_at = NOW 让下个 tick 立即拾取。不改 enabled / params_json / rule_type。 /// public async Task ResumeAsync(long id, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); var now = DateTime.Now; await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { PausedUntil = null, PauseReason = null, ConsecutiveFailureCount = 0, LastError = null, NextRunAt = now, UpdatedAt = now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return new { id, resumed = true, message = "已恢复,并将在下一轮调度中执行" }; } }