S8WatchRuleService.cs 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367
  1. using System.Text.Json;
  2. using Admin.NET.Plugin.AiDOP.Entity.S8;
  3. using Admin.NET.Plugin.AiDOP.Infrastructure;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure.S8;
  5. using Admin.NET.Plugin.AiDOP.Service.S8.Rules;
  6. using Admin.NET.Plugin.AiDOP.Service.S8.Rules.DataAccess;
  7. namespace Admin.NET.Plugin.AiDOP.Service.S8;
  8. public class S8WatchRuleService : ITransient
  9. {
  10. private readonly SqlSugarRepository<AdoS8WatchRule> _rep;
  11. private readonly SqlSugarRepository<AdoS8SceneConfig> _sceneRep;
  12. private readonly IS8DatasetCatalog _datasetCatalog;
  13. private readonly S8DatasetEnableGate _datasetEnableGate;
  14. public S8WatchRuleService(
  15. SqlSugarRepository<AdoS8WatchRule> rep,
  16. SqlSugarRepository<AdoS8SceneConfig> sceneRep,
  17. IS8DatasetCatalog datasetCatalog,
  18. S8DatasetEnableGate datasetEnableGate)
  19. {
  20. _rep = rep;
  21. _sceneRep = sceneRep;
  22. _datasetCatalog = datasetCatalog;
  23. _datasetEnableGate = datasetEnableGate;
  24. }
  25. public async Task<List<AdoS8WatchRule>> ListAsync(long tenantId, long factoryId) =>
  26. await _rep.AsQueryable()
  27. .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId)
  28. .ToListAsync();
  29. // S8-TENANT-FACTORY-P0-CLOSURE-1:归属一律由服务端可信作用域盖章,忽略 body.TenantId / body.FactoryId。
  30. //
  31. // S8-STANDARD-DATASET-HARD-CUTOVER-1:去掉 origin 形参与 S8RuleCreationPolicy。
  32. // 那套「按创建来源分治」的治理,存在意义是区分"调用方自带 SQL"与"服务端生成 SQL";
  33. // 两者现在都不存在了 —— 规则唯一的取数配置是 dataset_code,由 ValidateAsync 强制。
  34. public async Task<AdoS8WatchRule> CreateAsync(AdoS8WatchRule body, S8TrustedScope scope)
  35. {
  36. body.TenantId = scope.TenantId;
  37. body.FactoryId = scope.FactoryId;
  38. await ValidateAsync(body, scope);
  39. // S8-STEP6E:新建一律走 canonical 词表 + params schema,杜绝「存得下但运行时永远不生效」。
  40. ValidateVocabularyForCreate(body);
  41. if (!string.IsNullOrWhiteSpace(body.ParamsJson))
  42. ValidateParamsJsonByRuleType(body.RuleType, body.ParamsJson!.Trim());
  43. body.Id = 0;
  44. body.CreatedAt = DateTime.Now;
  45. // S8-STEP6E(与 CFG_DATASRC D-3 / CFG_ROLES 同源缺陷):回填自增主键。
  46. // 原 InsertAsync 只返回 bool,body.Id 保持 0,调用方随后 GET/PUT/DELETE 一律 404。
  47. body.Id = await _rep.AsInsertable(body).ExecuteReturnBigIdentityAsync();
  48. return body;
  49. }
  50. // ================================================================================
  51. // S8-STEP6E-CFG-WATCH-FIX-AND-SAFE-CERT-1:通用整实体 PUT 已退役。
  52. //
  53. // 原实现 `_rep.UpdateAsync(body)` 是**整列更新**,而 body 直接由客户端 JSON 绑定
  54. // (本实体即 DTO),服务端只重新盖章 5 个字段(Tenant/Factory/Id/CreatedAt/UpdatedAt)。
  55. // 实体与仓储层均无 UpdateIgnoreColumns / IsOnlyIgnoreUpdate 保护,因此调用方可写入
  56. // 全部 13 个 scheduler-owned 运行时列:
  57. // lock_token / locked_by / lock_until / running_started_at /
  58. // next_run_at / last_run_at / last_status / last_error / last_duration_ms / last_run_id /
  59. // consecutive_failure_count / paused_until / pause_reason
  60. //
  61. // 后果(按 S8WatchSchedulerService 的租约语义):
  62. // · 写 lock_token/lock_until → 窃取或作废活跃租约,令运行中实例的回写静默失败
  63. // · 清 lock → 第二实例重复拾取同一规则 → 重复建单
  64. // · lock_until 设远未来 → 该规则永不再被 PickReadyRulesAsync 选中(静默 DoS)
  65. // · paused_until 设远未来 → UI 仍显示「启用」但监控实际已停
  66. // · 写 last_status/last_run_id/last_error → 伪造调度审计轨迹
  67. // 且**无需恶意**:部分字段的 PUT body 会让这 13 列静默变 NULL(HTTP 200、无报错)。
  68. //
  69. // 退役而非改白名单,是因为该入口没有正式消费方(已穷举:前端 s8ConfigApi.watchRules
  70. // 无 update;e2e 只用 GET/POST/DELETE;服务端唯一引用是本 controller),
  71. // 而正式配置修改已有 UpdateParamsAsync / UpdateScheduleAsync 等窄入口。
  72. //
  73. // ⚠️ 抛出发生在**任何 DB 访问之前**:不做 LoadScopedAsync、不做重复性查询。
  74. // 这既保证零 DB 触碰,也保证任意 id(含越权 id)一律 410 而非 404
  75. // ——「这个能力没了」优先于「这条记录不属于你」,避免越权探测反推他租户数据是否存在。
  76. //
  77. // 保留方法签名是硬约束:S8TenantIsolationContractTests 用反射断言带 S8TrustedScope
  78. // 的写入口存在、且无 scope 的旧重载不存在。
  79. // ================================================================================
  80. public Task<AdoS8WatchRule> UpdateAsync(long id, AdoS8WatchRule body, S8TrustedScope scope) =>
  81. throw new S8WriteRetiredException(S8WriteRetiredException.WatchRuleUpdateMessage);
  82. // S8-TENANT-FACTORY-P0-CLOSURE-1:删除必须先按可信作用域绑行,禁止裸 DeleteByIdAsync(id)。
  83. public async Task DeleteAsync(long id, S8TrustedScope scope)
  84. {
  85. var e = await LoadScopedAsync(id, scope);
  86. await _rep.DeleteByIdAsync(e.Id);
  87. }
  88. /// <summary>按 Id + 可信作用域取行;不在作用域内一律按「不存在」处理,不泄露他租户资源是否存在。</summary>
  89. private async Task<AdoS8WatchRule> LoadScopedAsync(long id, S8TrustedScope scope) =>
  90. await _rep.AsQueryable()
  91. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  92. .FirstAsync() ?? throw new S8NotFoundException();
  93. /// <summary>
  94. /// R4 安全更新:只更新 params_json 与 enabled。expression / rule_code / data_source_id /
  95. /// scene_code / watch_object_type / rule_type / source_object_type 一律不通过此路径修改。
  96. /// 当 RuleType 非空时,按对应 evaluator 的 Params.Parse 进行 schema 校验,解析失败抛 S8BizException。
  97. /// </summary>
  98. public async Task<AdoS8WatchRule> UpdateParamsAsync(long id, S8WatchRuleParamsPayload payload, S8TrustedScope scope)
  99. {
  100. var entity = await LoadScopedAsync(id, scope);
  101. var paramsJson = payload.ParamsJson?.Trim();
  102. if (!string.IsNullOrEmpty(paramsJson))
  103. {
  104. ValidateParamsJsonByRuleType(entity.RuleType, paramsJson);
  105. }
  106. entity.ParamsJson = string.IsNullOrEmpty(paramsJson) ? null : paramsJson;
  107. // S8-DATASET-FOUNDATION-HARDENING-2:本方法是规则 enabled 的唯一切换入口,
  108. // 因此 Enable Gate 挂在此处。由 disabled → enabled 时必须满足完整运行条件;
  109. // 关闭规则不受 Gate 约束(否则数据集出问题后规则将无法被关停)。
  110. if (payload.Enabled && !entity.Enabled)
  111. S8WatchRuleDataAccessValidator.ValidateForEnable(entity, _datasetEnableGate, scope.TenantId, scope.FactoryId);
  112. entity.Enabled = payload.Enabled;
  113. // TASK-002-RESET-DIMENSION-MODEL-DEV-2B:维度归属 + 报警机制按 payload 原样落库(含 null 清空)。
  114. entity.StageCode = NormalizeOrNull(payload.StageCode);
  115. entity.OrderFlowCode = NormalizeOrNull(payload.OrderFlowCode);
  116. entity.RuleMechanism = NormalizeOrNull(payload.RuleMechanism);
  117. entity.UpdatedAt = DateTime.Now;
  118. await _rep.UpdateAsync(entity);
  119. return entity;
  120. }
  121. private static string? NormalizeOrNull(string? value)
  122. {
  123. if (string.IsNullOrWhiteSpace(value)) return null;
  124. var trimmed = value.Trim();
  125. return trimmed.Length == 0 ? null : trimmed;
  126. }
  127. /// <summary>
  128. /// S8-STEP6E:WATCH 侧 canonical 词表,**严格取自真实 evaluator dispatch**
  129. /// (S8WatchSchedulerService.RunSingleRuleAsync 的 switch,ordinal 大小写敏感),
  130. /// 不是从 NOTIFY 词表套用过来的。
  131. /// 注意 rule_type 为空/空白是**合法的历史态**(调度器按 rule_type_empty_skipped 跳过、
  132. /// 不算失败),但**不允许新建**——新建一条永远不会被任何 evaluator 承载的规则没有产品意义。
  133. /// </summary>
  134. private static readonly string[] CanonicalRuleTypes =
  135. {
  136. S8TimeoutRuleEvaluator.RuleTypeCode,
  137. S8ShortageRuleEvaluator.RuleTypeCode,
  138. S8OutOfRangeRuleEvaluator.RuleTypeCode,
  139. };
  140. /// <summary>
  141. /// S8-STEP6E:新建时的 canonical 词表校验(rule_type / scene / severity)。
  142. ///
  143. /// 只作用于 Create:
  144. /// · UpdateParamsAsync / UpdateScheduleAsync / Pause / Resume 都不修改这三个字段,无需重复校验;
  145. /// · TestAsync 是对**既有行**的探针,若在此加严会让 legacy 无效行连自检都跑不了。
  146. /// 即「legacy 无效数据允许读取与停用,但禁止继续创建」。
  147. /// </summary>
  148. private static void ValidateVocabularyForCreate(AdoS8WatchRule body)
  149. {
  150. // rule_type:必须是真实存在 evaluator 的三类之一。
  151. if (string.IsNullOrWhiteSpace(body.RuleType)
  152. || Array.IndexOf(CanonicalRuleTypes, body.RuleType) < 0)
  153. throw new S8BizException(
  154. "不支持的规则类型:" + (string.IsNullOrWhiteSpace(body.RuleType) ? "(空)" : body.RuleType)
  155. + ";当前仅支持 " + string.Join(" / ", CanonicalRuleTypes));
  156. // scene:S8SceneCode 本身没有 IsValid,canonical 单模块场景判定复用 S8ModuleCode.IsValid
  157. //(仓内唯一「严格 S1–S7、拒 legacy 复合场景」的现成实现,不另造第二套)。
  158. if (!S8ModuleCode.IsValid(body.SceneCode))
  159. throw new S8BizException(
  160. "不支持的场景编码:" + body.SceneCode + ";当前仅支持 " + string.Join(" / ", S8ModuleCode.All));
  161. // severity:必须在 S8SeverityCode.Normalize **之前**校验。
  162. // Normalize 的兜底分支是 `_ => Follow`,放在之后会把 HIGH / 拼写错误静默降级成 FOLLOW,
  163. // 门禁永远命中不了——DB 中 severity='HIGH' 的那一行正是这样进来的。
  164. // 也刻意不复用 S8SeverityCode.IsValid:那是宽松六值版(含 LOW/MEDIUM/HIGH/CRITICAL),
  165. // 供 legacy 查询参数兼容用,拿来当写入门禁会直接放行 legacy 值。
  166. if (!string.Equals(body.Severity, S8SeverityCode.Follow, StringComparison.Ordinal)
  167. && !string.Equals(body.Severity, S8SeverityCode.Serious, StringComparison.Ordinal))
  168. throw new S8BizException(
  169. "不支持的严重度:" + body.Severity + ";当前仅支持 "
  170. + S8SeverityCode.Follow + " / " + S8SeverityCode.Serious);
  171. }
  172. private static void ValidateParamsJsonByRuleType(string? ruleType, string paramsJson)
  173. {
  174. try
  175. {
  176. switch (ruleType)
  177. {
  178. case S8TimeoutRuleEvaluator.RuleTypeCode:
  179. {
  180. var p = S8TimeoutParams.Parse(paramsJson);
  181. if (string.IsNullOrWhiteSpace(p.DueAtField)
  182. || string.IsNullOrWhiteSpace(p.StatusField)
  183. || string.IsNullOrWhiteSpace(p.ExceptionTypeCode))
  184. throw new S8BizException("TIMEOUT params 缺少必填字段:dueAtField / statusField / exceptionTypeCode");
  185. break;
  186. }
  187. case S8ShortageRuleEvaluator.RuleTypeCode:
  188. {
  189. var p = S8ShortageParams.Parse(paramsJson);
  190. if (string.IsNullOrWhiteSpace(p.TargetQtyField)
  191. || string.IsNullOrWhiteSpace(p.ActualQtyField)
  192. || string.IsNullOrWhiteSpace(p.ExceptionTypeCode))
  193. throw new S8BizException("SHORTAGE params 缺少必填字段:targetQtyField / actualQtyField / exceptionTypeCode");
  194. break;
  195. }
  196. case S8OutOfRangeRuleEvaluator.RuleTypeCode:
  197. {
  198. var p = S8OutOfRangeParams.Parse(paramsJson);
  199. if (string.IsNullOrWhiteSpace(p.MeasuredValueField))
  200. throw new S8BizException("OUT_OF_RANGE params 缺少必填字段:measuredValueField");
  201. if (p.LowerBound == null && p.UpperBound == null
  202. && string.IsNullOrWhiteSpace(p.LowerBoundField)
  203. && string.IsNullOrWhiteSpace(p.UpperBoundField))
  204. throw new S8BizException("OUT_OF_RANGE params 必须提供 upperBound / lowerBound 或对应行内字段之一");
  205. break;
  206. }
  207. default:
  208. // RuleType 为空或非三类已知值:仅做 JSON 合法性校验,避免阻塞历史数据。
  209. using (JsonDocument.Parse(paramsJson)) { }
  210. break;
  211. }
  212. }
  213. catch (JsonException ex)
  214. {
  215. throw new S8BizException($"params_json 不是合法 JSON:{ex.Message}");
  216. }
  217. }
  218. public async Task<object> TestAsync(long id, S8TrustedScope scope)
  219. {
  220. var entity = await LoadScopedAsync(id, scope);
  221. await ValidateAsync(entity, scope, id);
  222. return new { id, success = true, message = "规则基础校验通过", pollIntervalSeconds = entity.PollIntervalSeconds };
  223. }
  224. // S8-SCHED-FRONTEND-1:远未来手工暂停哨兵值(与 SqlSugar DateTime 兼容;前端按 paused_until > now 判定)。
  225. private static readonly DateTime ManualPausedSentinel = new(9999, 12, 31, 23, 59, 59);
  226. /// <summary>
  227. /// S8-SCHED-FRONTEND-1:调度参数安全更新。仅修改 poll_interval_seconds / trigger_count_required /
  228. /// recover_count_required;不动 params_json / expression / rule_type / scene_code / data_source_id。
  229. /// </summary>
  230. public async Task<AdoS8WatchRule> UpdateScheduleAsync(long id, S8WatchRuleSchedulePayload payload, S8TrustedScope scope)
  231. {
  232. var entity = await LoadScopedAsync(id, scope);
  233. if (payload.PollIntervalSeconds < 60 || payload.PollIntervalSeconds > 86400)
  234. throw new S8BizException("poll_interval_seconds 必须在 60–86400 之间");
  235. if (payload.TriggerCountRequired < 1 || payload.TriggerCountRequired > 10)
  236. throw new S8BizException("trigger_count_required 必须在 1–10 之间");
  237. if (payload.RecoverCountRequired < 1 || payload.RecoverCountRequired > 10)
  238. throw new S8BizException("recover_count_required 必须在 1–10 之间");
  239. await _rep.Context.Updateable<AdoS8WatchRule>()
  240. .SetColumns(x => new AdoS8WatchRule
  241. {
  242. PollIntervalSeconds = payload.PollIntervalSeconds,
  243. TriggerCountRequired = payload.TriggerCountRequired,
  244. RecoverCountRequired = payload.RecoverCountRequired,
  245. UpdatedAt = DateTime.Now
  246. })
  247. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  248. .ExecuteCommandAsync();
  249. return await LoadScopedAsync(id, scope);
  250. }
  251. /// <summary>
  252. /// S8-SCHED-FRONTEND-1:立即执行一次。把 next_run_at 置为 NOW,让下个 tick 拾取。
  253. /// 不直接同步执行 evaluator;不阻塞请求;返回 200 + 提示。
  254. /// </summary>
  255. public async Task<object> RunNowAsync(long id, S8TrustedScope scope)
  256. {
  257. var entity = await LoadScopedAsync(id, scope);
  258. if (!entity.Enabled)
  259. throw new S8BizException("规则未启用,不能立即执行");
  260. var now = DateTime.Now;
  261. if (entity.PausedUntil.HasValue && entity.PausedUntil.Value > now)
  262. throw new S8BizException("规则已暂停,请先恢复");
  263. if (entity.LockUntil.HasValue && entity.LockUntil.Value > now)
  264. throw new S8BizException("规则正在执行中,请稍后再试");
  265. await _rep.Context.Updateable<AdoS8WatchRule>()
  266. .SetColumns(x => new AdoS8WatchRule
  267. {
  268. NextRunAt = now,
  269. UpdatedAt = now
  270. })
  271. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  272. .ExecuteCommandAsync();
  273. return new { id, queued = true, message = "已排队,最长 1 分钟内执行" };
  274. }
  275. /// <summary>
  276. /// S8-SCHED-FRONTEND-1:手工暂停。paused_until = 9999-12-31 哨兵 + pause_reason=MANUAL_PAUSED。
  277. /// 不强杀正在执行的 lease;当前运行完成后下一轮自然不被拾取。
  278. /// 不清 last_status / last_error。
  279. /// </summary>
  280. public async Task<object> PauseAsync(long id, S8TrustedScope scope)
  281. {
  282. var entity = await LoadScopedAsync(id, scope);
  283. await _rep.Context.Updateable<AdoS8WatchRule>()
  284. .SetColumns(x => new AdoS8WatchRule
  285. {
  286. PausedUntil = ManualPausedSentinel,
  287. PauseReason = "MANUAL_PAUSED",
  288. UpdatedAt = DateTime.Now
  289. })
  290. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  291. .ExecuteCommandAsync();
  292. return new { id, paused = true, message = "已暂停" };
  293. }
  294. /// <summary>
  295. /// S8-SCHED-FRONTEND-1:恢复。清 paused_until / pause_reason / last_error;归零 consecutive_failure_count;
  296. /// next_run_at = NOW 让下个 tick 立即拾取。不改 enabled / params_json / rule_type。
  297. /// </summary>
  298. public async Task<object> ResumeAsync(long id, S8TrustedScope scope)
  299. {
  300. var entity = await LoadScopedAsync(id, scope);
  301. var now = DateTime.Now;
  302. await _rep.Context.Updateable<AdoS8WatchRule>()
  303. .SetColumns(x => new AdoS8WatchRule
  304. {
  305. PausedUntil = null,
  306. PauseReason = null,
  307. ConsecutiveFailureCount = 0,
  308. LastError = null,
  309. NextRunAt = now,
  310. UpdatedAt = now
  311. })
  312. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  313. .ExecuteCommandAsync();
  314. return new { id, resumed = true, message = "已恢复,并将在下一轮调度中执行" };
  315. }
  316. private async Task ValidateAsync(AdoS8WatchRule body, S8TrustedScope scope, long? id = null)
  317. {
  318. if (string.IsNullOrWhiteSpace(body.RuleCode) || string.IsNullOrWhiteSpace(body.SceneCode))
  319. throw new S8BizException("规则编码和场景编码必填");
  320. var exists = await _rep.AsQueryable()
  321. .AnyAsync(x => x.Id != (id ?? 0) && x.TenantId == body.TenantId && x.FactoryId == body.FactoryId && x.RuleCode == body.RuleCode);
  322. if (exists) throw new S8BizException("监视规则编码已存在");
  323. // 保存级校验(dataset_code 结构 + 目录归属,不要求 Provider 已上线)。
  324. S8WatchRuleDataAccessValidator.ValidateForSave(body, _datasetCatalog);
  325. // 只有真正要启用时才要求完整运行条件;草稿 / disabled 规则允许 Provider 未上线。
  326. if (body.Enabled)
  327. S8WatchRuleDataAccessValidator.ValidateForEnable(body, _datasetEnableGate, scope.TenantId, scope.FactoryId);
  328. var scene = await _sceneRep.GetFirstAsync(x => x.TenantId == body.TenantId && x.FactoryId == body.FactoryId && x.SceneCode == body.SceneCode)
  329. ?? throw new S8BizException("关联场景不存在");
  330. if (!scene.Enabled) throw new S8BizException("关联场景未启用");
  331. if (body.PollIntervalSeconds <= 0) throw new S8BizException("轮询间隔必须大于 0");
  332. }
  333. }