using System.Security.Cryptography; using System.Text; using System.Text.Json; using Admin.NET.Plugin.AiDOP.DataPlatform.Executors; using Admin.NET.Plugin.AiDOP.DataPlatform.Sequence; using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; using Microsoft.Extensions.Logging; using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform.HotWatch; /// /// 热回读:按在途业务键对 165 做索引 seek 窄查询,变更落地到贴源层。 /// public sealed class MdpHotWatchService : ITransient { public const string SourceCode = NbrSequenceService.DefaultSourceCode; private const int MaxKeysPerBatch = 500; private readonly ISqlSugarClient _db; private readonly MdpSourceScopeFactory _scopeFactory; private readonly MdpStagingWriter _staging; private readonly ILogger _logger; public MdpHotWatchService( ISqlSugarClient db, MdpSourceScopeFactory scopeFactory, MdpStagingWriter staging, ILoggerFactory loggerFactory) { _db = db; _scopeFactory = scopeFactory; _staging = staging; _logger = loggerFactory.CreateLogger(nameof(MdpHotWatchService)); } /// Outbox 推送成功后登记在途关注。 public async Task EnrollAsync( string bizType, string bizKey, string domain, IEnumerable watchTables, long tenantId = 0, int pollIntervalSec = 5, CancellationToken ct = default) { bizType = (bizType ?? "").Trim(); bizKey = (bizKey ?? "").Trim(); domain = string.IsNullOrWhiteSpace(domain) ? "8010" : domain.Trim(); if (string.IsNullOrWhiteSpace(bizType) || string.IsNullOrWhiteSpace(bizKey)) return; if (tenantId <= 0 && string.Equals(bizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase)) tenantId = await ResolveLocalWorkOrdTenantAsync(bizKey, domain, ct); var exists = await _db.Queryable() .Where(x => x.BizType == bizType && x.BizKey == bizKey && x.Status == 0) .AnyAsync(ct); if (exists) return; var now = DateTime.Now; await _db.Insertable(new AdoMdpHotWatch { TenantId = tenantId, BizType = bizType, BizKey = bizKey, Domain = domain, WatchTables = JsonSerializer.Serialize(watchTables.ToArray()), EnrollTime = now, PollIntervalSec = pollIntervalSec, Status = 0, CreateTime = now, UpdateTime = now }).ExecuteCommandAsync(ct); _logger.LogInformation( "[MdpHotWatch] enrolled type={Type} key={Key} domain={Domain} tenant={Tenant}", bizType, bizKey, domain, tenantId); } /// /// Outbox 推送成功后按 action/idem 登记 WORK_ORDER 热关注。 /// idem 约定:wo|domain|workOrd|… / pick|domain|workOrd|… /// public Task TryEnrollFromOutboxSuccessAsync(MdpOutbox item, CancellationToken ct = default) { if (item == null) return Task.CompletedTask; if (!string.Equals(item.TargetSourceCode, SourceCode, StringComparison.OrdinalIgnoreCase)) return Task.CompletedTask; var action = (item.ActionCode ?? "").Trim().ToUpperInvariant(); var isWo = action.StartsWith("WO_MES_", StringComparison.Ordinal) || action is "PICK_WOM_UPSERT" or "PICK_WOR_STATUS"; if (!isWo) return Task.CompletedTask; if (!TryParseWoIdem(item.IdemKey, out var domain, out var workOrd)) return Task.CompletedTask; return EnrollAsync( "WORK_ORDER", workOrd, domain, new[] { "WorkOrdMaster", "WorkOrdRouting", "PeriodSequenceDet" }, item.TenantId, ct: ct); } /// 领料单写入 165 成功后登记 PICK_BILL(及关联工单)热关注。 public async Task EnrollPickBillAsync( string domain, string nbr, string? workOrd, long tenantId = 0, CancellationToken ct = default) { domain = string.IsNullOrWhiteSpace(domain) ? "8010" : domain.Trim(); nbr = (nbr ?? "").Trim(); if (string.IsNullOrWhiteSpace(nbr)) return; await EnrollAsync( "PICK_BILL", nbr, domain, new[] { "NbrMaster", "NbrDetail", "MissedPrint" }, tenantId, ct: ct); if (!string.IsNullOrWhiteSpace(workOrd)) { await EnrollAsync( "WORK_ORDER", workOrd.Trim(), domain, new[] { "WorkOrdMaster", "WorkOrdRouting", "PeriodSequenceDet" }, tenantId, ct: ct); } } /// /// 装箱标签推入 165 后登记采购单热关注,用于回读 WMS 扫码收货结果。 /// /// 挂在标签环节而非建单环节:标签生成意味着即将到货,此时开始轮询窗口最短。 /// public async Task EnrollPurOrderAsync( string domain, string purOrd, long tenantId = 0, CancellationToken ct = default) { purOrd = (purOrd ?? "").Trim(); if (string.IsNullOrWhiteSpace(purOrd)) return; domain = string.IsNullOrWhiteSpace(domain) ? "8010" : domain.Trim(); if (tenantId <= 0) tenantId = await ResolveLocalPurOrdTenantAsync(purOrd, domain, ct); await EnrollAsync( "PUR_ORDER", purOrd, domain, new[] { "PurOrdMaster", "PurOrdDetail", "MissedPrint" }, tenantId, ct: ct); } private static bool TryParseWoIdem(string? idem, out string domain, out string workOrd) { domain = "8010"; workOrd = ""; if (string.IsNullOrWhiteSpace(idem)) return false; var parts = idem.Split('|', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries); if (parts.Length < 3) return false; if (!parts[0].Equals("wo", StringComparison.OrdinalIgnoreCase) && !parts[0].Equals("pick", StringComparison.OrdinalIgnoreCase)) return false; domain = string.IsNullOrWhiteSpace(parts[1]) ? "8010" : parts[1]; workOrd = parts[2]; return !string.IsNullOrWhiteSpace(workOrd); } /// 取一批在途行并轮询 165。 public async Task<(int polled, int changed, int terminated)> PollOnceAsync( int take = 100, CancellationToken ct = default) { var due = await _db.Queryable() .Where(x => x.Status == 0) .OrderBy(x => x.LastPollTime ?? DateTime.MinValue) .Take(take) .ToListAsync(ct); if (due.Count == 0) return (0, 0, 0); MdpSource? source; ISqlSugarClient remote; try { source = await _db.Queryable() .Where(x => x.SourceCode == SourceCode && x.Status == 1) .FirstAsync(ct); if (source == null) { _logger.LogWarning("[MdpHotWatch] 源 {Source} 未启用,跳过本轮", SourceCode); return (0, 0, 0); } remote = await _scopeFactory.GetScopeAsync(SourceCode, ct); } catch (Exception ex) { _logger.LogWarning(ex, "[MdpHotWatch] 无法连接 165"); return (0, 0, 0); } var changed = 0; var terminated = 0; var now = DateTime.Now; foreach (var group in due.GroupBy(x => x.BizType)) { ct.ThrowIfCancellationRequested(); foreach (var watch in group.Take(MaxKeysPerBatch)) { var tables = ParseTables(watch.WatchTables); var sb = new StringBuilder(); foreach (var table in tables) { var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct); foreach (var row in rows) sb.Append(JsonSerializer.Serialize(row)); } var hash = Sha256(sb.ToString()); watch.LastPollTime = now; watch.UpdateTime = now; if (!string.Equals(hash, watch.LastSnapshotHash, StringComparison.Ordinal)) { watch.LastSnapshotHash = hash; changed++; var tableRows = new Dictionary>>(StringComparer.OrdinalIgnoreCase); // 变更落地:按表写 stg(执行侧字段快照) foreach (var table in tables) { var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct); tableRows[table] = rows; var entity = await ResolveEntityAsync(table, ct); if (entity == null) continue; foreach (var row in rows) { var dict = row.ToDictionary( kv => kv.Key, kv => (object?)kv.Value, StringComparer.OrdinalIgnoreCase); var rid = dict.TryGetValue("RecID", out var r) ? $"{r}" : watch.BizKey; var raw = JsonSerializer.Serialize(dict); await _staging.UpsertAsync( source, entity, table, dict, raw, rid, new MdpPullContext { TenantId = watch.TenantId, BatchId = $"hot-{now:yyyyMMddHHmmss}", FullRefresh = false }); } } // 工单执行量写回本库业务表,供看板直接读取 if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase)) { tableRows.TryGetValue("WorkOrdMaster", out var masters); tableRows.TryGetValue("WorkOrdRouting", out var routings); tableRows.TryGetValue("PeriodSequenceDet", out var periodDets); var effectiveTenantId = watch.TenantId; if (effectiveTenantId <= 0) effectiveTenantId = await ResolveLocalWorkOrdTenantAsync(watch.BizKey, watch.Domain, ct); await ApplyWorkOrderExecutionAsync( effectiveTenantId, watch.Domain, watch.BizKey, masters, routings, ct); await ApplyPeriodSequenceActualAsync( effectiveTenantId, watch.Domain, watch.BizKey, periodDets, routings, ct); } // 采购收货结果写回本库业务表,供发货单列表与齐套口径直接读取 if (string.Equals(watch.BizType, "PUR_ORDER", StringComparison.OrdinalIgnoreCase)) { tableRows.TryGetValue("PurOrdDetail", out var purDetails); tableRows.TryGetValue("MissedPrint", out var barcodes); var effectiveTenantId = watch.TenantId; if (effectiveTenantId <= 0) effectiveTenantId = await ResolveLocalPurOrdTenantAsync(watch.BizKey, watch.Domain, ct); await ApplyPurchaseReceiptAsync( effectiveTenantId, watch.Domain, watch.BizKey, purDetails, barcodes, ct); } if (await ShouldTerminateAsync(remote, watch, ct)) { watch.Status = 1; watch.TerminateTime = now; watch.TerminateReason = "auto"; terminated++; } } await _db.Updateable(watch) .UpdateColumns(x => new { x.LastPollTime, x.LastSnapshotHash, x.Status, x.TerminateTime, x.TerminateReason, x.UpdateTime }) .ExecuteCommandAsync(ct); } } return (due.Count, changed, terminated); } private async Task ResolveEntityAsync(string table, CancellationToken ct) { return await _db.Queryable() .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1) .FirstAsync(ct); } private static async Task>> QueryByBizKeyAsync( ISqlSugarClient remote, string table, string domain, string bizType, string bizKey, CancellationToken ct) { // 表白名单(防注入) if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$")) throw new InvalidOperationException($"非法表名:{table}"); string sql; SugarParameter[] pars; switch (table.ToUpperInvariant()) { case "NBRMASTER": sql = "SELECT * FROM NbrMaster WHERE Domain=@d AND Nbr=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "NBRDETAIL": sql = "SELECT * FROM NbrDetail WHERE Domain=@d AND Nbr=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "WORKORDMASTER": sql = "SELECT * FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "WORKORDROUTING": sql = "SELECT * FROM WorkOrdRouting WHERE Domain=@d AND WorkOrd=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "PERIODSEQUENCEDET": // 工序间衔接:按工单取全部行(含 MES 报工产生的 Period=0 影子行),由调用方按工序合并 sql = "SELECT * FROM PeriodSequenceDet WHERE Domain=@d AND WorkOrds=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "PURORDDETAIL": sql = "SELECT * FROM PurOrdDetail WHERE Domain=@d AND PurOrd=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "PURORDMASTER": sql = "SELECT * FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "LINESTATUSDET": sql = "SELECT * FROM LineStatusDet WHERE Domain=@d AND Line=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; case "MOBILETASK": // 堆表:按 TaskID sql = "SELECT * FROM MobileTask WHERE TaskID=@k"; pars = new[] { new SugarParameter("@k", bizKey) }; break; default: // MissedPrint 等:降级按 Domain + OrdNbr(可能扫表,见 WP8 E1) if (string.Equals(table, "MissedPrint", StringComparison.OrdinalIgnoreCase)) { sql = "SELECT TOP 200 * FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k"; pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) }; break; } return new List>(); } var dt = await remote.Ado.GetDataTableAsync(sql, pars); var list = new List>(); foreach (System.Data.DataRow row in dt.Rows) { var dict = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (System.Data.DataColumn col in dt.Columns) dict[col.ColumnName] = row[col] == DBNull.Value ? null! : row[col]; list.Add(dict); } return list; } /// /// 把 165 工单执行字段回写本库(MES 权威列:Status / QtyCompleted / QtyComplete / QtyReject)。 /// private async Task ApplyWorkOrderExecutionAsync( long tenantId, string domain, string workOrd, List>? masters, List>? routings, CancellationToken ct) { var now = DateTime.Now; var updatedRouting = 0; if (routings != null) { foreach (var row in routings) { if (!TryGetInt(row, "OP", out var op) && !TryGetInt(row, "Op", out op)) continue; var qtyComplete = GetDecimal(row, "QtyComplete"); var qtyReject = GetDecimal(row, "QtyReject"); var status = Trunc(GetString(row, "Status"), 1); updatedRouting += await _db.Ado.ExecuteCommandAsync( """ UPDATE WorkOrdRouting SET QtyComplete = @QtyComplete, QtyReject = @QtyReject, Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END, UpdateUser = 'MDP_HOT', UpdateTime = @Now WHERE WorkOrd = @WorkOrd AND OP = @Op AND IFNULL(Domain, '') = @Domain AND IFNULL(tenant_id, 0) = @TenantId """, new SugarParameter("@QtyComplete", qtyComplete), new SugarParameter("@QtyReject", qtyReject), new SugarParameter("@Status", status ?? ""), new SugarParameter("@Now", now), new SugarParameter("@WorkOrd", workOrd), new SugarParameter("@Op", op), new SugarParameter("@Domain", domain), new SugarParameter("@TenantId", tenantId)); } } if (masters is { Count: > 0 }) { var m = masters[0]; var qtyCompleted = GetDecimal(m, "QtyCompleted"); var status = Trunc(GetString(m, "Status"), 8); await _db.Ado.ExecuteCommandAsync( """ UPDATE WorkOrdMaster SET QtyCompleted = @QtyCompleted, Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END, UpdateUser = 'MDP_HOT', UpdateTime = @Now WHERE WorkOrd = @WorkOrd AND IFNULL(Domain, '') = @Domain AND IFNULL(tenant_id, 0) = @TenantId """, new SugarParameter("@QtyCompleted", qtyCompleted), new SugarParameter("@Status", status ?? ""), new SugarParameter("@Now", now), new SugarParameter("@WorkOrd", workOrd), new SugarParameter("@Domain", domain), new SugarParameter("@TenantId", tenantId)); } _logger.LogInformation( "[MdpHotWatch] applied WORK_ORDER execution wo={WorkOrd} domain={Domain} routingRows={Rows}", workOrd, domain, updatedRouting); } /// /// D-M01:把 165 的工序实绩按 (Domain, WorkOrds, Op) 合并后落到 ado_psd_op_actual。 /// /// 为什么按工序合并:MES APP 报工会新增 Period=0/IsActive=0 的影子行,实绩可能写在影子行、 /// 也可能写在 Period=1 计划行,且同一工序可能因重排跨多个 Line。只有工序级求和才与 /// WorkOrdRouting.QtyComplete(MES 自己的工序累计完成数)对齐。实测 Op501:3+5=8=QtyComplete。 /// /// 为什么不写 PeriodSequenceDet:本库 PSD 行会被排产整表重建(旧行置 IsActive=0), /// 实绩写进去会在下次全量重排时丢失。 /// /// 单位保持 165 原始口径:ActualTime 秒、SetupTime 小时(列名自带单位)。 /// private async Task ApplyPeriodSequenceActualAsync( long tenantId, string domain, string workOrd, List>? periodDets, List>? routings, CancellationToken ct) { if (periodDets is not { Count: > 0 }) return; // WorkOrdRouting.QtyComplete 作为交叉校验值(仅记录,不用于覆盖) var routingQty = new Dictionary(); if (routings != null) { foreach (var r in routings) { if (!TryGetInt(r, "OP", out var rop) && !TryGetInt(r, "Op", out rop)) continue; routingQty[rop] = GetDecimal(r, "QtyComplete"); } } var grouped = new Dictionary(); foreach (var row in periodDets) { if (!TryGetInt(row, "Op", out var op) && !TryGetInt(row, "OP", out op)) continue; if (!grouped.TryGetValue(op, out var acc)) { acc = new PsdOpActual { Op = op }; grouped[op] = acc; } acc.CompQty += GetDecimal(row, "CompQty"); acc.RejectQty += GetDecimal(row, "RejectQty"); acc.SetupTimeHour += GetDecimal(row, "SetupTime"); acc.ActualTimeSec += GetDecimal(row, "ActualTime"); acc.SrcRowCount++; acc.ItemNum ??= GetString(row, "ItemNum"); if (row.TryGetValue("UpdateTime", out var upd) && upd is DateTime dt && (acc.LastSrcUpdate == null || dt > acc.LastSrcUpdate)) acc.LastSrcUpdate = dt; } var now = DateTime.Now; var mismatches = new List(); foreach (var acc in grouped.Values) { var rq = routingQty.TryGetValue(acc.Op, out var q) ? (decimal?)q : null; if (rq.HasValue && rq.Value != acc.CompQty) mismatches.Add($"Op{acc.Op}: psdSum={acc.CompQty} routing={rq.Value}"); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO ado_psd_op_actual ( tenant_id, domain, work_ord, op, item_num, comp_qty, reject_qty, setup_time_hour, actual_time_sec, src_row_count, routing_qty, last_src_update, sync_time ) VALUES ( @TenantId, @Domain, @WorkOrd, @Op, @ItemNum, @CompQty, @RejectQty, @SetupTimeHour, @ActualTimeSec, @SrcRowCount, @RoutingQty, @LastSrcUpdate, @Now ) ON DUPLICATE KEY UPDATE item_num = VALUES(item_num), comp_qty = VALUES(comp_qty), reject_qty = VALUES(reject_qty), setup_time_hour = VALUES(setup_time_hour), actual_time_sec = VALUES(actual_time_sec), src_row_count = VALUES(src_row_count), routing_qty = VALUES(routing_qty), last_src_update = VALUES(last_src_update), sync_time = VALUES(sync_time) """, new SugarParameter("@TenantId", tenantId), new SugarParameter("@Domain", domain ?? ""), new SugarParameter("@WorkOrd", workOrd), new SugarParameter("@Op", acc.Op), new SugarParameter("@ItemNum", acc.ItemNum ?? (object)DBNull.Value), new SugarParameter("@CompQty", acc.CompQty), new SugarParameter("@RejectQty", acc.RejectQty), new SugarParameter("@SetupTimeHour", acc.SetupTimeHour), new SugarParameter("@ActualTimeSec", acc.ActualTimeSec), new SugarParameter("@SrcRowCount", acc.SrcRowCount), new SugarParameter("@RoutingQty", rq.HasValue ? rq.Value : (object)DBNull.Value), new SugarParameter("@LastSrcUpdate", acc.LastSrcUpdate ?? (object)DBNull.Value), new SugarParameter("@Now", now)); } if (mismatches.Count > 0) { // 不阻断落库:差异说明 165 侧口径需人工确认,先留证 _logger.LogWarning( "[MdpHotWatch] PSD 实绩与 WorkOrdRouting 不一致 wo={WorkOrd} domain={Domain} {Detail}", workOrd, domain, string.Join("; ", mismatches)); } _logger.LogInformation( "[MdpHotWatch] applied PSD actual wo={WorkOrd} domain={Domain} ops={Ops} srcRows={Rows}", workOrd, domain, grouped.Count, periodDets.Count); } /// /// 把 165 的采购收货结果回读本库:明细收货数、箱码状态,并据箱码推进送货单状态。 /// /// 口径(2026-08-11 定):WMS 扫码收货只入待检仓(InvTransHist 落 rct-po-ins / Loc=1000), /// 不等于合格入库,所以这里只推进 scm_shd/scm_shdzb.shzt,不写 rksl——合格入库数留给检验环节。 /// 165 的 scm_shdzb.rksl / shzt 在收货时并不更新,送货单进度只能由箱码状态反推。 /// private async Task ApplyPurchaseReceiptAsync( long tenantId, string domain, string purOrd, List>? purDetails, List>? barcodes, CancellationToken ct) { var now = DateTime.Now; var updatedLines = 0; var updatedBarcodes = 0; if (purDetails != null) { foreach (var row in purDetails) { if (!TryGetInt(row, "Line", out var line)) continue; // 自建单本库不填 Domain(165/MES 概念),所以空 Domain 也算命中;采购单号本库唯一 updatedLines += await _db.Ado.ExecuteCommandAsync( """ UPDATE PurOrdDetail SET RctQty = @RctQty, ReceiptQty = @ReceiptQty, QtyReturned = @QtyReturned, UpdateUser = 'MDP_HOT', UpdateTime = @Now WHERE PurOrd = @PurOrd AND Line = @Line AND IFNULL(Domain, '') IN ('', @Domain) AND (@TenantId = 0 OR IFNULL(tenant_id, 0) = @TenantId) """, new SugarParameter("@RctQty", GetDecimal(row, "RctQty")), new SugarParameter("@ReceiptQty", GetDecimal(row, "ReceiptQty")), new SugarParameter("@QtyReturned", GetDecimal(row, "QtyReturned")), new SugarParameter("@Now", now), new SugarParameter("@PurOrd", purOrd), new SugarParameter("@Line", line), new SugarParameter("@Domain", domain), new SugarParameter("@TenantId", tenantId)); } } // 本库 MissedPrint 无 tenant_id,按 Domain + BarCode 定位;作废行 BarCode 带 RecID 前缀不会误匹配 var shippers = new HashSet(StringComparer.OrdinalIgnoreCase); if (barcodes != null) { foreach (var row in barcodes) { var barCode = GetString(row, "BarCode")?.Trim(); if (string.IsNullOrWhiteSpace(barCode)) continue; updatedBarcodes += await _db.Ado.ExecuteCommandAsync( """ UPDATE MissedPrint SET Status = CASE WHEN IFNULL(@Status, '') = '' THEN Status ELSE @Status END, RctNbr = @RctNbr, Location = @Location, Shelf = @Shelf, InvStatus = @InvStatus, UpdateUser = 'MDP_HOT', UpdateTime = @Now WHERE BarCode = @BarCode AND IFNULL(Domain, '') IN ('', @Domain) """, new SugarParameter("@Status", Trunc(GetString(row, "Status"), 4) ?? ""), new SugarParameter("@RctNbr", Trunc(GetString(row, "RctNbr"), 48) ?? ""), new SugarParameter("@Location", Trunc(GetString(row, "Location"), 24) ?? ""), new SugarParameter("@Shelf", Trunc(GetString(row, "Shelf"), 40) ?? ""), new SugarParameter("@InvStatus", Trunc(GetString(row, "InvStatus"), 24) ?? ""), new SugarParameter("@Now", now), new SugarParameter("@BarCode", barCode), new SugarParameter("@Domain", domain)); var shipper = GetString(row, "ShipperNbr")?.Trim(); if (!string.IsNullOrWhiteSpace(shipper)) shippers.Add(shipper!); } } foreach (var shipper in shippers) await ApplyShipmentStatusAsync(tenantId, domain, shipper, ct); _logger.LogInformation( "[MdpHotWatch] applied PUR_ORDER receipt po={PurOrd} domain={Domain} lines={Lines} barcodes={Barcodes} shipments={Shipments}", purOrd, domain, updatedLines, updatedBarcodes, shippers.Count); } /// /// 按箱码待收数推进送货单状态:全部待收=待收,全部离开待收=完成,其余=收货中。 /// /// 完成判据取箱码而非数量(决策:少收/超收不影响状态),列表读的是 scm_shd.shzt。 /// private async Task ApplyShipmentStatusAsync( long tenantId, string domain, string shddh, CancellationToken ct) { var stats = await _db.Ado.SqlQueryAsync( """ SELECT COUNT(*) AS Total, SUM(CASE WHEN IFNULL(Status, '') = 'U' THEN 1 ELSE 0 END) AS Pending FROM MissedPrint WHERE ShipperNbr = @Shddh AND IFNULL(Domain, '') = @Domain AND IFNULL(PurOrd, '') NOT LIKE '作废%' """, new SugarParameter("@Shddh", shddh), new SugarParameter("@Domain", domain)); var stat = stats?.FirstOrDefault(); if (stat == null || stat.Total <= 0) return; var shzt = stat.Pending >= stat.Total ? "待收" : (stat.Pending == 0 ? "完成" : "收货中"); // 先定位表头再按主键更新:scm_shd.id 与 scm_shdzb.glid 排序规则不同(0900_ai_ci vs unicode_ci), // 直接 JOIN 会报 Illegal mix of collations,改由参数传 glid 规避 var heads = await _db.Ado.SqlQueryAsync( """ SELECT id AS Id FROM scm_shd WHERE shddh = @Shddh AND (@TenantId = 0 OR IFNULL(tenant_id, 0) = @TenantId) ORDER BY id LIMIT 1 """, new SugarParameter("@Shddh", shddh), new SugarParameter("@TenantId", tenantId)); var head = heads?.FirstOrDefault(); if (head == null) { _logger.LogWarning( "[MdpHotWatch] 送货单在本库缺失,跳过状态推进 shddh={Shddh} tenant={Tenant}", shddh, tenantId); return; } await _db.Ado.ExecuteCommandAsync( "UPDATE scm_shd SET shzt = @Shzt WHERE id = @Id", new SugarParameter("@Shzt", shzt), new SugarParameter("@Id", head.Id)); await _db.Ado.ExecuteCommandAsync( "UPDATE scm_shdzb SET shzt = @Shzt WHERE glid = @Glid", new SugarParameter("@Shzt", shzt), new SugarParameter("@Glid", head.Id.ToString())); } private sealed class ShipmentBarcodeStat { public int Total { get; set; } public int Pending { get; set; } } private sealed class ShipmentHeadRow { public long Id { get; set; } } private async Task ResolveLocalPurOrdTenantAsync(string purOrd, string domain, CancellationToken ct) { var tid = await _db.Ado.SqlQuerySingleAsync( """ SELECT IFNULL(tenant_id, 0) FROM PurOrdDetail WHERE PurOrd = @PurOrd AND IFNULL(Domain, '') IN ('', @Domain) ORDER BY Line LIMIT 1 """, new SugarParameter("@PurOrd", purOrd), new SugarParameter("@Domain", domain)); return tid ?? 0; } private sealed class PsdOpActual { public int Op { get; set; } public string? ItemNum { get; set; } public decimal CompQty { get; set; } public decimal RejectQty { get; set; } public decimal SetupTimeHour { get; set; } public decimal ActualTimeSec { get; set; } public int SrcRowCount { get; set; } public DateTime? LastSrcUpdate { get; set; } } private async Task ResolveLocalWorkOrdTenantAsync(string workOrd, string domain, CancellationToken ct) { var tid = await _db.Ado.SqlQuerySingleAsync( """ SELECT IFNULL(tenant_id, 0) FROM WorkOrdMaster WHERE WorkOrd = @WorkOrd AND IFNULL(Domain, '') = @Domain ORDER BY RecID DESC LIMIT 1 """, new SugarParameter("@WorkOrd", workOrd), new SugarParameter("@Domain", domain)); return tid ?? 0; } private static bool TryGetInt(Dictionary row, string key, out int value) { value = 0; if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return false; try { value = Convert.ToInt32(raw); return true; } catch { return false; } } private static decimal GetDecimal(Dictionary row, string key) { if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return 0m; try { return Convert.ToDecimal(raw); } catch { return 0m; } } private static string? GetString(Dictionary row, string key) { if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return null; return Convert.ToString(raw); } private static string? Trunc(string? s, int max) => string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]); private static async Task ShouldTerminateAsync( ISqlSugarClient remote, AdoMdpHotWatch watch, CancellationToken ct) { if (string.Equals(watch.BizType, "PICK_BILL", StringComparison.OrdinalIgnoreCase)) { var status = await remote.Ado.GetScalarAsync( "SELECT TOP 1 Status FROM NbrMaster WHERE Domain=@d AND Nbr=@k", new SugarParameter("@d", watch.Domain), new SugarParameter("@k", watch.BizKey)); var s = status?.ToString()?.Trim() ?? ""; // 终态取值待 WP4 固化;临时:C/Complete/关闭 等常见值 return s is "C" or "Complete" or "关闭" or "Y"; } if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase)) { var status = await remote.Ado.GetScalarAsync( "SELECT TOP 1 Status FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k", new SugarParameter("@d", watch.Domain), new SugarParameter("@k", watch.BizKey)); return string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase); } if (string.Equals(watch.BizType, "PUR_ORDER", StringComparison.OrdinalIgnoreCase)) { var status = await remote.Ado.GetScalarAsync( "SELECT TOP 1 Status FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k", new SugarParameter("@d", watch.Domain), new SugarParameter("@k", watch.BizKey)); if (string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase)) return true; // 箱码全部离开待收即收货结束;标签尚未推达(total=0)时不能终结 var total = Convert.ToInt32(await remote.Ado.GetScalarAsync( "SELECT COUNT(*) FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k", new SugarParameter("@d", watch.Domain), new SugarParameter("@k", watch.BizKey)) ?? 0); if (total > 0) { var pending = Convert.ToInt32(await remote.Ado.GetScalarAsync( "SELECT COUNT(*) FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k AND IsNull(Status,'')='U'", new SugarParameter("@d", watch.Domain), new SugarParameter("@k", watch.BizKey)) ?? 0); if (pending == 0) return true; } } // 兜底:超过 7 天强制终结 return watch.EnrollTime < DateTime.Now.AddDays(-7); } private static List ParseTables(string? json) { if (string.IsNullOrWhiteSpace(json)) return new List(); try { return JsonSerializer.Deserialize>(json!) ?? new List(); } catch { return new List(); } } private static string Sha256(string s) { var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s)); return Convert.ToHexString(bytes); } }