MdpOutboxEnqueueService.cs 3.1 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283
  1. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  3. /// <summary>Outbox 入队 + 立即唤醒推送。</summary>
  4. public sealed class MdpOutboxEnqueueService : ITransient
  5. {
  6. private readonly ISqlSugarClient _db;
  7. private readonly MdpOutboxWakeSignal _wake;
  8. public MdpOutboxEnqueueService(ISqlSugarClient db, MdpOutboxWakeSignal wake)
  9. {
  10. _db = db;
  11. _wake = wake;
  12. }
  13. /// <summary>
  14. /// 幂等入队:已存在 status∈(0,1) 的同 idem_key 则跳过。
  15. /// </summary>
  16. /// <param name="pulse">为 false 时不立即 Pulse(供调用方在本地事务提交后再唤醒,见 WP10 S1)。默认 true 保持既有行为。</param>
  17. public async Task<bool> TryEnqueueAsync(MdpOutbox item, CancellationToken ct = default, bool pulse = true)
  18. {
  19. if (item == null) throw new ArgumentNullException(nameof(item));
  20. if (string.IsNullOrWhiteSpace(item.IdemKey))
  21. throw new ArgumentException("idem_key 不能为空", nameof(item));
  22. var exists = await _db.Queryable<MdpOutbox>()
  23. .Where(x => x.IdemKey == item.IdemKey && (x.Status == 0 || x.Status == 1))
  24. .AnyAsync(ct);
  25. if (exists) return false;
  26. item.Status = 0;
  27. item.RetryCount = 0;
  28. item.CreateTime = DateTime.Now;
  29. item.UpdateTime = item.CreateTime;
  30. await _db.Insertable(item).ExecuteCommandAsync(ct);
  31. if (pulse) _wake.Pulse();
  32. return true;
  33. }
  34. /// <summary>
  35. /// 入队或刷新:同 idem_key 已存在则删除旧行、以新 id 重新入队(供工单快照重推)。
  36. /// 分发器按 id 升序消费,沿用旧 id 会让本轮的先后顺序被历史入队顺序倒置(P-033)。
  37. /// </summary>
  38. public async Task<bool> TryEnqueueOrRefreshAsync(MdpOutbox item, CancellationToken ct = default, bool pulse = true)
  39. {
  40. if (item == null) throw new ArgumentNullException(nameof(item));
  41. if (string.IsNullOrWhiteSpace(item.IdemKey))
  42. throw new ArgumentException("idem_key 不能为空", nameof(item));
  43. var now = DateTime.Now;
  44. item.Status = 0;
  45. item.RetryCount = 0;
  46. item.NextRetryTime = null;
  47. item.LastErrorCode = null;
  48. item.ErrorMsg = null;
  49. item.ResponseJson = null;
  50. item.CreateTime = now;
  51. item.UpdateTime = now;
  52. // uk_mdp_outbox_idem 为 idem_key 单列唯一索引,删除与插入须在同一事务内完成
  53. for (var attempt = 0; attempt < 2; attempt++)
  54. {
  55. item.Id = 0;
  56. var tran = await _db.Ado.UseTranAsync(async () =>
  57. {
  58. await _db.Ado.ExecuteCommandAsync(
  59. "DELETE FROM mdp_outbox WHERE idem_key = @IdemKey",
  60. new SugarParameter("@IdemKey", item.IdemKey));
  61. await _db.Insertable(item).ExecuteCommandAsync();
  62. });
  63. if (tran.IsSuccess)
  64. {
  65. if (pulse) _wake.Pulse();
  66. return true;
  67. }
  68. // 并发下另一线程可能已抢先插入同键行,重试一次
  69. }
  70. return false;
  71. }
  72. }