using System.Text.Json; using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; using Admin.NET.Plugin.AiDOP.DataPlatform.HotWatch; using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Wms; /// /// UAT-INT-001 · Mes 模式:工单物料变更后入队 Outbox,推送时在 165 本地事务内同步一张有效 SM。 /// Local 模式不入队,由调用方继续写本库 Nbr*。 /// public sealed class PickBillNbrSyncService : ITransient { public const string ActionCode = "PICK_NBR_SYNC"; public const string TargetSource = CreatePickBillService.TargetSource; private readonly ISqlSugarClient _db; private readonly MdpSourceScopeFactory _scopeFactory; private readonly MdpOutboxEnqueueService _enqueue; private readonly MdpHotWatchService _hotWatch; private readonly AidopPickBillOptions _opt; private readonly ILogger _logger; public PickBillNbrSyncService( ISqlSugarClient db, MdpSourceScopeFactory scopeFactory, MdpOutboxEnqueueService enqueue, MdpHotWatchService hotWatch, IOptions opt, ILogger logger) { _db = db; _scopeFactory = scopeFactory; _enqueue = enqueue; _hotWatch = hotWatch; _opt = opt.Value; _logger = logger; } public bool IsMes => _opt.IsMes; /// /// Mes 且工单已下达/投产/暂停时入队(同键刷新 payload)。Local 或未下达返回 false。 /// public async Task TryEnqueueAfterPlanChangeAsync( long tenantId, string workOrd, string? domain, string account, CancellationToken ct = default) { if (!_opt.IsMes) return false; workOrd = (workOrd ?? "").Trim(); if (string.IsNullOrWhiteSpace(workOrd)) return false; var status = await _db.Ado.GetStringAsync( """ SELECT IFNULL(LOWER(TRIM(Status)), '') FROM WorkOrdMaster WHERE tenant_id = @TenantId AND WorkOrd = @WorkOrd LIMIT 1 """, new SugarParameter("@TenantId", tenantId), new SugarParameter("@WorkOrd", workOrd)); if (status is not ("r" or "w" or "s")) return false; var mesDomain = ResolveMesDomain(domain); var idem = $"pick|{mesDomain}|{workOrd}|sm-sync"; if (idem.Length > 200) idem = idem[..200]; var item = new MdpOutbox { TenantId = tenantId, TargetSourceCode = TargetSource, ActionCode = ActionCode, IdemKey = idem, PayloadJson = JsonSerializer.Serialize(new PickNbrSyncPayload { Domain = mesDomain, WorkOrd = workOrd, Account = string.IsNullOrWhiteSpace(account) ? "aidop" : account.Trim() }) }; var ok = await _enqueue.TryEnqueueOrRefreshAsync(item, ct); _logger.LogInformation( "[PICK_NBR_SYNC] enqueue tenant={Tenant} wo={Wo} domain={Domain} refreshed={Ok}", tenantId, workOrd, mesDomain, ok); return true; } public async Task ApplyFromOutboxAsync(MdpOutbox item, CancellationToken ct) { PickNbrSyncPayload payload; try { payload = JsonSerializer.Deserialize(item.PayloadJson ?? "") ?? throw new InvalidOperationException("payload 为空"); } catch (Exception ex) { return MdpPushResult.Fail($"PAYLOAD_INVALID: {ex.Message}"); } var domain = ResolveMesDomain(payload.Domain); var workOrd = (payload.WorkOrd ?? "").Trim(); var account = Trunc(payload.Account, 24); if (string.IsNullOrWhiteSpace(workOrd)) return MdpPushResult.Fail("PAYLOAD_INVALID: workOrd 为空"); var plan = await LoadPlanLinesAsync(item.TenantId, workOrd); ISqlSugarClient ss; try { ss = await _scopeFactory.GetScopeAsync(TargetSource, ct); } catch (Exception ex) { return MdpPushResult.Fail($"MES_UNREACHABLE: {ex.Message}"); } try { var result = await ApplyOn165Async(ss, item.TenantId, domain, workOrd, account, plan, ct); return result; } catch (Exception ex) { _logger.LogWarning(ex, "[PICK_NBR_SYNC] apply failed wo={Wo} domain={Domain}", workOrd, domain); return MdpPushResult.Fail($"MES_WRITE_FAILED: {ex.Message}"); } } private async Task ApplyOn165Async( ISqlSugarClient ss, long tenantId, string domain, string workOrd, string account, List plan, CancellationToken ct) { var sm = (await ss.Ado.SqlQueryAsync( """ SELECT TOP 1 RecID, Nbr, Domain FROM NbrMaster WHERE Domain = @d AND WorkOrd = @w AND Type = 'SM' AND ISNULL(TransType, '') = '' AND ISNULL(IsActive, 1) = 1 ORDER BY RecID DESC """, new SugarParameter("@d", domain), new SugarParameter("@w", workOrd))).FirstOrDefault(); if (sm == null || string.IsNullOrWhiteSpace(sm.Nbr)) { return MdpPushResult.Skip(JsonSerializer.Serialize(new { idempotent = true, reason = "no_sm", workOrd, domain })); } var existing = await ss.Ado.SqlQueryAsync( """ SELECT RecID, ItemNum, QtyOrd, QtyRec, CurrQtyOpened, Line, Status FROM NbrDetail WHERE NbrRecID = @rid AND Type = 'SM' AND ISNULL(IsActive, 1) = 1 """, new SugarParameter("@rid", sm.RecID)); var planMap = plan .Where(p => !string.IsNullOrWhiteSpace(p.ItemNum)) .ToDictionary(p => p.ItemNum.Trim(), p => p, StringComparer.OrdinalIgnoreCase); var existingMap = existing .ToDictionary(x => (x.ItemNum ?? "").Trim(), x => x, StringComparer.OrdinalIgnoreCase); var now = DateTime.Now; var updated = 0; var closed = 0; var inserted = 0; await ss.Ado.BeginTranAsync(); try { foreach (var line in existing) { var key = (line.ItemNum ?? "").Trim(); var still = planMap.TryGetValue(key, out var p); var newQty = still ? p!.QtyRequired : 0m; var action = PickBillNbrSyncRules.DecideExisting(line.QtyRec, newQty, still); if (action == PickBillNbrSyncRules.LineAction.Close) { await ss.Ado.ExecuteCommandAsync( """ UPDATE NbrDetail SET Status = 'C', UpdateUser = @u, UpdateTime = @t WHERE RecID = @id """, new SugarParameter("@u", account), new SugarParameter("@t", now), new SugarParameter("@id", line.RecID)); closed++; continue; } if (action == PickBillNbrSyncRules.LineAction.UpdateQty && (line.QtyOrd != newQty || !string.Equals(line.Status ?? "", "", StringComparison.Ordinal))) { await ss.Ado.ExecuteCommandAsync( """ UPDATE NbrDetail SET QtyOrd = @q, CurrQtyOpened = @q, UM = @um, ItemName = @name, Status = '', UpdateUser = @u, UpdateTime = @t WHERE RecID = @id """, new SugarParameter("@q", newQty), new SugarParameter("@um", Trunc(p!.Unit, 8)), new SugarParameter("@name", Trunc(p.ItemName, 1000)), new SugarParameter("@u", account), new SugarParameter("@t", now), new SugarParameter("@id", line.RecID)); updated++; } } short nextLine = existing.Count > 0 ? (short)(existing.Max(x => x.Line) + 1) : (short)1; foreach (var p in plan) { var key = p.ItemNum.Trim(); if (existingMap.ContainsKey(key)) continue; await ss.Ado.ExecuteCommandAsync( """ INSERT INTO NbrDetail (Domain, Type, Nbr, Line, ItemNum, Dimension1, Dimension2, LocationFrom, LocationTo, QtyFrom, QtyTo, UM, [Print], Status, LotSerial, WorkOrd, QtyOrd, QtyRec, Address, BusinessID, CreateUser, UpdateUser, CreateTime, UpdateTime, IsActive, IsConfirm, QtyCache, CurrQtyOpened, IsChanged, NbrRecID, OrdNbr, ItemName, ERPfld1, ERPfld2, OrdLine, IsGP12Demand, IsGP12Checked, Material, SeqID) VALUES (@Domain, 'SM', @Nbr, @Line, @ItemNum, '', '', @LocationFrom, @LocationTo, 0, 0, @UM, 0, '', @LotSerial, @WorkOrd, @QtyOrd, 0, '', 0, @User, @User, @Now, @Now, 1, 0, 0, @QtyOrd, 1, @NbrRecID, '', @ItemName, '', '', 0, 0, 0, 0, 0) """, new SugarParameter("@Domain", domain), new SugarParameter("@Nbr", sm.Nbr), new SugarParameter("@Line", nextLine++), new SugarParameter("@ItemNum", Trunc(p.ItemNum, 24)), new SugarParameter("@LocationFrom", Trunc(p.LocationFrom, 8)), new SugarParameter("@LocationTo", Trunc(p.LocationTo, 8)), new SugarParameter("@UM", Trunc(p.Unit, 8)), new SugarParameter("@LotSerial", p.LotSerial ?? ""), new SugarParameter("@WorkOrd", workOrd), new SugarParameter("@QtyOrd", p.QtyRequired), new SugarParameter("@User", account), new SugarParameter("@Now", now), new SugarParameter("@NbrRecID", sm.RecID), new SugarParameter("@ItemName", Trunc(p.ItemName, 1000))); inserted++; } await ss.Ado.ExecuteCommandAsync( """ UPDATE NbrMaster SET QtyOrd = ( SELECT ISNULL(SUM(QtyOrd), 0) FROM NbrDetail WHERE NbrRecID = @rid AND Type = 'SM' AND ISNULL(IsActive, 1) = 1 AND ISNULL(Status, '') <> 'C' ), UpdateUser = @u, UpdateTime = @t WHERE RecID = @rid """, new SugarParameter("@rid", sm.RecID), new SugarParameter("@u", account), new SugarParameter("@t", now)); await ss.Ado.CommitTranAsync(); } catch { await ss.Ado.RollbackTranAsync(); throw; } try { await _hotWatch.EnrollPickBillAsync(domain, sm.Nbr!, workOrd, tenantId, ct); } catch (Exception ex) { _logger.LogWarning(ex, "[PICK_NBR_SYNC] hot-watch enroll failed nbr={Nbr} wo={Wo}", sm.Nbr, workOrd); } var response = JsonSerializer.Serialize(new { nbr = sm.Nbr, workOrd, domain, updated, closed, inserted }); _logger.LogInformation( "[PICK_NBR_SYNC] applied nbr={Nbr} wo={Wo} upd={U} close={C} ins={I}", sm.Nbr, workOrd, updated, closed, inserted); return MdpPushResult.Ok(updated + closed + inserted, response); } private async Task> LoadPlanLinesAsync(long tenantId, string workOrd) { return await _db.Ado.SqlQueryAsync( """ SELECT TRIM(d.ItemNum) AS ItemNum, SUM(d.QtyRequired) AS QtyRequired, MAX(IFNULL(d.UM, im.Um)) AS Unit, MAX(im.Descr) AS ItemName, MAX(IFNULL(NULLIF(TRIM(d.Location), ''), IFNULL(im.Location, ''))) AS LocationFrom, MAX(IFNULL(m.Location, '')) AS LocationTo, MAX(IFNULL(d.LotSerial, '')) AS LotSerial FROM WorkOrdDetail d LEFT JOIN WorkOrdMaster m ON m.WorkOrd = d.WorkOrd AND m.tenant_id = d.tenant_id LEFT JOIN ItemMaster im ON im.ItemNum = d.ItemNum AND im.tenant_id = d.tenant_id WHERE d.tenant_id = @TenantId AND d.WorkOrd = @WorkOrd AND IFNULL(d.IsActive, 0) = 1 GROUP BY TRIM(d.ItemNum) HAVING SUM(d.QtyRequired) > 0 ORDER BY TRIM(d.ItemNum) """, new SugarParameter("@TenantId", tenantId), new SugarParameter("@WorkOrd", workOrd)); } /// 165 Domain 最长 8 位;本库若把 tenant_id 写进 Domain 则回落到 8010。 public static string ResolveMesDomain(string? domain) { var d = (domain ?? "").Trim(); if (string.IsNullOrWhiteSpace(d) || d.Length > 8) return "8010"; return d; } private static string Trunc(string? s, int max) { s ??= ""; return s.Length <= max ? s : s[..max]; } private sealed class PickNbrSyncPayload { public string Domain { get; set; } = "8010"; public string WorkOrd { get; set; } = ""; public string Account { get; set; } = "aidop"; } private sealed class PlanLine { public string ItemNum { get; set; } = ""; public decimal QtyRequired { get; set; } public string? Unit { get; set; } public string? ItemName { get; set; } public string? LocationFrom { get; set; } public string? LocationTo { get; set; } public string? LotSerial { get; set; } } private sealed class SmHead { public int RecID { get; set; } public string? Nbr { get; set; } public string? Domain { get; set; } } private sealed class SmLine { public int RecID { get; set; } public string? ItemNum { get; set; } public decimal QtyOrd { get; set; } public decimal QtyRec { get; set; } public decimal CurrQtyOpened { get; set; } public short Line { get; set; } public string? Status { get; set; } } }