MdpHotWatchService.cs 38 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888
  1. using System.Security.Cryptography;
  2. using System.Text;
  3. using System.Text.Json;
  4. using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  5. using Admin.NET.Plugin.AiDOP.DataPlatform.Sequence;
  6. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  7. using Microsoft.Extensions.Logging;
  8. using SqlSugar;
  9. namespace Admin.NET.Plugin.AiDOP.DataPlatform.HotWatch;
  10. /// <summary>
  11. /// 热回读:按在途业务键对 165 做索引 seek 窄查询,变更落地到贴源层。
  12. /// </summary>
  13. public sealed class MdpHotWatchService : ITransient
  14. {
  15. public const string SourceCode = NbrSequenceService.DefaultSourceCode;
  16. private const int MaxKeysPerBatch = 500;
  17. private readonly ISqlSugarClient _db;
  18. private readonly MdpSourceScopeFactory _scopeFactory;
  19. private readonly MdpStagingWriter _staging;
  20. private readonly ILogger _logger;
  21. public MdpHotWatchService(
  22. ISqlSugarClient db,
  23. MdpSourceScopeFactory scopeFactory,
  24. MdpStagingWriter staging,
  25. ILoggerFactory loggerFactory)
  26. {
  27. _db = db;
  28. _scopeFactory = scopeFactory;
  29. _staging = staging;
  30. _logger = loggerFactory.CreateLogger(nameof(MdpHotWatchService));
  31. }
  32. /// <summary>Outbox 推送成功后登记在途关注。</summary>
  33. public async Task EnrollAsync(
  34. string bizType, string bizKey, string domain, IEnumerable<string> watchTables,
  35. long tenantId = 0, int pollIntervalSec = 5, CancellationToken ct = default)
  36. {
  37. bizType = (bizType ?? "").Trim();
  38. bizKey = (bizKey ?? "").Trim();
  39. domain = string.IsNullOrWhiteSpace(domain) ? "8010" : domain.Trim();
  40. if (string.IsNullOrWhiteSpace(bizType) || string.IsNullOrWhiteSpace(bizKey))
  41. return;
  42. if (tenantId <= 0 && string.Equals(bizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  43. tenantId = await ResolveLocalWorkOrdTenantAsync(bizKey, domain, ct);
  44. var exists = await _db.Queryable<AdoMdpHotWatch>()
  45. .Where(x => x.BizType == bizType && x.BizKey == bizKey && x.Status == 0)
  46. .AnyAsync(ct);
  47. if (exists) return;
  48. var now = DateTime.Now;
  49. await _db.Insertable(new AdoMdpHotWatch
  50. {
  51. TenantId = tenantId,
  52. BizType = bizType,
  53. BizKey = bizKey,
  54. Domain = domain,
  55. WatchTables = JsonSerializer.Serialize(watchTables.ToArray()),
  56. EnrollTime = now,
  57. PollIntervalSec = pollIntervalSec,
  58. Status = 0,
  59. CreateTime = now,
  60. UpdateTime = now
  61. }).ExecuteCommandAsync(ct);
  62. _logger.LogInformation(
  63. "[MdpHotWatch] enrolled type={Type} key={Key} domain={Domain} tenant={Tenant}",
  64. bizType, bizKey, domain, tenantId);
  65. }
  66. /// <summary>
  67. /// Outbox 推送成功后按 action/idem 登记 WORK_ORDER 热关注。
  68. /// idem 约定:<c>wo|domain|workOrd|…</c> / <c>pick|domain|workOrd|…</c>
  69. /// </summary>
  70. public Task TryEnrollFromOutboxSuccessAsync(MdpOutbox item, CancellationToken ct = default)
  71. {
  72. if (item == null) return Task.CompletedTask;
  73. if (!string.Equals(item.TargetSourceCode, SourceCode, StringComparison.OrdinalIgnoreCase))
  74. return Task.CompletedTask;
  75. var action = (item.ActionCode ?? "").Trim().ToUpperInvariant();
  76. var isWo =
  77. action.StartsWith("WO_MES_", StringComparison.Ordinal)
  78. || action is "PICK_WOM_UPSERT" or "PICK_WOR_STATUS";
  79. if (!isWo) return Task.CompletedTask;
  80. if (!TryParseWoIdem(item.IdemKey, out var domain, out var workOrd))
  81. return Task.CompletedTask;
  82. return EnrollAsync(
  83. "WORK_ORDER",
  84. workOrd,
  85. domain,
  86. new[] { "WorkOrdMaster", "WorkOrdRouting", "PeriodSequenceDet" },
  87. item.TenantId,
  88. ct: ct);
  89. }
  90. /// <summary>领料单写入 165 成功后登记 PICK_BILL(及关联工单)热关注。</summary>
  91. public async Task EnrollPickBillAsync(
  92. string domain, string nbr, string? workOrd, long tenantId = 0, CancellationToken ct = default)
  93. {
  94. domain = string.IsNullOrWhiteSpace(domain) ? "8010" : domain.Trim();
  95. nbr = (nbr ?? "").Trim();
  96. if (string.IsNullOrWhiteSpace(nbr)) return;
  97. await EnrollAsync(
  98. "PICK_BILL",
  99. nbr,
  100. domain,
  101. new[] { "NbrMaster", "NbrDetail", "MissedPrint" },
  102. tenantId,
  103. ct: ct);
  104. if (!string.IsNullOrWhiteSpace(workOrd))
  105. {
  106. await EnrollAsync(
  107. "WORK_ORDER",
  108. workOrd.Trim(),
  109. domain,
  110. new[] { "WorkOrdMaster", "WorkOrdRouting", "PeriodSequenceDet" },
  111. tenantId,
  112. ct: ct);
  113. }
  114. }
  115. /// <summary>
  116. /// 装箱标签推入 165 后登记采购单热关注,用于回读 WMS 扫码收货结果。
  117. ///
  118. /// 挂在标签环节而非建单环节:标签生成意味着即将到货,此时开始轮询窗口最短。
  119. /// </summary>
  120. public async Task EnrollPurOrderAsync(
  121. string domain, string purOrd, long tenantId = 0, CancellationToken ct = default)
  122. {
  123. purOrd = (purOrd ?? "").Trim();
  124. if (string.IsNullOrWhiteSpace(purOrd)) return;
  125. domain = string.IsNullOrWhiteSpace(domain) ? "8010" : domain.Trim();
  126. if (tenantId <= 0)
  127. tenantId = await ResolveLocalPurOrdTenantAsync(purOrd, domain, ct);
  128. await EnrollAsync(
  129. "PUR_ORDER",
  130. purOrd,
  131. domain,
  132. new[] { "PurOrdMaster", "PurOrdDetail", "MissedPrint" },
  133. tenantId,
  134. ct: ct);
  135. }
  136. private static bool TryParseWoIdem(string? idem, out string domain, out string workOrd)
  137. {
  138. domain = "8010";
  139. workOrd = "";
  140. if (string.IsNullOrWhiteSpace(idem)) return false;
  141. var parts = idem.Split('|', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
  142. if (parts.Length < 3) return false;
  143. if (!parts[0].Equals("wo", StringComparison.OrdinalIgnoreCase)
  144. && !parts[0].Equals("pick", StringComparison.OrdinalIgnoreCase))
  145. return false;
  146. domain = string.IsNullOrWhiteSpace(parts[1]) ? "8010" : parts[1];
  147. workOrd = parts[2];
  148. return !string.IsNullOrWhiteSpace(workOrd);
  149. }
  150. /// <summary>取一批在途行并轮询 165。</summary>
  151. public async Task<(int polled, int changed, int terminated)> PollOnceAsync(
  152. int take = 100, CancellationToken ct = default)
  153. {
  154. var due = await _db.Queryable<AdoMdpHotWatch>()
  155. .Where(x => x.Status == 0)
  156. .OrderBy(x => x.LastPollTime ?? DateTime.MinValue)
  157. .Take(take)
  158. .ToListAsync(ct);
  159. if (due.Count == 0) return (0, 0, 0);
  160. MdpSource? source;
  161. ISqlSugarClient remote;
  162. try
  163. {
  164. source = await _db.Queryable<MdpSource>()
  165. .Where(x => x.SourceCode == SourceCode && x.Status == 1)
  166. .FirstAsync(ct);
  167. if (source == null)
  168. {
  169. _logger.LogWarning("[MdpHotWatch] 源 {Source} 未启用,跳过本轮", SourceCode);
  170. return (0, 0, 0);
  171. }
  172. remote = await _scopeFactory.GetScopeAsync(SourceCode, ct);
  173. }
  174. catch (Exception ex)
  175. {
  176. _logger.LogWarning(ex, "[MdpHotWatch] 无法连接 165");
  177. return (0, 0, 0);
  178. }
  179. var changed = 0;
  180. var terminated = 0;
  181. var now = DateTime.Now;
  182. foreach (var group in due.GroupBy(x => x.BizType))
  183. {
  184. ct.ThrowIfCancellationRequested();
  185. foreach (var watch in group.Take(MaxKeysPerBatch))
  186. {
  187. var tables = ParseTables(watch.WatchTables);
  188. var sb = new StringBuilder();
  189. foreach (var table in tables)
  190. {
  191. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  192. foreach (var row in rows)
  193. sb.Append(JsonSerializer.Serialize(row));
  194. }
  195. var hash = Sha256(sb.ToString());
  196. watch.LastPollTime = now;
  197. watch.UpdateTime = now;
  198. if (!string.Equals(hash, watch.LastSnapshotHash, StringComparison.Ordinal))
  199. {
  200. watch.LastSnapshotHash = hash;
  201. changed++;
  202. var tableRows = new Dictionary<string, List<Dictionary<string, object>>>(StringComparer.OrdinalIgnoreCase);
  203. // 变更落地:按表写 stg(执行侧字段快照)
  204. foreach (var table in tables)
  205. {
  206. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  207. tableRows[table] = rows;
  208. var entity = await ResolveEntityAsync(table, ct);
  209. if (entity == null) continue;
  210. foreach (var row in rows)
  211. {
  212. var dict = row.ToDictionary(
  213. kv => kv.Key,
  214. kv => (object?)kv.Value,
  215. StringComparer.OrdinalIgnoreCase);
  216. var rid = dict.TryGetValue("RecID", out var r) ? $"{r}" : watch.BizKey;
  217. var raw = JsonSerializer.Serialize(dict);
  218. await _staging.UpsertAsync(
  219. source, entity, table, dict, raw, rid,
  220. new MdpPullContext
  221. {
  222. TenantId = watch.TenantId,
  223. BatchId = $"hot-{now:yyyyMMddHHmmss}",
  224. FullRefresh = false
  225. });
  226. }
  227. }
  228. // 工单执行量写回本库业务表,供看板直接读取
  229. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  230. {
  231. tableRows.TryGetValue("WorkOrdMaster", out var masters);
  232. tableRows.TryGetValue("WorkOrdRouting", out var routings);
  233. tableRows.TryGetValue("PeriodSequenceDet", out var periodDets);
  234. var effectiveTenantId = watch.TenantId;
  235. if (effectiveTenantId <= 0)
  236. effectiveTenantId = await ResolveLocalWorkOrdTenantAsync(watch.BizKey, watch.Domain, ct);
  237. await ApplyWorkOrderExecutionAsync(
  238. effectiveTenantId, watch.Domain, watch.BizKey, masters, routings, ct);
  239. await ApplyPeriodSequenceActualAsync(
  240. effectiveTenantId, watch.Domain, watch.BizKey, periodDets, routings, ct);
  241. }
  242. // 采购收货结果写回本库业务表,供发货单列表与齐套口径直接读取
  243. if (string.Equals(watch.BizType, "PUR_ORDER", StringComparison.OrdinalIgnoreCase))
  244. {
  245. tableRows.TryGetValue("PurOrdDetail", out var purDetails);
  246. tableRows.TryGetValue("MissedPrint", out var barcodes);
  247. var effectiveTenantId = watch.TenantId;
  248. if (effectiveTenantId <= 0)
  249. effectiveTenantId = await ResolveLocalPurOrdTenantAsync(watch.BizKey, watch.Domain, ct);
  250. await ApplyPurchaseReceiptAsync(
  251. effectiveTenantId, watch.Domain, watch.BizKey, purDetails, barcodes, ct);
  252. }
  253. if (await ShouldTerminateAsync(remote, watch, ct))
  254. {
  255. watch.Status = 1;
  256. watch.TerminateTime = now;
  257. watch.TerminateReason = "auto";
  258. terminated++;
  259. }
  260. }
  261. await _db.Updateable(watch)
  262. .UpdateColumns(x => new
  263. {
  264. x.LastPollTime, x.LastSnapshotHash, x.Status,
  265. x.TerminateTime, x.TerminateReason, x.UpdateTime
  266. })
  267. .ExecuteCommandAsync(ct);
  268. }
  269. }
  270. return (due.Count, changed, terminated);
  271. }
  272. private async Task<MdpEntity?> ResolveEntityAsync(string table, CancellationToken ct)
  273. {
  274. return await _db.Queryable<MdpEntity>()
  275. .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1)
  276. .FirstAsync(ct);
  277. }
  278. private static async Task<List<Dictionary<string, object>>> QueryByBizKeyAsync(
  279. ISqlSugarClient remote, string table, string domain, string bizType, string bizKey, CancellationToken ct)
  280. {
  281. // 表白名单(防注入)
  282. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$"))
  283. throw new InvalidOperationException($"非法表名:{table}");
  284. string sql;
  285. SugarParameter[] pars;
  286. switch (table.ToUpperInvariant())
  287. {
  288. case "NBRMASTER":
  289. sql = "SELECT * FROM NbrMaster WHERE Domain=@d AND Nbr=@k";
  290. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  291. break;
  292. case "NBRDETAIL":
  293. sql = "SELECT * FROM NbrDetail WHERE Domain=@d AND Nbr=@k";
  294. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  295. break;
  296. case "WORKORDMASTER":
  297. sql = "SELECT * FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k";
  298. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  299. break;
  300. case "WORKORDROUTING":
  301. sql = "SELECT * FROM WorkOrdRouting WHERE Domain=@d AND WorkOrd=@k";
  302. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  303. break;
  304. case "PERIODSEQUENCEDET":
  305. // 工序间衔接:按工单取全部行(含 MES 报工产生的 Period=0 影子行),由调用方按工序合并
  306. sql = "SELECT * FROM PeriodSequenceDet WHERE Domain=@d AND WorkOrds=@k";
  307. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  308. break;
  309. case "PURORDDETAIL":
  310. sql = "SELECT * FROM PurOrdDetail WHERE Domain=@d AND PurOrd=@k";
  311. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  312. break;
  313. case "PURORDMASTER":
  314. sql = "SELECT * FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k";
  315. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  316. break;
  317. case "LINESTATUSDET":
  318. sql = "SELECT * FROM LineStatusDet WHERE Domain=@d AND Line=@k";
  319. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  320. break;
  321. case "MOBILETASK":
  322. // 堆表:按 TaskID
  323. sql = "SELECT * FROM MobileTask WHERE TaskID=@k";
  324. pars = new[] { new SugarParameter("@k", bizKey) };
  325. break;
  326. default:
  327. // MissedPrint 等:降级按 Domain + OrdNbr(可能扫表,见 WP8 E1)
  328. if (string.Equals(table, "MissedPrint", StringComparison.OrdinalIgnoreCase))
  329. {
  330. sql = "SELECT TOP 200 * FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k";
  331. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  332. break;
  333. }
  334. return new List<Dictionary<string, object>>();
  335. }
  336. var dt = await remote.Ado.GetDataTableAsync(sql, pars);
  337. var list = new List<Dictionary<string, object>>();
  338. foreach (System.Data.DataRow row in dt.Rows)
  339. {
  340. var dict = new Dictionary<string, object>(StringComparer.OrdinalIgnoreCase);
  341. foreach (System.Data.DataColumn col in dt.Columns)
  342. dict[col.ColumnName] = row[col] == DBNull.Value ? null! : row[col];
  343. list.Add(dict);
  344. }
  345. return list;
  346. }
  347. /// <summary>
  348. /// 把 165 工单执行字段回写本库(MES 权威列:Status / QtyCompleted / QtyComplete / QtyReject)。
  349. /// </summary>
  350. private async Task ApplyWorkOrderExecutionAsync(
  351. long tenantId,
  352. string domain,
  353. string workOrd,
  354. List<Dictionary<string, object>>? masters,
  355. List<Dictionary<string, object>>? routings,
  356. CancellationToken ct)
  357. {
  358. var now = DateTime.Now;
  359. var updatedRouting = 0;
  360. if (routings != null)
  361. {
  362. foreach (var row in routings)
  363. {
  364. if (!TryGetInt(row, "OP", out var op) && !TryGetInt(row, "Op", out op))
  365. continue;
  366. var qtyComplete = GetDecimal(row, "QtyComplete");
  367. var qtyReject = GetDecimal(row, "QtyReject");
  368. var status = Trunc(GetString(row, "Status"), 1);
  369. updatedRouting += await _db.Ado.ExecuteCommandAsync(
  370. """
  371. UPDATE WorkOrdRouting
  372. SET QtyComplete = @QtyComplete,
  373. QtyReject = @QtyReject,
  374. Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END,
  375. UpdateUser = 'MDP_HOT',
  376. UpdateTime = @Now
  377. WHERE WorkOrd = @WorkOrd
  378. AND OP = @Op
  379. AND IFNULL(Domain, '') = @Domain
  380. AND IFNULL(tenant_id, 0) = @TenantId
  381. """,
  382. new SugarParameter("@QtyComplete", qtyComplete),
  383. new SugarParameter("@QtyReject", qtyReject),
  384. new SugarParameter("@Status", status ?? ""),
  385. new SugarParameter("@Now", now),
  386. new SugarParameter("@WorkOrd", workOrd),
  387. new SugarParameter("@Op", op),
  388. new SugarParameter("@Domain", domain),
  389. new SugarParameter("@TenantId", tenantId));
  390. }
  391. }
  392. if (masters is { Count: > 0 })
  393. {
  394. var m = masters[0];
  395. var qtyCompleted = GetDecimal(m, "QtyCompleted");
  396. var status = Trunc(GetString(m, "Status"), 8);
  397. await _db.Ado.ExecuteCommandAsync(
  398. """
  399. UPDATE WorkOrdMaster
  400. SET QtyCompleted = @QtyCompleted,
  401. Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END,
  402. UpdateUser = 'MDP_HOT',
  403. UpdateTime = @Now
  404. WHERE WorkOrd = @WorkOrd
  405. AND IFNULL(Domain, '') = @Domain
  406. AND IFNULL(tenant_id, 0) = @TenantId
  407. """,
  408. new SugarParameter("@QtyCompleted", qtyCompleted),
  409. new SugarParameter("@Status", status ?? ""),
  410. new SugarParameter("@Now", now),
  411. new SugarParameter("@WorkOrd", workOrd),
  412. new SugarParameter("@Domain", domain),
  413. new SugarParameter("@TenantId", tenantId));
  414. }
  415. _logger.LogInformation(
  416. "[MdpHotWatch] applied WORK_ORDER execution wo={WorkOrd} domain={Domain} routingRows={Rows}",
  417. workOrd, domain, updatedRouting);
  418. }
  419. /// <summary>
  420. /// D-M01:把 165 的工序实绩按 (Domain, WorkOrds, Op) 合并后落到 ado_psd_op_actual。
  421. ///
  422. /// 为什么按工序合并:MES APP 报工会新增 Period=0/IsActive=0 的影子行,实绩可能写在影子行、
  423. /// 也可能写在 Period=1 计划行,且同一工序可能因重排跨多个 Line。只有工序级求和才与
  424. /// WorkOrdRouting.QtyComplete(MES 自己的工序累计完成数)对齐。实测 Op501:3+5=8=QtyComplete。
  425. ///
  426. /// 为什么不写 PeriodSequenceDet:本库 PSD 行会被排产整表重建(旧行置 IsActive=0),
  427. /// 实绩写进去会在下次全量重排时丢失。
  428. ///
  429. /// 单位保持 165 原始口径:ActualTime 秒、SetupTime 小时(列名自带单位)。
  430. /// </summary>
  431. private async Task ApplyPeriodSequenceActualAsync(
  432. long tenantId,
  433. string domain,
  434. string workOrd,
  435. List<Dictionary<string, object>>? periodDets,
  436. List<Dictionary<string, object>>? routings,
  437. CancellationToken ct)
  438. {
  439. if (periodDets is not { Count: > 0 })
  440. return;
  441. // WorkOrdRouting.QtyComplete 作为交叉校验值(仅记录,不用于覆盖)
  442. var routingQty = new Dictionary<int, decimal>();
  443. if (routings != null)
  444. {
  445. foreach (var r in routings)
  446. {
  447. if (!TryGetInt(r, "OP", out var rop) && !TryGetInt(r, "Op", out rop))
  448. continue;
  449. routingQty[rop] = GetDecimal(r, "QtyComplete");
  450. }
  451. }
  452. var grouped = new Dictionary<int, PsdOpActual>();
  453. foreach (var row in periodDets)
  454. {
  455. if (!TryGetInt(row, "Op", out var op) && !TryGetInt(row, "OP", out op))
  456. continue;
  457. if (!grouped.TryGetValue(op, out var acc))
  458. {
  459. acc = new PsdOpActual { Op = op };
  460. grouped[op] = acc;
  461. }
  462. acc.CompQty += GetDecimal(row, "CompQty");
  463. acc.RejectQty += GetDecimal(row, "RejectQty");
  464. acc.SetupTimeHour += GetDecimal(row, "SetupTime");
  465. acc.ActualTimeSec += GetDecimal(row, "ActualTime");
  466. acc.SrcRowCount++;
  467. acc.ItemNum ??= GetString(row, "ItemNum");
  468. if (row.TryGetValue("UpdateTime", out var upd) && upd is DateTime dt
  469. && (acc.LastSrcUpdate == null || dt > acc.LastSrcUpdate))
  470. acc.LastSrcUpdate = dt;
  471. }
  472. var now = DateTime.Now;
  473. var mismatches = new List<string>();
  474. foreach (var acc in grouped.Values)
  475. {
  476. var rq = routingQty.TryGetValue(acc.Op, out var q) ? (decimal?)q : null;
  477. if (rq.HasValue && rq.Value != acc.CompQty)
  478. mismatches.Add($"Op{acc.Op}: psdSum={acc.CompQty} routing={rq.Value}");
  479. await _db.Ado.ExecuteCommandAsync(
  480. """
  481. INSERT INTO ado_psd_op_actual (
  482. tenant_id, domain, work_ord, op, item_num,
  483. comp_qty, reject_qty, setup_time_hour, actual_time_sec,
  484. src_row_count, routing_qty, last_src_update, sync_time
  485. ) VALUES (
  486. @TenantId, @Domain, @WorkOrd, @Op, @ItemNum,
  487. @CompQty, @RejectQty, @SetupTimeHour, @ActualTimeSec,
  488. @SrcRowCount, @RoutingQty, @LastSrcUpdate, @Now
  489. )
  490. ON DUPLICATE KEY UPDATE
  491. item_num = VALUES(item_num),
  492. comp_qty = VALUES(comp_qty),
  493. reject_qty = VALUES(reject_qty),
  494. setup_time_hour = VALUES(setup_time_hour),
  495. actual_time_sec = VALUES(actual_time_sec),
  496. src_row_count = VALUES(src_row_count),
  497. routing_qty = VALUES(routing_qty),
  498. last_src_update = VALUES(last_src_update),
  499. sync_time = VALUES(sync_time)
  500. """,
  501. new SugarParameter("@TenantId", tenantId),
  502. new SugarParameter("@Domain", domain ?? ""),
  503. new SugarParameter("@WorkOrd", workOrd),
  504. new SugarParameter("@Op", acc.Op),
  505. new SugarParameter("@ItemNum", acc.ItemNum ?? (object)DBNull.Value),
  506. new SugarParameter("@CompQty", acc.CompQty),
  507. new SugarParameter("@RejectQty", acc.RejectQty),
  508. new SugarParameter("@SetupTimeHour", acc.SetupTimeHour),
  509. new SugarParameter("@ActualTimeSec", acc.ActualTimeSec),
  510. new SugarParameter("@SrcRowCount", acc.SrcRowCount),
  511. new SugarParameter("@RoutingQty", rq.HasValue ? rq.Value : (object)DBNull.Value),
  512. new SugarParameter("@LastSrcUpdate", acc.LastSrcUpdate ?? (object)DBNull.Value),
  513. new SugarParameter("@Now", now));
  514. }
  515. if (mismatches.Count > 0)
  516. {
  517. // 不阻断落库:差异说明 165 侧口径需人工确认,先留证
  518. _logger.LogWarning(
  519. "[MdpHotWatch] PSD 实绩与 WorkOrdRouting 不一致 wo={WorkOrd} domain={Domain} {Detail}",
  520. workOrd, domain, string.Join("; ", mismatches));
  521. }
  522. _logger.LogInformation(
  523. "[MdpHotWatch] applied PSD actual wo={WorkOrd} domain={Domain} ops={Ops} srcRows={Rows}",
  524. workOrd, domain, grouped.Count, periodDets.Count);
  525. }
  526. /// <summary>
  527. /// 把 165 的采购收货结果回读本库:明细收货数、箱码状态,并据箱码推进送货单状态。
  528. ///
  529. /// 口径(2026-08-11 定):WMS 扫码收货只入待检仓(InvTransHist 落 rct-po-ins / Loc=1000),
  530. /// 不等于合格入库,所以这里只推进 scm_shd/scm_shdzb.shzt,不写 rksl——合格入库数留给检验环节。
  531. /// 165 的 scm_shdzb.rksl / shzt 在收货时并不更新,送货单进度只能由箱码状态反推。
  532. /// </summary>
  533. private async Task ApplyPurchaseReceiptAsync(
  534. long tenantId,
  535. string domain,
  536. string purOrd,
  537. List<Dictionary<string, object>>? purDetails,
  538. List<Dictionary<string, object>>? barcodes,
  539. CancellationToken ct)
  540. {
  541. var now = DateTime.Now;
  542. var updatedLines = 0;
  543. var updatedBarcodes = 0;
  544. if (purDetails != null)
  545. {
  546. foreach (var row in purDetails)
  547. {
  548. if (!TryGetInt(row, "Line", out var line))
  549. continue;
  550. // 自建单本库不填 Domain(165/MES 概念),所以空 Domain 也算命中;采购单号本库唯一
  551. updatedLines += await _db.Ado.ExecuteCommandAsync(
  552. """
  553. UPDATE PurOrdDetail
  554. SET RctQty = @RctQty,
  555. ReceiptQty = @ReceiptQty,
  556. QtyReturned = @QtyReturned,
  557. UpdateUser = 'MDP_HOT',
  558. UpdateTime = @Now
  559. WHERE PurOrd = @PurOrd
  560. AND Line = @Line
  561. AND IFNULL(Domain, '') IN ('', @Domain)
  562. AND (@TenantId = 0 OR IFNULL(tenant_id, 0) = @TenantId)
  563. """,
  564. new SugarParameter("@RctQty", GetDecimal(row, "RctQty")),
  565. new SugarParameter("@ReceiptQty", GetDecimal(row, "ReceiptQty")),
  566. new SugarParameter("@QtyReturned", GetDecimal(row, "QtyReturned")),
  567. new SugarParameter("@Now", now),
  568. new SugarParameter("@PurOrd", purOrd),
  569. new SugarParameter("@Line", line),
  570. new SugarParameter("@Domain", domain),
  571. new SugarParameter("@TenantId", tenantId));
  572. }
  573. }
  574. // 本库 MissedPrint 无 tenant_id,按 Domain + BarCode 定位;作废行 BarCode 带 RecID 前缀不会误匹配
  575. var shippers = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
  576. if (barcodes != null)
  577. {
  578. foreach (var row in barcodes)
  579. {
  580. var barCode = GetString(row, "BarCode")?.Trim();
  581. if (string.IsNullOrWhiteSpace(barCode))
  582. continue;
  583. updatedBarcodes += await _db.Ado.ExecuteCommandAsync(
  584. """
  585. UPDATE MissedPrint
  586. SET Status = CASE WHEN IFNULL(@Status, '') = '' THEN Status ELSE @Status END,
  587. RctNbr = @RctNbr,
  588. Location = @Location,
  589. Shelf = @Shelf,
  590. InvStatus = @InvStatus,
  591. UpdateUser = 'MDP_HOT',
  592. UpdateTime = @Now
  593. WHERE BarCode = @BarCode
  594. AND IFNULL(Domain, '') IN ('', @Domain)
  595. """,
  596. new SugarParameter("@Status", Trunc(GetString(row, "Status"), 4) ?? ""),
  597. new SugarParameter("@RctNbr", Trunc(GetString(row, "RctNbr"), 48) ?? ""),
  598. new SugarParameter("@Location", Trunc(GetString(row, "Location"), 24) ?? ""),
  599. new SugarParameter("@Shelf", Trunc(GetString(row, "Shelf"), 40) ?? ""),
  600. new SugarParameter("@InvStatus", Trunc(GetString(row, "InvStatus"), 24) ?? ""),
  601. new SugarParameter("@Now", now),
  602. new SugarParameter("@BarCode", barCode),
  603. new SugarParameter("@Domain", domain));
  604. var shipper = GetString(row, "ShipperNbr")?.Trim();
  605. if (!string.IsNullOrWhiteSpace(shipper))
  606. shippers.Add(shipper!);
  607. }
  608. }
  609. foreach (var shipper in shippers)
  610. await ApplyShipmentStatusAsync(tenantId, domain, shipper, ct);
  611. _logger.LogInformation(
  612. "[MdpHotWatch] applied PUR_ORDER receipt po={PurOrd} domain={Domain} lines={Lines} barcodes={Barcodes} shipments={Shipments}",
  613. purOrd, domain, updatedLines, updatedBarcodes, shippers.Count);
  614. }
  615. /// <summary>
  616. /// 按箱码待收数推进送货单状态:全部待收=待收,全部离开待收=完成,其余=收货中。
  617. ///
  618. /// 完成判据取箱码而非数量(决策:少收/超收不影响状态),列表读的是 scm_shd.shzt。
  619. /// </summary>
  620. private async Task ApplyShipmentStatusAsync(
  621. long tenantId, string domain, string shddh, CancellationToken ct)
  622. {
  623. var stats = await _db.Ado.SqlQueryAsync<ShipmentBarcodeStat>(
  624. """
  625. SELECT
  626. COUNT(*) AS Total,
  627. SUM(CASE WHEN IFNULL(Status, '') = 'U' THEN 1 ELSE 0 END) AS Pending
  628. FROM MissedPrint
  629. WHERE ShipperNbr = @Shddh
  630. AND IFNULL(Domain, '') = @Domain
  631. AND IFNULL(PurOrd, '') NOT LIKE '作废%'
  632. """,
  633. new SugarParameter("@Shddh", shddh),
  634. new SugarParameter("@Domain", domain));
  635. var stat = stats?.FirstOrDefault();
  636. if (stat == null || stat.Total <= 0)
  637. return;
  638. var shzt = stat.Pending >= stat.Total ? "待收" : (stat.Pending == 0 ? "完成" : "收货中");
  639. // 先定位表头再按主键更新:scm_shd.id 与 scm_shdzb.glid 排序规则不同(0900_ai_ci vs unicode_ci),
  640. // 直接 JOIN 会报 Illegal mix of collations,改由参数传 glid 规避
  641. var heads = await _db.Ado.SqlQueryAsync<ShipmentHeadRow>(
  642. """
  643. SELECT id AS Id
  644. FROM scm_shd
  645. WHERE shddh = @Shddh
  646. AND (@TenantId = 0 OR IFNULL(tenant_id, 0) = @TenantId)
  647. ORDER BY id
  648. LIMIT 1
  649. """,
  650. new SugarParameter("@Shddh", shddh),
  651. new SugarParameter("@TenantId", tenantId));
  652. var head = heads?.FirstOrDefault();
  653. if (head == null)
  654. {
  655. _logger.LogWarning(
  656. "[MdpHotWatch] 送货单在本库缺失,跳过状态推进 shddh={Shddh} tenant={Tenant}", shddh, tenantId);
  657. return;
  658. }
  659. await _db.Ado.ExecuteCommandAsync(
  660. "UPDATE scm_shd SET shzt = @Shzt WHERE id = @Id",
  661. new SugarParameter("@Shzt", shzt),
  662. new SugarParameter("@Id", head.Id));
  663. await _db.Ado.ExecuteCommandAsync(
  664. "UPDATE scm_shdzb SET shzt = @Shzt WHERE glid = @Glid",
  665. new SugarParameter("@Shzt", shzt),
  666. new SugarParameter("@Glid", head.Id.ToString()));
  667. }
  668. private sealed class ShipmentBarcodeStat
  669. {
  670. public int Total { get; set; }
  671. public int Pending { get; set; }
  672. }
  673. private sealed class ShipmentHeadRow
  674. {
  675. public long Id { get; set; }
  676. }
  677. private async Task<long> ResolveLocalPurOrdTenantAsync(string purOrd, string domain, CancellationToken ct)
  678. {
  679. var tid = await _db.Ado.SqlQuerySingleAsync<long?>(
  680. """
  681. SELECT IFNULL(tenant_id, 0)
  682. FROM PurOrdDetail
  683. WHERE PurOrd = @PurOrd AND IFNULL(Domain, '') IN ('', @Domain)
  684. ORDER BY Line
  685. LIMIT 1
  686. """,
  687. new SugarParameter("@PurOrd", purOrd),
  688. new SugarParameter("@Domain", domain));
  689. return tid ?? 0;
  690. }
  691. private sealed class PsdOpActual
  692. {
  693. public int Op { get; set; }
  694. public string? ItemNum { get; set; }
  695. public decimal CompQty { get; set; }
  696. public decimal RejectQty { get; set; }
  697. public decimal SetupTimeHour { get; set; }
  698. public decimal ActualTimeSec { get; set; }
  699. public int SrcRowCount { get; set; }
  700. public DateTime? LastSrcUpdate { get; set; }
  701. }
  702. private async Task<long> ResolveLocalWorkOrdTenantAsync(string workOrd, string domain, CancellationToken ct)
  703. {
  704. var tid = await _db.Ado.SqlQuerySingleAsync<long?>(
  705. """
  706. SELECT IFNULL(tenant_id, 0)
  707. FROM WorkOrdMaster
  708. WHERE WorkOrd = @WorkOrd AND IFNULL(Domain, '') = @Domain
  709. ORDER BY RecID DESC
  710. LIMIT 1
  711. """,
  712. new SugarParameter("@WorkOrd", workOrd),
  713. new SugarParameter("@Domain", domain));
  714. return tid ?? 0;
  715. }
  716. private static bool TryGetInt(Dictionary<string, object> row, string key, out int value)
  717. {
  718. value = 0;
  719. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return false;
  720. try
  721. {
  722. value = Convert.ToInt32(raw);
  723. return true;
  724. }
  725. catch
  726. {
  727. return false;
  728. }
  729. }
  730. private static decimal GetDecimal(Dictionary<string, object> row, string key)
  731. {
  732. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return 0m;
  733. try { return Convert.ToDecimal(raw); }
  734. catch { return 0m; }
  735. }
  736. private static string? GetString(Dictionary<string, object> row, string key)
  737. {
  738. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return null;
  739. return Convert.ToString(raw);
  740. }
  741. private static string? Trunc(string? s, int max) =>
  742. string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
  743. private static async Task<bool> ShouldTerminateAsync(
  744. ISqlSugarClient remote, AdoMdpHotWatch watch, CancellationToken ct)
  745. {
  746. if (string.Equals(watch.BizType, "PICK_BILL", StringComparison.OrdinalIgnoreCase))
  747. {
  748. var status = await remote.Ado.GetScalarAsync(
  749. "SELECT TOP 1 Status FROM NbrMaster WHERE Domain=@d AND Nbr=@k",
  750. new SugarParameter("@d", watch.Domain),
  751. new SugarParameter("@k", watch.BizKey));
  752. var s = status?.ToString()?.Trim() ?? "";
  753. // 终态取值待 WP4 固化;临时:C/Complete/关闭 等常见值
  754. return s is "C" or "Complete" or "关闭" or "Y";
  755. }
  756. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  757. {
  758. var status = await remote.Ado.GetScalarAsync(
  759. "SELECT TOP 1 Status FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k",
  760. new SugarParameter("@d", watch.Domain),
  761. new SugarParameter("@k", watch.BizKey));
  762. return string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase);
  763. }
  764. if (string.Equals(watch.BizType, "PUR_ORDER", StringComparison.OrdinalIgnoreCase))
  765. {
  766. var status = await remote.Ado.GetScalarAsync(
  767. "SELECT TOP 1 Status FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k",
  768. new SugarParameter("@d", watch.Domain),
  769. new SugarParameter("@k", watch.BizKey));
  770. if (string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase))
  771. return true;
  772. // 箱码全部离开待收即收货结束;标签尚未推达(total=0)时不能终结
  773. var total = Convert.ToInt32(await remote.Ado.GetScalarAsync(
  774. "SELECT COUNT(*) FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k",
  775. new SugarParameter("@d", watch.Domain),
  776. new SugarParameter("@k", watch.BizKey)) ?? 0);
  777. if (total > 0)
  778. {
  779. var pending = Convert.ToInt32(await remote.Ado.GetScalarAsync(
  780. "SELECT COUNT(*) FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k AND IsNull(Status,'')='U'",
  781. new SugarParameter("@d", watch.Domain),
  782. new SugarParameter("@k", watch.BizKey)) ?? 0);
  783. if (pending == 0) return true;
  784. }
  785. }
  786. // 兜底:超过 7 天强制终结
  787. return watch.EnrollTime < DateTime.Now.AddDays(-7);
  788. }
  789. private static List<string> ParseTables(string? json)
  790. {
  791. if (string.IsNullOrWhiteSpace(json)) return new List<string>();
  792. try
  793. {
  794. return JsonSerializer.Deserialize<List<string>>(json!) ?? new List<string>();
  795. }
  796. catch
  797. {
  798. return new List<string>();
  799. }
  800. }
  801. private static string Sha256(string s)
  802. {
  803. var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s));
  804. return Convert.ToHexString(bytes);
  805. }
  806. }