MdpHotWatchService.cs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286
  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. var exists = await _db.Queryable<AdoMdpHotWatch>()
  38. .Where(x => x.TenantId == tenantId && x.BizType == bizType && x.BizKey == bizKey && x.Status == 0)
  39. .AnyAsync(ct);
  40. if (exists) return;
  41. var now = DateTime.Now;
  42. await _db.Insertable(new AdoMdpHotWatch
  43. {
  44. TenantId = tenantId,
  45. BizType = bizType,
  46. BizKey = bizKey,
  47. Domain = domain,
  48. WatchTables = JsonSerializer.Serialize(watchTables.ToArray()),
  49. EnrollTime = now,
  50. PollIntervalSec = pollIntervalSec,
  51. Status = 0,
  52. CreateTime = now,
  53. UpdateTime = now
  54. }).ExecuteCommandAsync(ct);
  55. }
  56. /// <summary>取一批在途行并轮询 165。</summary>
  57. public async Task<(int polled, int changed, int terminated)> PollOnceAsync(
  58. int take = 100, CancellationToken ct = default)
  59. {
  60. var due = await _db.Queryable<AdoMdpHotWatch>()
  61. .Where(x => x.Status == 0)
  62. .OrderBy(x => x.LastPollTime ?? DateTime.MinValue)
  63. .Take(take)
  64. .ToListAsync(ct);
  65. if (due.Count == 0) return (0, 0, 0);
  66. MdpSource? source;
  67. ISqlSugarClient remote;
  68. try
  69. {
  70. source = await _db.Queryable<MdpSource>()
  71. .Where(x => x.SourceCode == SourceCode && x.Status == 1)
  72. .FirstAsync(ct);
  73. if (source == null)
  74. {
  75. _logger.LogWarning("[MdpHotWatch] 源 {Source} 未启用,跳过本轮", SourceCode);
  76. return (0, 0, 0);
  77. }
  78. remote = await _scopeFactory.GetScopeAsync(SourceCode, ct);
  79. }
  80. catch (Exception ex)
  81. {
  82. _logger.LogWarning(ex, "[MdpHotWatch] 无法连接 165");
  83. return (0, 0, 0);
  84. }
  85. var changed = 0;
  86. var terminated = 0;
  87. var now = DateTime.Now;
  88. foreach (var group in due.GroupBy(x => x.BizType))
  89. {
  90. ct.ThrowIfCancellationRequested();
  91. foreach (var watch in group.Take(MaxKeysPerBatch))
  92. {
  93. var tables = ParseTables(watch.WatchTables);
  94. var sb = new StringBuilder();
  95. foreach (var table in tables)
  96. {
  97. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  98. foreach (var row in rows)
  99. sb.Append(JsonSerializer.Serialize(row));
  100. }
  101. var hash = Sha256(sb.ToString());
  102. watch.LastPollTime = now;
  103. watch.UpdateTime = now;
  104. if (!string.Equals(hash, watch.LastSnapshotHash, StringComparison.Ordinal))
  105. {
  106. watch.LastSnapshotHash = hash;
  107. changed++;
  108. // 变更落地:按表写 stg(执行侧字段快照)
  109. foreach (var table in tables)
  110. {
  111. var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
  112. var entity = await ResolveEntityAsync(table, ct);
  113. if (entity == null) continue;
  114. foreach (var row in rows)
  115. {
  116. var dict = row.ToDictionary(
  117. kv => kv.Key,
  118. kv => (object?)kv.Value,
  119. StringComparer.OrdinalIgnoreCase);
  120. var rid = dict.TryGetValue("RecID", out var r) ? $"{r}" : watch.BizKey;
  121. var raw = JsonSerializer.Serialize(dict);
  122. await _staging.UpsertAsync(
  123. source, entity, table, dict, raw, rid,
  124. new MdpPullContext
  125. {
  126. TenantId = watch.TenantId,
  127. BatchId = $"hot-{now:yyyyMMddHHmmss}",
  128. FullRefresh = false
  129. });
  130. }
  131. }
  132. if (await ShouldTerminateAsync(remote, watch, ct))
  133. {
  134. watch.Status = 1;
  135. watch.TerminateTime = now;
  136. watch.TerminateReason = "auto";
  137. terminated++;
  138. }
  139. }
  140. await _db.Updateable(watch)
  141. .UpdateColumns(x => new
  142. {
  143. x.LastPollTime, x.LastSnapshotHash, x.Status,
  144. x.TerminateTime, x.TerminateReason, x.UpdateTime
  145. })
  146. .ExecuteCommandAsync(ct);
  147. }
  148. }
  149. return (due.Count, changed, terminated);
  150. }
  151. private async Task<MdpEntity?> ResolveEntityAsync(string table, CancellationToken ct)
  152. {
  153. return await _db.Queryable<MdpEntity>()
  154. .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1)
  155. .FirstAsync(ct);
  156. }
  157. private static async Task<List<Dictionary<string, object>>> QueryByBizKeyAsync(
  158. ISqlSugarClient remote, string table, string domain, string bizType, string bizKey, CancellationToken ct)
  159. {
  160. // 表白名单(防注入)
  161. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$"))
  162. throw new InvalidOperationException($"非法表名:{table}");
  163. string sql;
  164. SugarParameter[] pars;
  165. switch (table.ToUpperInvariant())
  166. {
  167. case "NBRMASTER":
  168. sql = "SELECT * FROM NbrMaster WHERE Domain=@d AND Nbr=@k";
  169. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  170. break;
  171. case "NBRDETAIL":
  172. sql = "SELECT * FROM NbrDetail WHERE Domain=@d AND Nbr=@k";
  173. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  174. break;
  175. case "WORKORDMASTER":
  176. sql = "SELECT * FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k";
  177. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  178. break;
  179. case "WORKORDROUTING":
  180. sql = "SELECT * FROM WorkOrdRouting WHERE Domain=@d AND WorkOrd=@k";
  181. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  182. break;
  183. case "PURORDDETAIL":
  184. sql = "SELECT * FROM PurOrdDetail WHERE Domain=@d AND PurOrd=@k";
  185. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  186. break;
  187. case "PURORDMASTER":
  188. sql = "SELECT * FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k";
  189. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  190. break;
  191. case "LINESTATUSDET":
  192. sql = "SELECT * FROM LineStatusDet WHERE Domain=@d AND Line=@k";
  193. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  194. break;
  195. case "MOBILETASK":
  196. // 堆表:按 TaskID
  197. sql = "SELECT * FROM MobileTask WHERE TaskID=@k";
  198. pars = new[] { new SugarParameter("@k", bizKey) };
  199. break;
  200. default:
  201. // MissedPrint 等:降级按 Domain + OrdNbr(可能扫表,见 WP8 E1)
  202. if (string.Equals(table, "MissedPrint", StringComparison.OrdinalIgnoreCase))
  203. {
  204. sql = "SELECT TOP 200 * FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k";
  205. pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
  206. break;
  207. }
  208. return new List<Dictionary<string, object>>();
  209. }
  210. var dt = await remote.Ado.GetDataTableAsync(sql, pars);
  211. var list = new List<Dictionary<string, object>>();
  212. foreach (System.Data.DataRow row in dt.Rows)
  213. {
  214. var dict = new Dictionary<string, object>(StringComparer.OrdinalIgnoreCase);
  215. foreach (System.Data.DataColumn col in dt.Columns)
  216. dict[col.ColumnName] = row[col] == DBNull.Value ? null! : row[col];
  217. list.Add(dict);
  218. }
  219. return list;
  220. }
  221. private static async Task<bool> ShouldTerminateAsync(
  222. ISqlSugarClient remote, AdoMdpHotWatch watch, CancellationToken ct)
  223. {
  224. if (string.Equals(watch.BizType, "PICK_BILL", StringComparison.OrdinalIgnoreCase))
  225. {
  226. var status = await remote.Ado.GetScalarAsync(
  227. "SELECT TOP 1 Status FROM NbrMaster WHERE Domain=@d AND Nbr=@k",
  228. new SugarParameter("@d", watch.Domain),
  229. new SugarParameter("@k", watch.BizKey));
  230. var s = status?.ToString()?.Trim() ?? "";
  231. // 终态取值待 WP4 固化;临时:C/Complete/关闭 等常见值
  232. return s is "C" or "Complete" or "关闭" or "Y";
  233. }
  234. if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
  235. {
  236. var status = await remote.Ado.GetScalarAsync(
  237. "SELECT TOP 1 Status FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k",
  238. new SugarParameter("@d", watch.Domain),
  239. new SugarParameter("@k", watch.BizKey));
  240. return string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase);
  241. }
  242. // 兜底:超过 7 天强制终结
  243. return watch.EnrollTime < DateTime.Now.AddDays(-7);
  244. }
  245. private static List<string> ParseTables(string? json)
  246. {
  247. if (string.IsNullOrWhiteSpace(json)) return new List<string>();
  248. try
  249. {
  250. return JsonSerializer.Deserialize<List<string>>(json!) ?? new List<string>();
  251. }
  252. catch
  253. {
  254. return new List<string>();
  255. }
  256. }
  257. private static string Sha256(string s)
  258. {
  259. var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s));
  260. return Convert.ToHexString(bytes);
  261. }
  262. }