MdpTargetPushDispatcher.cs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.HotWatch;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform.Wms;
  3. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  4. using Microsoft.Extensions.Logging;
  5. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  6. /// <summary>
  7. /// 出站分发:扫描 mdp_outbox,按 source_type 路由到 API / DB 执行器。
  8. /// WP10 S4b / D12:退避 1m/5m/15m/30m/60m/120m,MaxRetry=6;配置类失败不烧重试。
  9. /// P-029:MarkAsync 按 UpdateTime 乐观锁,避免与 TryEnqueueOrRefreshAsync 并发整列覆盖。
  10. /// </summary>
  11. public sealed class MdpTargetPushDispatcher : ITransient
  12. {
  13. public const int DefaultTake = 500;
  14. private const int MaxRetry = 6;
  15. /// <summary>D12:按 retry_count 下标取退避分钟数(clamp 到末项)。</summary>
  16. private static readonly int[] BackoffMinutes = [1, 5, 15, 30, 60, 120];
  17. private static readonly string[] PermanentErrorPrefixes =
  18. [
  19. "PAYLOAD_INVALID",
  20. "CONSTRAINT_VIOLATION",
  21. "SQL_CONSTRAINT",
  22. "UNSUPPORTED_OP",
  23. ];
  24. private readonly ISqlSugarClient _db;
  25. private readonly MdpApiPushExecutor _api;
  26. private readonly MdpDbPushExecutor _dbPush;
  27. private readonly MdpHotWatchService _hotWatch;
  28. private readonly PickBillNbrSyncService _pickNbrSync;
  29. private readonly MdpOutboundGate _gate;
  30. private readonly ILogger _logger;
  31. public MdpTargetPushDispatcher(
  32. ISqlSugarClient db,
  33. MdpApiPushExecutor api,
  34. MdpDbPushExecutor dbPush,
  35. MdpHotWatchService hotWatch,
  36. PickBillNbrSyncService pickNbrSync,
  37. MdpOutboundGate gate,
  38. ILoggerFactory loggerFactory)
  39. {
  40. _db = db;
  41. _api = api;
  42. _dbPush = dbPush;
  43. _hotWatch = hotWatch;
  44. _pickNbrSync = pickNbrSync;
  45. _gate = gate;
  46. _logger = loggerFactory.CreateLogger(nameof(MdpTargetPushDispatcher));
  47. }
  48. public async Task<(int success, int failed, int skipped)> PushPendingAsync(
  49. int take = DefaultTake, CancellationToken cancellationToken = default)
  50. {
  51. // 总开关关闭时不取消息。计数留给真正要发送的执行器(手工重投仍会走到那里)。
  52. _gate.AnnounceIfDisabled();
  53. if (!_gate.IsEnabled)
  54. return (0, 0, 0);
  55. var now = DateTime.Now;
  56. var pending = await _db.Queryable<MdpOutbox>()
  57. .Where(x => x.Status == 0
  58. && x.RetryCount < MaxRetry
  59. && (x.NextRetryTime == null || x.NextRetryTime <= now))
  60. .OrderBy(x => x.Id)
  61. .Take(take > 0 ? take : DefaultTake)
  62. .ToListAsync(cancellationToken);
  63. var success = 0;
  64. var failed = 0;
  65. var skipped = 0;
  66. foreach (var item in pending)
  67. {
  68. cancellationToken.ThrowIfCancellationRequested();
  69. // 乐观锁基准:仅当 UpdateTime 未变时写入(期间被 refresh 则放弃)
  70. var loadedUpdateTime = item.UpdateTime;
  71. try
  72. {
  73. var sources = await _db.Queryable<MdpSource>()
  74. .Where(x => x.SourceCode == item.TargetSourceCode && x.Status == 1)
  75. .Take(1)
  76. .ToListAsync(cancellationToken);
  77. var source = sources.FirstOrDefault();
  78. if (source == null)
  79. {
  80. // 配置类失败:保持 pending,不烧 retry_count,便于配置修复后被 60s 兜底作业捡起
  81. await MarkConfigIssueAsync(
  82. item,
  83. "SOURCE_DISABLED",
  84. $"目标源 {item.TargetSourceCode} 未启用",
  85. cancellationToken);
  86. skipped++;
  87. continue;
  88. }
  89. if (!TryResolve(source, out var executor, out var routeError))
  90. {
  91. await MarkConfigIssueAsync(item, "SOURCE_ROUTE_UNRESOLVED", routeError!, cancellationToken);
  92. skipped++;
  93. continue;
  94. }
  95. _logger.LogInformation(
  96. "[MdpTargetPushDispatcher] outbox id={Id} source={Source} type={SourceType} via={Executor} action={Action}",
  97. item.Id, source.SourceCode, source.SourceType, executor!.SupportedType, item.ActionCode);
  98. var result = string.Equals(item.ActionCode, PickBillNbrSyncService.ActionCode, StringComparison.OrdinalIgnoreCase)
  99. ? await _pickNbrSync.ApplyFromOutboxAsync(item, cancellationToken)
  100. : await executor.PushAsync(source, item, cancellationToken);
  101. if (result.Held)
  102. {
  103. skipped++;
  104. continue;
  105. }
  106. if (result.Success)
  107. {
  108. var marked = await MarkAsync(
  109. item, 1, result.ResponseJson, null, null, null, loadedUpdateTime, cancellationToken);
  110. if (!marked)
  111. {
  112. item.Status = 0;
  113. skipped++;
  114. continue;
  115. }
  116. try
  117. {
  118. await _hotWatch.TryEnrollFromOutboxSuccessAsync(item, cancellationToken);
  119. }
  120. catch (Exception enrollEx)
  121. {
  122. _logger.LogWarning(
  123. enrollEx,
  124. "[MdpTargetPushDispatcher] hot-watch enroll failed outbox id={Id} action={Action}",
  125. item.Id, item.ActionCode);
  126. }
  127. success++;
  128. }
  129. else
  130. {
  131. await ApplyFailureAsync(
  132. item, result.ResponseJson, result.ErrorMessage, loadedUpdateTime, cancellationToken);
  133. if (item.Status == 2) failed++; else skipped++;
  134. }
  135. }
  136. catch (Exception ex)
  137. {
  138. await ApplyFailureAsync(item, null, ex.Message, loadedUpdateTime, cancellationToken);
  139. if (item.Status == 2) failed++; else skipped++;
  140. _logger.LogWarning(ex, "[MdpTargetPushDispatcher] outbox id={Id} failed", item.Id);
  141. }
  142. }
  143. return (success, failed, skipped);
  144. }
  145. private async Task ApplyFailureAsync(
  146. MdpOutbox item,
  147. string? responseJson,
  148. string? errorMessage,
  149. DateTime loadedUpdateTime,
  150. CancellationToken ct)
  151. {
  152. var msg = Truncate(errorMessage, 900);
  153. var (errorCode, permanent) = ClassifyFailure(msg);
  154. if (permanent)
  155. {
  156. // 立即死信,不递增 retry_count
  157. item.NextRetryTime = null;
  158. item.LastErrorCode = errorCode;
  159. var marked = await MarkAsync(item, 2, responseJson, msg, null, errorCode, loadedUpdateTime, ct);
  160. if (!marked) item.Status = 0;
  161. return;
  162. }
  163. item.RetryCount++;
  164. item.LastErrorCode = errorCode;
  165. if (item.RetryCount >= MaxRetry)
  166. {
  167. item.NextRetryTime = null;
  168. var marked = await MarkAsync(item, 2, responseJson, msg, null, errorCode, loadedUpdateTime, ct);
  169. if (!marked) item.Status = 0;
  170. }
  171. else
  172. {
  173. item.NextRetryTime = DateTime.Now.AddMinutes(GetBackoffMinutes(item.RetryCount));
  174. var marked = await MarkAsync(
  175. item, 0, responseJson, msg, item.NextRetryTime, errorCode, loadedUpdateTime, ct);
  176. if (!marked) item.Status = 0;
  177. }
  178. }
  179. /// <summary>
  180. /// 可重试 vs 不可重试。配置类走 <see cref="MarkConfigIssueAsync"/>,不进本方法。
  181. /// </summary>
  182. private static (string? errorCode, bool permanent) ClassifyFailure(string? error)
  183. {
  184. if (string.IsNullOrWhiteSpace(error))
  185. return ("UNKNOWN", false);
  186. var trimmed = error.Trim();
  187. foreach (var prefix in PermanentErrorPrefixes)
  188. {
  189. if (trimmed.StartsWith(prefix, StringComparison.OrdinalIgnoreCase))
  190. return (prefix, true);
  191. }
  192. // payload 反序列化 / 解析失败 → 永久
  193. if (trimmed.Contains("payload 解析失败", StringComparison.OrdinalIgnoreCase)
  194. || trimmed.Contains("payload parse", StringComparison.OrdinalIgnoreCase)
  195. || trimmed.Contains("JsonException", StringComparison.OrdinalIgnoreCase)
  196. || trimmed.Contains("反序列化", StringComparison.OrdinalIgnoreCase))
  197. return ("PAYLOAD_INVALID", true);
  198. // SQL 约束 / 唯一冲突 → 永久
  199. if (ContainsAny(trimmed,
  200. "unique constraint", "unique key", "duplicate key", "duplicate entry",
  201. "primary key", "violation of unique", "violation of primary",
  202. "cannot insert duplicate", "约束", "唯一索引", "主键冲突"))
  203. return ("CONSTRAINT_VIOLATION", true);
  204. if (trimmed.Contains("不支持的 op", StringComparison.OrdinalIgnoreCase))
  205. return ("UNSUPPORTED_OP", true);
  206. // 连接 / 超时 → 可重试
  207. if (ContainsAny(trimmed,
  208. "timeout", "timed out", "connection", "network", "could not open",
  209. "unable to connect", "transport-level", "broken pipe", "socket",
  210. "连接", "超时", "无法连接"))
  211. return ("TRANSIENT", false);
  212. return ("RETRYABLE", false);
  213. }
  214. private static bool ContainsAny(string haystack, params string[] needles)
  215. {
  216. foreach (var n in needles)
  217. {
  218. if (haystack.Contains(n, StringComparison.OrdinalIgnoreCase))
  219. return true;
  220. }
  221. return false;
  222. }
  223. /// <summary>retry_count 从 1 起对应 Backoff[0];越界 clamp 到末项。</summary>
  224. private static int GetBackoffMinutes(int retryCount)
  225. {
  226. var idx = Math.Max(0, retryCount - 1);
  227. if (idx >= BackoffMinutes.Length)
  228. idx = BackoffMinutes.Length - 1;
  229. return BackoffMinutes[idx];
  230. }
  231. /// <summary>
  232. /// 按 source_type / 连接信息解析执行器;无法解析时返回可诊断错误(不再静默落到 API)。
  233. /// </summary>
  234. private bool TryResolve(MdpSource source, out IMdpTargetPushExecutor? executor, out string? errorMessage)
  235. {
  236. var t = (source.SourceType ?? "").Trim().ToUpperInvariant();
  237. if (t is "DB" or "SQLSERVER" or "MYSQL" or "DATABASE")
  238. {
  239. executor = _dbPush;
  240. errorMessage = null;
  241. return true;
  242. }
  243. if (t == "API" || t == "HTTP")
  244. {
  245. executor = _api;
  246. errorMessage = null;
  247. return true;
  248. }
  249. // 未标注类型时:有库连接信息则走 DB,否则不再猜 API
  250. if (!string.IsNullOrWhiteSpace(source.DbHost) && !string.IsNullOrWhiteSpace(source.DbName))
  251. {
  252. executor = _dbPush;
  253. errorMessage = null;
  254. return true;
  255. }
  256. _logger.LogWarning(
  257. "[MdpTargetPushDispatcher] 路由未解析 source={Code} source_type='{Type}' DbHost='{Host}' DbName='{Name}' ApiBaseUrl='{Api}'",
  258. source.SourceCode, source.SourceType, source.DbHost, source.DbName, source.ApiBaseUrl);
  259. executor = null;
  260. errorMessage =
  261. $"SOURCE_ROUTE_UNRESOLVED: 源 {source.SourceCode} 既未标注 source_type,也无 DbHost/DbName";
  262. return false;
  263. }
  264. /// <summary>配置类问题:仅写 ErrorMsg/LastErrorCode/UpdateTime,status 保持 0,retry_count 不变。</summary>
  265. private async Task MarkConfigIssueAsync(MdpOutbox item, string errorCode, string error, CancellationToken ct)
  266. {
  267. item.ErrorMsg = Truncate(error, 900);
  268. item.LastErrorCode = Truncate(errorCode, 64);
  269. item.UpdateTime = DateTime.Now;
  270. await _db.Updateable(item)
  271. .UpdateColumns(x => new { x.ErrorMsg, x.LastErrorCode, x.UpdateTime })
  272. .ExecuteCommandAsync(ct);
  273. }
  274. /// <summary>
  275. /// 条件更新:仅当 UpdateTime 与加载时一致才写入,避免与入队 refresh 丢失更新。
  276. /// </summary>
  277. private async Task<bool> MarkAsync(
  278. MdpOutbox item,
  279. int status,
  280. string? responseJson,
  281. string? error,
  282. DateTime? nextRetryTime,
  283. string? lastErrorCode,
  284. DateTime loadedUpdateTime,
  285. CancellationToken ct)
  286. {
  287. var now = DateTime.Now;
  288. var resp = Truncate(responseJson, 4000);
  289. var code = Truncate(lastErrorCode, 64);
  290. var affected = await _db.Updateable<MdpOutbox>()
  291. .SetColumns(x => x.Status == status)
  292. .SetColumns(x => x.RetryCount == item.RetryCount)
  293. .SetColumns(x => x.ResponseJson == resp)
  294. .SetColumns(x => x.ErrorMsg == error)
  295. .SetColumns(x => x.NextRetryTime == nextRetryTime)
  296. .SetColumns(x => x.LastErrorCode == code)
  297. .SetColumns(x => x.UpdateTime == now)
  298. .Where(x => x.Id == item.Id && x.UpdateTime == loadedUpdateTime)
  299. .ExecuteCommandAsync(ct);
  300. if (affected <= 0)
  301. {
  302. _logger.LogInformation(
  303. "[MdpTargetPushDispatcher] 标记被跳过(期间已被重新入队)outbox id={Id} idem={Idem} wantStatus={Status}",
  304. item.Id, item.IdemKey, status);
  305. return false;
  306. }
  307. item.Status = status;
  308. item.ResponseJson = resp;
  309. item.ErrorMsg = error;
  310. item.NextRetryTime = nextRetryTime;
  311. item.LastErrorCode = code;
  312. item.UpdateTime = now;
  313. return true;
  314. }
  315. private static string? Truncate(string? s, int max) =>
  316. string.IsNullOrEmpty(s) ? s : (s.Length <= max ? s : s[..max]);
  317. }