MdpRebuildContracts.cs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465
  1. using System.Text.RegularExpressions;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  3. public sealed record MdpRebuildScope(string ModuleCode, long TenantId, long FactoryId)
  4. {
  5. public static readonly HashSet<string> RebuildModules = new(StringComparer.OrdinalIgnoreCase)
  6. {
  7. "S1", "S2", "S3", "S4", "S5", "S6", "S7"
  8. };
  9. public static readonly HashSet<string> SnapshotModules = new(StringComparer.OrdinalIgnoreCase)
  10. {
  11. "S1", "S2", "S3", "S4", "S5", "S6", "S7", "S9"
  12. };
  13. public static MdpRebuildScope Create(string moduleCode, long tenantId, long factoryId)
  14. {
  15. var mc = NormalizeModule(moduleCode);
  16. if (!RebuildModules.Contains(mc))
  17. throw new InvalidOperationException($"模块 {mc} 不支持全量数据重算");
  18. if (tenantId <= 0)
  19. throw new InvalidOperationException($"{mc} 全量重算必须指定有效 tenantId,禁止 tenantId=0");
  20. if (factoryId <= 0)
  21. throw new InvalidOperationException($"{mc} 全量重算必须指定有效 factoryId");
  22. if (factoryId == tenantId)
  23. throw new InvalidOperationException($"{mc} 全量重算的 factoryId 不能等于 tenantId");
  24. return new MdpRebuildScope(mc, tenantId, factoryId);
  25. }
  26. public static string NormalizeModule(string? moduleCode)
  27. {
  28. var mc = (moduleCode ?? string.Empty).Trim().ToUpperInvariant();
  29. if (mc.Length == 0) throw new InvalidOperationException("moduleCode 不能为空");
  30. return mc;
  31. }
  32. public string LockName => $"{ModuleCode}_MDP_FULL:{TenantId}:{FactoryId}";
  33. public string ScopeKey => $"{ModuleCode}:{TenantId}:{FactoryId}";
  34. }
  35. public static class ModuleRebuildStatus
  36. {
  37. public const string Queued = "QUEUED";
  38. public const string Running = "RUNNING";
  39. public const string Success = "SUCCESS";
  40. public const string Failed = "FAILED";
  41. public const string Cancelled = "CANCELLED";
  42. }
  43. /// <summary>
  44. /// <c>ado_module_dashboard_rebuild_job.trigger_type</c> 取值。
  45. ///
  46. /// <para><b>新增自动类触发器时必须同步检查领取过滤</b>:
  47. /// <c>ModuleRebuildJobStore.ClaimNextQueuedAsync</c> 对非执行机采用**白名单**——
  48. /// 只领 <see cref="Manual"/> 或带 <c>requested_by</c> 的行。
  49. /// 因此新加的自动触发类型天然不会被非执行机领走,这是刻意设计,勿改成黑名单。</para>
  50. /// </summary>
  51. public static class ModuleRebuildTriggerType
  52. {
  53. /// <summary>页面「数据重算」按钮。非执行机也可领取。</summary>
  54. public const string Manual = "MANUAL";
  55. /// <summary>定时任务入队。只有执行机可领取。</summary>
  56. public const string Auto = "AUTO";
  57. /// <summary>首次建库/冷启动灌数。只有执行机可领取。</summary>
  58. public const string Bootstrap = "BOOTSTRAP";
  59. /// <summary>每日 03:40 夜间串行全量兜底。穿透 AUTO 冷却窗口,只有执行机可领取。</summary>
  60. public const string AutoNightly = "AUTO_NIGHTLY";
  61. /// <summary>
  62. /// 只有小时级 AUTO 走增量拉数;夜间兜底、启动引导、人工重算一律全量。
  63. /// 用于把 <c>triggerType</c> 翻译成 <c>MdpPullContext.FullRefresh</c>。
  64. /// </summary>
  65. public static bool ShouldPullFull(string triggerType)
  66. => !string.Equals(triggerType, Auto, StringComparison.OrdinalIgnoreCase);
  67. /// <summary>
  68. /// 非执行机能否领取该任务——即「这是不是一个人工触发的任务」。
  69. ///
  70. /// <para>这是该语义的**规范定义**。<c>ModuleRebuildJobStore.ClaimNextQueuedAsync</c> 里
  71. /// 的 <c>WhereIF</c> 条件必须与本函数等价(表达式树要翻译成 SQL,无法直接复用本函数)。
  72. /// 改动任一处都必须同步另一处,<c>ModuleRebuildClaimFilterTests</c> 会同时守住两边。</para>
  73. /// </summary>
  74. public static bool IsManualTrigger(string triggerType, long? requestedBy)
  75. => string.Equals(triggerType, Manual, StringComparison.OrdinalIgnoreCase) || requestedBy != null;
  76. }
  77. public static class ModuleRebuildStages
  78. {
  79. public const string Queued = "QUEUED";
  80. public const string AcquiringLock = "ACQUIRING_LOCK";
  81. public const string Preparing = "PREPARING";
  82. public const string Staging = "STAGING";
  83. public const string Standard = "STANDARD";
  84. public const string Dwd = "DWD";
  85. public const string Kpi = "KPI";
  86. public const string Atomic = "ATOMIC";
  87. public const string T8Inbound = "T8_INBOUND";
  88. public const string KpiPreparing = "KPI_PREPARING";
  89. public const string KpiCalculating = "KPI_CALCULATING";
  90. public const string Finalizing = "FINALIZING";
  91. public const string Success = "SUCCESS";
  92. public const string Failed = "FAILED";
  93. public const string Cancelled = "CANCELLED";
  94. public static int StageTotal(string moduleCode) => Normalize(moduleCode) switch
  95. {
  96. "S4" => 4,
  97. "S5" or "S6" or "S7" => 4,
  98. _ => 5
  99. };
  100. public static int ToStageIndex(string moduleCode, string stage)
  101. {
  102. var mc = Normalize(moduleCode);
  103. if (mc is "S5" or "S6" or "S7")
  104. {
  105. return stage switch
  106. {
  107. T8Inbound => 1,
  108. KpiPreparing => 2,
  109. KpiCalculating => 3,
  110. Atomic or Finalizing or Success => 4,
  111. _ => 0
  112. };
  113. }
  114. if (mc == "S4")
  115. {
  116. return stage switch
  117. {
  118. Staging => 1,
  119. Standard => 2,
  120. Dwd => 3,
  121. Kpi or Finalizing or Success => 4,
  122. _ => 0
  123. };
  124. }
  125. return stage switch
  126. {
  127. Staging => 1,
  128. Standard => 2,
  129. Dwd => 3,
  130. Kpi => 4,
  131. Atomic or Finalizing or Success => 5,
  132. _ => 0
  133. };
  134. }
  135. public static string ToChinese(string? stage) => stage switch
  136. {
  137. Queued => "排队中",
  138. AcquiringLock => "等待资源",
  139. Preparing => "准备运行环境",
  140. Staging => "拉取源数据",
  141. Standard => "标准化数据",
  142. Dwd => "生成 DWD 明细",
  143. Kpi => "重算 KPI",
  144. Atomic => "重建原子数据",
  145. T8Inbound => "T8 基表入站",
  146. KpiPreparing => "准备 KPI 计算",
  147. KpiCalculating => "计算 KPI",
  148. Finalizing => "收尾",
  149. Success => "已完成",
  150. Failed => "失败",
  151. Cancelled => "已取消",
  152. _ => stage ?? ""
  153. };
  154. public static string ConfirmTitle(string moduleCode) => $"确认重算 {Normalize(moduleCode)} 数据?";
  155. private static string Normalize(string? moduleCode) =>
  156. string.IsNullOrWhiteSpace(moduleCode) ? "" : moduleCode.Trim().ToUpperInvariant();
  157. }
  158. public sealed record ModuleProgressUpdate(
  159. string Stage,
  160. int StageIndex,
  161. int ProgressPercent,
  162. string Message,
  163. int? Rows = null,
  164. string? CompletedStage = null,
  165. string? DetailJson = null);
  166. public sealed class ModuleRebuildResult
  167. {
  168. public string BatchId { get; set; } = string.Empty;
  169. public long RunLogId { get; set; }
  170. public int StageRows { get; set; }
  171. public int StandardRows { get; set; }
  172. public int DwdRows { get; set; }
  173. public int KpiRows { get; set; }
  174. public int AtomicRows { get; set; }
  175. public string? DetailJson { get; set; }
  176. }
  177. /// <summary>管理员在阶段边界请求取消。由 <c>RunClaimedAsync</c> 捕获并收口为 CANCELLED。</summary>
  178. public sealed class ModuleRebuildCancelledException : Exception
  179. {
  180. public ModuleRebuildCancelledException() : base("已由管理员取消") { }
  181. /// <summary>
  182. /// 已开跑的那一轮 transform run log。进度回调拿不到它(它在各模块服务内部),
  183. /// 由 <c>RunClaimedAsync</c> 在异常冒泡时按 batchId 回填,收口时写成 ABORTED。
  184. /// </summary>
  185. public long RunLogId { get; set; }
  186. }
  187. public sealed class ModuleRebuildAlreadyRunningException : Exception
  188. {
  189. public ModuleRebuildAlreadyRunningException(string moduleCode)
  190. : base($"{moduleCode} 数据重算正在执行,请勿重复提交")
  191. {
  192. ModuleCode = moduleCode;
  193. }
  194. public string ModuleCode { get; }
  195. }
  196. public static class MdpSqlScope
  197. {
  198. private static readonly Regex FromAlias = new(
  199. @"\bFROM\s+(?<table>[`.\w]+)\s+(?:AS\s+)?(?<alias>[`\w]+)",
  200. RegexOptions.IgnoreCase | RegexOptions.Compiled);
  201. /// <summary>
  202. /// 不带租户列的全局字典表,注入器必须原样放过。
  203. /// <c>mdp_source_table_registry</c> 只有 (id, source_table, source_system, remark),
  204. /// 它由 <c>MdpSourceIdentity.Resolve</c> 拼成的标量子查询会出现在 S2/S3/S4 的每条 STD 语句里;
  205. /// 给它补上 <c>r.tenant_id=@TenantId</c> 就是 Unknown column,整轮跑批 FAILED。
  206. /// </summary>
  207. private static readonly HashSet<string> ScopelessTables = new(StringComparer.OrdinalIgnoreCase)
  208. {
  209. "mdp_source_table_registry"
  210. };
  211. /// <summary>扫到这些词就算当前 WHERE 子句结束。</summary>
  212. private static readonly string[] ClauseEnders =
  213. ["GROUP", "HAVING", "ORDER", "LIMIT", "UNION", "ON"];
  214. public static string InjectTenantFactory(string sql)
  215. => Inject(sql, includeFactory: true);
  216. public static string InjectTenantOnly(string sql)
  217. => Inject(sql, includeFactory: false);
  218. private static string Inject(string sql, bool includeFactory)
  219. {
  220. if (string.IsNullOrWhiteSpace(sql)) return sql;
  221. return Regex.Replace(sql, @"\bWHERE\b", m =>
  222. {
  223. if (ClauseHasTenantParam(sql, m.Index + m.Length)) return m.Value;
  224. var (prefix, scopeless) = ResolveSource(sql, m.Index);
  225. if (scopeless) return m.Value;
  226. var factoryFilter = includeFactory
  227. ? $"COALESCE(NULLIF({prefix}factory_id,0),1)=@FactoryId AND "
  228. : string.Empty;
  229. return $"WHERE {prefix}tenant_id=@TenantId AND {factoryFilter}";
  230. }, RegexOptions.IgnoreCase);
  231. }
  232. /// <summary>
  233. /// 本条 WHERE 子句自己是否已经写了租户条件。只认同层级的 <c>@TenantId</c>:
  234. /// 嵌套子查询自带的租户条件不能替外层背书,否则外层就漏过滤了。
  235. ///
  236. /// <para>这里原先是「往后数 240 个字符里有没有 @TenantId」。那个窗口既够不到写在后面的
  237. /// 条件(该跳过的没跳过),又会把邻近语句里不相干的 <c>@TenantId</c> 算进来(不该跳过的跳过了),
  238. /// 同一段 SQL 加不加租户条件取决于文本长度。<c>mdp_source_table_registry</c> 的子查询
  239. /// 在 S4 正好被窗口盖住而幸免,在 S3 就被注入成 <c>r.tenant_id=@TenantId</c> 并报 Unknown column。</para>
  240. /// </summary>
  241. private static bool ClauseHasTenantParam(string sql, int start)
  242. {
  243. var depth = 0;
  244. for (var i = start; i < sql.Length; i++)
  245. {
  246. var c = sql[i];
  247. if (c == '\'')
  248. {
  249. i = SkipQuoted(sql, i);
  250. continue;
  251. }
  252. if (c == '(')
  253. {
  254. depth++;
  255. continue;
  256. }
  257. if (c == ')')
  258. {
  259. if (depth == 0) return false; // 退出了所在子查询,子句到此为止
  260. depth--;
  261. continue;
  262. }
  263. if (depth > 0) continue;
  264. if (c == '@' && IsWord(sql, i, "@TenantId")) return true;
  265. foreach (var ender in ClauseEnders)
  266. if (IsWord(sql, i, ender)) return false;
  267. }
  268. return false;
  269. }
  270. /// <summary>从 <paramref name="openIndex"/> 处的单引号跳到配对的结束引号;支持 '' 与反斜杠转义。</summary>
  271. private static int SkipQuoted(string sql, int openIndex)
  272. {
  273. for (var i = openIndex + 1; i < sql.Length; i++)
  274. {
  275. if (sql[i] == '\\') { i++; continue; }
  276. if (sql[i] != '\'') continue;
  277. if (i + 1 < sql.Length && sql[i + 1] == '\'') { i++; continue; }
  278. return i;
  279. }
  280. return sql.Length;
  281. }
  282. /// <summary><paramref name="sql"/> 的 <paramref name="at"/> 处是否是一个完整的 <paramref name="word"/> 而非某个标识符的片段。</summary>
  283. private static bool IsWord(string sql, int at, string word)
  284. {
  285. if (!sql.AsSpan(at).StartsWith(word, StringComparison.OrdinalIgnoreCase)) return false;
  286. if (at > 0 && (char.IsLetterOrDigit(sql[at - 1]) || sql[at - 1] is '_' or '.' or '@' or '`')) return false;
  287. var after = at + word.Length;
  288. return after >= sql.Length || !(char.IsLetterOrDigit(sql[after]) || sql[after] == '_');
  289. }
  290. private static (string Prefix, bool Scopeless) ResolveSource(string sql, int whereIndex)
  291. {
  292. var depth = 0;
  293. for (var i = whereIndex - 1; i >= 0; i--)
  294. {
  295. if (sql[i] == ')')
  296. {
  297. depth++;
  298. continue;
  299. }
  300. if (sql[i] == '(')
  301. {
  302. if (depth > 0) depth--;
  303. continue;
  304. }
  305. if (depth != 0 || i < 3) continue;
  306. if (!sql.AsSpan(i - 3, 4).Equals("FROM", StringComparison.OrdinalIgnoreCase)) continue;
  307. if (i > 3 && (char.IsLetterOrDigit(sql[i - 4]) || sql[i - 4] == '_')) continue;
  308. var match = FromAlias.Match(sql, i - 3);
  309. if (!match.Success || match.Index != i - 3) return (string.Empty, false);
  310. var table = match.Groups["table"].Value.Trim('`');
  311. var scopeless = ScopelessTables.Contains(table)
  312. || table.StartsWith("information_schema.", StringComparison.OrdinalIgnoreCase);
  313. var alias = match.Groups["alias"].Value.Trim('`');
  314. var prefix = alias is "LEFT" or "RIGHT" or "INNER" or "CROSS" or "JOIN" or "ON" or "AS" or "WHERE"
  315. ? string.Empty
  316. : alias + ".";
  317. return (prefix, scopeless);
  318. }
  319. return (string.Empty, false);
  320. }
  321. }
  322. public interface IModuleRebuildCapability
  323. {
  324. bool IsEnabled(string moduleCode);
  325. /// <summary>单进程并发上限:本实例的 <c>ModuleRebuildWorker</c> 最多同时跑几个 scope。</summary>
  326. int MaxParallelScopes { get; }
  327. /// <summary>
  328. /// 全库并发上限:**所有**连同一个库的实例加起来最多同时跑几个 scope,
  329. /// 由 <c>ModuleRebuildJobStore.ClaimNextQueuedAsync</c> 在抢占时收口。
  330. /// <para>只有 <see cref="MaxParallelScopes"/> 时它约束的是单进程,多实例共库会成倍放大:
  331. /// 2026-09-24 取证里配置值是 1,而实测同一瞬间有 2–3 个重算在跑。</para>
  332. /// </summary>
  333. int GlobalMaxParallelScopes { get; }
  334. /// <summary>
  335. /// AUTO / BOOTSTRAP 入队的最小间隔小时数。MANUAL 与 AUTO_NIGHTLY 不受此限。
  336. /// 取 0 表示关闭冷却,仅供排障时临时恢复旧行为。
  337. /// </summary>
  338. int AutoMinIntervalHours { get; }
  339. }
  340. public sealed class ModuleRebuildCapability : IModuleRebuildCapability, ISingleton
  341. {
  342. public bool IsEnabled(string moduleCode)
  343. {
  344. var mc = MdpRebuildScope.NormalizeModule(moduleCode);
  345. try
  346. {
  347. return Furion.App.GetConfig<bool?>($"AiDOP:MdpRebuild:Modules:{mc}:Enabled", true) ?? false;
  348. }
  349. catch
  350. {
  351. return false;
  352. }
  353. }
  354. public int MaxParallelScopes
  355. {
  356. get
  357. {
  358. try
  359. {
  360. return Math.Clamp(Furion.App.GetConfig<int?>("AiDOP:MdpRebuild:MaxParallelScopes", true) ?? 2, 1, 8);
  361. }
  362. catch
  363. {
  364. return 2;
  365. }
  366. }
  367. }
  368. public int GlobalMaxParallelScopes
  369. {
  370. get
  371. {
  372. try
  373. {
  374. // 缺省回落到单进程上限:未显式配置时行为等价于「只有一个实例在跑」,
  375. // 而不是悄悄放开成无上限。
  376. var v = Furion.App.GetConfig<int?>("AiDOP:MdpRebuild:GlobalMaxParallelScopes", true);
  377. return Math.Clamp(v ?? MaxParallelScopes, 1, 16);
  378. }
  379. catch
  380. {
  381. return MaxParallelScopes;
  382. }
  383. }
  384. }
  385. public int AutoMinIntervalHours
  386. {
  387. get
  388. {
  389. try
  390. {
  391. return Math.Clamp(Furion.App.GetConfig<int?>("AiDOP:MdpRebuild:AutoMinIntervalHours", true) ?? 6, 0, 168);
  392. }
  393. catch
  394. {
  395. return 6;
  396. }
  397. }
  398. }
  399. }
  400. public sealed class AlwaysOnModuleRebuildCapability : IModuleRebuildCapability
  401. {
  402. public bool IsEnabled(string moduleCode) => true;
  403. public int MaxParallelScopes => 2;
  404. public int GlobalMaxParallelScopes => 2;
  405. public int AutoMinIntervalHours => 6;
  406. }