using Admin.NET.Plugin.AiDOP.DataPlatform;
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
using Admin.NET.Plugin.AiDOP.Entity.S0.Sales;
using Admin.NET.Plugin.AiDOP.Entity.S0.Supply;
using Admin.NET.Plugin.AiDOP.Entity.S0.Warehouse;
using Admin.NET.Plugin.AiDOP.Infrastructure;
using Microsoft.Extensions.Logging;
using Yitter.IdGenerator;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.Inbound;
/// 主数据 stg 批 COMMITTED 后回灌 S0 镜像。失败不阻断接收。
public sealed class MdmMirrorUpsertService : ITransient
{
private readonly ISqlSugarClient _db;
private readonly AdoS0ReferenceChecker _refs;
private readonly ILogger _logger;
private readonly MdpNeutralSourceGate _neutralGate;
private readonly EmployeePositionMapService _positions;
public MdmMirrorUpsertService(
ISqlSugarClient db,
AdoS0ReferenceChecker refs,
ILogger logger,
MdpNeutralSourceGate neutralGate,
EmployeePositionMapService positions)
{
_db = db;
_refs = refs;
_logger = logger;
_neutralGate = neutralGate;
_positions = positions;
}
public async Task MirrorCommittedAsync(
string entityCode,
long tenantId,
long? factoryId,
string sourceCode,
IReadOnlyList rows,
CancellationToken ct)
{
if (rows == null || rows.Count == 0)
return;
using var _ = MdmMirrorWriteScope.Enter();
var factory = factoryId is > 0 ? factoryId.Value : 1L;
var code0 = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
var parentRoute = MdmMirrorCatalog.Find(code0);
if (!string.IsNullOrWhiteSpace(parentRoute?.ParentColumn))
{
await MirrorParentReplaceAsync(tenantId, factory, sourceCode, parentRoute!, rows, ct);
return;
}
foreach (var row in rows)
{
try
{
var code = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
var route = MdmMirrorCatalog.Find(code);
if (route == null)
continue;
if (code is "MDM_ITEM")
await MirrorItemAsync(tenantId, factory, sourceCode, row, ct);
else if (code is "MDM_CUSTOMER")
await MirrorCustomerAsync(tenantId, factory, sourceCode, row, ct);
else if (code is "MDM_SUPPLIER")
await MirrorSupplierAsync(tenantId, factory, sourceCode, row, ct);
else if (code is "MDM_LOCATION")
await MirrorLocationAsync(tenantId, factory, sourceCode, row, ct);
else if (code is "MDM_EMPLOYEE_HEADCOUNT")
await MirrorEmployeeAsync(tenantId, factory, sourceCode, row, ct);
else
await MirrorGenericAsync(tenantId, factory, sourceCode, route, row, ct);
}
catch (Exception ex)
{
_logger.LogError(ex, "inbound mirror failed entity={Entity} biz={Biz}", entityCode, row.BizKey);
}
}
}
private async Task MirrorParentReplaceAsync(
long tenantId, long factoryId, string sourceCode, MdmMirrorCatalog.Route route,
IReadOnlyList rows, CancellationToken ct)
{
var columns = await LoadColumnsAsync(route.Table, ct);
var parentCol = MatchColumn(columns, route.ParentColumn!);
var tenantCol = MatchColumn(columns, "TenantId");
var sourceCol = MatchColumn(columns, "SourceSystem");
if (parentCol == null || tenantCol == null)
return;
foreach (var group in rows.GroupBy(r => Str(r.Dict, route.ParentColumn!) ?? ""))
{
if (string.IsNullOrWhiteSpace(group.Key))
continue;
if (sourceCol != null)
{
await _db.Ado.ExecuteCommandAsync(
$"DELETE FROM `{route.Table}` WHERE `{tenantCol}`=@tenant AND `{parentCol}`=@parent AND `{sourceCol}` IS NOT NULL",
new SugarParameter("@tenant", tenantId),
new SugarParameter("@parent", group.Key));
}
foreach (var row in group)
await InsertAlignedAsync(route.Table, columns, tenantId, factoryId, sourceCode, row, ct);
}
}
private async Task MirrorGenericAsync(
long tenantId, long factoryId, string sourceCode, MdmMirrorCatalog.Route route,
MdpInboundPreparedRow row, CancellationToken ct)
{
var columns = await LoadColumnsAsync(route.Table, ct);
var keyCol = MatchColumn(columns, route.KeyColumn);
var tenantCol = MatchColumn(columns, "TenantId");
if (keyCol == null || tenantCol == null)
return;
var key = Str(row.Dict, route.KeyColumn);
if (string.IsNullOrWhiteSpace(key))
return;
var assignments = AlignedAssignments(columns, tenantId, factoryId, sourceCode, row)
.Where(x => !string.Equals(x.Column, keyCol, StringComparison.OrdinalIgnoreCase)
&& !string.Equals(x.Column, tenantCol, StringComparison.OrdinalIgnoreCase)
&& !string.Equals(x.Column, "Id", StringComparison.OrdinalIgnoreCase))
.ToList();
if (assignments.Count == 0)
return;
var setSql = string.Join(", ", assignments.Select(x => $"`{x.Column}`=@{x.Name}"));
var pars = assignments.Select(x => new SugarParameter("@" + x.Name, x.Value)).ToList();
pars.Add(new SugarParameter("@tenant", tenantId));
pars.Add(new SugarParameter("@key", key));
var updated = await _db.Ado.ExecuteCommandAsync(
$"UPDATE `{route.Table}` SET {setSql} WHERE `{tenantCol}`=@tenant AND `{keyCol}`=@key",
pars);
if (updated == 0)
await InsertAlignedAsync(route.Table, columns, tenantId, factoryId, sourceCode, row, ct);
}
private async Task InsertAlignedAsync(
string table, IReadOnlyList columns, long tenantId, long factoryId, string sourceCode,
MdpInboundPreparedRow row, CancellationToken ct)
{
ct.ThrowIfCancellationRequested();
var fields = AlignedAssignments(columns, tenantId, factoryId, sourceCode, row).ToList();
if (MatchColumn(columns, "Id") is { } idCol && fields.All(x => !string.Equals(x.Column, idCol, StringComparison.OrdinalIgnoreCase)))
fields.Add(("Id", "id", YitIdHelper.NextId()));
if (fields.Count == 0)
return;
var cols = string.Join(", ", fields.Select(x => $"`{x.Column}`"));
var vals = string.Join(", ", fields.Select(x => "@" + x.Name));
var pars = fields.Select(x => new SugarParameter("@" + x.Name, x.Value)).ToArray();
await _db.Ado.ExecuteCommandAsync($"INSERT INTO `{table}` ({cols}) VALUES ({vals})", pars);
}
private static List<(string Column, string Name, object? Value)> AlignedAssignments(
IReadOnlyList columns, long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row)
{
var list = new List<(string, string, object?)>();
var used = new HashSet(StringComparer.OrdinalIgnoreCase);
void Add(string? column, object? value)
{
if (column == null || !used.Add(column))
return;
list.Add((column, "p" + list.Count, value));
}
foreach (var kv in row.Dict)
Add(MatchColumn(columns, kv.Key), kv.Value);
Add(MatchColumn(columns, "TenantId"), tenantId);
Add(MatchColumn(columns, "FactoryId"), factoryId);
Add(MatchColumn(columns, "SourceSystem"), sourceCode);
Add(MatchColumn(columns, "SourceUpdatedAt"), row.SourceUpdatedAt);
return list;
}
private async Task> LoadColumnsAsync(string table, CancellationToken ct)
{
var rows = await _db.Ado.SqlQueryAsync(
"""
SELECT COLUMN_NAME FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = @t
""",
new SugarParameter("@t", table));
ct.ThrowIfCancellationRequested();
return rows;
}
private static string? MatchColumn(IReadOnlyList columns, string name) =>
columns.FirstOrDefault(c => string.Equals(c, name, StringComparison.OrdinalIgnoreCase));
private async Task MirrorItemAsync(
long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
{
var itemNum = Str(row.Dict, "ItemNum") ?? Str(row.Dict, "number");
if (string.IsNullOrWhiteSpace(itemNum))
return;
var incomingAt = ParseTime(row.SourceUpdatedAt);
var existing = await _db.Queryable()
.Where(x => x.TenantId == tenantId && x.ItemNum == itemNum)
.FirstAsync(ct);
if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
{
await InsertConflictAsync(tenantId, "MDM_ITEM", row.BizKey, sourceCode, row.RawJson, "ItemMaster", ct);
return;
}
if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
{
_logger.LogInformation("inbound mirror skipped_stale item={Item}", itemNum);
return;
}
if (!string.IsNullOrWhiteSpace(Str(row.Dict, "Location")))
await _refs.LocationExistsAsync(tenantId, Str(row.Dict, "Location"));
var now = DateTime.Now;
if (existing == null)
{
await _db.Insertable(new AdoS0ItemMaster
{
TenantId = tenantId,
FactoryRefId = factoryId,
DomainCode = Str(row.Dict, "Domain"),
ItemNum = itemNum,
Descr = Str(row.Dict, "Descr") ?? Str(row.Dict, "name") ?? itemNum,
Drawing = Str(row.Dict, "Drawing") ?? Str(row.Dict, "model"),
UM = Str(row.Dict, "UM") ?? Str(row.Dict, "unit"),
ItemType = Str(row.Dict, "ItemType"),
Status = Str(row.Dict, "Status") ?? "normal",
IsActive = !IsInactive(Str(row.Dict, "Status") ?? Str(row.Dict, "is_active")),
SourceSystem = sourceCode,
SourceUpdatedAt = incomingAt,
CreateTime = now,
UpdateTime = now,
UpdateUser = "API_INBOUND"
}).ExecuteCommandAsync(ct);
return;
}
existing.Descr = Str(row.Dict, "Descr") ?? Str(row.Dict, "name") ?? existing.Descr;
existing.Drawing = Str(row.Dict, "Drawing") ?? Str(row.Dict, "model") ?? existing.Drawing;
existing.UM = Str(row.Dict, "UM") ?? Str(row.Dict, "unit") ?? existing.UM;
existing.ItemType = Str(row.Dict, "ItemType") ?? existing.ItemType;
existing.DomainCode = Str(row.Dict, "Domain") ?? existing.DomainCode;
if (!string.IsNullOrWhiteSpace(Str(row.Dict, "Status")))
{
existing.Status = Str(row.Dict, "Status");
existing.IsActive = !IsInactive(existing.Status);
}
existing.SourceSystem = sourceCode;
existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
existing.UpdateTime = now;
existing.UpdateUser = "API_INBOUND";
await _db.Updateable(existing)
.IgnoreColumns(x => new { x.Location, x.DefaultShelf, x.SafetyStk, x.LotSerialControl, x.AllocateSingleLot })
.ExecuteCommandAsync(ct);
}
private async Task MirrorCustomerAsync(
long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
{
var cust = Str(row.Dict, "Cust") ?? Str(row.Dict, "custom_no");
if (string.IsNullOrWhiteSpace(cust))
return;
var incomingAt = ParseTime(row.SourceUpdatedAt);
var existing = await _db.Queryable()
.Where(x => x.TenantId == tenantId && x.Cust == cust)
.FirstAsync(ct);
if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
{
await InsertConflictAsync(tenantId, "MDM_CUSTOMER", row.BizKey, sourceCode, row.RawJson, "CustMaster", ct);
return;
}
if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
return;
var now = DateTime.Now;
if (existing == null)
{
await _db.Insertable(new AdoS0CustMaster
{
TenantId = tenantId,
FactoryRefId = factoryId,
Cust = cust,
SortName = Str(row.Dict, "SortName") ?? Str(row.Dict, "custom_name"),
SourceSystem = sourceCode,
SourceUpdatedAt = incomingAt,
CreateTime = now,
UpdateTime = now
}).ExecuteCommandAsync(ct);
return;
}
existing.SortName = Str(row.Dict, "SortName") ?? existing.SortName;
existing.SourceSystem = sourceCode;
existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
existing.UpdateTime = now;
await _db.Updateable(existing).ExecuteCommandAsync(ct);
}
private async Task MirrorSupplierAsync(
long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
{
var supp = Str(row.Dict, "Supp") ?? Str(row.Dict, "supplier_number");
if (string.IsNullOrWhiteSpace(supp))
return;
var incomingAt = ParseTime(row.SourceUpdatedAt);
var existing = await _db.Queryable()
.Where(x => x.TenantId == tenantId && x.Supp == supp)
.FirstAsync(ct);
if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
{
await InsertConflictAsync(tenantId, "MDM_SUPPLIER", row.BizKey, sourceCode, row.RawJson, "SuppMaster", ct);
return;
}
if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
return;
var now = DateTime.Now;
if (existing == null)
{
await _db.Insertable(new AdoS0SuppMaster
{
TenantId = tenantId,
FactoryRefId = factoryId,
Supp = supp,
SortName = Str(row.Dict, "SortName"),
SourceSystem = sourceCode,
SourceUpdatedAt = incomingAt,
CreateTime = now,
UpdateTime = now
}).ExecuteCommandAsync(ct);
return;
}
existing.SortName = Str(row.Dict, "SortName") ?? existing.SortName;
existing.SourceSystem = sourceCode;
existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
existing.UpdateTime = now;
await _db.Updateable(existing).ExecuteCommandAsync(ct);
}
private async Task MirrorLocationAsync(
long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
{
var loc = Str(row.Dict, "location") ?? Str(row.Dict, "Location");
if (string.IsNullOrWhiteSpace(loc))
return;
var incomingAt = ParseTime(row.SourceUpdatedAt);
var existing = await _db.Queryable()
.Where(x => x.TenantId == tenantId && x.Location == loc)
.FirstAsync(ct);
if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
{
await InsertConflictAsync(tenantId, "MDM_LOCATION", row.BizKey, sourceCode, row.RawJson, "LocationMaster", ct);
return;
}
if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
return;
var now = DateTime.Now;
if (existing == null)
{
await _db.Insertable(new AdoS0LocationMaster
{
TenantId = tenantId,
FactoryRefId = factoryId,
DomainCode = Str(row.Dict, "Domain") ?? "",
Location = loc,
Descr = Str(row.Dict, "descr") ?? Str(row.Dict, "Descr"),
Typed = Str(row.Dict, "Typed"),
SourceSystem = sourceCode,
SourceUpdatedAt = incomingAt,
CreateTime = now,
UpdateTime = now
}).ExecuteCommandAsync(ct);
await UpsertLocationRoleAsync(tenantId, Str(row.Dict, "Domain"), loc, Str(row.Dict, "Typed"), Str(row.Dict, "descr") ?? Str(row.Dict, "Descr"), ct);
return;
}
existing.Descr = Str(row.Dict, "descr") ?? existing.Descr;
existing.Typed = Str(row.Dict, "Typed") ?? existing.Typed;
existing.SourceSystem = sourceCode;
existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
existing.UpdateTime = now;
await _db.Updateable(existing).ExecuteCommandAsync(ct);
await UpsertLocationRoleAsync(tenantId, existing.DomainCode, loc, existing.Typed, existing.Descr, ct);
}
///
/// 库位写入后按现行关键词规则增量分类。MANUAL 与 RATIFIED 不覆盖。
///
private async Task UpsertLocationRoleAsync(
long tenantId, string? domain, string location, string? typed, string? descr, CancellationToken ct)
{
ct.ThrowIfCancellationRequested();
var role = NeutralTransTypeCodes.ClassifyLocationRole(typed, descr);
if (string.IsNullOrWhiteSpace(role))
role = "OTHER";
try
{
await _db.Ado.ExecuteCommandAsync(
"""
INSERT INTO mdp_location_role
(tenant_id, domain, location, location_role, role_source, remark, create_time)
VALUES (@tid, @domain, @loc, @role, 'AUTO_DESCR', @remark, NOW())
ON DUPLICATE KEY UPDATE
location_role = IF(role_source IN ('MANUAL','RATIFIED'), location_role, VALUES(location_role)),
remark = IF(role_source IN ('MANUAL','RATIFIED'), remark, VALUES(remark))
""",
new SugarParameter("@tid", tenantId),
new SugarParameter("@domain", domain ?? ""),
new SugarParameter("@loc", location),
new SugarParameter("@role", role),
new SugarParameter("@remark", "入库镜像按库位说明分类"));
}
catch (Exception ex)
{
_logger.LogWarning(ex, "location role classify skipped loc={Location}", location);
}
}
private async Task MirrorEmployeeAsync(
long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
{
var emp = Str(row.Dict, "Employee") ?? Str(row.Dict, "employee");
if (string.IsNullOrWhiteSpace(emp))
return;
if (!string.IsNullOrWhiteSpace(Str(row.Dict, "Department")))
await _refs.DepartmentExistsAsync(tenantId, Str(row.Dict, "Department"));
var incomingAt = ParseTime(row.SourceUpdatedAt);
var existing = await _db.Queryable()
.Where(x => x.TenantId == tenantId && x.Employee == emp)
.FirstAsync(ct);
if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
{
await InsertConflictAsync(tenantId, "MDM_EMPLOYEE_HEADCOUNT", row.BizKey, sourceCode, row.RawJson, "EmployeeMaster", ct);
return;
}
if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
return;
var now = DateTime.Now;
if (existing == null)
{
await _db.Insertable(new AdoS0EmployeeMaster
{
TenantId = tenantId,
FactoryRefId = factoryId,
DomainCode = Str(row.Dict, "Domain") ?? "",
Employee = emp,
Name = Str(row.Dict, "Name"),
SourceSystem = sourceCode,
SourceUpdatedAt = incomingAt,
CreateTime = now,
UpdateTime = now
}).ExecuteCommandAsync(ct);
}
else
{
existing.Name = Str(row.Dict, "Name") ?? existing.Name;
existing.SourceSystem = sourceCode;
existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
existing.UpdateTime = now;
await _db.Updateable(existing).ExecuteCommandAsync(ct);
}
try
{
var empSource = string.IsNullOrWhiteSpace(sourceCode) ? "API" : sourceCode;
if (!await _neutralGate.AllowsAsync(tenantId, "EMPLOYEE", empSource))
return;
await _db.Ado.ExecuteCommandAsync(
"""
INSERT INTO mdp_std_employee
(tenant_id, factory_id, source_system, domain, employee_no, employee_name, department_code,
position_code, src_position_raw, employment_status, src_employment_status_raw,
source_row_id, source_biz_key, sync_batch_id, sync_time)
VALUES
(@tid, @factory, @src, @domain, @emp, @name, @dept,
@pos, @posRaw, @status, @statusRaw,
@emp, @biz, @batch, @now)
ON DUPLICATE KEY UPDATE
employee_name=VALUES(employee_name), department_code=VALUES(department_code),
position_code=VALUES(position_code), src_position_raw=VALUES(src_position_raw),
employment_status=VALUES(employment_status), src_employment_status_raw=VALUES(src_employment_status_raw),
sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
""",
new SugarParameter("@tid", tenantId),
new SugarParameter("@factory", factoryId),
new SugarParameter("@src", empSource),
new SugarParameter("@domain", Str(row.Dict, "Domain") ?? ""),
new SugarParameter("@emp", emp),
new SugarParameter("@name", Str(row.Dict, "Name")),
new SugarParameter("@dept", Str(row.Dict, "Department")),
new SugarParameter("@pos", await ResolvePositionCodeAsync(tenantId, empSource, Str(row.Dict, "Domain") ?? "", Str(row.Dict, "Position"), ct)),
new SugarParameter("@posRaw", Str(row.Dict, "Position")),
new SugarParameter("@status", NormalizeEmployment(Str(row.Dict, "EmploymentStatus"))),
new SugarParameter("@statusRaw", Str(row.Dict, "EmploymentStatus")),
new SugarParameter("@biz", $"{Str(row.Dict, "Domain")}:{emp}"),
new SugarParameter("@batch", "INBOUND_EMP"),
new SugarParameter("@now", now));
}
catch (Exception ex)
{
_logger.LogWarning(ex, "雇员中立层未就绪,已跳过 mdp_std_employee,镜像主数据仍已写入");
}
}
private async Task ResolvePositionCodeAsync(long tenantId, string sourceSystem, string domain, string? raw, CancellationToken ct)
{
if (string.IsNullOrWhiteSpace(raw))
return EmployeePositionRules.Unknown;
await _positions.UpsertAsync(tenantId, sourceSystem, "Position",
[new EmployeePositionMapService.RawPosition { Domain = domain, Raw = raw }], ct);
var codes = await _db.Ado.SqlQueryAsync(
"""
SELECT position_code FROM mdp_employee_position_map
WHERE tenant_id=@tid AND source_system=@src AND domain=@domain AND src_position_raw=@raw
LIMIT 1
""",
new { tid = tenantId, src = sourceSystem, domain, raw });
return codes.FirstOrDefault() ?? EmployeePositionRules.Classify(raw);
}
private static string NormalizeEmployment(string? raw)
{
var v = (raw ?? string.Empty).Trim();
var u = v.ToUpperInvariant();
if (v is "在职" || u is "ACTIVE" or "ONJOB" or "ON_JOB") return "ACTIVE";
if (v is "离职" or "辞退" || u is "LEFT" or "LEAVE" or "RESIGNED" or "TERMINATED" or "QUIT") return "LEFT";
if (v is "停用" || u is "INACTIVE" or "DISABLED" or "SUSPENDED") return "INACTIVE";
return "UNKNOWN";
}
private async Task InsertConflictAsync(
long tenantId, string entityCode, string bizKey, string sourceCode, string raw, string table, CancellationToken ct)
{
var now = DateTime.Now;
await _db.Insertable(new MdpInboundConflict
{
TenantId = tenantId,
EntityCode = entityCode,
BizKey = bizKey ?? string.Empty,
SourceSystem = sourceCode,
IncomingRaw = raw ?? string.Empty,
MirrorTable = table,
ConflictType = "MANUAL_ROW_EXISTS",
Status = "PENDING",
CreateTime = now,
UpdateTime = now
}).ExecuteCommandAsync(ct);
}
private static string Str(IDictionary row, string key)
{
foreach (var kv in row)
{
if (string.Equals(kv.Key, key, StringComparison.OrdinalIgnoreCase))
return kv.Value?.ToString();
}
return null;
}
private static DateTime? ParseTime(string raw) =>
DateTimeOffset.TryParse(raw, out var dto) ? dto.LocalDateTime : null;
private static bool IsInactive(string status) =>
string.Equals(status, "INACTIVE", StringComparison.OrdinalIgnoreCase)
|| status == "0"
|| string.Equals(status, "false", StringComparison.OrdinalIgnoreCase);
}