| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286 |
- using System.Security.Cryptography;
- using System.Text;
- using System.Text.Json;
- using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
- using Admin.NET.Plugin.AiDOP.DataPlatform.Sequence;
- using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
- using Microsoft.Extensions.Logging;
- using SqlSugar;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.HotWatch;
- /// <summary>
- /// 热回读:按在途业务键对 165 做索引 seek 窄查询,变更落地到贴源层。
- /// </summary>
- public sealed class MdpHotWatchService : ITransient
- {
- public const string SourceCode = NbrSequenceService.DefaultSourceCode;
- private const int MaxKeysPerBatch = 500;
- private readonly ISqlSugarClient _db;
- private readonly MdpSourceScopeFactory _scopeFactory;
- private readonly MdpStagingWriter _staging;
- private readonly ILogger _logger;
- public MdpHotWatchService(
- ISqlSugarClient db,
- MdpSourceScopeFactory scopeFactory,
- MdpStagingWriter staging,
- ILoggerFactory loggerFactory)
- {
- _db = db;
- _scopeFactory = scopeFactory;
- _staging = staging;
- _logger = loggerFactory.CreateLogger(nameof(MdpHotWatchService));
- }
- /// <summary>Outbox 推送成功后登记在途关注。</summary>
- public async Task EnrollAsync(
- string bizType, string bizKey, string domain, IEnumerable<string> watchTables,
- long tenantId = 0, int pollIntervalSec = 5, CancellationToken ct = default)
- {
- var exists = await _db.Queryable<AdoMdpHotWatch>()
- .Where(x => x.TenantId == tenantId && x.BizType == bizType && x.BizKey == bizKey && x.Status == 0)
- .AnyAsync(ct);
- if (exists) return;
- var now = DateTime.Now;
- await _db.Insertable(new AdoMdpHotWatch
- {
- TenantId = tenantId,
- BizType = bizType,
- BizKey = bizKey,
- Domain = domain,
- WatchTables = JsonSerializer.Serialize(watchTables.ToArray()),
- EnrollTime = now,
- PollIntervalSec = pollIntervalSec,
- Status = 0,
- CreateTime = now,
- UpdateTime = now
- }).ExecuteCommandAsync(ct);
- }
- /// <summary>取一批在途行并轮询 165。</summary>
- public async Task<(int polled, int changed, int terminated)> PollOnceAsync(
- int take = 100, CancellationToken ct = default)
- {
- var due = await _db.Queryable<AdoMdpHotWatch>()
- .Where(x => x.Status == 0)
- .OrderBy(x => x.LastPollTime ?? DateTime.MinValue)
- .Take(take)
- .ToListAsync(ct);
- if (due.Count == 0) return (0, 0, 0);
- MdpSource? source;
- ISqlSugarClient remote;
- try
- {
- source = await _db.Queryable<MdpSource>()
- .Where(x => x.SourceCode == SourceCode && x.Status == 1)
- .FirstAsync(ct);
- if (source == null)
- {
- _logger.LogWarning("[MdpHotWatch] 源 {Source} 未启用,跳过本轮", SourceCode);
- return (0, 0, 0);
- }
- remote = await _scopeFactory.GetScopeAsync(SourceCode, ct);
- }
- catch (Exception ex)
- {
- _logger.LogWarning(ex, "[MdpHotWatch] 无法连接 165");
- return (0, 0, 0);
- }
- var changed = 0;
- var terminated = 0;
- var now = DateTime.Now;
- foreach (var group in due.GroupBy(x => x.BizType))
- {
- ct.ThrowIfCancellationRequested();
- foreach (var watch in group.Take(MaxKeysPerBatch))
- {
- var tables = ParseTables(watch.WatchTables);
- var sb = new StringBuilder();
- foreach (var table in tables)
- {
- var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
- foreach (var row in rows)
- sb.Append(JsonSerializer.Serialize(row));
- }
- var hash = Sha256(sb.ToString());
- watch.LastPollTime = now;
- watch.UpdateTime = now;
- if (!string.Equals(hash, watch.LastSnapshotHash, StringComparison.Ordinal))
- {
- watch.LastSnapshotHash = hash;
- changed++;
- // 变更落地:按表写 stg(执行侧字段快照)
- foreach (var table in tables)
- {
- var rows = await QueryByBizKeyAsync(remote, table, watch.Domain, watch.BizType, watch.BizKey, ct);
- var entity = await ResolveEntityAsync(table, ct);
- if (entity == null) continue;
- foreach (var row in rows)
- {
- var dict = row.ToDictionary(
- kv => kv.Key,
- kv => (object?)kv.Value,
- StringComparer.OrdinalIgnoreCase);
- var rid = dict.TryGetValue("RecID", out var r) ? $"{r}" : watch.BizKey;
- var raw = JsonSerializer.Serialize(dict);
- await _staging.UpsertAsync(
- source, entity, table, dict, raw, rid,
- new MdpPullContext
- {
- TenantId = watch.TenantId,
- BatchId = $"hot-{now:yyyyMMddHHmmss}",
- FullRefresh = false
- });
- }
- }
- if (await ShouldTerminateAsync(remote, watch, ct))
- {
- watch.Status = 1;
- watch.TerminateTime = now;
- watch.TerminateReason = "auto";
- terminated++;
- }
- }
- await _db.Updateable(watch)
- .UpdateColumns(x => new
- {
- x.LastPollTime, x.LastSnapshotHash, x.Status,
- x.TerminateTime, x.TerminateReason, x.UpdateTime
- })
- .ExecuteCommandAsync(ct);
- }
- }
- return (due.Count, changed, terminated);
- }
- private async Task<MdpEntity?> ResolveEntityAsync(string table, CancellationToken ct)
- {
- return await _db.Queryable<MdpEntity>()
- .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1)
- .FirstAsync(ct);
- }
- private static async Task<List<Dictionary<string, object>>> QueryByBizKeyAsync(
- ISqlSugarClient remote, string table, string domain, string bizType, string bizKey, CancellationToken ct)
- {
- // 表白名单(防注入)
- if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_]+$"))
- throw new InvalidOperationException($"非法表名:{table}");
- string sql;
- SugarParameter[] pars;
- switch (table.ToUpperInvariant())
- {
- case "NBRMASTER":
- sql = "SELECT * FROM NbrMaster WHERE Domain=@d AND Nbr=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "NBRDETAIL":
- sql = "SELECT * FROM NbrDetail WHERE Domain=@d AND Nbr=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "WORKORDMASTER":
- sql = "SELECT * FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "WORKORDROUTING":
- sql = "SELECT * FROM WorkOrdRouting WHERE Domain=@d AND WorkOrd=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "PURORDDETAIL":
- sql = "SELECT * FROM PurOrdDetail WHERE Domain=@d AND PurOrd=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "PURORDMASTER":
- sql = "SELECT * FROM PurOrdMaster WHERE Domain=@d AND PurOrd=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "LINESTATUSDET":
- sql = "SELECT * FROM LineStatusDet WHERE Domain=@d AND Line=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- case "MOBILETASK":
- // 堆表:按 TaskID
- sql = "SELECT * FROM MobileTask WHERE TaskID=@k";
- pars = new[] { new SugarParameter("@k", bizKey) };
- break;
- default:
- // MissedPrint 等:降级按 Domain + OrdNbr(可能扫表,见 WP8 E1)
- if (string.Equals(table, "MissedPrint", StringComparison.OrdinalIgnoreCase))
- {
- sql = "SELECT TOP 200 * FROM MissedPrint WHERE Domain=@d AND OrdNbr=@k";
- pars = new[] { new SugarParameter("@d", domain), new SugarParameter("@k", bizKey) };
- break;
- }
- return new List<Dictionary<string, object>>();
- }
- var dt = await remote.Ado.GetDataTableAsync(sql, pars);
- var list = new List<Dictionary<string, object>>();
- foreach (System.Data.DataRow row in dt.Rows)
- {
- var dict = new Dictionary<string, object>(StringComparer.OrdinalIgnoreCase);
- foreach (System.Data.DataColumn col in dt.Columns)
- dict[col.ColumnName] = row[col] == DBNull.Value ? null! : row[col];
- list.Add(dict);
- }
- return list;
- }
- private static async Task<bool> ShouldTerminateAsync(
- ISqlSugarClient remote, AdoMdpHotWatch watch, CancellationToken ct)
- {
- if (string.Equals(watch.BizType, "PICK_BILL", StringComparison.OrdinalIgnoreCase))
- {
- var status = await remote.Ado.GetScalarAsync(
- "SELECT TOP 1 Status FROM NbrMaster WHERE Domain=@d AND Nbr=@k",
- new SugarParameter("@d", watch.Domain),
- new SugarParameter("@k", watch.BizKey));
- var s = status?.ToString()?.Trim() ?? "";
- // 终态取值待 WP4 固化;临时:C/Complete/关闭 等常见值
- return s is "C" or "Complete" or "关闭" or "Y";
- }
- if (string.Equals(watch.BizType, "WORK_ORDER", StringComparison.OrdinalIgnoreCase))
- {
- var status = await remote.Ado.GetScalarAsync(
- "SELECT TOP 1 Status FROM WorkOrdMaster WHERE Domain=@d AND WorkOrd=@k",
- new SugarParameter("@d", watch.Domain),
- new SugarParameter("@k", watch.BizKey));
- return string.Equals(status?.ToString()?.Trim(), "C", StringComparison.OrdinalIgnoreCase);
- }
- // 兜底:超过 7 天强制终结
- return watch.EnrollTime < DateTime.Now.AddDays(-7);
- }
- private static List<string> ParseTables(string? json)
- {
- if (string.IsNullOrWhiteSpace(json)) return new List<string>();
- try
- {
- return JsonSerializer.Deserialize<List<string>>(json!) ?? new List<string>();
- }
- catch
- {
- return new List<string>();
- }
- }
- private static string Sha256(string s)
- {
- var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s));
- return Convert.ToHexString(bytes);
- }
- }
|