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