| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465 |
- 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<string> RebuildModules = new(StringComparer.OrdinalIgnoreCase)
- {
- "S1", "S2", "S3", "S4", "S5", "S6", "S7"
- };
- public static readonly HashSet<string> 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";
- }
- /// <summary>
- /// <c>ado_module_dashboard_rebuild_job.trigger_type</c> 取值。
- ///
- /// <para><b>新增自动类触发器时必须同步检查领取过滤</b>:
- /// <c>ModuleRebuildJobStore.ClaimNextQueuedAsync</c> 对非执行机采用**白名单**——
- /// 只领 <see cref="Manual"/> 或带 <c>requested_by</c> 的行。
- /// 因此新加的自动触发类型天然不会被非执行机领走,这是刻意设计,勿改成黑名单。</para>
- /// </summary>
- public static class ModuleRebuildTriggerType
- {
- /// <summary>页面「数据重算」按钮。非执行机也可领取。</summary>
- public const string Manual = "MANUAL";
- /// <summary>定时任务入队。只有执行机可领取。</summary>
- public const string Auto = "AUTO";
- /// <summary>首次建库/冷启动灌数。只有执行机可领取。</summary>
- public const string Bootstrap = "BOOTSTRAP";
- /// <summary>每日 03:40 夜间串行全量兜底。穿透 AUTO 冷却窗口,只有执行机可领取。</summary>
- public const string AutoNightly = "AUTO_NIGHTLY";
- /// <summary>
- /// 只有小时级 AUTO 走增量拉数;夜间兜底、启动引导、人工重算一律全量。
- /// 用于把 <c>triggerType</c> 翻译成 <c>MdpPullContext.FullRefresh</c>。
- /// </summary>
- public static bool ShouldPullFull(string triggerType)
- => !string.Equals(triggerType, Auto, StringComparison.OrdinalIgnoreCase);
- /// <summary>
- /// 非执行机能否领取该任务——即「这是不是一个人工触发的任务」。
- ///
- /// <para>这是该语义的**规范定义**。<c>ModuleRebuildJobStore.ClaimNextQueuedAsync</c> 里
- /// 的 <c>WhereIF</c> 条件必须与本函数等价(表达式树要翻译成 SQL,无法直接复用本函数)。
- /// 改动任一处都必须同步另一处,<c>ModuleRebuildClaimFilterTests</c> 会同时守住两边。</para>
- /// </summary>
- 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; }
- }
- /// <summary>管理员在阶段边界请求取消。由 <c>RunClaimedAsync</c> 捕获并收口为 CANCELLED。</summary>
- public sealed class ModuleRebuildCancelledException : Exception
- {
- public ModuleRebuildCancelledException() : base("已由管理员取消") { }
- /// <summary>
- /// 已开跑的那一轮 transform run log。进度回调拿不到它(它在各模块服务内部),
- /// 由 <c>RunClaimedAsync</c> 在异常冒泡时按 batchId 回填,收口时写成 ABORTED。
- /// </summary>
- 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+(?<table>[`.\w]+)\s+(?:AS\s+)?(?<alias>[`\w]+)",
- RegexOptions.IgnoreCase | RegexOptions.Compiled);
- /// <summary>
- /// 不带租户列的全局字典表,注入器必须原样放过。
- /// <c>mdp_source_table_registry</c> 只有 (id, source_table, source_system, remark),
- /// 它由 <c>MdpSourceIdentity.Resolve</c> 拼成的标量子查询会出现在 S2/S3/S4 的每条 STD 语句里;
- /// 给它补上 <c>r.tenant_id=@TenantId</c> 就是 Unknown column,整轮跑批 FAILED。
- /// </summary>
- private static readonly HashSet<string> ScopelessTables = new(StringComparer.OrdinalIgnoreCase)
- {
- "mdp_source_table_registry"
- };
- /// <summary>扫到这些词就算当前 WHERE 子句结束。</summary>
- 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);
- }
- /// <summary>
- /// 本条 WHERE 子句自己是否已经写了租户条件。只认同层级的 <c>@TenantId</c>:
- /// 嵌套子查询自带的租户条件不能替外层背书,否则外层就漏过滤了。
- ///
- /// <para>这里原先是「往后数 240 个字符里有没有 @TenantId」。那个窗口既够不到写在后面的
- /// 条件(该跳过的没跳过),又会把邻近语句里不相干的 <c>@TenantId</c> 算进来(不该跳过的跳过了),
- /// 同一段 SQL 加不加租户条件取决于文本长度。<c>mdp_source_table_registry</c> 的子查询
- /// 在 S4 正好被窗口盖住而幸免,在 S3 就被注入成 <c>r.tenant_id=@TenantId</c> 并报 Unknown column。</para>
- /// </summary>
- 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;
- }
- /// <summary>从 <paramref name="openIndex"/> 处的单引号跳到配对的结束引号;支持 '' 与反斜杠转义。</summary>
- 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;
- }
- /// <summary><paramref name="sql"/> 的 <paramref name="at"/> 处是否是一个完整的 <paramref name="word"/> 而非某个标识符的片段。</summary>
- 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);
- /// <summary>单进程并发上限:本实例的 <c>ModuleRebuildWorker</c> 最多同时跑几个 scope。</summary>
- int MaxParallelScopes { get; }
- /// <summary>
- /// 全库并发上限:**所有**连同一个库的实例加起来最多同时跑几个 scope,
- /// 由 <c>ModuleRebuildJobStore.ClaimNextQueuedAsync</c> 在抢占时收口。
- /// <para>只有 <see cref="MaxParallelScopes"/> 时它约束的是单进程,多实例共库会成倍放大:
- /// 2026-09-24 取证里配置值是 1,而实测同一瞬间有 2–3 个重算在跑。</para>
- /// </summary>
- int GlobalMaxParallelScopes { get; }
- /// <summary>
- /// AUTO / BOOTSTRAP 入队的最小间隔小时数。MANUAL 与 AUTO_NIGHTLY 不受此限。
- /// 取 0 表示关闭冷却,仅供排障时临时恢复旧行为。
- /// </summary>
- int AutoMinIntervalHours { get; }
- }
- public sealed class ModuleRebuildCapability : IModuleRebuildCapability, ISingleton
- {
- public bool IsEnabled(string moduleCode)
- {
- var mc = MdpRebuildScope.NormalizeModule(moduleCode);
- try
- {
- return Furion.App.GetConfig<bool?>($"AiDOP:MdpRebuild:Modules:{mc}:Enabled", true) ?? false;
- }
- catch
- {
- return false;
- }
- }
- public int MaxParallelScopes
- {
- get
- {
- try
- {
- return Math.Clamp(Furion.App.GetConfig<int?>("AiDOP:MdpRebuild:MaxParallelScopes", true) ?? 2, 1, 8);
- }
- catch
- {
- return 2;
- }
- }
- }
- public int GlobalMaxParallelScopes
- {
- get
- {
- try
- {
- // 缺省回落到单进程上限:未显式配置时行为等价于「只有一个实例在跑」,
- // 而不是悄悄放开成无上限。
- var v = Furion.App.GetConfig<int?>("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<int?>("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;
- }
|