MdpOutboxEnqueueService.cs 3.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100
  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 已存在则覆盖 payload 并复位为待推(供工单快照重推)。
  36. /// 返回 true 表示新建或已刷新;false 仅在参数非法时不应出现。
  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 existingList = await _db.Queryable<MdpOutbox>()
  44. .Where(x => x.IdemKey == item.IdemKey)
  45. .OrderByDescending(x => x.Id)
  46. .Take(1)
  47. .ToListAsync(ct);
  48. var existing = existingList.FirstOrDefault();
  49. var now = DateTime.Now;
  50. if (existing != null)
  51. {
  52. existing.TargetSourceCode = item.TargetSourceCode;
  53. existing.ActionCode = item.ActionCode;
  54. existing.PayloadJson = item.PayloadJson;
  55. existing.TenantId = item.TenantId;
  56. existing.Status = 0;
  57. existing.RetryCount = 0;
  58. existing.NextRetryTime = null;
  59. existing.LastErrorCode = null;
  60. existing.ErrorMsg = null;
  61. existing.ResponseJson = null;
  62. existing.UpdateTime = now;
  63. await _db.Updateable(existing)
  64. .UpdateColumns(x => new
  65. {
  66. x.TargetSourceCode,
  67. x.ActionCode,
  68. x.PayloadJson,
  69. x.TenantId,
  70. x.Status,
  71. x.RetryCount,
  72. x.NextRetryTime,
  73. x.LastErrorCode,
  74. x.ErrorMsg,
  75. x.ResponseJson,
  76. x.UpdateTime
  77. })
  78. .ExecuteCommandAsync(ct);
  79. if (pulse) _wake.Pulse();
  80. return true;
  81. }
  82. item.Status = 0;
  83. item.RetryCount = 0;
  84. item.CreateTime = now;
  85. item.UpdateTime = now;
  86. await _db.Insertable(item).ExecuteCommandAsync(ct);
  87. if (pulse) _wake.Pulse();
  88. return true;
  89. }
  90. }