MdmMirrorUpsertService.cs 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567
  1. using Admin.NET.Plugin.AiDOP.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.Entity.S0.Sales;
  4. using Admin.NET.Plugin.AiDOP.Entity.S0.Supply;
  5. using Admin.NET.Plugin.AiDOP.Entity.S0.Warehouse;
  6. using Admin.NET.Plugin.AiDOP.Infrastructure;
  7. using Microsoft.Extensions.Logging;
  8. using Yitter.IdGenerator;
  9. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Inbound;
  10. /// <summary>主数据 stg 批 COMMITTED 后回灌 S0 镜像。失败不阻断接收。</summary>
  11. public sealed class MdmMirrorUpsertService : ITransient
  12. {
  13. private readonly ISqlSugarClient _db;
  14. private readonly AdoS0ReferenceChecker _refs;
  15. private readonly ILogger<MdmMirrorUpsertService> _logger;
  16. private readonly MdpNeutralSourceGate _neutralGate;
  17. private readonly EmployeePositionMapService _positions;
  18. public MdmMirrorUpsertService(
  19. ISqlSugarClient db,
  20. AdoS0ReferenceChecker refs,
  21. ILogger<MdmMirrorUpsertService> logger,
  22. MdpNeutralSourceGate neutralGate,
  23. EmployeePositionMapService positions)
  24. {
  25. _db = db;
  26. _refs = refs;
  27. _logger = logger;
  28. _neutralGate = neutralGate;
  29. _positions = positions;
  30. }
  31. public async Task MirrorCommittedAsync(
  32. string entityCode,
  33. long tenantId,
  34. long? factoryId,
  35. string sourceCode,
  36. IReadOnlyList<MdpInboundPreparedRow> rows,
  37. CancellationToken ct)
  38. {
  39. if (rows == null || rows.Count == 0)
  40. return;
  41. using var _ = MdmMirrorWriteScope.Enter();
  42. var factory = factoryId is > 0 ? factoryId.Value : 1L;
  43. var code0 = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
  44. var parentRoute = MdmMirrorCatalog.Find(code0);
  45. if (!string.IsNullOrWhiteSpace(parentRoute?.ParentColumn))
  46. {
  47. await MirrorParentReplaceAsync(tenantId, factory, sourceCode, parentRoute!, rows, ct);
  48. return;
  49. }
  50. foreach (var row in rows)
  51. {
  52. try
  53. {
  54. var code = (entityCode ?? string.Empty).Trim().ToUpperInvariant();
  55. var route = MdmMirrorCatalog.Find(code);
  56. if (route == null)
  57. continue;
  58. if (code is "MDM_ITEM")
  59. await MirrorItemAsync(tenantId, factory, sourceCode, row, ct);
  60. else if (code is "MDM_CUSTOMER")
  61. await MirrorCustomerAsync(tenantId, factory, sourceCode, row, ct);
  62. else if (code is "MDM_SUPPLIER")
  63. await MirrorSupplierAsync(tenantId, factory, sourceCode, row, ct);
  64. else if (code is "MDM_LOCATION")
  65. await MirrorLocationAsync(tenantId, factory, sourceCode, row, ct);
  66. else if (code is "MDM_EMPLOYEE_HEADCOUNT")
  67. await MirrorEmployeeAsync(tenantId, factory, sourceCode, row, ct);
  68. else
  69. await MirrorGenericAsync(tenantId, factory, sourceCode, route, row, ct);
  70. }
  71. catch (Exception ex)
  72. {
  73. _logger.LogError(ex, "inbound mirror failed entity={Entity} biz={Biz}", entityCode, row.BizKey);
  74. }
  75. }
  76. }
  77. private async Task MirrorParentReplaceAsync(
  78. long tenantId, long factoryId, string sourceCode, MdmMirrorCatalog.Route route,
  79. IReadOnlyList<MdpInboundPreparedRow> rows, CancellationToken ct)
  80. {
  81. var columns = await LoadColumnsAsync(route.Table, ct);
  82. var parentCol = MatchColumn(columns, route.ParentColumn!);
  83. var tenantCol = MatchColumn(columns, "TenantId");
  84. var sourceCol = MatchColumn(columns, "SourceSystem");
  85. if (parentCol == null || tenantCol == null)
  86. return;
  87. foreach (var group in rows.GroupBy(r => Str(r.Dict, route.ParentColumn!) ?? ""))
  88. {
  89. if (string.IsNullOrWhiteSpace(group.Key))
  90. continue;
  91. if (sourceCol != null)
  92. {
  93. await _db.Ado.ExecuteCommandAsync(
  94. $"DELETE FROM `{route.Table}` WHERE `{tenantCol}`=@tenant AND `{parentCol}`=@parent AND `{sourceCol}` IS NOT NULL",
  95. new SugarParameter("@tenant", tenantId),
  96. new SugarParameter("@parent", group.Key));
  97. }
  98. foreach (var row in group)
  99. await InsertAlignedAsync(route.Table, columns, tenantId, factoryId, sourceCode, row, ct);
  100. }
  101. }
  102. private async Task MirrorGenericAsync(
  103. long tenantId, long factoryId, string sourceCode, MdmMirrorCatalog.Route route,
  104. MdpInboundPreparedRow row, CancellationToken ct)
  105. {
  106. var columns = await LoadColumnsAsync(route.Table, ct);
  107. var keyCol = MatchColumn(columns, route.KeyColumn);
  108. var tenantCol = MatchColumn(columns, "TenantId");
  109. if (keyCol == null || tenantCol == null)
  110. return;
  111. var key = Str(row.Dict, route.KeyColumn);
  112. if (string.IsNullOrWhiteSpace(key))
  113. return;
  114. var assignments = AlignedAssignments(columns, tenantId, factoryId, sourceCode, row)
  115. .Where(x => !string.Equals(x.Column, keyCol, StringComparison.OrdinalIgnoreCase)
  116. && !string.Equals(x.Column, tenantCol, StringComparison.OrdinalIgnoreCase)
  117. && !string.Equals(x.Column, "Id", StringComparison.OrdinalIgnoreCase))
  118. .ToList();
  119. if (assignments.Count == 0)
  120. return;
  121. var setSql = string.Join(", ", assignments.Select(x => $"`{x.Column}`=@{x.Name}"));
  122. var pars = assignments.Select(x => new SugarParameter("@" + x.Name, x.Value)).ToList();
  123. pars.Add(new SugarParameter("@tenant", tenantId));
  124. pars.Add(new SugarParameter("@key", key));
  125. var updated = await _db.Ado.ExecuteCommandAsync(
  126. $"UPDATE `{route.Table}` SET {setSql} WHERE `{tenantCol}`=@tenant AND `{keyCol}`=@key",
  127. pars);
  128. if (updated == 0)
  129. await InsertAlignedAsync(route.Table, columns, tenantId, factoryId, sourceCode, row, ct);
  130. }
  131. private async Task InsertAlignedAsync(
  132. string table, IReadOnlyList<string> columns, long tenantId, long factoryId, string sourceCode,
  133. MdpInboundPreparedRow row, CancellationToken ct)
  134. {
  135. ct.ThrowIfCancellationRequested();
  136. var fields = AlignedAssignments(columns, tenantId, factoryId, sourceCode, row).ToList();
  137. if (MatchColumn(columns, "Id") is { } idCol && fields.All(x => !string.Equals(x.Column, idCol, StringComparison.OrdinalIgnoreCase)))
  138. fields.Add(("Id", "id", YitIdHelper.NextId()));
  139. if (fields.Count == 0)
  140. return;
  141. var cols = string.Join(", ", fields.Select(x => $"`{x.Column}`"));
  142. var vals = string.Join(", ", fields.Select(x => "@" + x.Name));
  143. var pars = fields.Select(x => new SugarParameter("@" + x.Name, x.Value)).ToArray();
  144. await _db.Ado.ExecuteCommandAsync($"INSERT INTO `{table}` ({cols}) VALUES ({vals})", pars);
  145. }
  146. private static List<(string Column, string Name, object? Value)> AlignedAssignments(
  147. IReadOnlyList<string> columns, long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row)
  148. {
  149. var list = new List<(string, string, object?)>();
  150. var used = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
  151. void Add(string? column, object? value)
  152. {
  153. if (column == null || !used.Add(column))
  154. return;
  155. list.Add((column, "p" + list.Count, value));
  156. }
  157. foreach (var kv in row.Dict)
  158. Add(MatchColumn(columns, kv.Key), kv.Value);
  159. Add(MatchColumn(columns, "TenantId"), tenantId);
  160. Add(MatchColumn(columns, "FactoryId"), factoryId);
  161. Add(MatchColumn(columns, "SourceSystem"), sourceCode);
  162. Add(MatchColumn(columns, "SourceUpdatedAt"), row.SourceUpdatedAt);
  163. return list;
  164. }
  165. private async Task<IReadOnlyList<string>> LoadColumnsAsync(string table, CancellationToken ct)
  166. {
  167. var rows = await _db.Ado.SqlQueryAsync<string>(
  168. """
  169. SELECT COLUMN_NAME FROM information_schema.COLUMNS
  170. WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = @t
  171. """,
  172. new SugarParameter("@t", table));
  173. ct.ThrowIfCancellationRequested();
  174. return rows;
  175. }
  176. private static string? MatchColumn(IReadOnlyList<string> columns, string name) =>
  177. columns.FirstOrDefault(c => string.Equals(c, name, StringComparison.OrdinalIgnoreCase));
  178. private async Task MirrorItemAsync(
  179. long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
  180. {
  181. var itemNum = Str(row.Dict, "ItemNum") ?? Str(row.Dict, "number");
  182. if (string.IsNullOrWhiteSpace(itemNum))
  183. return;
  184. var incomingAt = ParseTime(row.SourceUpdatedAt);
  185. var existing = await _db.Queryable<AdoS0ItemMaster>()
  186. .Where(x => x.TenantId == tenantId && x.ItemNum == itemNum)
  187. .FirstAsync(ct);
  188. if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
  189. {
  190. await InsertConflictAsync(tenantId, "MDM_ITEM", row.BizKey, sourceCode, row.RawJson, "ItemMaster", ct);
  191. return;
  192. }
  193. if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
  194. {
  195. _logger.LogInformation("inbound mirror skipped_stale item={Item}", itemNum);
  196. return;
  197. }
  198. if (!string.IsNullOrWhiteSpace(Str(row.Dict, "Location")))
  199. await _refs.LocationExistsAsync(tenantId, Str(row.Dict, "Location"));
  200. var now = DateTime.Now;
  201. if (existing == null)
  202. {
  203. await _db.Insertable(new AdoS0ItemMaster
  204. {
  205. TenantId = tenantId,
  206. FactoryRefId = factoryId,
  207. DomainCode = Str(row.Dict, "Domain"),
  208. ItemNum = itemNum,
  209. Descr = Str(row.Dict, "Descr") ?? Str(row.Dict, "name") ?? itemNum,
  210. Drawing = Str(row.Dict, "Drawing") ?? Str(row.Dict, "model"),
  211. UM = Str(row.Dict, "UM") ?? Str(row.Dict, "unit"),
  212. ItemType = Str(row.Dict, "ItemType"),
  213. Status = Str(row.Dict, "Status") ?? "normal",
  214. IsActive = !IsInactive(Str(row.Dict, "Status") ?? Str(row.Dict, "is_active")),
  215. SourceSystem = sourceCode,
  216. SourceUpdatedAt = incomingAt,
  217. CreateTime = now,
  218. UpdateTime = now,
  219. UpdateUser = "API_INBOUND"
  220. }).ExecuteCommandAsync(ct);
  221. return;
  222. }
  223. existing.Descr = Str(row.Dict, "Descr") ?? Str(row.Dict, "name") ?? existing.Descr;
  224. existing.Drawing = Str(row.Dict, "Drawing") ?? Str(row.Dict, "model") ?? existing.Drawing;
  225. existing.UM = Str(row.Dict, "UM") ?? Str(row.Dict, "unit") ?? existing.UM;
  226. existing.ItemType = Str(row.Dict, "ItemType") ?? existing.ItemType;
  227. existing.DomainCode = Str(row.Dict, "Domain") ?? existing.DomainCode;
  228. if (!string.IsNullOrWhiteSpace(Str(row.Dict, "Status")))
  229. {
  230. existing.Status = Str(row.Dict, "Status");
  231. existing.IsActive = !IsInactive(existing.Status);
  232. }
  233. existing.SourceSystem = sourceCode;
  234. existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
  235. existing.UpdateTime = now;
  236. existing.UpdateUser = "API_INBOUND";
  237. await _db.Updateable(existing)
  238. .IgnoreColumns(x => new { x.Location, x.DefaultShelf, x.SafetyStk, x.LotSerialControl, x.AllocateSingleLot })
  239. .ExecuteCommandAsync(ct);
  240. }
  241. private async Task MirrorCustomerAsync(
  242. long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
  243. {
  244. var cust = Str(row.Dict, "Cust") ?? Str(row.Dict, "custom_no");
  245. if (string.IsNullOrWhiteSpace(cust))
  246. return;
  247. var incomingAt = ParseTime(row.SourceUpdatedAt);
  248. var existing = await _db.Queryable<AdoS0CustMaster>()
  249. .Where(x => x.TenantId == tenantId && x.Cust == cust)
  250. .FirstAsync(ct);
  251. if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
  252. {
  253. await InsertConflictAsync(tenantId, "MDM_CUSTOMER", row.BizKey, sourceCode, row.RawJson, "CustMaster", ct);
  254. return;
  255. }
  256. if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
  257. return;
  258. var now = DateTime.Now;
  259. if (existing == null)
  260. {
  261. await _db.Insertable(new AdoS0CustMaster
  262. {
  263. TenantId = tenantId,
  264. FactoryRefId = factoryId,
  265. Cust = cust,
  266. SortName = Str(row.Dict, "SortName") ?? Str(row.Dict, "custom_name"),
  267. SourceSystem = sourceCode,
  268. SourceUpdatedAt = incomingAt,
  269. CreateTime = now,
  270. UpdateTime = now
  271. }).ExecuteCommandAsync(ct);
  272. return;
  273. }
  274. existing.SortName = Str(row.Dict, "SortName") ?? existing.SortName;
  275. existing.SourceSystem = sourceCode;
  276. existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
  277. existing.UpdateTime = now;
  278. await _db.Updateable(existing).ExecuteCommandAsync(ct);
  279. }
  280. private async Task MirrorSupplierAsync(
  281. long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
  282. {
  283. var supp = Str(row.Dict, "Supp") ?? Str(row.Dict, "supplier_number");
  284. if (string.IsNullOrWhiteSpace(supp))
  285. return;
  286. var incomingAt = ParseTime(row.SourceUpdatedAt);
  287. var existing = await _db.Queryable<AdoS0SuppMaster>()
  288. .Where(x => x.TenantId == tenantId && x.Supp == supp)
  289. .FirstAsync(ct);
  290. if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
  291. {
  292. await InsertConflictAsync(tenantId, "MDM_SUPPLIER", row.BizKey, sourceCode, row.RawJson, "SuppMaster", ct);
  293. return;
  294. }
  295. if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
  296. return;
  297. var now = DateTime.Now;
  298. if (existing == null)
  299. {
  300. await _db.Insertable(new AdoS0SuppMaster
  301. {
  302. TenantId = tenantId,
  303. FactoryRefId = factoryId,
  304. Supp = supp,
  305. SortName = Str(row.Dict, "SortName"),
  306. SourceSystem = sourceCode,
  307. SourceUpdatedAt = incomingAt,
  308. CreateTime = now,
  309. UpdateTime = now
  310. }).ExecuteCommandAsync(ct);
  311. return;
  312. }
  313. existing.SortName = Str(row.Dict, "SortName") ?? existing.SortName;
  314. existing.SourceSystem = sourceCode;
  315. existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
  316. existing.UpdateTime = now;
  317. await _db.Updateable(existing).ExecuteCommandAsync(ct);
  318. }
  319. private async Task MirrorLocationAsync(
  320. long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
  321. {
  322. var loc = Str(row.Dict, "location") ?? Str(row.Dict, "Location");
  323. if (string.IsNullOrWhiteSpace(loc))
  324. return;
  325. var incomingAt = ParseTime(row.SourceUpdatedAt);
  326. var existing = await _db.Queryable<AdoS0LocationMaster>()
  327. .Where(x => x.TenantId == tenantId && x.Location == loc)
  328. .FirstAsync(ct);
  329. if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
  330. {
  331. await InsertConflictAsync(tenantId, "MDM_LOCATION", row.BizKey, sourceCode, row.RawJson, "LocationMaster", ct);
  332. return;
  333. }
  334. if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
  335. return;
  336. var now = DateTime.Now;
  337. if (existing == null)
  338. {
  339. await _db.Insertable(new AdoS0LocationMaster
  340. {
  341. TenantId = tenantId,
  342. FactoryRefId = factoryId,
  343. DomainCode = Str(row.Dict, "Domain") ?? "",
  344. Location = loc,
  345. Descr = Str(row.Dict, "descr") ?? Str(row.Dict, "Descr"),
  346. Typed = Str(row.Dict, "Typed"),
  347. SourceSystem = sourceCode,
  348. SourceUpdatedAt = incomingAt,
  349. CreateTime = now,
  350. UpdateTime = now
  351. }).ExecuteCommandAsync(ct);
  352. await UpsertLocationRoleAsync(tenantId, Str(row.Dict, "Domain"), loc, Str(row.Dict, "Typed"), Str(row.Dict, "descr") ?? Str(row.Dict, "Descr"), ct);
  353. return;
  354. }
  355. existing.Descr = Str(row.Dict, "descr") ?? existing.Descr;
  356. existing.Typed = Str(row.Dict, "Typed") ?? existing.Typed;
  357. existing.SourceSystem = sourceCode;
  358. existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
  359. existing.UpdateTime = now;
  360. await _db.Updateable(existing).ExecuteCommandAsync(ct);
  361. await UpsertLocationRoleAsync(tenantId, existing.DomainCode, loc, existing.Typed, existing.Descr, ct);
  362. }
  363. /// <summary>
  364. /// 库位写入后按现行关键词规则增量分类。MANUAL 与 RATIFIED 不覆盖。
  365. /// </summary>
  366. private async Task UpsertLocationRoleAsync(
  367. long tenantId, string? domain, string location, string? typed, string? descr, CancellationToken ct)
  368. {
  369. ct.ThrowIfCancellationRequested();
  370. var role = NeutralTransTypeCodes.ClassifyLocationRole(typed, descr);
  371. if (string.IsNullOrWhiteSpace(role))
  372. role = "OTHER";
  373. try
  374. {
  375. await _db.Ado.ExecuteCommandAsync(
  376. """
  377. INSERT INTO mdp_location_role
  378. (tenant_id, domain, location, location_role, role_source, remark, create_time)
  379. VALUES (@tid, @domain, @loc, @role, 'AUTO_DESCR', @remark, NOW())
  380. ON DUPLICATE KEY UPDATE
  381. location_role = IF(role_source IN ('MANUAL','RATIFIED'), location_role, VALUES(location_role)),
  382. remark = IF(role_source IN ('MANUAL','RATIFIED'), remark, VALUES(remark))
  383. """,
  384. new SugarParameter("@tid", tenantId),
  385. new SugarParameter("@domain", domain ?? ""),
  386. new SugarParameter("@loc", location),
  387. new SugarParameter("@role", role),
  388. new SugarParameter("@remark", "入库镜像按库位说明分类"));
  389. }
  390. catch (Exception ex)
  391. {
  392. _logger.LogWarning(ex, "location role classify skipped loc={Location}", location);
  393. }
  394. }
  395. private async Task MirrorEmployeeAsync(
  396. long tenantId, long factoryId, string sourceCode, MdpInboundPreparedRow row, CancellationToken ct)
  397. {
  398. var emp = Str(row.Dict, "Employee") ?? Str(row.Dict, "employee");
  399. if (string.IsNullOrWhiteSpace(emp))
  400. return;
  401. if (!string.IsNullOrWhiteSpace(Str(row.Dict, "Department")))
  402. await _refs.DepartmentExistsAsync(tenantId, Str(row.Dict, "Department"));
  403. var incomingAt = ParseTime(row.SourceUpdatedAt);
  404. var existing = await _db.Queryable<AdoS0EmployeeMaster>()
  405. .Where(x => x.TenantId == tenantId && x.Employee == emp)
  406. .FirstAsync(ct);
  407. if (existing != null && string.IsNullOrWhiteSpace(existing.SourceSystem))
  408. {
  409. await InsertConflictAsync(tenantId, "MDM_EMPLOYEE_HEADCOUNT", row.BizKey, sourceCode, row.RawJson, "EmployeeMaster", ct);
  410. return;
  411. }
  412. if (existing?.SourceUpdatedAt is { } oldAt && incomingAt is { } neu && oldAt > neu)
  413. return;
  414. var now = DateTime.Now;
  415. if (existing == null)
  416. {
  417. await _db.Insertable(new AdoS0EmployeeMaster
  418. {
  419. TenantId = tenantId,
  420. FactoryRefId = factoryId,
  421. DomainCode = Str(row.Dict, "Domain") ?? "",
  422. Employee = emp,
  423. Name = Str(row.Dict, "Name"),
  424. SourceSystem = sourceCode,
  425. SourceUpdatedAt = incomingAt,
  426. CreateTime = now,
  427. UpdateTime = now
  428. }).ExecuteCommandAsync(ct);
  429. }
  430. else
  431. {
  432. existing.Name = Str(row.Dict, "Name") ?? existing.Name;
  433. existing.SourceSystem = sourceCode;
  434. existing.SourceUpdatedAt = incomingAt ?? existing.SourceUpdatedAt;
  435. existing.UpdateTime = now;
  436. await _db.Updateable(existing).ExecuteCommandAsync(ct);
  437. }
  438. try
  439. {
  440. var empSource = string.IsNullOrWhiteSpace(sourceCode) ? "API" : sourceCode;
  441. if (!await _neutralGate.AllowsAsync(tenantId, "EMPLOYEE", empSource))
  442. return;
  443. await _db.Ado.ExecuteCommandAsync(
  444. """
  445. INSERT INTO mdp_std_employee
  446. (tenant_id, factory_id, source_system, domain, employee_no, employee_name, department_code,
  447. position_code, src_position_raw, employment_status, src_employment_status_raw,
  448. source_row_id, source_biz_key, sync_batch_id, sync_time)
  449. VALUES
  450. (@tid, @factory, @src, @domain, @emp, @name, @dept,
  451. @pos, @posRaw, @status, @statusRaw,
  452. @emp, @biz, @batch, @now)
  453. ON DUPLICATE KEY UPDATE
  454. employee_name=VALUES(employee_name), department_code=VALUES(department_code),
  455. position_code=VALUES(position_code), src_position_raw=VALUES(src_position_raw),
  456. employment_status=VALUES(employment_status), src_employment_status_raw=VALUES(src_employment_status_raw),
  457. sync_batch_id=VALUES(sync_batch_id), sync_time=VALUES(sync_time)
  458. """,
  459. new SugarParameter("@tid", tenantId),
  460. new SugarParameter("@factory", factoryId),
  461. new SugarParameter("@src", empSource),
  462. new SugarParameter("@domain", Str(row.Dict, "Domain") ?? ""),
  463. new SugarParameter("@emp", emp),
  464. new SugarParameter("@name", Str(row.Dict, "Name")),
  465. new SugarParameter("@dept", Str(row.Dict, "Department")),
  466. new SugarParameter("@pos", await ResolvePositionCodeAsync(tenantId, empSource, Str(row.Dict, "Domain") ?? "", Str(row.Dict, "Position"), ct)),
  467. new SugarParameter("@posRaw", Str(row.Dict, "Position")),
  468. new SugarParameter("@status", NormalizeEmployment(Str(row.Dict, "EmploymentStatus"))),
  469. new SugarParameter("@statusRaw", Str(row.Dict, "EmploymentStatus")),
  470. new SugarParameter("@biz", $"{Str(row.Dict, "Domain")}:{emp}"),
  471. new SugarParameter("@batch", "INBOUND_EMP"),
  472. new SugarParameter("@now", now));
  473. }
  474. catch (Exception ex)
  475. {
  476. _logger.LogWarning(ex, "雇员中立层未就绪,已跳过 mdp_std_employee,镜像主数据仍已写入");
  477. }
  478. }
  479. private async Task<string> ResolvePositionCodeAsync(long tenantId, string sourceSystem, string domain, string? raw, CancellationToken ct)
  480. {
  481. if (string.IsNullOrWhiteSpace(raw))
  482. return EmployeePositionRules.Unknown;
  483. await _positions.UpsertAsync(tenantId, sourceSystem, "Position",
  484. [new EmployeePositionMapService.RawPosition { Domain = domain, Raw = raw }], ct);
  485. var codes = await _db.Ado.SqlQueryAsync<string>(
  486. """
  487. SELECT position_code FROM mdp_employee_position_map
  488. WHERE tenant_id=@tid AND source_system=@src AND domain=@domain AND src_position_raw=@raw
  489. LIMIT 1
  490. """,
  491. new { tid = tenantId, src = sourceSystem, domain, raw });
  492. return codes.FirstOrDefault() ?? EmployeePositionRules.Classify(raw);
  493. }
  494. private static string NormalizeEmployment(string? raw)
  495. {
  496. var v = (raw ?? string.Empty).Trim();
  497. var u = v.ToUpperInvariant();
  498. if (v is "在职" || u is "ACTIVE" or "ONJOB" or "ON_JOB") return "ACTIVE";
  499. if (v is "离职" or "辞退" || u is "LEFT" or "LEAVE" or "RESIGNED" or "TERMINATED" or "QUIT") return "LEFT";
  500. if (v is "停用" || u is "INACTIVE" or "DISABLED" or "SUSPENDED") return "INACTIVE";
  501. return "UNKNOWN";
  502. }
  503. private async Task InsertConflictAsync(
  504. long tenantId, string entityCode, string bizKey, string sourceCode, string raw, string table, CancellationToken ct)
  505. {
  506. var now = DateTime.Now;
  507. await _db.Insertable(new MdpInboundConflict
  508. {
  509. TenantId = tenantId,
  510. EntityCode = entityCode,
  511. BizKey = bizKey ?? string.Empty,
  512. SourceSystem = sourceCode,
  513. IncomingRaw = raw ?? string.Empty,
  514. MirrorTable = table,
  515. ConflictType = "MANUAL_ROW_EXISTS",
  516. Status = "PENDING",
  517. CreateTime = now,
  518. UpdateTime = now
  519. }).ExecuteCommandAsync(ct);
  520. }
  521. private static string Str(IDictionary<string, object?> row, string key)
  522. {
  523. foreach (var kv in row)
  524. {
  525. if (string.Equals(kv.Key, key, StringComparison.OrdinalIgnoreCase))
  526. return kv.Value?.ToString();
  527. }
  528. return null;
  529. }
  530. private static DateTime? ParseTime(string raw) =>
  531. DateTimeOffset.TryParse(raw, out var dto) ? dto.LocalDateTime : null;
  532. private static bool IsInactive(string status) =>
  533. string.Equals(status, "INACTIVE", StringComparison.OrdinalIgnoreCase)
  534. || status == "0"
  535. || string.Equals(status, "false", StringComparison.OrdinalIgnoreCase);
  536. }