MdpHotWatchService.cs 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646
  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. private static bool TryParseWoIdem(string? idem, out string domain, out string workOrd)
  116. {
  117. domain = "8010";
  118. workOrd = "";
  119. if (string.IsNullOrWhiteSpace(idem)) return false;
  120. var parts = idem.Split('|', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
  121. if (parts.Length < 3) return false;
  122. if (!parts[0].Equals("wo", StringComparison.OrdinalIgnoreCase)
  123. && !parts[0].Equals("pick", StringComparison.OrdinalIgnoreCase))
  124. return false;
  125. domain = string.IsNullOrWhiteSpace(parts[1]) ? "8010" : parts[1];
  126. workOrd = parts[2];
  127. return !string.IsNullOrWhiteSpace(workOrd);
  128. }
  129. /// <summary>取一批在途行并轮询 165。</summary>
  130. public async Task<(int polled, int changed, int terminated)> PollOnceAsync(
  131. int take = 100, CancellationToken ct = default)
  132. {
  133. var due = await _db.Queryable<AdoMdpHotWatch>()
  134. .Where(x => x.Status == 0)
  135. .OrderBy(x => x.LastPollTime ?? DateTime.MinValue)
  136. .Take(take)
  137. .ToListAsync(ct);
  138. if (due.Count == 0) return (0, 0, 0);
  139. MdpSource? source;
  140. ISqlSugarClient remote;
  141. try
  142. {
  143. source = await _db.Queryable<MdpSource>()
  144. .Where(x => x.SourceCode == SourceCode && x.Status == 1)
  145. .FirstAsync(ct);
  146. if (source == null)
  147. {
  148. _logger.LogWarning("[MdpHotWatch] 源 {Source} 未启用,跳过本轮", SourceCode);
  149. return (0, 0, 0);
  150. }
  151. remote = await _scopeFactory.GetScopeAsync(SourceCode, ct);
  152. }
  153. catch (Exception ex)
  154. {
  155. _logger.LogWarning(ex, "[MdpHotWatch] 无法连接 165");
  156. return (0, 0, 0);
  157. }
  158. var changed = 0;
  159. var terminated = 0;
  160. var now = DateTime.Now;
  161. foreach (var group in due.GroupBy(x => x.BizType))
  162. {
  163. ct.ThrowIfCancellationRequested();
  164. foreach (var watch in group.Take(MaxKeysPerBatch))
  165. {
  166. var tables = ParseTables(watch.WatchTables);
  167. var sb = new StringBuilder();
  168. foreach (var table in tables)
  169. {
  170. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  171. foreach (var row in rows)
  172. sb.Append(JsonSerializer.Serialize(row));
  173. }
  174. var hash = Sha256(sb.ToString());
  175. watch.LastPollTime = now;
  176. watch.UpdateTime = now;
  177. if (!string.Equals(hash, watch.LastSnapshotHash, StringComparison.Ordinal))
  178. {
  179. watch.LastSnapshotHash = hash;
  180. changed++;
  181. var tableRows = new Dictionary<string, List<Dictionary<string, object>>>(StringComparer.OrdinalIgnoreCase);
  182. // 变更落地:按表写 stg(执行侧字段快照)
  183. foreach (var table in tables)
  184. {
  185. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  186. tableRows[table] = rows;
  187. var entity = await ResolveEntityAsync(table, ct);
  188. if (entity == null) continue;
  189. foreach (var row in rows)
  190. {
  191. var dict = row.ToDictionary(
  192. kv => kv.Key,
  193. kv => (object?)kv.Value,
  194. StringComparer.OrdinalIgnoreCase);
  195. var rid = dict.TryGetValue("RecID", out var r) ? $"{r}" : watch.BizKey;
  196. var raw = JsonSerializer.Serialize(dict);
  197. await _staging.UpsertAsync(
  198. source, entity, table, dict, raw, rid,
  199. new MdpPullContext
  200. {
  201. TenantId = watch.TenantId,
  202. BatchId = $"hot-{now:yyyyMMddHHmmss}",
  203. FullRefresh = false
  204. });
  205. }
  206. }
  207. // 工单执行量写回本库业务表,供看板直接读取
  208. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  209. {
  210. tableRows.TryGetValue("WorkOrdMaster", out var masters);
  211. tableRows.TryGetValue("WorkOrdRouting", out var routings);
  212. tableRows.TryGetValue("PeriodSequenceDet", out var periodDets);
  213. var effectiveTenantId = watch.TenantId;
  214. if (effectiveTenantId <= 0)
  215. effectiveTenantId = await ResolveLocalWorkOrdTenantAsync(watch.BizKey, watch.Domain, ct);
  216. await ApplyWorkOrderExecutionAsync(
  217. effectiveTenantId, watch.Domain, watch.BizKey, masters, routings, ct);
  218. await ApplyPeriodSequenceActualAsync(
  219. effectiveTenantId, watch.Domain, watch.BizKey, periodDets, routings, ct);
  220. }
  221. if (await ShouldTerminateAsync(remote, watch, ct))
  222. {
  223. watch.Status = 1;
  224. watch.TerminateTime = now;
  225. watch.TerminateReason = "auto";
  226. terminated++;
  227. }
  228. }
  229. await _db.Updateable(watch)
  230. .UpdateColumns(x => new
  231. {
  232. x.LastPollTime, x.LastSnapshotHash, x.Status,
  233. x.TerminateTime, x.TerminateReason, x.UpdateTime
  234. })
  235. .ExecuteCommandAsync(ct);
  236. }
  237. }
  238. return (due.Count, changed, terminated);
  239. }
  240. private async Task<MdpEntity?> ResolveEntityAsync(string table, CancellationToken ct)
  241. {
  242. return await _db.Queryable<MdpEntity>()
  243. .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1)
  244. .FirstAsync(ct);
  245. }
  246. private static async Task<List<Dictionary<string, object>>> QueryByBizKeyAsync(
  247. ISqlSugarClient remote, string table, string domain, string bizType, string bizKey, CancellationToken ct)
  248. {
  249. // 表白名单(防注入)
  250. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$"))
  251. throw new InvalidOperationException($"非法表名:{table}");
  252. string sql;
  253. SugarParameter[] pars;
  254. switch (table.ToUpperInvariant())
  255. {
  256. case "NBRMASTER":
  257. sql = "SELECT * FROM NbrMaster WHERE Domain=@d AND Nbr=@k";
  258. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  259. break;
  260. case "NBRDETAIL":
  261. sql = "SELECT * FROM NbrDetail WHERE Domain=@d AND Nbr=@k";
  262. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  263. break;
  264. case "WORKORDMASTER":
  265. sql = "SELECT * FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k";
  266. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  267. break;
  268. case "WORKORDROUTING":
  269. sql = "SELECT * FROM WorkOrdRouting WHERE Domain=@d AND WorkOrd=@k";
  270. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  271. break;
  272. case "PERIODSEQUENCEDET":
  273. // 工序间衔接:按工单取全部行(含 MES 报工产生的 Period=0 影子行),由调用方按工序合并
  274. sql = "SELECT * FROM PeriodSequenceDet WHERE Domain=@d AND WorkOrds=@k";
  275. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  276. break;
  277. case "PURORDDETAIL":
  278. sql = "SELECT * FROM PurOrdDetail WHERE Domain=@d AND PurOrd=@k";
  279. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  280. break;
  281. case "PURORDMASTER":
  282. sql = "SELECT * FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k";
  283. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  284. break;
  285. case "LINESTATUSDET":
  286. sql = "SELECT * FROM LineStatusDet WHERE Domain=@d AND Line=@k";
  287. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  288. break;
  289. case "MOBILETASK":
  290. // 堆表:按 TaskID
  291. sql = "SELECT * FROM MobileTask WHERE TaskID=@k";
  292. pars = new[] { new SugarParameter("@k", bizKey) };
  293. break;
  294. default:
  295. // MissedPrint 等:降级按 Domain + OrdNbr(可能扫表,见 WP8 E1)
  296. if (string.Equals(table, "MissedPrint", StringComparison.OrdinalIgnoreCase))
  297. {
  298. sql = "SELECT TOP 200 * FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k";
  299. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  300. break;
  301. }
  302. return new List<Dictionary<string, object>>();
  303. }
  304. var dt = await remote.Ado.GetDataTableAsync(sql, pars);
  305. var list = new List<Dictionary<string, object>>();
  306. foreach (System.Data.DataRow row in dt.Rows)
  307. {
  308. var dict = new Dictionary<string, object>(StringComparer.OrdinalIgnoreCase);
  309. foreach (System.Data.DataColumn col in dt.Columns)
  310. dict[col.ColumnName] = row[col] == DBNull.Value ? null! : row[col];
  311. list.Add(dict);
  312. }
  313. return list;
  314. }
  315. /// <summary>
  316. /// 把 165 工单执行字段回写本库(MES 权威列:Status / QtyCompleted / QtyComplete / QtyReject)。
  317. /// </summary>
  318. private async Task ApplyWorkOrderExecutionAsync(
  319. long tenantId,
  320. string domain,
  321. string workOrd,
  322. List<Dictionary<string, object>>? masters,
  323. List<Dictionary<string, object>>? routings,
  324. CancellationToken ct)
  325. {
  326. var now = DateTime.Now;
  327. var updatedRouting = 0;
  328. if (routings != null)
  329. {
  330. foreach (var row in routings)
  331. {
  332. if (!TryGetInt(row, "OP", out var op) && !TryGetInt(row, "Op", out op))
  333. continue;
  334. var qtyComplete = GetDecimal(row, "QtyComplete");
  335. var qtyReject = GetDecimal(row, "QtyReject");
  336. var status = Trunc(GetString(row, "Status"), 1);
  337. updatedRouting += await _db.Ado.ExecuteCommandAsync(
  338. """
  339. UPDATE WorkOrdRouting
  340. SET QtyComplete = @QtyComplete,
  341. QtyReject = @QtyReject,
  342. Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END,
  343. UpdateUser = 'MDP_HOT',
  344. UpdateTime = @Now
  345. WHERE WorkOrd = @WorkOrd
  346. AND OP = @Op
  347. AND IFNULL(Domain, '') = @Domain
  348. AND IFNULL(tenant_id, 0) = @TenantId
  349. """,
  350. new SugarParameter("@QtyComplete", qtyComplete),
  351. new SugarParameter("@QtyReject", qtyReject),
  352. new SugarParameter("@Status", status ?? ""),
  353. new SugarParameter("@Now", now),
  354. new SugarParameter("@WorkOrd", workOrd),
  355. new SugarParameter("@Op", op),
  356. new SugarParameter("@Domain", domain),
  357. new SugarParameter("@TenantId", tenantId));
  358. }
  359. }
  360. if (masters is { Count: > 0 })
  361. {
  362. var m = masters[0];
  363. var qtyCompleted = GetDecimal(m, "QtyCompleted");
  364. var status = Trunc(GetString(m, "Status"), 8);
  365. await _db.Ado.ExecuteCommandAsync(
  366. """
  367. UPDATE WorkOrdMaster
  368. SET QtyCompleted = @QtyCompleted,
  369. Status = CASE WHEN IFNULL(@Status,'') = '' THEN Status ELSE @Status END,
  370. UpdateUser = 'MDP_HOT',
  371. UpdateTime = @Now
  372. WHERE WorkOrd = @WorkOrd
  373. AND IFNULL(Domain, '') = @Domain
  374. AND IFNULL(tenant_id, 0) = @TenantId
  375. """,
  376. new SugarParameter("@QtyCompleted", qtyCompleted),
  377. new SugarParameter("@Status", status ?? ""),
  378. new SugarParameter("@Now", now),
  379. new SugarParameter("@WorkOrd", workOrd),
  380. new SugarParameter("@Domain", domain),
  381. new SugarParameter("@TenantId", tenantId));
  382. }
  383. _logger.LogInformation(
  384. "[MdpHotWatch] applied WORK_ORDER execution wo={WorkOrd} domain={Domain} routingRows={Rows}",
  385. workOrd, domain, updatedRouting);
  386. }
  387. /// <summary>
  388. /// D-M01:把 165 的工序实绩按 (Domain, WorkOrds, Op) 合并后落到 ado_psd_op_actual。
  389. ///
  390. /// 为什么按工序合并:MES APP 报工会新增 Period=0/IsActive=0 的影子行,实绩可能写在影子行、
  391. /// 也可能写在 Period=1 计划行,且同一工序可能因重排跨多个 Line。只有工序级求和才与
  392. /// WorkOrdRouting.QtyComplete(MES 自己的工序累计完成数)对齐。实测 Op501:3+5=8=QtyComplete。
  393. ///
  394. /// 为什么不写 PeriodSequenceDet:本库 PSD 行会被排产整表重建(旧行置 IsActive=0),
  395. /// 实绩写进去会在下次全量重排时丢失。
  396. ///
  397. /// 单位保持 165 原始口径:ActualTime 秒、SetupTime 小时(列名自带单位)。
  398. /// </summary>
  399. private async Task ApplyPeriodSequenceActualAsync(
  400. long tenantId,
  401. string domain,
  402. string workOrd,
  403. List<Dictionary<string, object>>? periodDets,
  404. List<Dictionary<string, object>>? routings,
  405. CancellationToken ct)
  406. {
  407. if (periodDets is not { Count: > 0 })
  408. return;
  409. // WorkOrdRouting.QtyComplete 作为交叉校验值(仅记录,不用于覆盖)
  410. var routingQty = new Dictionary<int, decimal>();
  411. if (routings != null)
  412. {
  413. foreach (var r in routings)
  414. {
  415. if (!TryGetInt(r, "OP", out var rop) && !TryGetInt(r, "Op", out rop))
  416. continue;
  417. routingQty[rop] = GetDecimal(r, "QtyComplete");
  418. }
  419. }
  420. var grouped = new Dictionary<int, PsdOpActual>();
  421. foreach (var row in periodDets)
  422. {
  423. if (!TryGetInt(row, "Op", out var op) && !TryGetInt(row, "OP", out op))
  424. continue;
  425. if (!grouped.TryGetValue(op, out var acc))
  426. {
  427. acc = new PsdOpActual { Op = op };
  428. grouped[op] = acc;
  429. }
  430. acc.CompQty += GetDecimal(row, "CompQty");
  431. acc.RejectQty += GetDecimal(row, "RejectQty");
  432. acc.SetupTimeHour += GetDecimal(row, "SetupTime");
  433. acc.ActualTimeSec += GetDecimal(row, "ActualTime");
  434. acc.SrcRowCount++;
  435. acc.ItemNum ??= GetString(row, "ItemNum");
  436. if (row.TryGetValue("UpdateTime", out var upd) && upd is DateTime dt
  437. && (acc.LastSrcUpdate == null || dt > acc.LastSrcUpdate))
  438. acc.LastSrcUpdate = dt;
  439. }
  440. var now = DateTime.Now;
  441. var mismatches = new List<string>();
  442. foreach (var acc in grouped.Values)
  443. {
  444. var rq = routingQty.TryGetValue(acc.Op, out var q) ? (decimal?)q : null;
  445. if (rq.HasValue && rq.Value != acc.CompQty)
  446. mismatches.Add($"Op{acc.Op}: psdSum={acc.CompQty} routing={rq.Value}");
  447. await _db.Ado.ExecuteCommandAsync(
  448. """
  449. INSERT INTO ado_psd_op_actual (
  450. tenant_id, domain, work_ord, op, item_num,
  451. comp_qty, reject_qty, setup_time_hour, actual_time_sec,
  452. src_row_count, routing_qty, last_src_update, sync_time
  453. ) VALUES (
  454. @TenantId, @Domain, @WorkOrd, @Op, @ItemNum,
  455. @CompQty, @RejectQty, @SetupTimeHour, @ActualTimeSec,
  456. @SrcRowCount, @RoutingQty, @LastSrcUpdate, @Now
  457. )
  458. ON DUPLICATE KEY UPDATE
  459. item_num = VALUES(item_num),
  460. comp_qty = VALUES(comp_qty),
  461. reject_qty = VALUES(reject_qty),
  462. setup_time_hour = VALUES(setup_time_hour),
  463. actual_time_sec = VALUES(actual_time_sec),
  464. src_row_count = VALUES(src_row_count),
  465. routing_qty = VALUES(routing_qty),
  466. last_src_update = VALUES(last_src_update),
  467. sync_time = VALUES(sync_time)
  468. """,
  469. new SugarParameter("@TenantId", tenantId),
  470. new SugarParameter("@Domain", domain ?? ""),
  471. new SugarParameter("@WorkOrd", workOrd),
  472. new SugarParameter("@Op", acc.Op),
  473. new SugarParameter("@ItemNum", acc.ItemNum ?? (object)DBNull.Value),
  474. new SugarParameter("@CompQty", acc.CompQty),
  475. new SugarParameter("@RejectQty", acc.RejectQty),
  476. new SugarParameter("@SetupTimeHour", acc.SetupTimeHour),
  477. new SugarParameter("@ActualTimeSec", acc.ActualTimeSec),
  478. new SugarParameter("@SrcRowCount", acc.SrcRowCount),
  479. new SugarParameter("@RoutingQty", rq.HasValue ? rq.Value : (object)DBNull.Value),
  480. new SugarParameter("@LastSrcUpdate", acc.LastSrcUpdate ?? (object)DBNull.Value),
  481. new SugarParameter("@Now", now));
  482. }
  483. if (mismatches.Count > 0)
  484. {
  485. // 不阻断落库:差异说明 165 侧口径需人工确认,先留证
  486. _logger.LogWarning(
  487. "[MdpHotWatch] PSD 实绩与 WorkOrdRouting 不一致 wo={WorkOrd} domain={Domain} {Detail}",
  488. workOrd, domain, string.Join("; ", mismatches));
  489. }
  490. _logger.LogInformation(
  491. "[MdpHotWatch] applied PSD actual wo={WorkOrd} domain={Domain} ops={Ops} srcRows={Rows}",
  492. workOrd, domain, grouped.Count, periodDets.Count);
  493. }
  494. private sealed class PsdOpActual
  495. {
  496. public int Op { get; set; }
  497. public string? ItemNum { get; set; }
  498. public decimal CompQty { get; set; }
  499. public decimal RejectQty { get; set; }
  500. public decimal SetupTimeHour { get; set; }
  501. public decimal ActualTimeSec { get; set; }
  502. public int SrcRowCount { get; set; }
  503. public DateTime? LastSrcUpdate { get; set; }
  504. }
  505. private async Task<long> ResolveLocalWorkOrdTenantAsync(string workOrd, string domain, CancellationToken ct)
  506. {
  507. var tid = await _db.Ado.SqlQuerySingleAsync<long?>(
  508. """
  509. SELECT IFNULL(tenant_id, 0)
  510. FROM WorkOrdMaster
  511. WHERE WorkOrd = @WorkOrd AND IFNULL(Domain, '') = @Domain
  512. ORDER BY RecID DESC
  513. LIMIT 1
  514. """,
  515. new SugarParameter("@WorkOrd", workOrd),
  516. new SugarParameter("@Domain", domain));
  517. return tid ?? 0;
  518. }
  519. private static bool TryGetInt(Dictionary<string, object> row, string key, out int value)
  520. {
  521. value = 0;
  522. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return false;
  523. try
  524. {
  525. value = Convert.ToInt32(raw);
  526. return true;
  527. }
  528. catch
  529. {
  530. return false;
  531. }
  532. }
  533. private static decimal GetDecimal(Dictionary<string, object> row, string key)
  534. {
  535. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return 0m;
  536. try { return Convert.ToDecimal(raw); }
  537. catch { return 0m; }
  538. }
  539. private static string? GetString(Dictionary<string, object> row, string key)
  540. {
  541. if (!row.TryGetValue(key, out var raw) || raw == null || raw is DBNull) return null;
  542. return Convert.ToString(raw);
  543. }
  544. private static string? Trunc(string? s, int max) =>
  545. string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
  546. private static async Task<bool> ShouldTerminateAsync(
  547. ISqlSugarClient remote, AdoMdpHotWatch watch, CancellationToken ct)
  548. {
  549. if (string.Equals(watch.BizType, "PICK_BILL", StringComparison.OrdinalIgnoreCase))
  550. {
  551. var status = await remote.Ado.GetScalarAsync(
  552. "SELECT TOP 1 Status FROM NbrMaster WHERE Domain=@d AND Nbr=@k",
  553. new SugarParameter("@d", watch.Domain),
  554. new SugarParameter("@k", watch.BizKey));
  555. var s = status?.ToString()?.Trim() ?? "";
  556. // 终态取值待 WP4 固化;临时:C/Complete/关闭 等常见值
  557. return s is "C" or "Complete" or "关闭" or "Y";
  558. }
  559. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  560. {
  561. var status = await remote.Ado.GetScalarAsync(
  562. "SELECT TOP 1 Status FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k",
  563. new SugarParameter("@d", watch.Domain),
  564. new SugarParameter("@k", watch.BizKey));
  565. return string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase);
  566. }
  567. // 兜底:超过 7 天强制终结
  568. return watch.EnrollTime < DateTime.Now.AddDays(-7);
  569. }
  570. private static List<string> ParseTables(string? json)
  571. {
  572. if (string.IsNullOrWhiteSpace(json)) return new List<string>();
  573. try
  574. {
  575. return JsonSerializer.Deserialize<List<string>>(json!) ?? new List<string>();
  576. }
  577. catch
  578. {
  579. return new List<string>();
  580. }
  581. }
  582. private static string Sha256(string s)
  583. {
  584. var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s));
  585. return Convert.ToHexString(bytes);
  586. }
  587. }