| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100 |
- using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
- /// <summary>Outbox 入队 + 立即唤醒推送。</summary>
- public sealed class MdpOutboxEnqueueService : ITransient
- {
- private readonly ISqlSugarClient _db;
- private readonly MdpOutboxWakeSignal _wake;
- public MdpOutboxEnqueueService(ISqlSugarClient db, MdpOutboxWakeSignal wake)
- {
- _db = db;
- _wake = wake;
- }
- /// <summary>
- /// 幂等入队:已存在 status∈(0,1) 的同 idem_key 则跳过。
- /// </summary>
- /// <param name="pulse">为 false 时不立即 Pulse(供调用方在本地事务提交后再唤醒,见 WP10 S1)。默认 true 保持既有行为。</param>
- public async Task<bool> 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<MdpOutbox>()
- .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;
- }
- /// <summary>
- /// 入队或刷新:同 idem_key 已存在则覆盖 payload 并复位为待推(供工单快照重推)。
- /// 返回 true 表示新建或已刷新;false 仅在参数非法时不应出现。
- /// </summary>
- public async Task<bool> 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<MdpOutbox>()
- .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;
- }
- }
|