using System.Text.RegularExpressions; namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild; public sealed record MdpRebuildScope(string ModuleCode, long TenantId, long FactoryId) { public static readonly HashSet RebuildModules = new(StringComparer.OrdinalIgnoreCase) { "S1", "S2", "S3", "S4", "S5", "S6", "S7" }; public static readonly HashSet SnapshotModules = new(StringComparer.OrdinalIgnoreCase) { "S1", "S2", "S3", "S4", "S5", "S6", "S7", "S9" }; public static MdpRebuildScope Create(string moduleCode, long tenantId, long factoryId) { var mc = NormalizeModule(moduleCode); if (!RebuildModules.Contains(mc)) throw new InvalidOperationException($"模块 {mc} 不支持全量数据重算"); if (tenantId <= 0) throw new InvalidOperationException($"{mc} 全量重算必须指定有效 tenantId,禁止 tenantId=0"); if (factoryId <= 0) throw new InvalidOperationException($"{mc} 全量重算必须指定有效 factoryId"); if (factoryId == tenantId) throw new InvalidOperationException($"{mc} 全量重算的 factoryId 不能等于 tenantId"); return new MdpRebuildScope(mc, tenantId, factoryId); } public static string NormalizeModule(string? moduleCode) { var mc = (moduleCode ?? string.Empty).Trim().ToUpperInvariant(); if (mc.Length == 0) throw new InvalidOperationException("moduleCode 不能为空"); return mc; } public string LockName => $"{ModuleCode}_MDP_FULL:{TenantId}:{FactoryId}"; public string ScopeKey => $"{ModuleCode}:{TenantId}:{FactoryId}"; } public static class ModuleRebuildStatus { public const string Queued = "QUEUED"; public const string Running = "RUNNING"; public const string Success = "SUCCESS"; public const string Failed = "FAILED"; public const string Cancelled = "CANCELLED"; } /// /// ado_module_dashboard_rebuild_job.trigger_type 取值。 /// /// 新增自动类触发器时必须同步检查领取过滤: /// ModuleRebuildJobStore.ClaimNextQueuedAsync 对非执行机采用**白名单**—— /// 只领 或带 requested_by 的行。 /// 因此新加的自动触发类型天然不会被非执行机领走,这是刻意设计,勿改成黑名单。 /// public static class ModuleRebuildTriggerType { /// 页面「数据重算」按钮。非执行机也可领取。 public const string Manual = "MANUAL"; /// 定时任务入队。只有执行机可领取。 public const string Auto = "AUTO"; /// 首次建库/冷启动灌数。只有执行机可领取。 public const string Bootstrap = "BOOTSTRAP"; /// 每日 03:40 夜间串行全量兜底。穿透 AUTO 冷却窗口,只有执行机可领取。 public const string AutoNightly = "AUTO_NIGHTLY"; /// /// 只有小时级 AUTO 走增量拉数;夜间兜底、启动引导、人工重算一律全量。 /// 用于把 triggerType 翻译成 MdpPullContext.FullRefresh。 /// public static bool ShouldPullFull(string triggerType) => !string.Equals(triggerType, Auto, StringComparison.OrdinalIgnoreCase); /// /// 非执行机能否领取该任务——即「这是不是一个人工触发的任务」。 /// /// 这是该语义的**规范定义**。ModuleRebuildJobStore.ClaimNextQueuedAsync 里 /// 的 WhereIF 条件必须与本函数等价(表达式树要翻译成 SQL,无法直接复用本函数)。 /// 改动任一处都必须同步另一处,ModuleRebuildClaimFilterTests 会同时守住两边。 /// public static bool IsManualTrigger(string triggerType, long? requestedBy) => string.Equals(triggerType, Manual, StringComparison.OrdinalIgnoreCase) || requestedBy != null; } public static class ModuleRebuildStages { public const string Queued = "QUEUED"; public const string AcquiringLock = "ACQUIRING_LOCK"; public const string Preparing = "PREPARING"; public const string Staging = "STAGING"; public const string Standard = "STANDARD"; public const string Dwd = "DWD"; public const string Kpi = "KPI"; public const string Atomic = "ATOMIC"; public const string T8Inbound = "T8_INBOUND"; public const string KpiPreparing = "KPI_PREPARING"; public const string KpiCalculating = "KPI_CALCULATING"; public const string Finalizing = "FINALIZING"; public const string Success = "SUCCESS"; public const string Failed = "FAILED"; public const string Cancelled = "CANCELLED"; public static int StageTotal(string moduleCode) => Normalize(moduleCode) switch { "S4" => 4, "S5" or "S6" or "S7" => 4, _ => 5 }; public static int ToStageIndex(string moduleCode, string stage) { var mc = Normalize(moduleCode); if (mc is "S5" or "S6" or "S7") { return stage switch { T8Inbound => 1, KpiPreparing => 2, KpiCalculating => 3, Atomic or Finalizing or Success => 4, _ => 0 }; } if (mc == "S4") { return stage switch { Staging => 1, Standard => 2, Dwd => 3, Kpi or Finalizing or Success => 4, _ => 0 }; } return stage switch { Staging => 1, Standard => 2, Dwd => 3, Kpi => 4, Atomic or Finalizing or Success => 5, _ => 0 }; } public static string ToChinese(string? stage) => stage switch { Queued => "排队中", AcquiringLock => "等待资源", Preparing => "准备运行环境", Staging => "拉取源数据", Standard => "标准化数据", Dwd => "生成 DWD 明细", Kpi => "重算 KPI", Atomic => "重建原子数据", T8Inbound => "T8 基表入站", KpiPreparing => "准备 KPI 计算", KpiCalculating => "计算 KPI", Finalizing => "收尾", Success => "已完成", Failed => "失败", Cancelled => "已取消", _ => stage ?? "" }; public static string ConfirmTitle(string moduleCode) => $"确认重算 {Normalize(moduleCode)} 数据?"; private static string Normalize(string? moduleCode) => string.IsNullOrWhiteSpace(moduleCode) ? "" : moduleCode.Trim().ToUpperInvariant(); } public sealed record ModuleProgressUpdate( string Stage, int StageIndex, int ProgressPercent, string Message, int? Rows = null, string? CompletedStage = null, string? DetailJson = null); public sealed class ModuleRebuildResult { public string BatchId { get; set; } = string.Empty; public long RunLogId { get; set; } public int StageRows { get; set; } public int StandardRows { get; set; } public int DwdRows { get; set; } public int KpiRows { get; set; } public int AtomicRows { get; set; } public string? DetailJson { get; set; } } /// 管理员在阶段边界请求取消。由 RunClaimedAsync 捕获并收口为 CANCELLED。 public sealed class ModuleRebuildCancelledException : Exception { public ModuleRebuildCancelledException() : base("已由管理员取消") { } /// /// 已开跑的那一轮 transform run log。进度回调拿不到它(它在各模块服务内部), /// 由 RunClaimedAsync 在异常冒泡时按 batchId 回填,收口时写成 ABORTED。 /// public long RunLogId { get; set; } } public sealed class ModuleRebuildAlreadyRunningException : Exception { public ModuleRebuildAlreadyRunningException(string moduleCode) : base($"{moduleCode} 数据重算正在执行,请勿重复提交") { ModuleCode = moduleCode; } public string ModuleCode { get; } } public static class MdpSqlScope { private static readonly Regex FromAlias = new( @"\bFROM\s+(?[`.\w]+)\s+(?:AS\s+)?(?[`\w]+)", RegexOptions.IgnoreCase | RegexOptions.Compiled); /// /// 不带租户列的全局字典表,注入器必须原样放过。 /// mdp_source_table_registry 只有 (id, source_table, source_system, remark), /// 它由 MdpSourceIdentity.Resolve 拼成的标量子查询会出现在 S2/S3/S4 的每条 STD 语句里; /// 给它补上 r.tenant_id=@TenantId 就是 Unknown column,整轮跑批 FAILED。 /// private static readonly HashSet ScopelessTables = new(StringComparer.OrdinalIgnoreCase) { "mdp_source_table_registry" }; /// 扫到这些词就算当前 WHERE 子句结束。 private static readonly string[] ClauseEnders = ["GROUP", "HAVING", "ORDER", "LIMIT", "UNION", "ON"]; public static string InjectTenantFactory(string sql) => Inject(sql, includeFactory: true); public static string InjectTenantOnly(string sql) => Inject(sql, includeFactory: false); private static string Inject(string sql, bool includeFactory) { if (string.IsNullOrWhiteSpace(sql)) return sql; return Regex.Replace(sql, @"\bWHERE\b", m => { if (ClauseHasTenantParam(sql, m.Index + m.Length)) return m.Value; var (prefix, scopeless) = ResolveSource(sql, m.Index); if (scopeless) return m.Value; var factoryFilter = includeFactory ? $"COALESCE(NULLIF({prefix}factory_id,0),1)=@FactoryId AND " : string.Empty; return $"WHERE {prefix}tenant_id=@TenantId AND {factoryFilter}"; }, RegexOptions.IgnoreCase); } /// /// 本条 WHERE 子句自己是否已经写了租户条件。只认同层级的 @TenantId: /// 嵌套子查询自带的租户条件不能替外层背书,否则外层就漏过滤了。 /// /// 这里原先是「往后数 240 个字符里有没有 @TenantId」。那个窗口既够不到写在后面的 /// 条件(该跳过的没跳过),又会把邻近语句里不相干的 @TenantId 算进来(不该跳过的跳过了), /// 同一段 SQL 加不加租户条件取决于文本长度。mdp_source_table_registry 的子查询 /// 在 S4 正好被窗口盖住而幸免,在 S3 就被注入成 r.tenant_id=@TenantId 并报 Unknown column。 /// private static bool ClauseHasTenantParam(string sql, int start) { var depth = 0; for (var i = start; i < sql.Length; i++) { var c = sql[i]; if (c == '\'') { i = SkipQuoted(sql, i); continue; } if (c == '(') { depth++; continue; } if (c == ')') { if (depth == 0) return false; // 退出了所在子查询,子句到此为止 depth--; continue; } if (depth > 0) continue; if (c == '@' && IsWord(sql, i, "@TenantId")) return true; foreach (var ender in ClauseEnders) if (IsWord(sql, i, ender)) return false; } return false; } /// 从 处的单引号跳到配对的结束引号;支持 '' 与反斜杠转义。 private static int SkipQuoted(string sql, int openIndex) { for (var i = openIndex + 1; i < sql.Length; i++) { if (sql[i] == '\\') { i++; continue; } if (sql[i] != '\'') continue; if (i + 1 < sql.Length && sql[i + 1] == '\'') { i++; continue; } return i; } return sql.Length; } /// 的 处是否是一个完整的 而非某个标识符的片段。 private static bool IsWord(string sql, int at, string word) { if (!sql.AsSpan(at).StartsWith(word, StringComparison.OrdinalIgnoreCase)) return false; if (at > 0 && (char.IsLetterOrDigit(sql[at - 1]) || sql[at - 1] is '_' or '.' or '@' or '`')) return false; var after = at + word.Length; return after >= sql.Length || !(char.IsLetterOrDigit(sql[after]) || sql[after] == '_'); } private static (string Prefix, bool Scopeless) ResolveSource(string sql, int whereIndex) { var depth = 0; for (var i = whereIndex - 1; i >= 0; i--) { if (sql[i] == ')') { depth++; continue; } if (sql[i] == '(') { if (depth > 0) depth--; continue; } if (depth != 0 || i < 3) continue; if (!sql.AsSpan(i - 3, 4).Equals("FROM", StringComparison.OrdinalIgnoreCase)) continue; if (i > 3 && (char.IsLetterOrDigit(sql[i - 4]) || sql[i - 4] == '_')) continue; var match = FromAlias.Match(sql, i - 3); if (!match.Success || match.Index != i - 3) return (string.Empty, false); var table = match.Groups["table"].Value.Trim('`'); var scopeless = ScopelessTables.Contains(table) || table.StartsWith("information_schema.", StringComparison.OrdinalIgnoreCase); var alias = match.Groups["alias"].Value.Trim('`'); var prefix = alias is "LEFT" or "RIGHT" or "INNER" or "CROSS" or "JOIN" or "ON" or "AS" or "WHERE" ? string.Empty : alias + "."; return (prefix, scopeless); } return (string.Empty, false); } } public interface IModuleRebuildCapability { bool IsEnabled(string moduleCode); /// 单进程并发上限:本实例的 ModuleRebuildWorker 最多同时跑几个 scope。 int MaxParallelScopes { get; } /// /// 全库并发上限:**所有**连同一个库的实例加起来最多同时跑几个 scope, /// 由 ModuleRebuildJobStore.ClaimNextQueuedAsync 在抢占时收口。 /// 只有 时它约束的是单进程,多实例共库会成倍放大: /// 2026-09-24 取证里配置值是 1,而实测同一瞬间有 2–3 个重算在跑。 /// int GlobalMaxParallelScopes { get; } /// /// AUTO / BOOTSTRAP 入队的最小间隔小时数。MANUAL 与 AUTO_NIGHTLY 不受此限。 /// 取 0 表示关闭冷却,仅供排障时临时恢复旧行为。 /// int AutoMinIntervalHours { get; } } public sealed class ModuleRebuildCapability : IModuleRebuildCapability, ISingleton { public bool IsEnabled(string moduleCode) { var mc = MdpRebuildScope.NormalizeModule(moduleCode); try { return Furion.App.GetConfig($"AiDOP:MdpRebuild:Modules:{mc}:Enabled", true) ?? false; } catch { return false; } } public int MaxParallelScopes { get { try { return Math.Clamp(Furion.App.GetConfig("AiDOP:MdpRebuild:MaxParallelScopes", true) ?? 2, 1, 8); } catch { return 2; } } } public int GlobalMaxParallelScopes { get { try { // 缺省回落到单进程上限:未显式配置时行为等价于「只有一个实例在跑」, // 而不是悄悄放开成无上限。 var v = Furion.App.GetConfig("AiDOP:MdpRebuild:GlobalMaxParallelScopes", true); return Math.Clamp(v ?? MaxParallelScopes, 1, 16); } catch { return MaxParallelScopes; } } } public int AutoMinIntervalHours { get { try { return Math.Clamp(Furion.App.GetConfig("AiDOP:MdpRebuild:AutoMinIntervalHours", true) ?? 6, 0, 168); } catch { return 6; } } } } public sealed class AlwaysOnModuleRebuildCapability : IModuleRebuildCapability { public bool IsEnabled(string moduleCode) => true; public int MaxParallelScopes => 2; public int GlobalMaxParallelScopes => 2; public int AutoMinIntervalHours => 6; }