namespace JobEngineDAL.Query { public static class JobQueueQB { public const string ENQUEUE = @" IF NOT EXISTS (SELECT 1 FROM TJOBQUEUE WHERE MESSAGEID = @MessageId AND TENANTID = @TenantId) BEGIN INSERT INTO TJOBQUEUE (TENANTID, MESSAGETYPE, PAYLOAD, PRIORITY, STATUS, SCHEDULEDFOR, MAXRETRIES, MESSAGEID, CORRELATIONID, SOURCESERVICE, CREATEDDATE) VALUES (@TenantId, @MessageType, @Payload, @Priority, 'PENDING', @ScheduledFor, @MaxRetries, @MessageId, @CorrelationId, @SourceService, GETUTCDATE()); SELECT SCOPE_IDENTITY(); END ELSE SELECT -1;"; // Atomic claim: UPDATE + OUTPUT in a single statement prevents double-claiming public const string CLAIM_BATCH = @" UPDATE TOP(@BatchSize) Q SET STATUS = 'PROCESSING', PROCESSINGSTARTED = GETUTCDATE(), UPDATEDDATE = GETUTCDATE() OUTPUT INSERTED.QUEUEID, INSERTED.TENANTID, INSERTED.MESSAGETYPE, INSERTED.PAYLOAD, INSERTED.PRIORITY, INSERTED.STATUS, INSERTED.SCHEDULEDFOR, INSERTED.MAXRETRIES, INSERTED.RETRYATTEMPT, INSERTED.NEXTATTEMPTAFTER, INSERTED.PROCESSINGSTARTED, INSERTED.MESSAGEID, INSERTED.CORRELATIONID, INSERTED.SOURCESERVICE, INSERTED.CREATEDDATE, INSERTED.ERRORDETAILS, INSERTED.UPDATEDDATE FROM TJOBQUEUE Q WITH (ROWLOCK, READPAST) WHERE Q.STATUS IN ('PENDING', 'FAILED') AND Q.TENANTID = @TenantId AND Q.MESSAGETYPE IN @MessageTypes AND Q.SCHEDULEDFOR <= GETUTCDATE() AND (Q.NEXTATTEMPTAFTER IS NULL OR Q.NEXTATTEMPTAFTER <= GETUTCDATE())"; public const string COMPLETE_ITEM = @" UPDATE TJOBQUEUE SET STATUS = 'COMPLETED', UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId"; public const string RETRY_ITEM = @" UPDATE TJOBQUEUE SET STATUS = 'FAILED', RETRYATTEMPT = RETRYATTEMPT + 1, NEXTATTEMPTAFTER = @NextAttemptAfter, ERRORDETAILS = @ErrorDetails, UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId"; public const string MOVE_TO_DLQ = @" UPDATE TJOBQUEUE SET STATUS = 'DLQ', ERRORDETAILS = @ErrorDetails, UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId"; public const string REQUEUE_ITEM = @" UPDATE TJOBQUEUE SET STATUS = 'PENDING', RETRYATTEMPT = 0, NEXTATTEMPTAFTER = NULL, ERRORDETAILS = NULL, SCHEDULEDFOR = GETUTCDATE(), UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId AND STATUS = 'DLQ'"; public const string CANCEL_ITEM = @" UPDATE TJOBQUEUE SET STATUS = 'CANCELLED', UPDATEDDATE = GETUTCDATE() WHERE QUEUEID = @QueueId AND TENANTID = @TenantId AND STATUS = 'PENDING'"; public const string PURGE_DLQ = @" DELETE FROM TJOBQUEUE WHERE STATUS = 'DLQ' AND TENANTID = @TenantId AND CREATEDDATE < DATEADD(DAY, -@DaysOld, GETUTCDATE())"; public const string GET_QUEUE_PAGED = @" SELECT QUEUEID, TENANTID, MESSAGETYPE, PRIORITY, STATUS, SCHEDULEDFOR, MAXRETRIES, RETRYATTEMPT, NEXTATTEMPTAFTER, PROCESSINGSTARTED, ERRORDETAILS, MESSAGEID, CORRELATIONID, SOURCESERVICE, CREATEDDATE, UPDATEDDATE FROM TJOBQUEUE WHERE TENANTID = @TenantId {WHERE_CLAUSE} ORDER BY CREATEDDATE DESC OFFSET @Offset ROWS FETCH NEXT @PageSize ROWS ONLY"; public const string GET_QUEUE_COUNT = @" SELECT COUNT(1) FROM TJOBQUEUE WHERE TENANTID = @TenantId {WHERE_CLAUSE}"; public const string GET_QUEUE_STATS = @" SELECT MESSAGETYPE, SUM(CASE WHEN STATUS = 'PENDING' THEN 1 ELSE 0 END) AS Pending, SUM(CASE WHEN STATUS = 'PROCESSING' THEN 1 ELSE 0 END) AS Processing, SUM(CASE WHEN STATUS = 'DLQ' THEN 1 ELSE 0 END) AS Dlq, SUM(CASE WHEN STATUS = 'FAILED' THEN 1 ELSE 0 END) AS Failed FROM TJOBQUEUE WHERE TENANTID = @TenantId GROUP BY MESSAGETYPE ORDER BY MESSAGETYPE"; public const string GET_DLQ_TOTAL = @" SELECT COUNT(1) FROM TJOBQUEUE WHERE STATUS = 'DLQ' AND TENANTID = @TenantId"; public const string GET_PENDING_TOTAL = @" SELECT COUNT(1) FROM TJOBQUEUE WHERE STATUS IN ('PENDING','FAILED') AND TENANTID = @TenantId"; } }