MdpHotWatchService.cs 56 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215
  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. // S5 IQC:在采购收货热关注上叠加报检/检验单三表,供申请列表与任务列表出数
  129. var tables = new[]
  130. {
  131. "PurOrdMaster", "PurOrdDetail", "MissedPrint",
  132. "qms_qcp_inspecapplyn", "qms_qcp_insappnentry", "qms_qcp_inspbill"
  133. };
  134. var existing = await _db.Queryable<AdoMdpHotWatch>()
  135. .Where(x => x.BizType == "PUR_ORDER" && x.BizKey == purOrd && x.Status == 0)
  136. .FirstAsync(ct);
  137. if (existing != null)
  138. {
  139. // 已登记的关注行补齐 qms 三表(旧行只有采购三表)
  140. var cur = ParseTables(existing.WatchTables);
  141. var needUpdate = tables.Any(t => !cur.Contains(t, StringComparer.OrdinalIgnoreCase));
  142. if (needUpdate)
  143. {
  144. existing.WatchTables = JsonSerializer.Serialize(tables);
  145. existing.UpdateTime = DateTime.Now;
  146. if (tenantId > 0 && existing.TenantId <= 0) existing.TenantId = tenantId;
  147. await _db.Updateable(existing)
  148. .UpdateColumns(x => new { x.WatchTables, x.TenantId, x.UpdateTime })
  149. .ExecuteCommandAsync(ct);
  150. }
  151. return;
  152. }
  153. await EnrollAsync("PUR_ORDER", purOrd, domain, tables, tenantId, ct: ct);
  154. }
  155. private static bool TryParseWoIdem(string? idem, out string domain, out string workOrd)
  156. {
  157. domain = "8010";
  158. workOrd = "";
  159. if (string.IsNullOrWhiteSpace(idem)) return false;
  160. var parts = idem.Split('|', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
  161. if (parts.Length < 3) return false;
  162. if (!parts[0].Equals("wo", StringComparison.OrdinalIgnoreCase)
  163. && !parts[0].Equals("pick", StringComparison.OrdinalIgnoreCase))
  164. return false;
  165. domain = string.IsNullOrWhiteSpace(parts[1]) ? "8010" : parts[1];
  166. workOrd = parts[2];
  167. return !string.IsNullOrWhiteSpace(workOrd);
  168. }
  169. /// <summary>取一批在途行并轮询 165。</summary>
  170. public async Task<(int polled, int changed, int terminated)> PollOnceAsync(
  171. int take = 100, CancellationToken ct = default)
  172. {
  173. var due = await _db.Queryable<AdoMdpHotWatch>()
  174. .Where(x => x.Status == 0)
  175. .OrderBy(x => x.LastPollTime ?? DateTime.MinValue)
  176. .Take(take)
  177. .ToListAsync(ct);
  178. if (due.Count == 0) return (0, 0, 0);
  179. MdpSource? source;
  180. ISqlSugarClient remote;
  181. try
  182. {
  183. source = await _db.Queryable<MdpSource>()
  184. .Where(x => x.SourceCode == SourceCode && x.Status == 1)
  185. .FirstAsync(ct);
  186. if (source == null)
  187. {
  188. _logger.LogWarning("[MdpHotWatch] 源 {Source} 未启用,跳过本轮", SourceCode);
  189. return (0, 0, 0);
  190. }
  191. remote = await _scopeFactory.GetScopeAsync(SourceCode, ct);
  192. }
  193. catch (Exception ex)
  194. {
  195. _logger.LogWarning(ex, "[MdpHotWatch] 无法连接 165");
  196. return (0, 0, 0);
  197. }
  198. var changed = 0;
  199. var terminated = 0;
  200. var now = DateTime.Now;
  201. foreach (var group in due.GroupBy(x => x.BizType))
  202. {
  203. ct.ThrowIfCancellationRequested();
  204. foreach (var watch in group.Take(MaxKeysPerBatch))
  205. {
  206. var tables = ParseTables(watch.WatchTables);
  207. // 序无关哈希:QueryByBizKey 多段无 ORDER BY,行序抖动会导致指纹每轮变化并无谓重放 UPSERT
  208. var hashParts = new List<string>();
  209. foreach (var table in tables)
  210. {
  211. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  212. foreach (var row in rows)
  213. hashParts.Add(JsonSerializer.Serialize(row));
  214. }
  215. hashParts.Sort(StringComparer.Ordinal);
  216. var hash = Sha256(string.Concat(hashParts));
  217. watch.LastPollTime = now;
  218. watch.UpdateTime = now;
  219. if (!string.Equals(hash, watch.LastSnapshotHash, StringComparison.Ordinal))
  220. {
  221. watch.LastSnapshotHash = hash;
  222. changed++;
  223. var tableRows = new Dictionary<string, List<Dictionary<string, object>>>(StringComparer.OrdinalIgnoreCase);
  224. // 变更落地:按表写 stg(执行侧字段快照)
  225. foreach (var table in tables)
  226. {
  227. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  228. tableRows[table] = rows;
  229. var entity = await ResolveEntityAsync(table, ct);
  230. if (entity == null) continue;
  231. foreach (var row in rows)
  232. {
  233. var dict = row.ToDictionary(
  234. kv => kv.Key,
  235. kv => (object?)kv.Value,
  236. StringComparer.OrdinalIgnoreCase);
  237. var rid = dict.TryGetValue("RecID", out var r) ? $"{r}" : watch.BizKey;
  238. var raw = JsonSerializer.Serialize(dict);
  239. await _staging.UpsertAsync(
  240. source, entity, table, dict, raw, rid,
  241. new MdpPullContext
  242. {
  243. TenantId = watch.TenantId,
  244. BatchId = $"hot-{now:yyyyMMddHHmmss}",
  245. FullRefresh = false
  246. });
  247. }
  248. }
  249. // 工单执行量写回本库业务表,供看板直接读取
  250. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  251. {
  252. tableRows.TryGetValue("WorkOrdMaster", out var masters);
  253. tableRows.TryGetValue("WorkOrdRouting", out var routings);
  254. tableRows.TryGetValue("PeriodSequenceDet", out var periodDets);
  255. var effectiveTenantId = watch.TenantId;
  256. if (effectiveTenantId <= 0)
  257. effectiveTenantId = await ResolveLocalWorkOrdTenantAsync(watch.BizKey, watch.Domain, ct);
  258. await ApplyWorkOrderExecutionAsync(
  259. effectiveTenantId, watch.Domain, watch.BizKey, masters, routings, ct);
  260. await ApplyPeriodSequenceActualAsync(
  261. effectiveTenantId, watch.Domain, watch.BizKey, periodDets, routings, ct);
  262. }
  263. // 采购收货结果写回本库业务表,供发货单列表与齐套口径直接读取
  264. if (string.Equals(watch.BizType, "PUR_ORDER", StringComparison.OrdinalIgnoreCase))
  265. {
  266. tableRows.TryGetValue("PurOrdDetail", out var purDetails);
  267. tableRows.TryGetValue("MissedPrint", out var barcodes);
  268. tableRows.TryGetValue("qms_qcp_inspecapplyn", out var iqcApplies);
  269. tableRows.TryGetValue("qms_qcp_insappnentry", out var iqcEntries);
  270. tableRows.TryGetValue("qms_qcp_inspbill", out var iqcBills);
  271. var effectiveTenantId = watch.TenantId;
  272. if (effectiveTenantId <= 0)
  273. effectiveTenantId = await ResolveLocalPurOrdTenantAsync(watch.BizKey, watch.Domain, ct);
  274. await ApplyPurchaseReceiptAsync(
  275. effectiveTenantId, watch.Domain, watch.BizKey, purDetails, barcodes, ct);
  276. await ApplyIqcAsync(
  277. effectiveTenantId, watch.Domain, watch.BizKey,
  278. iqcApplies, iqcEntries, iqcBills, ct);
  279. }
  280. if (await ShouldTerminateAsync(remote, watch, ct))
  281. {
  282. watch.Status = 1;
  283. watch.TerminateTime = now;
  284. watch.TerminateReason = "auto";
  285. terminated++;
  286. }
  287. }
  288. await _db.Updateable(watch)
  289. .UpdateColumns(x => new
  290. {
  291. x.LastPollTime, x.LastSnapshotHash, x.Status,
  292. x.TerminateTime, x.TerminateReason, x.UpdateTime
  293. })
  294. .ExecuteCommandAsync(ct);
  295. }
  296. }
  297. return (due.Count, changed, terminated);
  298. }
  299. private async Task<MdpEntity?> ResolveEntityAsync(string table, CancellationToken ct)
  300. {
  301. return await _db.Queryable<MdpEntity>()
  302. .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1)
  303. .FirstAsync(ct);
  304. }
  305. private static async Task<List<Dictionary<string, object>>> QueryByBizKeyAsync(
  306. ISqlSugarClient remote, string table, string domain, string bizType, string bizKey, CancellationToken ct)
  307. {
  308. // 表白名单(防注入)
  309. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$"))
  310. throw new InvalidOperationException($"非法表名:{table}");
  311. string sql;
  312. SugarParameter[] pars;
  313. switch (table.ToUpperInvariant())
  314. {
  315. case "NBRMASTER":
  316. sql = "SELECT * FROM NbrMaster WHERE Domain=@d AND Nbr=@k";
  317. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  318. break;
  319. case "NBRDETAIL":
  320. sql = "SELECT * FROM NbrDetail WHERE Domain=@d AND Nbr=@k";
  321. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  322. break;
  323. case "WORKORDMASTER":
  324. sql = "SELECT * FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k";
  325. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  326. break;
  327. case "WORKORDROUTING":
  328. sql = "SELECT * FROM WorkOrdRouting WHERE Domain=@d AND WorkOrd=@k";
  329. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  330. break;
  331. case "PERIODSEQUENCEDET":
  332. // 工序间衔接:按工单取全部行(含 MES 报工产生的 Period=0 影子行),由调用方按工序合并
  333. sql = "SELECT * FROM PeriodSequenceDet WHERE Domain=@d AND WorkOrds=@k";
  334. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  335. break;
  336. case "PURORDDETAIL":
  337. sql = "SELECT * FROM PurOrdDetail WHERE Domain=@d AND PurOrd=@k";
  338. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  339. break;
  340. case "PURORDMASTER":
  341. sql = "SELECT * FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k";
  342. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  343. break;
  344. case "LINESTATUSDET":
  345. sql = "SELECT * FROM LineStatusDet WHERE Domain=@d AND Line=@k";
  346. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  347. break;
  348. case "MOBILETASK":
  349. // 堆表:按 TaskID
  350. sql = "SELECT * FROM MobileTask WHERE TaskID=@k";
  351. pars = new[] { new SugarParameter("@k", bizKey) };
  352. break;
  353. // S5 IQC:按采购单 → 收货单 → 报检分录/主表/检验单(业务键窄查)
  354. case "QMS_QCP_INSAPPNENTRY":
  355. sql = """
  356. SELECT TOP 200 e.*
  357. FROM qms_qcp_insappnentry e WITH (NOLOCK)
  358. INNER JOIN PurOrdRctDetail d WITH (NOLOCK)
  359. ON d.Receiver = e.FSRCORDERNUM AND d.ItemNum = e.FMATERIALCFG
  360. WHERE d.Domain = @d AND d.OrdNbr = @k
  361. """;
  362. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  363. break;
  364. case "QMS_QCP_INSPECAPPLYN":
  365. sql = """
  366. SELECT DISTINCT TOP 200 a.*
  367. FROM qms_qcp_inspecapplyn a WITH (NOLOCK)
  368. INNER JOIN qms_qcp_insappnentry e WITH (NOLOCK) ON a.id = e.glid
  369. INNER JOIN PurOrdRctDetail d WITH (NOLOCK)
  370. ON d.Receiver = e.FSRCORDERNUM AND d.ItemNum = e.FMATERIALCFG
  371. WHERE d.Domain = @d AND d.OrdNbr = @k
  372. """;
  373. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  374. break;
  375. case "QMS_QCP_INSPBILL":
  376. sql = """
  377. SELECT DISTINCT TOP 200 b.*
  378. FROM qms_qcp_inspbill b WITH (NOLOCK)
  379. INNER JOIN qms_qcp_inspecapplyn a WITH (NOLOCK) ON a.FBILLNO = b.lydjbh
  380. INNER JOIN qms_qcp_insappnentry e WITH (NOLOCK) ON a.id = e.glid
  381. INNER JOIN PurOrdRctDetail d WITH (NOLOCK)
  382. ON d.Receiver = e.FSRCORDERNUM AND d.ItemNum = e.FMATERIALCFG
  383. WHERE d.Domain = @d AND d.OrdNbr = @k
  384. """;
  385. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  386. break;
  387. default:
  388. // MissedPrint 等:降级按 Domain + OrdNbr(可能扫表,见 WP8 E1)
  389. if (string.Equals(table, "MissedPrint", StringComparison.OrdinalIgnoreCase))
  390. {
  391. sql = "SELECT TOP 200 * FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k";
  392. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  393. break;
  394. }
  395. return new List<Dictionary<string, object>>();
  396. }
  397. var dt = await remote.Ado.GetDataTableAsync(sql, pars);
  398. var list = new List<Dictionary<string, object>>();
  399. foreach (System.Data.DataRow row in dt.Rows)
  400. {
  401. var dict = new Dictionary<string, object>(StringComparer.OrdinalIgnoreCase);
  402. foreach (System.Data.DataColumn col in dt.Columns)
  403. dict[col.ColumnName] = row[col] == DBNull.Value ? null! : row[col];
  404. list.Add(dict);
  405. }
  406. return list;
  407. }
  408. /// <summary>
  409. /// 把 165 工单执行字段回写本库(MES 权威列:Status / QtyCompleted / QtyComplete / QtyReject / QtyScrap)。
  410. /// </summary>
  411. private async Task ApplyWorkOrderExecutionAsync(
  412. long tenantId,
  413. string domain,
  414. string workOrd,
  415. List<Dictionary<string, object>>? masters,
  416. List<Dictionary<string, object>>? routings,
  417. CancellationToken ct)
  418. {
  419. var now = DateTime.Now;
  420. var updatedRouting = 0;
  421. if (routings != null)
  422. {
  423. foreach (var row in routings)
  424. {
  425. if (!TryGetInt(row, "OP", out var op) && !TryGetInt(row, "Op", out op))
  426. continue;
  427. var qtyComplete = GetDecimal(row, "QtyComplete");
  428. var qtyReject = GetDecimal(row, "QtyReject");
  429. var qtyScrap = GetNullableDecimal(row, "QtyScrap");
  430. var status = Trunc(GetString(row, "Status"), 1);
  431. updatedRouting += await _db.Ado.ExecuteCommandAsync(
  432. """
  433. UPDATE WorkOrdRouting
  434. SET QtyComplete = @QtyComplete,
  435. QtyReject = @QtyReject,
  436. QtyScrap = COALESCE(@QtyScrap, QtyScrap),
  437. Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END,
  438. UpdateUser = 'MDP_HOT',
  439. UpdateTime = @Now
  440. WHERE WorkOrd = @WorkOrd
  441. AND OP = @Op
  442. AND IFNULL(Domain, '') = @Domain
  443. AND IFNULL(tenant_id, 0) = @TenantId
  444. """,
  445. new SugarParameter("@QtyComplete", qtyComplete),
  446. new SugarParameter("@QtyReject", qtyReject),
  447. new SugarParameter("@QtyScrap", qtyScrap ?? (object)DBNull.Value),
  448. new SugarParameter("@Status", status ?? ""),
  449. new SugarParameter("@Now", now),
  450. new SugarParameter("@WorkOrd", workOrd),
  451. new SugarParameter("@Op", op),
  452. new SugarParameter("@Domain", domain),
  453. new SugarParameter("@TenantId", tenantId));
  454. }
  455. }
  456. if (masters is { Count: > 0 })
  457. {
  458. var m = masters[0];
  459. var qtyCompleted = GetDecimal(m, "QtyCompleted");
  460. var status = Trunc(GetString(m, "Status"), 8);
  461. await _db.Ado.ExecuteCommandAsync(
  462. """
  463. UPDATE WorkOrdMaster
  464. SET QtyCompleted = @QtyCompleted,
  465. Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END,
  466. UpdateUser = 'MDP_HOT',
  467. UpdateTime = @Now
  468. WHERE WorkOrd = @WorkOrd
  469. AND IFNULL(Domain, '') = @Domain
  470. AND IFNULL(tenant_id, 0) = @TenantId
  471. """,
  472. new SugarParameter("@QtyCompleted", qtyCompleted),
  473. new SugarParameter("@Status", status ?? ""),
  474. new SugarParameter("@Now", now),
  475. new SugarParameter("@WorkOrd", workOrd),
  476. new SugarParameter("@Domain", domain),
  477. new SugarParameter("@TenantId", tenantId));
  478. }
  479. _logger.LogInformation(
  480. "[MdpHotWatch] applied WORK_ORDER execution wo={WorkOrd} domain={Domain} routingRows={Rows}",
  481. workOrd, domain, updatedRouting);
  482. }
  483. /// <summary>
  484. /// D-M01:把 165 的工序实绩按 (Domain, WorkOrds, Op) 合并后落到 ado_psd_op_actual。
  485. ///
  486. /// 为什么按工序合并:MES APP 报工会新增 Period=0/IsActive=0 的影子行,实绩可能写在影子行、
  487. /// 也可能写在 Period=1 计划行,且同一工序可能因重排跨多个 Line。只有工序级求和才与
  488. /// WorkOrdRouting.QtyComplete(MES 自己的工序累计完成数)对齐。实测 Op501:3+5=8=QtyComplete。
  489. ///
  490. /// 为什么不写 PeriodSequenceDet:本库 PSD 行会被排产整表重建(旧行置 IsActive=0),
  491. /// 实绩写进去会在下次全量重排时丢失。
  492. ///
  493. /// 单位保持 165 原始口径:ActualTime 秒、SetupTime 小时(列名自带单位)。
  494. /// </summary>
  495. private async Task ApplyPeriodSequenceActualAsync(
  496. long tenantId,
  497. string domain,
  498. string workOrd,
  499. List<Dictionary<string, object>>? periodDets,
  500. List<Dictionary<string, object>>? routings,
  501. CancellationToken ct)
  502. {
  503. if (periodDets is not { Count: > 0 })
  504. return;
  505. // WorkOrdRouting.QtyComplete 作为交叉校验值(仅记录,不用于覆盖)
  506. var routingQty = new Dictionary<int, decimal>();
  507. if (routings != null)
  508. {
  509. foreach (var r in routings)
  510. {
  511. if (!TryGetInt(r, "OP", out var rop) && !TryGetInt(r, "Op", out rop))
  512. continue;
  513. routingQty[rop] = GetDecimal(r, "QtyComplete");
  514. }
  515. }
  516. var grouped = new Dictionary<int, PsdOpActual>();
  517. foreach (var row in periodDets)
  518. {
  519. if (!TryGetInt(row, "Op", out var op) && !TryGetInt(row, "OP", out op))
  520. continue;
  521. if (!grouped.TryGetValue(op, out var acc))
  522. {
  523. acc = new PsdOpActual { Op = op };
  524. grouped[op] = acc;
  525. }
  526. acc.CompQty += GetDecimal(row, "CompQty");
  527. acc.RejectQty += GetDecimal(row, "RejectQty");
  528. acc.SetupTimeHour += GetDecimal(row, "SetupTime");
  529. acc.ActualTimeSec += GetDecimal(row, "ActualTime");
  530. acc.SrcRowCount++;
  531. acc.ItemNum ??= GetString(row, "ItemNum");
  532. if (row.TryGetValue("UpdateTime", out var upd) && upd is DateTime dt
  533. && (acc.LastSrcUpdate == null || dt > acc.LastSrcUpdate))
  534. acc.LastSrcUpdate = dt;
  535. }
  536. var now = DateTime.Now;
  537. var mismatches = new List<string>();
  538. foreach (var acc in grouped.Values)
  539. {
  540. var rq = routingQty.TryGetValue(acc.Op, out var q) ? (decimal?)q : null;
  541. if (rq.HasValue && rq.Value != acc.CompQty)
  542. mismatches.Add($"Op{acc.Op}: psdSum={acc.CompQty} routing={rq.Value}");
  543. await _db.Ado.ExecuteCommandAsync(
  544. """
  545. INSERT INTO ado_psd_op_actual (
  546. tenant_id, domain, work_ord, op, item_num,
  547. comp_qty, reject_qty, setup_time_hour, actual_time_sec,
  548. src_row_count, routing_qty, last_src_update, sync_time
  549. ) VALUES (
  550. @TenantId, @Domain, @WorkOrd, @Op, @ItemNum,
  551. @CompQty, @RejectQty, @SetupTimeHour, @ActualTimeSec,
  552. @SrcRowCount, @RoutingQty, @LastSrcUpdate, @Now
  553. )
  554. ON DUPLICATE KEY UPDATE
  555. item_num = VALUES(item_num),
  556. comp_qty = VALUES(comp_qty),
  557. reject_qty = VALUES(reject_qty),
  558. setup_time_hour = VALUES(setup_time_hour),
  559. actual_time_sec = VALUES(actual_time_sec),
  560. src_row_count = VALUES(src_row_count),
  561. routing_qty = VALUES(routing_qty),
  562. last_src_update = VALUES(last_src_update),
  563. sync_time = VALUES(sync_time)
  564. """,
  565. new SugarParameter("@TenantId", tenantId),
  566. new SugarParameter("@Domain", domain ?? ""),
  567. new SugarParameter("@WorkOrd", workOrd),
  568. new SugarParameter("@Op", acc.Op),
  569. new SugarParameter("@ItemNum", acc.ItemNum ?? (object)DBNull.Value),
  570. new SugarParameter("@CompQty", acc.CompQty),
  571. new SugarParameter("@RejectQty", acc.RejectQty),
  572. new SugarParameter("@SetupTimeHour", acc.SetupTimeHour),
  573. new SugarParameter("@ActualTimeSec", acc.ActualTimeSec),
  574. new SugarParameter("@SrcRowCount", acc.SrcRowCount),
  575. new SugarParameter("@RoutingQty", rq.HasValue ? rq.Value : (object)DBNull.Value),
  576. new SugarParameter("@LastSrcUpdate", acc.LastSrcUpdate ?? (object)DBNull.Value),
  577. new SugarParameter("@Now", now));
  578. }
  579. if (mismatches.Count > 0)
  580. {
  581. // 不阻断落库:差异说明 165 侧口径需人工确认,先留证
  582. _logger.LogWarning(
  583. "[MdpHotWatch] PSD 实绩与 WorkOrdRouting 不一致 wo={WorkOrd} domain={Domain} {Detail}",
  584. workOrd, domain, string.Join("; ", mismatches));
  585. }
  586. _logger.LogInformation(
  587. "[MdpHotWatch] applied PSD actual wo={WorkOrd} domain={Domain} ops={Ops} srcRows={Rows}",
  588. workOrd, domain, grouped.Count, periodDets.Count);
  589. }
  590. /// <summary>
  591. /// 把 165 的报检单/分录/检验单 UPSERT 进本库业务表(S5 IQC P1)。
  592. /// 同名列 1:1 + 补 tenant_id;主键沿用 165 的 id(本库 id 非自增)。
  593. /// </summary>
  594. private async Task ApplyIqcAsync(
  595. long tenantId,
  596. string domain,
  597. string purOrd,
  598. List<Dictionary<string, object>>? applies,
  599. List<Dictionary<string, object>>? entries,
  600. List<Dictionary<string, object>>? bills,
  601. CancellationToken ct)
  602. {
  603. if (tenantId <= 0)
  604. {
  605. _logger.LogWarning(
  606. "[MdpHotWatch] ApplyIqc 跳过:无法解析租户 purOrd={PurOrd} domain={Domain}", purOrd, domain);
  607. return;
  608. }
  609. var nApply = 0;
  610. var nEntry = 0;
  611. var nBill = 0;
  612. if (applies != null)
  613. {
  614. foreach (var row in applies)
  615. {
  616. if (!TryGetInt64(row, "id", out var id) || id <= 0) continue;
  617. var billNo = GetString(row, "FBILLNO");
  618. if (string.IsNullOrWhiteSpace(billNo)) continue;
  619. nApply += await _db.Ado.ExecuteCommandAsync(
  620. """
  621. INSERT INTO qms_qcp_inspecapplyn (
  622. id, tenant_id, FBILLNO, FBILLTYPE, FBIZTYPE, FAPPLYTIME, FCOMMENT
  623. ) VALUES (
  624. @Id, @TenantId, @FBILLNO, @FBILLTYPE, @FBIZTYPE, @FAPPLYTIME, @FCOMMENT
  625. )
  626. ON DUPLICATE KEY UPDATE
  627. FBILLNO = VALUES(FBILLNO),
  628. FBILLTYPE = VALUES(FBILLTYPE),
  629. FBIZTYPE = VALUES(FBIZTYPE),
  630. FAPPLYTIME = VALUES(FAPPLYTIME),
  631. FCOMMENT = VALUES(FCOMMENT),
  632. tenant_id = IF(IFNULL(tenant_id,0)=0, VALUES(tenant_id), tenant_id)
  633. """,
  634. new SugarParameter("@Id", id),
  635. new SugarParameter("@TenantId", tenantId),
  636. new SugarParameter("@FBILLNO", billNo),
  637. new SugarParameter("@FBILLTYPE", GetString(row, "FBILLTYPE") ?? (object)DBNull.Value),
  638. new SugarParameter("@FBIZTYPE", GetString(row, "FBIZTYPE") ?? (object)DBNull.Value),
  639. new SugarParameter("@FAPPLYTIME", GetDateTime(row, "FAPPLYTIME") ?? (object)DBNull.Value),
  640. new SugarParameter("@FCOMMENT", GetString(row, "FCOMMENT") ?? (object)DBNull.Value));
  641. }
  642. }
  643. if (entries != null)
  644. {
  645. foreach (var row in entries)
  646. {
  647. if (!TryGetInt64(row, "id", out var id) || id <= 0) continue;
  648. // jyfzr/yxj 为 Ai-DOP 独占故不回读;状态只许前进,165 的旧值/空值不得打回
  649. nEntry += await _db.Ado.ExecuteCommandAsync(
  650. """
  651. INSERT INTO qms_qcp_insappnentry (
  652. id, tenant_id, glid, FSEQ, FMATERIALCFG, wlmc, ggxh, FLOTNUMBER,
  653. FSRCORDERTYPE, FSRCORDERNUM, FAPPLYQTY, FINSPECTSTATUS,
  654. jykssj, jywcsj, FWAREHOUSEID, FLOCATIONID, FSUPPLIER, gysbm, gysmc, shdh
  655. ) VALUES (
  656. @Id, @TenantId, @Glid, @FSEQ, @FMATERIALCFG, @Wlmc, @Ggxh, @FLOTNUMBER,
  657. @FSRCORDERTYPE, @FSRCORDERNUM, @FAPPLYQTY, @FINSPECTSTATUS,
  658. @Jykssj, @Jywcsj, @FWAREHOUSEID, @FLOCATIONID, @FSUPPLIER, @Gysbm, @Gysmc, @Shdh
  659. )
  660. ON DUPLICATE KEY UPDATE
  661. glid = VALUES(glid),
  662. FSEQ = VALUES(FSEQ),
  663. FMATERIALCFG = VALUES(FMATERIALCFG),
  664. wlmc = VALUES(wlmc),
  665. ggxh = VALUES(ggxh),
  666. FLOTNUMBER = VALUES(FLOTNUMBER),
  667. FSRCORDERTYPE = VALUES(FSRCORDERTYPE),
  668. FSRCORDERNUM = VALUES(FSRCORDERNUM),
  669. FAPPLYQTY = VALUES(FAPPLYQTY),
  670. FINSPECTSTATUS = IF(
  671. FIELD(VALUES(FINSPECTSTATUS), '未检验', '检验中', '检验完成') > 0
  672. AND FIELD(VALUES(FINSPECTSTATUS), '未检验', '检验中', '检验完成')
  673. >= FIELD(FINSPECTSTATUS, '未检验', '检验中', '检验完成'),
  674. VALUES(FINSPECTSTATUS), FINSPECTSTATUS),
  675. jykssj = IF(IFNULL(VALUES(jykssj),'')='', jykssj, VALUES(jykssj)),
  676. jywcsj = IF(IFNULL(VALUES(jywcsj),'')='', jywcsj, VALUES(jywcsj)),
  677. FWAREHOUSEID = VALUES(FWAREHOUSEID),
  678. FLOCATIONID = VALUES(FLOCATIONID),
  679. FSUPPLIER = VALUES(FSUPPLIER),
  680. gysbm = VALUES(gysbm),
  681. gysmc = VALUES(gysmc),
  682. shdh = VALUES(shdh),
  683. tenant_id = IF(IFNULL(tenant_id,0)=0, VALUES(tenant_id), tenant_id)
  684. """,
  685. new SugarParameter("@Id", id),
  686. new SugarParameter("@TenantId", tenantId),
  687. new SugarParameter("@Glid", GetInt64OrNull(row, "glid") ?? (object)DBNull.Value),
  688. new SugarParameter("@FSEQ", GetIntOrNull(row, "FSEQ") ?? (object)DBNull.Value),
  689. new SugarParameter("@FMATERIALCFG", GetString(row, "FMATERIALCFG") ?? (object)DBNull.Value),
  690. new SugarParameter("@Wlmc", GetString(row, "wlmc") ?? (object)DBNull.Value),
  691. new SugarParameter("@Ggxh", GetString(row, "ggxh") ?? (object)DBNull.Value),
  692. new SugarParameter("@FLOTNUMBER", GetString(row, "FLOTNUMBER") ?? (object)DBNull.Value),
  693. new SugarParameter("@FSRCORDERTYPE", GetString(row, "FSRCORDERTYPE") ?? (object)DBNull.Value),
  694. new SugarParameter("@FSRCORDERNUM", GetString(row, "FSRCORDERNUM") ?? (object)DBNull.Value),
  695. new SugarParameter("@FAPPLYQTY", GetDecimal(row, "FAPPLYQTY")),
  696. new SugarParameter("@FINSPECTSTATUS", GetString(row, "FINSPECTSTATUS") ?? (object)DBNull.Value),
  697. new SugarParameter("@Jykssj", GetString(row, "jykssj") ?? (object)DBNull.Value),
  698. new SugarParameter("@Jywcsj", GetString(row, "jywcsj") ?? (object)DBNull.Value),
  699. new SugarParameter("@FWAREHOUSEID", GetString(row, "FWAREHOUSEID") ?? (object)DBNull.Value),
  700. new SugarParameter("@FLOCATIONID", GetString(row, "FLOCATIONID") ?? (object)DBNull.Value),
  701. new SugarParameter("@FSUPPLIER", GetString(row, "FSUPPLIER") ?? (object)DBNull.Value),
  702. new SugarParameter("@Gysbm", GetString(row, "gysbm") ?? (object)DBNull.Value),
  703. new SugarParameter("@Gysmc", GetString(row, "gysmc") ?? (object)DBNull.Value),
  704. new SugarParameter("@Shdh", GetString(row, "shdh") ?? (object)DBNull.Value));
  705. }
  706. }
  707. if (bills != null)
  708. {
  709. foreach (var row in bills)
  710. {
  711. if (!TryGetInt64(row, "id", out var id) || id <= 0) continue;
  712. var billNo = GetString(row, "FBILLNO");
  713. if (string.IsNullOrWhiteSpace(billNo)) continue;
  714. nBill += await _db.Ado.ExecuteCommandAsync(
  715. """
  716. INSERT INTO qms_qcp_inspbill (
  717. id, tenant_id, FBILLNO, lydjbh, hid, FMATERIALCFG, wlmc, ggxh, gysmc, pch,
  718. FRINSQTY, jysl, dhsl, bhgsl, pd, clfs, FBILLSTATUS, jyr, FINSPECTORID,
  719. FINSPESTARTDATE, FCREATETIME
  720. ) VALUES (
  721. @Id, @TenantId, @FBILLNO, @Lydjbh, @Hid, @FMATERIALCFG, @Wlmc, @Ggxh, @Gysmc, @Pch,
  722. @FRINSQTY, @Jysl, @Dhsl, @Bhgsl, @Pd, @Clfs, @FBILLSTATUS, @Jyr, @FINSPECTORID,
  723. @FINSPESTARTDATE, @FCREATETIME
  724. )
  725. ON DUPLICATE KEY UPDATE
  726. FBILLNO = VALUES(FBILLNO),
  727. lydjbh = VALUES(lydjbh),
  728. hid = VALUES(hid),
  729. FMATERIALCFG = VALUES(FMATERIALCFG),
  730. wlmc = VALUES(wlmc),
  731. ggxh = VALUES(ggxh),
  732. gysmc = VALUES(gysmc),
  733. pch = VALUES(pch),
  734. FRINSQTY = VALUES(FRINSQTY),
  735. jysl = VALUES(jysl),
  736. dhsl = IF(VALUES(dhsl) IS NULL, dhsl, VALUES(dhsl)),
  737. bhgsl = IF(VALUES(bhgsl) IS NULL, bhgsl, VALUES(bhgsl)),
  738. pd = IF(VALUES(pd) IS NULL, pd, VALUES(pd)),
  739. clfs = IF(VALUES(clfs) IS NULL, clfs, VALUES(clfs)),
  740. FBILLSTATUS = IF(IFNULL(VALUES(FBILLSTATUS),'')='', FBILLSTATUS, VALUES(FBILLSTATUS)),
  741. jyr = IF(IFNULL(VALUES(jyr),'')='', jyr, VALUES(jyr)),
  742. FINSPECTORID = IF(IFNULL(VALUES(FINSPECTORID),'')='', FINSPECTORID, VALUES(FINSPECTORID)),
  743. FINSPESTARTDATE = IF(VALUES(FINSPESTARTDATE) IS NULL, FINSPESTARTDATE, VALUES(FINSPESTARTDATE)),
  744. tenant_id = IF(IFNULL(tenant_id,0)=0, VALUES(tenant_id), tenant_id)
  745. """,
  746. new SugarParameter("@Id", id),
  747. new SugarParameter("@TenantId", tenantId),
  748. new SugarParameter("@FBILLNO", billNo),
  749. new SugarParameter("@Lydjbh", GetString(row, "lydjbh") ?? (object)DBNull.Value),
  750. new SugarParameter("@Hid", GetInt64OrNull(row, "hid") ?? (object)DBNull.Value),
  751. new SugarParameter("@FMATERIALCFG", GetString(row, "FMATERIALCFG") ?? (object)DBNull.Value),
  752. new SugarParameter("@Wlmc", GetString(row, "wlmc") ?? (object)DBNull.Value),
  753. new SugarParameter("@Ggxh", GetString(row, "ggxh") ?? (object)DBNull.Value),
  754. new SugarParameter("@Gysmc", GetString(row, "gysmc") ?? (object)DBNull.Value),
  755. new SugarParameter("@Pch", GetString(row, "pch") ?? (object)DBNull.Value),
  756. new SugarParameter("@FRINSQTY", GetDecimal(row, "FRINSQTY")),
  757. new SugarParameter("@Jysl", GetDecimal(row, "jysl")),
  758. new SugarParameter("@Dhsl", GetDecimal(row, "dhsl")),
  759. new SugarParameter("@Bhgsl", GetDecimal(row, "bhgsl")),
  760. new SugarParameter("@Pd", GetInt64OrNull(row, "pd") ?? (object)DBNull.Value),
  761. new SugarParameter("@Clfs", GetInt64OrNull(row, "clfs") ?? (object)DBNull.Value),
  762. new SugarParameter("@FBILLSTATUS", GetString(row, "FBILLSTATUS") ?? (object)DBNull.Value),
  763. new SugarParameter("@Jyr", GetString(row, "jyr") ?? (object)DBNull.Value),
  764. new SugarParameter("@FINSPECTORID", GetString(row, "FINSPECTORID") ?? (object)DBNull.Value),
  765. new SugarParameter("@FINSPESTARTDATE", GetDateTime(row, "FINSPESTARTDATE") ?? (object)DBNull.Value),
  766. new SugarParameter("@FCREATETIME", GetDateTime(row, "FCREATETIME") ?? (object)DBNull.Value));
  767. }
  768. }
  769. if (nApply + nEntry + nBill > 0)
  770. {
  771. _logger.LogInformation(
  772. "[MdpHotWatch] applied IQC purOrd={PurOrd} domain={Domain} apply={A} entry={E} bill={B}",
  773. purOrd, domain, nApply, nEntry, nBill);
  774. }
  775. }
  776. /// <summary>
  777. /// 把 165 的采购收货结果回读本库:明细收货数、箱码状态,并据箱码推进送货单状态。
  778. ///
  779. /// 口径(2026-08-11 定):WMS 扫码收货只入待检仓(InvTransHist 落 rct-po-ins / Loc=1000),
  780. /// 不等于合格入库,所以这里只推进 scm_shd/scm_shdzb.shzt,不写 rksl——合格入库数留给检验环节。
  781. /// 165 的 scm_shdzb.rksl / shzt 在收货时并不更新,送货单进度只能由箱码状态反推。
  782. /// </summary>
  783. private async Task ApplyPurchaseReceiptAsync(
  784. long tenantId,
  785. string domain,
  786. string purOrd,
  787. List<Dictionary<string, object>>? purDetails,
  788. List<Dictionary<string, object>>? barcodes,
  789. CancellationToken ct)
  790. {
  791. var now = DateTime.Now;
  792. var updatedLines = 0;
  793. var updatedBarcodes = 0;
  794. if (purDetails != null)
  795. {
  796. foreach (var row in purDetails)
  797. {
  798. if (!TryGetInt(row, "Line", out var line))
  799. continue;
  800. // 自建单本库不填 Domain(165/MES 概念),所以空 Domain 也算命中;采购单号本库唯一
  801. updatedLines += await _db.Ado.ExecuteCommandAsync(
  802. """
  803. UPDATE PurOrdDetail
  804. SET RctQty = @RctQty,
  805. ReceiptQty = @ReceiptQty,
  806. QtyReturned = @QtyReturned,
  807. UpdateUser = 'MDP_HOT',
  808. UpdateTime = @Now
  809. WHERE PurOrd = @PurOrd
  810. AND Line = @Line
  811. AND IFNULL(Domain, '') IN ('', @Domain)
  812. AND (@TenantId = 0 OR IFNULL(tenant_id, 0) = @TenantId)
  813. """,
  814. new SugarParameter("@RctQty", GetDecimal(row, "RctQty")),
  815. new SugarParameter("@ReceiptQty", GetDecimal(row, "ReceiptQty")),
  816. new SugarParameter("@QtyReturned", GetDecimal(row, "QtyReturned")),
  817. new SugarParameter("@Now", now),
  818. new SugarParameter("@PurOrd", purOrd),
  819. new SugarParameter("@Line", line),
  820. new SugarParameter("@Domain", domain),
  821. new SugarParameter("@TenantId", tenantId));
  822. }
  823. }
  824. // 本库 MissedPrint 无 tenant_id,按 Domain + BarCode 定位;作废行 BarCode 带 RecID 前缀不会误匹配
  825. var shippers = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
  826. if (barcodes != null)
  827. {
  828. foreach (var row in barcodes)
  829. {
  830. var barCode = GetString(row, "BarCode")?.Trim();
  831. if (string.IsNullOrWhiteSpace(barCode))
  832. continue;
  833. updatedBarcodes += await _db.Ado.ExecuteCommandAsync(
  834. """
  835. UPDATE MissedPrint
  836. SET Status = CASE WHEN IFNULL(@Status, '') = '' THEN Status ELSE @Status END,
  837. RctNbr = @RctNbr,
  838. Location = @Location,
  839. Shelf = @Shelf,
  840. InvStatus = @InvStatus,
  841. UpdateUser = 'MDP_HOT',
  842. UpdateTime = @Now
  843. WHERE BarCode = @BarCode
  844. AND IFNULL(Domain, '') IN ('', @Domain)
  845. """,
  846. new SugarParameter("@Status", Trunc(GetString(row, "Status"), 4) ?? ""),
  847. new SugarParameter("@RctNbr", Trunc(GetString(row, "RctNbr"), 48) ?? ""),
  848. new SugarParameter("@Location", Trunc(GetString(row, "Location"), 24) ?? ""),
  849. new SugarParameter("@Shelf", Trunc(GetString(row, "Shelf"), 40) ?? ""),
  850. new SugarParameter("@InvStatus", Trunc(GetString(row, "InvStatus"), 24) ?? ""),
  851. new SugarParameter("@Now", now),
  852. new SugarParameter("@BarCode", barCode),
  853. new SugarParameter("@Domain", domain));
  854. var shipper = GetString(row, "ShipperNbr")?.Trim();
  855. if (!string.IsNullOrWhiteSpace(shipper))
  856. shippers.Add(shipper!);
  857. }
  858. }
  859. foreach (var shipper in shippers)
  860. await ApplyShipmentStatusAsync(tenantId, domain, shipper, ct);
  861. _logger.LogInformation(
  862. "[MdpHotWatch] applied PUR_ORDER receipt po={PurOrd} domain={Domain} lines={Lines} barcodes={Barcodes} shipments={Shipments}",
  863. purOrd, domain, updatedLines, updatedBarcodes, shippers.Count);
  864. }
  865. /// <summary>
  866. /// 按箱码待收数推进送货单状态:全部待收=待收,全部离开待收=完成,其余=收货中。
  867. ///
  868. /// 完成判据取箱码而非数量(决策:少收/超收不影响状态),列表读的是 scm_shd.shzt。
  869. /// </summary>
  870. private async Task ApplyShipmentStatusAsync(
  871. long tenantId, string domain, string shddh, CancellationToken ct)
  872. {
  873. var stats = await _db.Ado.SqlQueryAsync<ShipmentBarcodeStat>(
  874. """
  875. SELECT
  876. COUNT(*) AS Total,
  877. SUM(CASE WHEN IFNULL(Status, '') = 'U' THEN 1 ELSE 0 END) AS Pending
  878. FROM MissedPrint
  879. WHERE ShipperNbr = @Shddh
  880. AND IFNULL(Domain, '') = @Domain
  881. AND IFNULL(PurOrd, '') NOT LIKE '作废%'
  882. """,
  883. new SugarParameter("@Shddh", shddh),
  884. new SugarParameter("@Domain", domain));
  885. var stat = stats?.FirstOrDefault();
  886. if (stat == null || stat.Total <= 0)
  887. return;
  888. var shzt = stat.Pending >= stat.Total ? "待收" : (stat.Pending == 0 ? "完成" : "收货中");
  889. // 先定位表头再按主键更新:scm_shd.id 与 scm_shdzb.glid 排序规则不同(0900_ai_ci vs unicode_ci),
  890. // 直接 JOIN 会报 Illegal mix of collations,改由参数传 glid 规避
  891. var heads = await _db.Ado.SqlQueryAsync<ShipmentHeadRow>(
  892. """
  893. SELECT id AS Id
  894. FROM scm_shd
  895. WHERE shddh = @Shddh
  896. AND (@TenantId = 0 OR IFNULL(tenant_id, 0) = @TenantId)
  897. ORDER BY id
  898. LIMIT 1
  899. """,
  900. new SugarParameter("@Shddh", shddh),
  901. new SugarParameter("@TenantId", tenantId));
  902. var head = heads?.FirstOrDefault();
  903. if (head == null)
  904. {
  905. _logger.LogWarning(
  906. "[MdpHotWatch] 送货单在本库缺失,跳过状态推进 shddh={Shddh} tenant={Tenant}", shddh, tenantId);
  907. return;
  908. }
  909. await _db.Ado.ExecuteCommandAsync(
  910. "UPDATE scm_shd SET shzt = @Shzt WHERE id = @Id",
  911. new SugarParameter("@Shzt", shzt),
  912. new SugarParameter("@Id", head.Id));
  913. await _db.Ado.ExecuteCommandAsync(
  914. "UPDATE scm_shdzb SET shzt = @Shzt WHERE glid = @Glid",
  915. new SugarParameter("@Shzt", shzt),
  916. new SugarParameter("@Glid", head.Id.ToString()));
  917. }
  918. private sealed class ShipmentBarcodeStat
  919. {
  920. public int Total { get; set; }
  921. public int Pending { get; set; }
  922. }
  923. private sealed class ShipmentHeadRow
  924. {
  925. public long Id { get; set; }
  926. }
  927. private async Task<long> ResolveLocalPurOrdTenantAsync(string purOrd, string domain, CancellationToken ct)
  928. {
  929. var tid = await _db.Ado.SqlQuerySingleAsync<long?>(
  930. """
  931. SELECT IFNULL(tenant_id, 0)
  932. FROM PurOrdDetail
  933. WHERE PurOrd = @PurOrd AND IFNULL(Domain, '') IN ('', @Domain)
  934. ORDER BY Line
  935. LIMIT 1
  936. """,
  937. new SugarParameter("@PurOrd", purOrd),
  938. new SugarParameter("@Domain", domain));
  939. return tid ?? 0;
  940. }
  941. private sealed class PsdOpActual
  942. {
  943. public int Op { get; set; }
  944. public string? ItemNum { get; set; }
  945. public decimal CompQty { get; set; }
  946. public decimal RejectQty { get; set; }
  947. public decimal SetupTimeHour { get; set; }
  948. public decimal ActualTimeSec { get; set; }
  949. public int SrcRowCount { get; set; }
  950. public DateTime? LastSrcUpdate { get; set; }
  951. }
  952. private async Task<long> ResolveLocalWorkOrdTenantAsync(string workOrd, string domain, CancellationToken ct)
  953. {
  954. var tid = await _db.Ado.SqlQuerySingleAsync<long?>(
  955. """
  956. SELECT IFNULL(tenant_id, 0)
  957. FROM WorkOrdMaster
  958. WHERE WorkOrd = @WorkOrd AND IFNULL(Domain, '') = @Domain
  959. ORDER BY RecID DESC
  960. LIMIT 1
  961. """,
  962. new SugarParameter("@WorkOrd", workOrd),
  963. new SugarParameter("@Domain", domain));
  964. return tid ?? 0;
  965. }
  966. private static bool TryGetInt(Dictionary<string, object> row, string key, out int value)
  967. {
  968. value = 0;
  969. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return false;
  970. try
  971. {
  972. value = Convert.ToInt32(raw);
  973. return true;
  974. }
  975. catch
  976. {
  977. return false;
  978. }
  979. }
  980. private static bool TryGetInt64(Dictionary<string, object> row, string key, out long value)
  981. {
  982. value = 0;
  983. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return false;
  984. try
  985. {
  986. value = Convert.ToInt64(raw);
  987. return true;
  988. }
  989. catch
  990. {
  991. return false;
  992. }
  993. }
  994. private static long? GetInt64OrNull(Dictionary<string, object> row, string key)
  995. {
  996. if (!TryGetInt64(row, key, out var v)) return null;
  997. return v;
  998. }
  999. private static int? GetIntOrNull(Dictionary<string, object> row, string key)
  1000. {
  1001. if (!TryGetInt(row, key, out var v)) return null;
  1002. return v;
  1003. }
  1004. private static DateTime? GetDateTime(Dictionary<string, object> row, string key)
  1005. {
  1006. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return null;
  1007. if (raw is DateTime dt) return dt;
  1008. if (DateTime.TryParse(Convert.ToString(raw), out var parsed)) return parsed;
  1009. return null;
  1010. }
  1011. private static decimal GetDecimal(Dictionary<string, object> row, string key)
  1012. {
  1013. return GetNullableDecimal(row, key) ?? 0m;
  1014. }
  1015. private static decimal? GetNullableDecimal(Dictionary<string, object> row, string key)
  1016. {
  1017. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return null;
  1018. try { return Convert.ToDecimal(raw); }
  1019. catch { return null; }
  1020. }
  1021. private static string? GetString(Dictionary<string, object> row, string key)
  1022. {
  1023. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return null;
  1024. return Convert.ToString(raw);
  1025. }
  1026. private static string? Trunc(string? s, int max) =>
  1027. string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
  1028. private static async Task<bool> ShouldTerminateAsync(
  1029. ISqlSugarClient remote, AdoMdpHotWatch watch, CancellationToken ct)
  1030. {
  1031. if (string.Equals(watch.BizType, "PICK_BILL", StringComparison.OrdinalIgnoreCase))
  1032. {
  1033. var status = await remote.Ado.GetScalarAsync(
  1034. "SELECT TOP 1 Status FROM NbrMaster WHERE Domain=@d AND Nbr=@k",
  1035. new SugarParameter("@d", watch.Domain),
  1036. new SugarParameter("@k", watch.BizKey));
  1037. var s = status?.ToString()?.Trim() ?? "";
  1038. // 终态取值待 WP4 固化;临时:C/Complete/关闭 等常见值
  1039. return s is "C" or "Complete" or "关闭" or "Y";
  1040. }
  1041. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  1042. {
  1043. var status = await remote.Ado.GetScalarAsync(
  1044. "SELECT TOP 1 Status FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k",
  1045. new SugarParameter("@d", watch.Domain),
  1046. new SugarParameter("@k", watch.BizKey));
  1047. return string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase);
  1048. }
  1049. if (string.Equals(watch.BizType, "PUR_ORDER", StringComparison.OrdinalIgnoreCase))
  1050. {
  1051. var status = await remote.Ado.GetScalarAsync(
  1052. "SELECT TOP 1 Status FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k",
  1053. new SugarParameter("@d", watch.Domain),
  1054. new SugarParameter("@k", watch.BizKey));
  1055. if (string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase))
  1056. return true;
  1057. // S5 IQC:收货完成后还要等报检分录全部「检验完成」才终结(否则 IQC 回读窗口被掐断)
  1058. var total = Convert.ToInt32(await remote.Ado.GetScalarAsync(
  1059. "SELECT COUNT(*) FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k",
  1060. new SugarParameter("@d", watch.Domain),
  1061. new SugarParameter("@k", watch.BizKey)) ?? 0);
  1062. if (total > 0)
  1063. {
  1064. var pending = Convert.ToInt32(await remote.Ado.GetScalarAsync(
  1065. "SELECT COUNT(*) FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k AND IsNull(Status,'')='U'",
  1066. new SugarParameter("@d", watch.Domain),
  1067. new SugarParameter("@k", watch.BizKey)) ?? 0);
  1068. if (pending == 0)
  1069. {
  1070. var entryTotal = Convert.ToInt32(await remote.Ado.GetScalarAsync(
  1071. """
  1072. SELECT COUNT(*)
  1073. FROM qms_qcp_insappnentry e WITH (NOLOCK)
  1074. INNER JOIN PurOrdRctDetail d WITH (NOLOCK)
  1075. ON d.Receiver = e.FSRCORDERNUM AND d.ItemNum = e.FMATERIALCFG
  1076. WHERE d.Domain = @d AND d.OrdNbr = @k
  1077. """,
  1078. new SugarParameter("@d", watch.Domain),
  1079. new SugarParameter("@k", watch.BizKey)) ?? 0);
  1080. if (entryTotal <= 0)
  1081. return false; // 已收完但报检尚未出现:继续等
  1082. var entryPending = Convert.ToInt32(await remote.Ado.GetScalarAsync(
  1083. """
  1084. SELECT COUNT(*)
  1085. FROM qms_qcp_insappnentry e WITH (NOLOCK)
  1086. INNER JOIN PurOrdRctDetail d WITH (NOLOCK)
  1087. ON d.Receiver = e.FSRCORDERNUM AND d.ItemNum = e.FMATERIALCFG
  1088. WHERE d.Domain = @d AND d.OrdNbr = @k
  1089. AND IsNull(e.FINSPECTSTATUS,'') <> N'检验完成'
  1090. """,
  1091. new SugarParameter("@d", watch.Domain),
  1092. new SugarParameter("@k", watch.BizKey)) ?? 0);
  1093. if (entryPending == 0) return true;
  1094. }
  1095. }
  1096. }
  1097. // 兜底:超过 7 天强制终结
  1098. return watch.EnrollTime < DateTime.Now.AddDays(-7);
  1099. }
  1100. private static List<string> ParseTables(string? json)
  1101. {
  1102. if (string.IsNullOrWhiteSpace(json)) return new List<string>();
  1103. try
  1104. {
  1105. return JsonSerializer.Deserialize<List<string>>(json!) ?? new List<string>();
  1106. }
  1107. catch
  1108. {
  1109. return new List<string>();
  1110. }
  1111. }
  1112. private static string Sha256(string s)
  1113. {
  1114. var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s));
  1115. return Convert.ToHexString(bytes);
  1116. }
  1117. }