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 已存在则覆盖 payload 并复位为待推(供工单快照重推)。 /// 返回 true 表示新建或已刷新;false 仅在参数非法时不应出现。 /// 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 existingList = await _db.Queryable() .Where(x => x.IdemKey == item.IdemKey) .OrderByDescending(x => x.Id) .Take(1) .ToListAsync(ct); var existing = existingList.FirstOrDefault(); var now = DateTime.Now; if (existing != null) { existing.TargetSourceCode = item.TargetSourceCode; existing.ActionCode = item.ActionCode; existing.PayloadJson = item.PayloadJson; existing.TenantId = item.TenantId; existing.Status = 0; existing.RetryCount = 0; existing.NextRetryTime = null; existing.LastErrorCode = null; existing.ErrorMsg = null; existing.ResponseJson = null; existing.UpdateTime = now; await _db.Updateable(existing) .UpdateColumns(x => new { x.TargetSourceCode, x.ActionCode, x.PayloadJson, x.TenantId, x.Status, x.RetryCount, x.NextRetryTime, x.LastErrorCode, x.ErrorMsg, x.ResponseJson, x.UpdateTime }) .ExecuteCommandAsync(ct); if (pulse) _wake.Pulse(); return true; } item.Status = 0; item.RetryCount = 0; item.CreateTime = now; item.UpdateTime = now; await _db.Insertable(item).ExecuteCommandAsync(ct); if (pulse) _wake.Pulse(); return true; } }