using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; namespace GB5Shared.QueueReader { public class JobQueueEnqueuer : IJobQueueEnqueuer { private readonly IQueryExecutor _qe; private const string EnqueueSql = @" IF NOT EXISTS (SELECT 1 FROM TJOBQUEUE WHERE MESSAGEID = @MessageId AND TENANTID = @TenantId) BEGIN INSERT INTO TJOBQUEUE (TENANTID, JOBEXECUTIONID, MESSAGETYPE, PAYLOAD, PRIORITY, SCHEDULEDFOR, MAXRETRIES, MESSAGEID, CORRELATIONID, SOURCESERVICE, STATUS, CREATEDDATE) VALUES (@TenantId, @JobExecutionId, @MessageType, @Payload, @Priority, @ScheduledFor, @MaxRetries, @MessageId, @CorrelationId, @SourceService, 'PENDING', GETUTCDATE()) SELECT CAST(SCOPE_IDENTITY() AS BIGINT) END ELSE SELECT CAST(0 AS BIGINT)"; public JobQueueEnqueuer(IQueryExecutor queryExecutor) => _qe = queryExecutor; public async Task EnqueueAsync( string messageType, string payload, LoginDTO login, string? messageId = null, byte priority = 5, DateTime? scheduledFor = null, int maxRetries = 3, string? correlationId = null, string? sourceService = null, long? jobExecutionId = null, CancellationToken ct = default) { var effectiveMessageId = messageId ?? Guid.NewGuid().ToString("N"); return await _qe.ExecuteScalarAsync(login, EnqueueSql, new { TenantId = login.ClientId, JobExecutionId = jobExecutionId, MessageType = messageType, Payload = payload, Priority = priority, ScheduledFor = scheduledFor ?? DateTime.UtcNow, MaxRetries = maxRetries, MessageId = effectiveMessageId, CorrelationId = correlationId, SourceService = sourceService ?? "JobEngine" }, cancellationToken: ct).ConfigureAwait(false); } } }