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; namespace Admin.NET.Plugin.AiDOP.Service.S8; public class S8WatchRuleService : ITransient { private readonly SqlSugarRepository _rep; private readonly SqlSugarRepository _dataSourceRep; private readonly SqlSugarRepository _sceneRep; public S8WatchRuleService( SqlSugarRepository rep, SqlSugarRepository dataSourceRep, SqlSugarRepository sceneRep) { _rep = rep; _dataSourceRep = dataSourceRep; _sceneRep = sceneRep; } public async Task> ListAsync(long tenantId, long factoryId) => await _rep.AsQueryable() .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId) .ToListAsync(); // S8-TENANT-FACTORY-P0-CLOSURE-1:归属一律由服务端可信作用域盖章,忽略 body.TenantId / body.FactoryId。 public async Task CreateAsync(AdoS8WatchRule body, S8TrustedScope scope) { body.TenantId = scope.TenantId; body.FactoryId = scope.FactoryId; await ValidateAsync(body, scope); // S8-STEP6E:新建一律走 canonical 词表 + params schema,杜绝「存得下但运行时永远不生效」。 ValidateVocabularyForCreate(body); if (!string.IsNullOrWhiteSpace(body.ParamsJson)) ValidateParamsJsonByRuleType(body.RuleType, body.ParamsJson!.Trim()); body.Id = 0; body.CreatedAt = DateTime.Now; // S8-STEP6E(与 CFG_DATASRC D-3 / CFG_ROLES 同源缺陷):回填自增主键。 // 原 InsertAsync 只返回 bool,body.Id 保持 0,调用方随后 GET/PUT/DELETE 一律 404。 body.Id = await _rep.AsInsertable(body).ExecuteReturnBigIdentityAsync(); return body; } // ================================================================================ // 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-TENANT-FACTORY-P0-CLOSURE-1:删除必须先按可信作用域绑行,禁止裸 DeleteByIdAsync(id)。 public async Task DeleteAsync(long id, S8TrustedScope scope) { var e = await LoadScopedAsync(id, scope); await _rep.DeleteByIdAsync(e.Id); } /// 按 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(); /// /// R4 安全更新:只更新 params_json 与 enabled。expression / rule_code / data_source_id / /// scene_code / watch_object_type / rule_type / source_object_type 一律不通过此路径修改。 /// 当 RuleType 非空时,按对应 evaluator 的 Params.Parse 进行 schema 校验,解析失败抛 S8BizException。 /// public async Task UpdateParamsAsync(long id, S8WatchRuleParamsPayload payload, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); var paramsJson = payload.ParamsJson?.Trim(); if (!string.IsNullOrEmpty(paramsJson)) { ValidateParamsJsonByRuleType(entity.RuleType, paramsJson); } entity.ParamsJson = string.IsNullOrEmpty(paramsJson) ? null : paramsJson; entity.Enabled = payload.Enabled; // TASK-002-RESET-DIMENSION-MODEL-DEV-2B:维度归属 + 报警机制按 payload 原样落库(含 null 清空)。 entity.StageCode = NormalizeOrNull(payload.StageCode); entity.OrderFlowCode = NormalizeOrNull(payload.OrderFlowCode); entity.RuleMechanism = NormalizeOrNull(payload.RuleMechanism); entity.UpdatedAt = DateTime.Now; await _rep.UpdateAsync(entity); return entity; } private static string? NormalizeOrNull(string? value) { if (string.IsNullOrWhiteSpace(value)) return null; var trimmed = value.Trim(); return trimmed.Length == 0 ? null : trimmed; } /// /// S8-STEP6E:WATCH 侧 canonical 词表,**严格取自真实 evaluator dispatch** /// (S8WatchSchedulerService.RunSingleRuleAsync 的 switch,ordinal 大小写敏感), /// 不是从 NOTIFY 词表套用过来的。 /// 注意 rule_type 为空/空白是**合法的历史态**(调度器按 rule_type_empty_skipped 跳过、 /// 不算失败),但**不允许新建**——新建一条永远不会被任何 evaluator 承载的规则没有产品意义。 /// private static readonly string[] CanonicalRuleTypes = { S8TimeoutRuleEvaluator.RuleTypeCode, S8ShortageRuleEvaluator.RuleTypeCode, S8OutOfRangeRuleEvaluator.RuleTypeCode, }; /// /// S8-STEP6E:新建时的 canonical 词表校验(rule_type / scene / severity)。 /// /// 只作用于 Create: /// · UpdateParamsAsync / UpdateScheduleAsync / Pause / Resume 都不修改这三个字段,无需重复校验; /// · TestAsync 是对**既有行**的探针,若在此加严会让 legacy 无效行连自检都跑不了。 /// 即「legacy 无效数据允许读取与停用,但禁止继续创建」。 /// private static void ValidateVocabularyForCreate(AdoS8WatchRule body) { // rule_type:必须是真实存在 evaluator 的三类之一。 if (string.IsNullOrWhiteSpace(body.RuleType) || Array.IndexOf(CanonicalRuleTypes, body.RuleType) < 0) throw new S8BizException( "不支持的规则类型:" + (string.IsNullOrWhiteSpace(body.RuleType) ? "(空)" : body.RuleType) + ";当前仅支持 " + string.Join(" / ", CanonicalRuleTypes)); // scene:S8SceneCode 本身没有 IsValid,canonical 单模块场景判定复用 S8ModuleCode.IsValid //(仓内唯一「严格 S1–S7、拒 legacy 复合场景」的现成实现,不另造第二套)。 if (!S8ModuleCode.IsValid(body.SceneCode)) throw new S8BizException( "不支持的场景编码:" + body.SceneCode + ";当前仅支持 " + string.Join(" / ", S8ModuleCode.All)); // severity:必须在 S8SeverityCode.Normalize **之前**校验。 // Normalize 的兜底分支是 `_ => Follow`,放在之后会把 HIGH / 拼写错误静默降级成 FOLLOW, // 门禁永远命中不了——DB 中 severity='HIGH' 的那一行正是这样进来的。 // 也刻意不复用 S8SeverityCode.IsValid:那是宽松六值版(含 LOW/MEDIUM/HIGH/CRITICAL), // 供 legacy 查询参数兼容用,拿来当写入门禁会直接放行 legacy 值。 if (!string.Equals(body.Severity, S8SeverityCode.Follow, StringComparison.Ordinal) && !string.Equals(body.Severity, S8SeverityCode.Serious, StringComparison.Ordinal)) throw new S8BizException( "不支持的严重度:" + body.Severity + ";当前仅支持 " + S8SeverityCode.Follow + " / " + S8SeverityCode.Serious); } private static void ValidateParamsJsonByRuleType(string? ruleType, string paramsJson) { try { switch (ruleType) { case S8TimeoutRuleEvaluator.RuleTypeCode: { var p = S8TimeoutParams.Parse(paramsJson); if (string.IsNullOrWhiteSpace(p.DueAtField) || string.IsNullOrWhiteSpace(p.StatusField) || string.IsNullOrWhiteSpace(p.ExceptionTypeCode)) throw new S8BizException("TIMEOUT params 缺少必填字段:dueAtField / statusField / exceptionTypeCode"); break; } case S8ShortageRuleEvaluator.RuleTypeCode: { var p = S8ShortageParams.Parse(paramsJson); if (string.IsNullOrWhiteSpace(p.TargetQtyField) || string.IsNullOrWhiteSpace(p.ActualQtyField) || string.IsNullOrWhiteSpace(p.ExceptionTypeCode)) throw new S8BizException("SHORTAGE params 缺少必填字段:targetQtyField / actualQtyField / exceptionTypeCode"); break; } case S8OutOfRangeRuleEvaluator.RuleTypeCode: { var p = S8OutOfRangeParams.Parse(paramsJson); if (string.IsNullOrWhiteSpace(p.MeasuredValueField)) throw new S8BizException("OUT_OF_RANGE params 缺少必填字段:measuredValueField"); if (p.LowerBound == null && p.UpperBound == null && string.IsNullOrWhiteSpace(p.LowerBoundField) && string.IsNullOrWhiteSpace(p.UpperBoundField)) throw new S8BizException("OUT_OF_RANGE params 必须提供 upperBound / lowerBound 或对应行内字段之一"); break; } default: // RuleType 为空或非三类已知值:仅做 JSON 合法性校验,避免阻塞历史数据。 using (JsonDocument.Parse(paramsJson)) { } break; } } catch (JsonException ex) { throw new S8BizException($"params_json 不是合法 JSON:{ex.Message}"); } } public async Task TestAsync(long id, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); await ValidateAsync(entity, scope, id); return new { id, success = true, message = "规则基础校验通过", pollIntervalSeconds = entity.PollIntervalSeconds }; } // 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 async Task UpdateScheduleAsync(long id, S8WatchRuleSchedulePayload payload, S8TrustedScope scope) { var entity = await LoadScopedAsync(id, scope); if (payload.PollIntervalSeconds < 60 || payload.PollIntervalSeconds > 86400) throw new S8BizException("poll_interval_seconds 必须在 60–86400 之间"); if (payload.TriggerCountRequired < 1 || payload.TriggerCountRequired > 10) throw new S8BizException("trigger_count_required 必须在 1–10 之间"); if (payload.RecoverCountRequired < 1 || payload.RecoverCountRequired > 10) throw new S8BizException("recover_count_required 必须在 1–10 之间"); await _rep.Context.Updateable() .SetColumns(x => new AdoS8WatchRule { PollIntervalSeconds = payload.PollIntervalSeconds, TriggerCountRequired = payload.TriggerCountRequired, RecoverCountRequired = payload.RecoverCountRequired, UpdatedAt = DateTime.Now }) .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) .ExecuteCommandAsync(); return await LoadScopedAsync(id, 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 = "已恢复,并将在下一轮调度中执行" }; } private async Task ValidateAsync(AdoS8WatchRule body, S8TrustedScope scope, long? id = null) { if (string.IsNullOrWhiteSpace(body.RuleCode) || string.IsNullOrWhiteSpace(body.SceneCode)) throw new S8BizException("规则编码和场景编码必填"); var exists = await _rep.AsQueryable() .AnyAsync(x => x.Id != (id ?? 0) && x.TenantId == body.TenantId && x.FactoryId == body.FactoryId && x.RuleCode == body.RuleCode); if (exists) throw new S8BizException("监视规则编码已存在"); // S8-TENANT-FACTORY-P0-CLOSURE-1:关联数据源必须同属可信作用域,禁止引用他租户数据源。 var dataSource = await _dataSourceRep.GetFirstAsync( x => x.Id == body.DataSourceId && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId) ?? throw new S8BizException("关联数据源不存在"); if (!dataSource.Enabled) throw new S8BizException("关联数据源未启用"); var scene = await _sceneRep.GetFirstAsync(x => x.TenantId == body.TenantId && x.FactoryId == body.FactoryId && x.SceneCode == body.SceneCode) ?? throw new S8BizException("关联场景不存在"); if (!scene.Enabled) throw new S8BizException("关联场景未启用"); if (body.PollIntervalSeconds <= 0) throw new S8BizException("轮询间隔必须大于 0"); } }