using GB5Shared.DTO.Framework.Login; using GB5Shared.Resource.Response; using GB5Shared.Telemetry; using JobEngineDAL.DTOs; using JobEngineDAL.Interfaces; using JobEngineBLL.Interfaces; using Microsoft.Extensions.Logging; using Newtonsoft.Json; namespace JobEngineBLL.Implementations { public class JobQueueBLL : IJobQueueBLL { private readonly IJobQueueDAL _dal; private readonly ILogger _logger; public JobQueueBLL(IJobQueueDAL dal, ILogger logger) { _dal = dal; _logger = logger; } public async Task EnqueueAsync(JobQueueItemDTO item, LoginDTO login, CancellationToken ct) { try { GB5Trace.Step("validate-enqueue", new { item.MessageType, item.MessageId }); if (string.IsNullOrWhiteSpace(item.MessageType)) throw new ArgumentException("MessageType is required"); if (string.IsNullOrWhiteSpace(item.Payload)) throw new ArgumentException("Payload is required"); if (string.IsNullOrWhiteSpace(item.MessageId)) throw new ArgumentException("MessageId (idempotency key) is required"); GB5Trace.Step("enqueue-item", new { item.MessageType, item.Priority }); var queueId = await _dal.EnqueueAsync(item, login, ct).ConfigureAwait(false); return $"{SuccessResponse.SaveSuccessMessage} {queueId}"; } catch (Exception ex) { GB5Trace.MarkFailed("enqueue-failed", ex); _logger.LogError(ex, "Enqueue failed for MessageType {Type} MessageId {Id}", item.MessageType, item.MessageId); throw; } } public async Task GetQueueAsync( string? status, string? messageType, int page, int pageSize, LoginDTO login, CancellationToken ct) { try { var (items, total) = await _dal.GetQueuePagedAsync( login.ClientId, status, messageType, page, pageSize, login, ct).ConfigureAwait(false); return JsonConvert.SerializeObject(new { Items = items, TotalCount = total, Page = page, PageSize = pageSize }); } catch (Exception ex) { GB5Trace.MarkFailed("get-queue-failed", ex); _logger.LogError(ex, "GetQueue failed for tenant {TenantId}", login.ClientId); throw; } } public async Task GetQueueStatsAsync(LoginDTO login, CancellationToken ct) { try { var stats = await _dal.GetQueueStatsAsync(login, ct).ConfigureAwait(false); return JsonConvert.SerializeObject(stats); } catch (Exception ex) { GB5Trace.MarkFailed("get-queue-stats-failed", ex); _logger.LogError(ex, "GetQueueStats failed for tenant {TenantId}", login.ClientId); throw; } } public async Task RequeueItemAsync(long queueId, LoginDTO login, CancellationToken ct) { try { GB5Trace.Step("requeue-item", new { queueId }); await _dal.RequeueItemAsync(queueId, login, ct).ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { GB5Trace.MarkFailed("requeue-item-failed", ex); _logger.LogError(ex, "RequeueItem failed for QueueId {Id}", queueId); throw; } } public async Task CancelItemAsync(long queueId, LoginDTO login, CancellationToken ct) { try { GB5Trace.Step("cancel-queue-item", new { queueId }); await _dal.CancelItemAsync(queueId, login, ct).ConfigureAwait(false); return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { GB5Trace.MarkFailed("cancel-queue-item-failed", ex); _logger.LogError(ex, "CancelItem failed for QueueId {Id}", queueId); throw; } } public async Task PurgeDLQAsync(int daysOld, LoginDTO login, CancellationToken ct) { try { GB5Trace.Step("purge-dlq", new { daysOld, TenantId = login.ClientId }); var purged = await _dal.PurgeDLQAsync(daysOld, login, ct).ConfigureAwait(false); _logger.LogInformation("PurgeDLQ removed {Count} items for tenant {TenantId}", purged, login.ClientId); return SuccessResponse.DeleteSuccessMessage; } catch (Exception ex) { GB5Trace.MarkFailed("purge-dlq-failed", ex); _logger.LogError(ex, "PurgeDLQ failed for tenant {TenantId}", login.ClientId); throw; } } } }