using System.Text.Json;
using Microsoft.Extensions.Logging;
using SqlSugar;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
///
/// S0 维度刷新编排:批次号、逐维度物化、run log 汇总。
///
/// 一次刷新 = 一个 batch_id(同时是 staging / dim 的 sync_batch_id 与
/// mdp_transform_run_log.batch_id;后者是 UNIQUE,故只写一条汇总行,
/// 逐维度明细进 summary_json)。
///
public sealed class S0DimRefreshService : ITransient
{
/// run log 的作业编码。
public const string JobCode = "S0_DIM_REFRESH";
private static readonly JsonSerializerOptions JsonOpts = new()
{
Encoder = System.Text.Encodings.Web.JavaScriptEncoder.UnsafeRelaxedJsonEscaping
};
private readonly ISqlSugarClient _db;
private readonly S0DimMaterializer _materializer;
private readonly S0DimReconciler _reconciler;
private readonly ILogger _logger;
/// 构造。
public S0DimRefreshService(
ISqlSugarClient db,
S0DimMaterializer materializer,
S0DimReconciler reconciler,
ILogger logger)
{
_db = db;
_materializer = materializer;
_reconciler = reconciler;
_logger = logger;
}
///
/// 刷新指定租户的维度。 为空表示全部(按 的依赖顺序)。
///
public async Task RefreshAsync(
long tenantId, string? key = null, string triggerType = "MANUAL", CancellationToken ct = default)
{
if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}");
S0DimCatalog.ValidateAll();
var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)];
var startedAt = DateTime.Now;
var batchId = $"S0DIM_{tenantId}_{startedAt:yyyyMMddHHmmssfff}";
var result = new S0DimRefreshResult { BatchId = batchId, TenantId = tenantId };
result.RunLogId = await InsertRunLogAsync(tenantId, batchId, startedAt, triggerType);
try
{
foreach (var def in defs)
{
ct.ThrowIfCancellationRequested();
result.Items.Add(await _materializer.RunAsync(def, tenantId, batchId, ct));
}
}
catch (OperationCanceledException)
{
// run log 没有 stale 回收机制,而 LastRunAsync 会把 RUNNING 行当成「最近一次结果」返回。
// 若取消时直接抛出,这条记录会永久停在 RUNNING,读者无法区分「正在跑」与「早就断了」。
// 故先收口再重抛。CompleteRunLogAsync 不接收 ct —— 否则它自己也会被同一个 token 取消。
result.Status = "CANCELED";
result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds;
await CompleteRunLogAsync(result, startedAt);
_logger.LogWarning("[S0Dim] refresh canceled tenant={Tenant} batch={Batch} 已完成 {Done}/{Total} 个维度",
tenantId, batchId, result.Items.Count, defs.Count);
throw;
}
result.Status = result.Items.Any(i => i.Status == "FAILED")
? "FAILED"
: result.Items.Any(i => i.Reconcile?.Warnings.Count > 0)
? "SUCCESS_WITH_WARNING"
: "SUCCESS";
result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds;
await CompleteRunLogAsync(result, startedAt);
_logger.LogInformation("[S0Dim] refresh done tenant={Tenant} batch={Batch} status={Status}",
tenantId, batchId, result.Status);
return result;
}
/// 只读对账(不改数据)。 为空表示全部。
public async Task> ReconcileAsync(
long tenantId, string? key = null, CancellationToken ct = default)
{
if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}");
S0DimCatalog.ValidateAll();
var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)];
var list = new List();
foreach (var def in defs)
{
ct.ThrowIfCancellationRequested();
// batchId 传 null:按「当前状态」口径对账,不限定某一批
list.Add(await _reconciler.ReconcileAsync(def, tenantId, null, ct));
}
return list;
}
/// 最近一次运行记录(按租户)。
public async Task