using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; namespace GB5Shared.ActionProcessor { public interface IActionOutboxDAL { Task InsertAsync(ActionOutboxDTO dto, LoginDTO login, CancellationToken ct = default); Task> GetPendingBatchAsync( int batchSize, LoginDTO login, CancellationToken ct = default); Task MarkSentAsync(int outboxId, LoginDTO login, CancellationToken ct = default); Task MarkFailedAsync(int outboxId, LoginDTO login, CancellationToken ct = default); } public sealed class ActionOutboxDAL : IActionOutboxDAL { private readonly IQueryExecutor _qe; public ActionOutboxDAL(IQueryExecutor qe) => _qe = qe; public async Task InsertAsync(ActionOutboxDTO dto, LoginDTO login, CancellationToken ct = default) { await _qe.ExecuteAsync( login, ActionOutboxQB.INSERT_OUTBOX, new { dto.ActionRunId, dto.DestinationTopic, dto.Payload, dto.CorrelationKey, CreatedById = login.UserId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } public async Task> GetPendingBatchAsync( int batchSize, LoginDTO login, CancellationToken ct = default) { return (await _qe.QueryAsync( login, ActionOutboxQB.GET_PENDING_BATCH, new { BatchSize = batchSize, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false)).ToList(); } public async Task MarkSentAsync(int outboxId, LoginDTO login, CancellationToken ct = default) { await _qe.ExecuteAsync( login, ActionOutboxQB.UPDATE_SENT, new { OutboxId = outboxId, ModifiedById = login.UserId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } public async Task MarkFailedAsync(int outboxId, LoginDTO login, CancellationToken ct = default) { await _qe.ExecuteAsync( login, ActionOutboxQB.UPDATE_ATTEMPTS_FAILED, new { OutboxId = outboxId, ModifiedById = login.UserId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } } }