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);
}
}