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; 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; public MdmMirrorUpsertService( ISqlSugarClient db, AdoS0ReferenceChecker refs, ILogger logger) { _db = db; _refs = refs; _logger = logger; } public async Task MirrorCommittedAsync( string entityCode, long tenantId, long? factoryId, string sourceCode, IReadOnlyList rows, CancellationToken ct) { if (rows == null || rows.Count == 0) return; var factory = factoryId is > 0 ? factoryId.Value : 1L; foreach (var row in rows) { try { var code = (entityCode ?? string.Empty).Trim().ToUpperInvariant(); 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); } catch (Exception ex) { _logger.LogError(ex, "inbound mirror failed entity={Entity} biz={Biz}", entityCode, row.BizKey); } } } 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"), SourceSystem = sourceCode, SourceUpdatedAt = incomingAt, CreateTime = now, UpdateTime = now }).ExecuteCommandAsync(ct); return; } existing.Descr = Str(row.Dict, "descr") ?? existing.Descr; existing.SourceSystem = sourceCode; existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt; existing.UpdateTime = now; await _db.Updateable(existing).ExecuteCommandAsync(ct); } 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 { 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", string.IsNullOrWhiteSpace(sourceCode) ? "API" : sourceCode), 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", NormalizePosition(Str(row.Dict, "Position"))), 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 static string NormalizePosition(string? raw) { var v = (raw ?? string.Empty).Trim().ToUpperInvariant(); return v is "PLANNER" or "WAREHOUSE" or "PRODUCTION" ? v : "OTHER"; } 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); }