using Microsoft.Extensions.Logging;
namespace Admin.NET.Plugin.AiDOP.DataPlatform;
///
/// 删除某租户某标准对象里来源与登记不一致的中立行。不删贴源,不删目录之外的表。
///
public sealed class MdpNeutralSourceCleanup : ITransient
{
private readonly ISqlSugarClient _db;
private readonly ILogger _logger;
public MdpNeutralSourceCleanup(ISqlSugarClient db, ILoggerFactory loggerFactory)
{
_db = db;
_logger = loggerFactory.CreateLogger(nameof(MdpNeutralSourceCleanup));
}
public async Task PurgeModuleAsync(long tenantId, string moduleCode, CancellationToken cancellationToken = default)
{
var total = 0;
foreach (var def in MdpStdObjectCatalog.All.Where(d => d.Module == moduleCode))
{
cancellationToken.ThrowIfCancellationRequested();
total += await PurgeObjectAsync(tenantId, def.Code, cancellationToken);
}
return total;
}
public async Task PurgeObjectAsync(long tenantId, string stdObject, CancellationToken cancellationToken = default)
{
var def = MdpStdObjectCatalog.Find(stdObject);
if (def == null || tenantId <= 0)
return 0;
var total = 0;
foreach (var table in def.Tables)
{
cancellationToken.ThrowIfCancellationRequested();
if (!await TableReadyAsync(table))
continue;
var docFilter = string.IsNullOrEmpty(def.DocType) || !await ColumnExistsAsync(table, "doc_type")
? ""
: " AND t.doc_type=@doc";
if (!await RegisteredSourceHasRowsAsync(tenantId, def.Code, table, docFilter, cancellationToken))
continue;
var sql = $"""
DELETE t FROM `{table}` t
INNER JOIN mdp_tenant_std_source s
ON s.tenant_id=t.tenant_id AND s.std_object=@obj
WHERE t.tenant_id=@tenant
AND t.source_system<>s.source_system
{docFilter}
""";
var n = await _db.Ado.ExecuteCommandAsync(sql, new { tenant = tenantId, obj = def.Code, doc = def.DocType });
if (n > 0)
_logger.LogWarning(
"中立层来源清理 tenant={Tenant} object={Object} table={Table} deleted={Deleted}",
tenantId, def.Code, table, n);
total += n;
}
return total;
}
///
/// 登记来源在该表一行都没有时,不执行清理。
///
/// 清理是「按登记收口」,不是「清空对象」。若登记来源尚未产出任何行,删掉非登记来源就等于
/// 把该对象整片清空:指标随即算出 NO_DATA 且全程不报错,看板一夜之间没数(S6/S7 就是这么掉的)。
/// 这种情况下真正的问题是登记与数据不一致,须人工核对登记,不能靠重算去「修正」。
/// 故此处保留存量并按 PURGE_SKIPPED_EMPTY_SOURCE 告警,宁可留旧数据也不静默清空。
///
///
private async Task RegisteredSourceHasRowsAsync(
long tenantId, string stdObject, string table, string docFilter, CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
var rows = await _db.Ado.SqlQueryAsync(
$"""
SELECT MAX(s.source_system) AS RegisteredSource,
SUM(t.source_system=s.source_system) AS RegisteredRows,
SUM(t.source_system<>s.source_system) AS ForeignRows
FROM `{table}` t
INNER JOIN mdp_tenant_std_source s
ON s.tenant_id=t.tenant_id AND s.std_object=@obj
WHERE t.tenant_id=@tenant
{docFilter}
""",
new { tenant = tenantId, obj = stdObject, doc = MdpStdObjectCatalog.Find(stdObject)?.DocType });
var stat = rows.FirstOrDefault();
if (stat == null || stat.ForeignRows <= 0)
return true;
if (stat.RegisteredRows > 0)
return true;
_logger.LogError(
"中立层来源清理已跳过:登记来源无数据 tenant={Tenant} object={Object} table={Table} 登记来源={Source} 待清理行={Foreign}",
tenantId, stdObject, table, stat.RegisteredSource, stat.ForeignRows);
await _db.Ado.ExecuteCommandAsync(
"""
INSERT INTO mdp_source_gate_log
(tenant_id, std_object, source_system, gate_reason, row_count, sample_keys, sync_batch_id)
VALUES (@tenant, @obj, @src, 'PURGE_SKIPPED_EMPTY_SOURCE', @rows, @sample, '')
""",
new
{
tenant = tenantId,
obj = stdObject,
src = stat.RegisteredSource ?? "",
rows = stat.ForeignRows,
sample = table
});
return false;
}
private sealed class SourceRowCount
{
public string? RegisteredSource { get; set; }
public long RegisteredRows { get; set; }
public long ForeignRows { get; set; }
}
private async Task TableReadyAsync(string table)
{
var n = await _db.Ado.SqlQueryAsync(
"""
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table AND COLUMN_NAME='source_system'
""",
new { table });
return n.FirstOrDefault() > 0;
}
private async Task ColumnExistsAsync(string table, string column)
{
var n = await _db.Ado.SqlQueryAsync(
"""
SELECT COUNT(*) FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@table AND COLUMN_NAME=@column
""",
new { table, column });
return n.FirstOrDefault() > 0;
}
}