S8WatchRuleService.cs 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469
  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. using Admin.NET.Plugin.AiDOP.Service.S8.Rules.Definitions;
  8. namespace Admin.NET.Plugin.AiDOP.Service.S8;
  9. public class S8WatchRuleService : ITransient
  10. {
  11. private readonly SqlSugarRepository<AdoS8WatchRule> _rep;
  12. private readonly SqlSugarRepository<AdoS8SceneConfig> _sceneRep;
  13. private readonly IS8DatasetCatalog _datasetCatalog;
  14. private readonly S8DatasetEnableGate _datasetEnableGate;
  15. private readonly IS8RuleCatalog _ruleCatalog;
  16. public S8WatchRuleService(
  17. SqlSugarRepository<AdoS8WatchRule> rep,
  18. SqlSugarRepository<AdoS8SceneConfig> sceneRep,
  19. IS8DatasetCatalog datasetCatalog,
  20. S8DatasetEnableGate datasetEnableGate,
  21. IS8RuleCatalog ruleCatalog)
  22. {
  23. _rep = rep;
  24. _sceneRep = sceneRep;
  25. _datasetCatalog = datasetCatalog;
  26. _datasetEnableGate = datasetEnableGate;
  27. _ruleCatalog = ruleCatalog;
  28. }
  29. public async Task<List<AdoS8WatchRule>> ListAsync(long tenantId, long factoryId) =>
  30. await _rep.AsQueryable()
  31. .Where(x => x.TenantId == tenantId && x.FactoryId == factoryId)
  32. .ToListAsync();
  33. // S8-TENANT-FACTORY-P0-CLOSURE-1:归属一律由服务端可信作用域盖章,忽略 body.TenantId / body.FactoryId。
  34. //
  35. // S8-STANDARD-DATASET-HARD-CUTOVER-1:去掉 origin 形参与 S8RuleCreationPolicy。
  36. // 那套「按创建来源分治」的治理,存在意义是区分"调用方自带 SQL"与"服务端生成 SQL";
  37. // 两者现在都不存在了 —— 规则唯一的取数配置是 dataset_code,由 ValidateAsync 强制。
  38. public async Task<AdoS8WatchRule> CreateAsync(AdoS8WatchRule body, S8TrustedScope scope)
  39. {
  40. body.TenantId = scope.TenantId;
  41. body.FactoryId = scope.FactoryId;
  42. await ValidateAsync(body, scope);
  43. // S8-STEP6E:新建一律走 canonical 词表 + params schema,杜绝「存得下但运行时永远不生效」。
  44. ValidateVocabularyForCreate(body);
  45. if (!string.IsNullOrWhiteSpace(body.ParamsJson))
  46. ValidateParamsJsonShape(body.ParamsJson!.Trim());
  47. body.Id = 0;
  48. body.CreatedAt = DateTime.Now;
  49. // S8-STEP6E(与 CFG_DATASRC D-3 / CFG_ROLES 同源缺陷):回填自增主键。
  50. // 原 InsertAsync 只返回 bool,body.Id 保持 0,调用方随后 GET/PUT/DELETE 一律 404。
  51. body.Id = await _rep.AsInsertable(body).ExecuteReturnBigIdentityAsync();
  52. return body;
  53. }
  54. // ================================================================================
  55. // S8-STEP6E-CFG-WATCH-FIX-AND-SAFE-CERT-1:通用整实体 PUT 已退役。
  56. //
  57. // 原实现 `_rep.UpdateAsync(body)` 是**整列更新**,而 body 直接由客户端 JSON 绑定
  58. // (本实体即 DTO),服务端只重新盖章 5 个字段(Tenant/Factory/Id/CreatedAt/UpdatedAt)。
  59. // 实体与仓储层均无 UpdateIgnoreColumns / IsOnlyIgnoreUpdate 保护,因此调用方可写入
  60. // 全部 13 个 scheduler-owned 运行时列:
  61. // lock_token / locked_by / lock_until / running_started_at /
  62. // next_run_at / last_run_at / last_status / last_error / last_duration_ms / last_run_id /
  63. // consecutive_failure_count / paused_until / pause_reason
  64. //
  65. // 后果(按 S8WatchSchedulerService 的租约语义):
  66. // · 写 lock_token/lock_until → 窃取或作废活跃租约,令运行中实例的回写静默失败
  67. // · 清 lock → 第二实例重复拾取同一规则 → 重复建单
  68. // · lock_until 设远未来 → 该规则永不再被 PickReadyRulesAsync 选中(静默 DoS)
  69. // · paused_until 设远未来 → UI 仍显示「启用」但监控实际已停
  70. // · 写 last_status/last_run_id/last_error → 伪造调度审计轨迹
  71. // 且**无需恶意**:部分字段的 PUT body 会让这 13 列静默变 NULL(HTTP 200、无报错)。
  72. //
  73. // 退役而非改白名单,是因为该入口没有正式消费方(已穷举:前端 s8ConfigApi.watchRules
  74. // 无 update;e2e 只用 GET/POST/DELETE;服务端唯一引用是本 controller),
  75. // 而正式配置修改已有 UpdateParamsAsync / UpdateScheduleAsync 等窄入口。
  76. //
  77. // ⚠️ 抛出发生在**任何 DB 访问之前**:不做 LoadScopedAsync、不做重复性查询。
  78. // 这既保证零 DB 触碰,也保证任意 id(含越权 id)一律 410 而非 404
  79. // ——「这个能力没了」优先于「这条记录不属于你」,避免越权探测反推他租户数据是否存在。
  80. //
  81. // 保留方法签名是硬约束:S8TenantIsolationContractTests 用反射断言带 S8TrustedScope
  82. // 的写入口存在、且无 scope 的旧重载不存在。
  83. // ================================================================================
  84. public Task<AdoS8WatchRule> UpdateAsync(long id, AdoS8WatchRule body, S8TrustedScope scope) =>
  85. throw new S8WriteRetiredException(S8WriteRetiredException.WatchRuleUpdateMessage);
  86. // S8-TENANT-FACTORY-P0-CLOSURE-1:删除必须先按可信作用域绑行,禁止裸 DeleteByIdAsync(id)。
  87. public async Task DeleteAsync(long id, S8TrustedScope scope)
  88. {
  89. var e = await LoadScopedAsync(id, scope);
  90. await _rep.DeleteByIdAsync(e.Id);
  91. }
  92. /// <summary>按 Id + 可信作用域取行;不在作用域内一律按「不存在」处理,不泄露他租户资源是否存在。</summary>
  93. private async Task<AdoS8WatchRule> LoadScopedAsync(long id, S8TrustedScope scope) =>
  94. await _rep.AsQueryable()
  95. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  96. .FirstAsync() ?? throw new S8NotFoundException();
  97. /// <summary>
  98. /// S8-RULE-GOVERNANCE-BATCH1:运行参数的**部分更新**。
  99. ///
  100. /// <para>与被它取代的 <c>UpdateParamsAsync</c> 的三处根本差异:</para>
  101. /// <list type="number">
  102. /// <item><b>不再接受 params_json 原文</b>。判定语义(dueAtField / statusField /
  103. /// completedStates / objectIdField / exceptionTypeCode)现在只存在于代码定义里,
  104. /// 没有任何 API 能改到它们;</item>
  105. /// <item><b>不再承担启停</b>。enabled 走 <see cref="EnableAsync"/> / <see cref="DisableAsync"/>;</item>
  106. /// <item><b>PATCH 而非整块覆盖</b>。未提供的字段保持原值 —— 这是 G1 的直接修复:
  107. /// 旧实现里一次 <c>{"enabled":false}</c> 就会把 params_json 抹成 NULL。</item>
  108. /// </list>
  109. ///
  110. /// <para>写入用 <c>UpdateColumns</c> 白名单,物理上无法触碰 Definition 投影列与 13 个调度运行态列。</para>
  111. /// </summary>
  112. public async Task<AdoS8WatchRule> UpdateParametersAsync(long id, S8RuleParametersPayload payload, S8TrustedScope scope)
  113. {
  114. // 前置拒绝:在任何 DB 访问之前判掉非法载荷,越权 id 也不会被用来探测记录是否存在。
  115. // 同一道守卫在 ApplyParameters 内再做一次 —— 那里才是所有调用方的必经之处。
  116. EnsureNoDefinitionFields(payload);
  117. var entity = await LoadScopedAsync(id, scope);
  118. var definition = _ruleCatalog.GetRequired(entity.RuleCode);
  119. var next = ApplyParameters(entity, payload, definition);
  120. await _rep.Context.Updateable<AdoS8WatchRule>()
  121. .SetColumns(x => new AdoS8WatchRule
  122. {
  123. PollIntervalSeconds = next.PollIntervalSeconds,
  124. TriggerCountRequired = next.TriggerCountRequired,
  125. RecoverCountRequired = next.RecoverCountRequired,
  126. Severity = next.Severity,
  127. ParamsJson = next.ToParamsJson(),
  128. UpdatedAt = DateTime.Now
  129. })
  130. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  131. .ExecuteCommandAsync();
  132. return await LoadScopedAsync(id, scope);
  133. }
  134. /// <summary>
  135. /// 纯函数:把 PATCH 载荷叠加到「当前生效参数」上,并按 Definition 声明的取值域校验。
  136. ///
  137. /// <para>抽成 static 是为了让 G1 的核心不变量(未提供 = 不变)能在**不接数据库**的情况下被测试
  138. /// 逐字段断言。之前那条缺陷之所以能活到生产,正是因为它藏在一个必须有仓储才能跑的方法里。</para>
  139. /// </summary>
  140. internal static S8RuleRuntimeParameters ApplyParameters(
  141. AdoS8WatchRule entity, S8RuleParametersPayload payload, S8RuleDefinition definition)
  142. {
  143. EnsureNoDefinitionFields(payload);
  144. var current = S8RuleRuntimeParameters.Resolve(entity, definition);
  145. var policy = definition.Parameters ?? new S8RuleParameterPolicy();
  146. var poll = payload.PollIntervalSeconds ?? current.PollIntervalSeconds;
  147. var trigger = payload.TriggerCountRequired ?? current.TriggerCountRequired;
  148. var recover = payload.RecoverCountRequired ?? current.RecoverCountRequired;
  149. var grace = payload.GraceMinutes ?? current.GraceMinutes;
  150. var severity = payload.Severity?.Trim() ?? current.Severity;
  151. EnsureInRange("轮询间隔(秒)", poll, policy.PollIntervalSecondsMin, policy.PollIntervalSecondsMax);
  152. EnsureInRange("连续命中建单次数", trigger, policy.TriggerCountRequiredMin, policy.TriggerCountRequiredMax);
  153. EnsureInRange("连续未命中恢复次数", recover, policy.RecoverCountRequiredMin, policy.RecoverCountRequiredMax);
  154. EnsureInRange("宽限分钟", grace, policy.GraceMinutesMin, policy.GraceMinutesMax);
  155. if (!policy.AllowedSeverities.Contains(severity, StringComparer.Ordinal))
  156. throw new S8BizException(
  157. $"不支持的严重度:{severity};当前规则仅支持 {string.Join(" / ", policy.AllowedSeverities)}");
  158. var occurrenceDept = payload.DefaultOccurrenceDeptId ?? current.DefaultOccurrenceDeptId;
  159. var responsibleDept = payload.DefaultResponsibleDeptId ?? current.DefaultResponsibleDeptId;
  160. if (!policy.AllowsDepartmentDefaults && (occurrenceDept.HasValue || responsibleDept.HasValue))
  161. throw new S8BizException("当前规则不支持配置部门兜底");
  162. return new S8RuleRuntimeParameters
  163. {
  164. PollIntervalSeconds = poll,
  165. TriggerCountRequired = trigger,
  166. RecoverCountRequired = recover,
  167. Severity = severity,
  168. GraceMinutes = grace,
  169. DefaultOccurrenceDeptId = occurrenceDept,
  170. DefaultResponsibleDeptId = responsibleDept
  171. };
  172. }
  173. /// <summary>
  174. /// 载荷里出现 Definition 字段(或已迁走的 enabled)一律**显式拒绝**。
  175. ///
  176. /// <para>静默忽略等于告诉调用方"改成功了",而实际什么都没发生 —— 那比报错更危险,
  177. /// 因为调用方会据此认为规则已经按新口径运行。</para>
  178. ///
  179. /// <para>放在 <see cref="ApplyParameters"/> 内而不是只放在 API 层:
  180. /// 纯函数是所有写入路径的必经之处,守卫挂在这里才不会被下一个调用方绕过。</para>
  181. /// </summary>
  182. private static void EnsureNoDefinitionFields(S8RuleParametersPayload payload)
  183. {
  184. if (payload == null) throw new S8BizException("请求体不能为空");
  185. var rejected = payload.RejectedDefinitionFields;
  186. if (rejected.Count > 0)
  187. throw new S8BizException(
  188. "以下字段由代码定义,不能通过参数接口修改:" + string.Join(" / ", rejected)
  189. + ";启停请使用 /enable 与 /disable");
  190. }
  191. private static void EnsureInRange(string label, int value, int min, int max)
  192. {
  193. if (value < min || value > max)
  194. throw new S8BizException($"{label} 必须在 {min}–{max} 之间,当前值 {value}");
  195. }
  196. /// <summary>
  197. /// 启用规则。<b>只写 enabled 与调度触发时间,绝不触碰任何参数列或 Definition 投影列。</b>
  198. ///
  199. /// <para>幂等:已启用时直接返回,不产生写入 —— 重复调用不会重排下次执行时间,
  200. /// 也就不会被用来变相"插队"调度。</para>
  201. /// </summary>
  202. public async Task<AdoS8WatchRule> EnableAsync(long id, S8TrustedScope scope)
  203. {
  204. var entity = await LoadScopedAsync(id, scope);
  205. if (entity.Enabled) return entity;
  206. // ① 没有代码定义的规则不得启用。这是 Create API 仍然存在期间的安全过渡:
  207. // 业务即使造出一条任意 rule_code 的规则,也无法让调度器替它跑。
  208. var definition = _ruleCatalog.GetRequired(entity.RuleCode);
  209. // ② 数据集侧完整运行条件。按 Definition 的 dataset_code / rule_type 判定,
  210. // 而不是 DB 上那两列 —— 后者是投影,可能被人为改过。
  211. _datasetEnableGate.EnsureCanEnable(
  212. definition.DatasetCode, definition.RuleType, definition.RuleCode, scope.TenantId, scope.FactoryId);
  213. var now = DateTime.Now;
  214. await _rep.Context.Updateable<AdoS8WatchRule>()
  215. .SetColumns(x => new AdoS8WatchRule
  216. {
  217. Enabled = true,
  218. NextRunAt = now,
  219. UpdatedAt = now
  220. })
  221. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  222. .ExecuteCommandAsync();
  223. return await LoadScopedAsync(id, scope);
  224. }
  225. /// <summary>
  226. /// 停用规则。<b>只写 enabled。</b>
  227. ///
  228. /// <para>幂等:已停用时直接返回。</para>
  229. /// <para>不动 lease / next_run_at:<c>PickReadyRulesAsync</c> 的候选谓词第一条就是
  230. /// <c>x.Enabled</c>,停用后自然不会再被拾取;正在执行中的那一轮由
  231. /// <c>ResetExpiredLeasesAsync</c> 按既有租约语义收尾。强行清租约反而会与正在跑的实例撕扯。</para>
  232. /// </summary>
  233. public async Task<AdoS8WatchRule> DisableAsync(long id, S8TrustedScope scope)
  234. {
  235. var entity = await LoadScopedAsync(id, scope);
  236. if (!entity.Enabled) return entity;
  237. // 停用**不过** Enable Gate:数据集出问题之后仍然必须能把规则关掉。
  238. await _rep.Context.Updateable<AdoS8WatchRule>()
  239. .SetColumns(x => new AdoS8WatchRule
  240. {
  241. Enabled = false,
  242. UpdatedAt = DateTime.Now
  243. })
  244. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  245. .ExecuteCommandAsync();
  246. return await LoadScopedAsync(id, scope);
  247. }
  248. /// <summary>
  249. /// S8-STEP6E:WATCH 侧 canonical 词表,**严格取自真实 evaluator dispatch**
  250. /// (S8WatchSchedulerService.RunSingleRuleAsync 的 switch,ordinal 大小写敏感),
  251. /// 不是从 NOTIFY 词表套用过来的。
  252. /// 注意 rule_type 为空/空白是**合法的历史态**(调度器按 rule_type_empty_skipped 跳过、
  253. /// 不算失败),但**不允许新建**——新建一条永远不会被任何 evaluator 承载的规则没有产品意义。
  254. /// </summary>
  255. private static readonly string[] CanonicalRuleTypes =
  256. {
  257. S8TimeoutRuleEvaluator.RuleTypeCode,
  258. S8ShortageRuleEvaluator.RuleTypeCode,
  259. S8OutOfRangeRuleEvaluator.RuleTypeCode,
  260. };
  261. /// <summary>
  262. /// S8-STEP6E:新建时的 canonical 词表校验(rule_type / scene / severity)。
  263. ///
  264. /// 只作用于 Create:
  265. /// · UpdateParamsAsync / UpdateScheduleAsync / Pause / Resume 都不修改这三个字段,无需重复校验;
  266. /// · TestAsync 是对**既有行**的探针,若在此加严会让 legacy 无效行连自检都跑不了。
  267. /// 即「legacy 无效数据允许读取与停用,但禁止继续创建」。
  268. /// </summary>
  269. private static void ValidateVocabularyForCreate(AdoS8WatchRule body)
  270. {
  271. // rule_type:必须是真实存在 evaluator 的三类之一。
  272. if (string.IsNullOrWhiteSpace(body.RuleType)
  273. || Array.IndexOf(CanonicalRuleTypes, body.RuleType) < 0)
  274. throw new S8BizException(
  275. "不支持的规则类型:" + (string.IsNullOrWhiteSpace(body.RuleType) ? "(空)" : body.RuleType)
  276. + ";当前仅支持 " + string.Join(" / ", CanonicalRuleTypes));
  277. // scene:S8SceneCode 本身没有 IsValid,canonical 单模块场景判定复用 S8ModuleCode.IsValid
  278. //(仓内唯一「严格 S1–S7、拒 legacy 复合场景」的现成实现,不另造第二套)。
  279. if (!S8ModuleCode.IsValid(body.SceneCode))
  280. throw new S8BizException(
  281. "不支持的场景编码:" + body.SceneCode + ";当前仅支持 " + string.Join(" / ", S8ModuleCode.All));
  282. // severity:必须在 S8SeverityCode.Normalize **之前**校验。
  283. // Normalize 的兜底分支是 `_ => Follow`,放在之后会把 HIGH / 拼写错误静默降级成 FOLLOW,
  284. // 门禁永远命中不了——DB 中 severity='HIGH' 的那一行正是这样进来的。
  285. // 也刻意不复用 S8SeverityCode.IsValid:那是宽松六值版(含 LOW/MEDIUM/HIGH/CRITICAL),
  286. // 供 legacy 查询参数兼容用,拿来当写入门禁会直接放行 legacy 值。
  287. if (!string.Equals(body.Severity, S8SeverityCode.Follow, StringComparison.Ordinal)
  288. && !string.Equals(body.Severity, S8SeverityCode.Serious, StringComparison.Ordinal))
  289. throw new S8BizException(
  290. "不支持的严重度:" + body.Severity + ";当前仅支持 "
  291. + S8SeverityCode.Follow + " / " + S8SeverityCode.Serious);
  292. }
  293. /// <summary>
  294. /// S8-RULE-GOVERNANCE-BATCH1:params_json 的按类型 schema 校验已退役。
  295. ///
  296. /// <para>它校验的是 <c>dueAtField</c> / <c>statusField</c> / <c>exceptionTypeCode</c> 这些
  297. /// **A 类判定语义**是否齐备 —— 而这些字段现在只存在于代码定义里,params_json 根本不再承载它们。
  298. /// 继续校验等于要求调用方提交一份已经没有意义的 JSON。</para>
  299. ///
  300. /// <para>保留 JSON 合法性检查:脏文本存进去会让运行期解析失败,
  301. /// 虽然 <see cref="S8RuleRuntimeParameters.Resolve"/> 会回落默认值而不报错,
  302. /// 但让一份读不懂的文本静静躺在库里没有任何好处。</para>
  303. /// </summary>
  304. private static void ValidateParamsJsonShape(string paramsJson)
  305. {
  306. try { using (JsonDocument.Parse(paramsJson)) { } }
  307. catch (JsonException ex)
  308. {
  309. throw new S8BizException($"params_json 不是合法 JSON:{ex.Message}");
  310. }
  311. }
  312. public async Task<object> TestAsync(long id, S8TrustedScope scope)
  313. {
  314. var entity = await LoadScopedAsync(id, scope);
  315. await ValidateAsync(entity, scope, id);
  316. return new { id, success = true, message = "规则基础校验通过", pollIntervalSeconds = entity.PollIntervalSeconds };
  317. }
  318. // S8-SCHED-FRONTEND-1:远未来手工暂停哨兵值(与 SqlSugar DateTime 兼容;前端按 paused_until > now 判定)。
  319. private static readonly DateTime ManualPausedSentinel = new(9999, 12, 31, 23, 59, 59);
  320. /// <summary>
  321. /// S8-SCHED-FRONTEND-1:调度参数安全更新。仅修改 poll_interval_seconds / trigger_count_required /
  322. /// recover_count_required;不动 params_json / expression / rule_type / scene_code / data_source_id。
  323. /// </summary>
  324. public Task<AdoS8WatchRule> UpdateScheduleAsync(long id, S8WatchRuleSchedulePayload payload, S8TrustedScope scope) =>
  325. // S8-RULE-GOVERNANCE-BATCH1:三个调度字段已并入统一参数白名单,本入口只做形状转换后委托,
  326. // 不再各自维护一份取值域 —— 两处取值域一旦漂移,就会出现「A 接口存得下、B 接口存不下」。
  327. UpdateParametersAsync(id, new S8RuleParametersPayload
  328. {
  329. PollIntervalSeconds = payload?.PollIntervalSeconds,
  330. TriggerCountRequired = payload?.TriggerCountRequired,
  331. RecoverCountRequired = payload?.RecoverCountRequired
  332. }, scope);
  333. /// <summary>
  334. /// S8-SCHED-FRONTEND-1:立即执行一次。把 next_run_at 置为 NOW,让下个 tick 拾取。
  335. /// 不直接同步执行 evaluator;不阻塞请求;返回 200 + 提示。
  336. /// </summary>
  337. public async Task<object> RunNowAsync(long id, S8TrustedScope scope)
  338. {
  339. var entity = await LoadScopedAsync(id, scope);
  340. if (!entity.Enabled)
  341. throw new S8BizException("规则未启用,不能立即执行");
  342. var now = DateTime.Now;
  343. if (entity.PausedUntil.HasValue && entity.PausedUntil.Value > now)
  344. throw new S8BizException("规则已暂停,请先恢复");
  345. if (entity.LockUntil.HasValue && entity.LockUntil.Value > now)
  346. throw new S8BizException("规则正在执行中,请稍后再试");
  347. await _rep.Context.Updateable<AdoS8WatchRule>()
  348. .SetColumns(x => new AdoS8WatchRule
  349. {
  350. NextRunAt = now,
  351. UpdatedAt = now
  352. })
  353. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  354. .ExecuteCommandAsync();
  355. return new { id, queued = true, message = "已排队,最长 1 分钟内执行" };
  356. }
  357. /// <summary>
  358. /// S8-SCHED-FRONTEND-1:手工暂停。paused_until = 9999-12-31 哨兵 + pause_reason=MANUAL_PAUSED。
  359. /// 不强杀正在执行的 lease;当前运行完成后下一轮自然不被拾取。
  360. /// 不清 last_status / last_error。
  361. /// </summary>
  362. public async Task<object> PauseAsync(long id, S8TrustedScope scope)
  363. {
  364. var entity = await LoadScopedAsync(id, scope);
  365. await _rep.Context.Updateable<AdoS8WatchRule>()
  366. .SetColumns(x => new AdoS8WatchRule
  367. {
  368. PausedUntil = ManualPausedSentinel,
  369. PauseReason = "MANUAL_PAUSED",
  370. UpdatedAt = DateTime.Now
  371. })
  372. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  373. .ExecuteCommandAsync();
  374. return new { id, paused = true, message = "已暂停" };
  375. }
  376. /// <summary>
  377. /// S8-SCHED-FRONTEND-1:恢复。清 paused_until / pause_reason / last_error;归零 consecutive_failure_count;
  378. /// next_run_at = NOW 让下个 tick 立即拾取。不改 enabled / params_json / rule_type。
  379. /// </summary>
  380. public async Task<object> ResumeAsync(long id, S8TrustedScope scope)
  381. {
  382. var entity = await LoadScopedAsync(id, scope);
  383. var now = DateTime.Now;
  384. await _rep.Context.Updateable<AdoS8WatchRule>()
  385. .SetColumns(x => new AdoS8WatchRule
  386. {
  387. PausedUntil = null,
  388. PauseReason = null,
  389. ConsecutiveFailureCount = 0,
  390. LastError = null,
  391. NextRunAt = now,
  392. UpdatedAt = now
  393. })
  394. .Where(x => x.Id == id && x.TenantId == scope.TenantId && x.FactoryId == scope.FactoryId)
  395. .ExecuteCommandAsync();
  396. return new { id, resumed = true, message = "已恢复,并将在下一轮调度中执行" };
  397. }
  398. private async Task ValidateAsync(AdoS8WatchRule body, S8TrustedScope scope, long? id = null)
  399. {
  400. if (string.IsNullOrWhiteSpace(body.RuleCode) || string.IsNullOrWhiteSpace(body.SceneCode))
  401. throw new S8BizException("规则编码和场景编码必填");
  402. var exists = await _rep.AsQueryable()
  403. .AnyAsync(x => x.Id != (id ?? 0) && x.TenantId == body.TenantId && x.FactoryId == body.FactoryId && x.RuleCode == body.RuleCode);
  404. if (exists) throw new S8BizException("监视规则编码已存在");
  405. // 保存级校验(dataset_code 结构 + 目录归属,不要求 Provider 已上线)。
  406. S8WatchRuleDataAccessValidator.ValidateForSave(body, _datasetCatalog);
  407. // 只有真正要启用时才要求完整运行条件;草稿 / disabled 规则允许 Provider 未上线。
  408. // S8-RULE-GOVERNANCE-BATCH1:启用条件按**代码定义**判定,不按 body 上的投影列。
  409. if (body.Enabled)
  410. {
  411. var definition = _ruleCatalog.GetRequired(body.RuleCode);
  412. _datasetEnableGate.EnsureCanEnable(
  413. definition.DatasetCode, definition.RuleType, definition.RuleCode, scope.TenantId, scope.FactoryId);
  414. }
  415. var scene = await _sceneRep.GetFirstAsync(x => x.TenantId == body.TenantId && x.FactoryId == body.FactoryId && x.SceneCode == body.SceneCode)
  416. ?? throw new S8BizException("关联场景不存在");
  417. if (!scene.Enabled) throw new S8BizException("关联场景未启用");
  418. if (body.PollIntervalSeconds <= 0) throw new S8BizException("轮询间隔必须大于 0");
  419. }
  420. }