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); }