using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors; /// Outbox 入队 + 立即唤醒推送。 public sealed class MdpOutboxEnqueueService : ITransient { private readonly ISqlSugarClient _db; private readonly MdpOutboxWakeSignal _wake; public MdpOutboxEnqueueService(ISqlSugarClient db, MdpOutboxWakeSignal wake) { _db = db; _wake = wake; } /// /// 幂等入队:已存在 status∈(0,1) 的同 idem_key 则跳过。 /// /// 为 false 时不立即 Pulse(供调用方在本地事务提交后再唤醒,见 WP10 S1)。默认 true 保持既有行为。 public async Task TryEnqueueAsync(MdpOutbox item, CancellationToken ct = default, bool pulse = true) { if (item == null) throw new ArgumentNullException(nameof(item)); if (string.IsNullOrWhiteSpace(item.IdemKey)) throw new ArgumentException("idem_key 不能为空", nameof(item)); var exists = await _db.Queryable() .Where(x => x.IdemKey == item.IdemKey && (x.Status == 0 || x.Status == 1)) .AnyAsync(ct); if (exists) return false; item.Status = 0; item.RetryCount = 0; item.CreateTime = DateTime.Now; item.UpdateTime = item.CreateTime; await _db.Insertable(item).ExecuteCommandAsync(ct); if (pulse) _wake.Pulse(); return true; } /// /// 入队或刷新:同 idem_key 已存在则删除旧行、以新 id 重新入队(供工单快照重推)。 /// 分发器按 id 升序消费,沿用旧 id 会让本轮的先后顺序被历史入队顺序倒置(P-033)。 /// public async Task TryEnqueueOrRefreshAsync(MdpOutbox item, CancellationToken ct = default, bool pulse = true) { if (item == null) throw new ArgumentNullException(nameof(item)); if (string.IsNullOrWhiteSpace(item.IdemKey)) throw new ArgumentException("idem_key 不能为空", nameof(item)); var now = DateTime.Now; item.Status = 0; item.RetryCount = 0; item.NextRetryTime = null; item.LastErrorCode = null; item.ErrorMsg = null; item.ResponseJson = null; item.CreateTime = now; item.UpdateTime = now; // uk_mdp_outbox_idem 为 idem_key 单列唯一索引,删除与插入须在同一事务内完成 for (var attempt = 0; attempt < 2; attempt++) { item.Id = 0; var tran = await _db.Ado.UseTranAsync(async () => { await _db.Ado.ExecuteCommandAsync( "DELETE FROM mdp_outbox WHERE idem_key = @IdemKey", new SugarParameter("@IdemKey", item.IdemKey)); await _db.Insertable(item).ExecuteCommandAsync(); }); if (tran.IsSuccess) { if (pulse) _wake.Pulse(); return true; } // 并发下另一线程可能已抢先插入同键行,重试一次 } return false; } }