using Dapper; using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using JobEngineDAL.DTOs; using JobEngineDAL.Interfaces; using JobEngineDAL.Query; namespace JobEngineDAL.Implementations { public class JobQueueDAL : IJobQueueDAL { private readonly IQueryExecutor _qe; public JobQueueDAL(IQueryExecutor queryExecutor) => _qe = queryExecutor; public async Task EnqueueAsync(JobQueueItemDTO item, LoginDTO login, CancellationToken ct) => await _qe.ExecuteScalarAsync(login, JobQueueQB.ENQUEUE, new { TenantId = login.ClientId, item.MessageType, item.Payload, item.Priority, item.ScheduledFor, item.MaxRetries, item.MessageId, item.CorrelationId, item.SourceService }, cancellationToken: ct).ConfigureAwait(false); public async Task> ClaimBatchAsync( int tenantId, IReadOnlyList messageTypes, int batchSize, LoginDTO login, CancellationToken ct) { // Dapper needs an IEnumerable for IN clause — use DynamicParameters var param = new DynamicParameters(); param.Add("BatchSize", batchSize); param.Add("TenantId", tenantId); param.Add("MessageTypes", messageTypes); return await _qe.QueryAsync(login, JobQueueQB.CLAIM_BATCH, param, cancellationToken: ct).ConfigureAwait(false); } public async Task CompleteAsync(long queueId, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobQueueQB.COMPLETE_ITEM, new { QueueId = queueId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task RetryAsync(long queueId, TimeSpan delay, string errorDetails, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobQueueQB.RETRY_ITEM, new { QueueId = queueId, NextAttemptAfter = DateTime.UtcNow.Add(delay), ErrorDetails = errorDetails, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task MoveToDLQAsync(long queueId, string errorDetails, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobQueueQB.MOVE_TO_DLQ, new { QueueId = queueId, ErrorDetails = errorDetails, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task RequeueItemAsync(long queueId, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobQueueQB.REQUEUE_ITEM, new { QueueId = queueId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task CancelItemAsync(long queueId, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobQueueQB.CANCEL_ITEM, new { QueueId = queueId, TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task PurgeDLQAsync(int daysOld, LoginDTO login, CancellationToken ct) => await _qe.ExecuteAsync(login, JobQueueQB.PURGE_DLQ, new { TenantId = login.ClientId, DaysOld = daysOld }, cancellationToken: ct).ConfigureAwait(false); public async Task<(IEnumerable Items, int TotalCount)> GetQueuePagedAsync( int tenantId, string? status, string? messageType, int page, int pageSize, LoginDTO login, CancellationToken ct) { var whereClause = ""; if (!string.IsNullOrWhiteSpace(status)) whereClause += " AND STATUS = @Status"; if (!string.IsNullOrWhiteSpace(messageType)) whereClause += " AND MESSAGETYPE = @MessageType"; var sql = JobQueueQB.GET_QUEUE_PAGED.Replace("{WHERE_CLAUSE}", whereClause); var countSql = JobQueueQB.GET_QUEUE_COUNT.Replace("{WHERE_CLAUSE}", whereClause); var param = new { TenantId = tenantId, Status = status, MessageType = messageType, Offset = (page - 1) * pageSize, PageSize = pageSize }; var items = await _qe.QueryAsync(login, sql, param, cancellationToken: ct).ConfigureAwait(false); var total = await _qe.ExecuteScalarAsync(login, countSql, new { TenantId = tenantId, Status = status, MessageType = messageType }, cancellationToken: ct).ConfigureAwait(false); return (items, total); } public async Task> GetQueueStatsAsync(LoginDTO login, CancellationToken ct) => await _qe.QueryAsync(login, JobQueueQB.GET_QUEUE_STATS, new { TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task GetDLQTotalAsync(LoginDTO login, CancellationToken ct) => await _qe.ExecuteScalarAsync(login, JobQueueQB.GET_DLQ_TOTAL, new { TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); public async Task GetPendingTotalAsync(LoginDTO login, CancellationToken ct) => await _qe.ExecuteScalarAsync(login, JobQueueQB.GET_PENDING_TOTAL, new { TenantId = login.ClientId }, cancellationToken: ct).ConfigureAwait(false); } }