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;
}
}