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; /// /// 热回读:按在途业务键对 165 做索引 seek 窄查询,变更落地到贴源层。 /// 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)); } /// Outbox 推送成功后登记在途关注。 public async Task EnrollAsync( string bizType, string bizKey, string domain, IEnumerable watchTables, long tenantId = 0, int pollIntervalSec = 5, CancellationToken ct = default) { var exists = await _db.Queryable() .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); } /// 取一批在途行并轮询 165。 public async Task<(int polled, int changed, int terminated)> PollOnceAsync( int take = 100, CancellationToken ct = default) { var due = await _db.Queryable() .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() .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 ResolveEntityAsync(string table, CancellationToken ct) { return await _db.Queryable() .Where(x => x.SourceTableName == table && x.EntityCode.EndsWith("_SQLSERVER") && x.Status == 1) .FirstAsync(ct); } private static async Task>> 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>(); } var dt = await remote.Ado.GetDataTableAsync(sql, pars); var list = new List>(); foreach (System.Data.DataRow row in dt.Rows) { var dict = new Dictionary(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 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 ParseTables(string? json) { if (string.IsNullOrWhiteSpace(json)) return new List(); try { return JsonSerializer.Deserialize>(json!) ?? new List(); } catch { return new List(); } } private static string Sha256(string s) { var bytes = SHA256.HashData(Encoding.UTF8.GetBytes(s)); return Convert.ToHexString(bytes); } }