diff --git a/RedisMessagePipeline.Demo/Program.cs b/RedisMessagePipeline.Demo/Program.cs index 74b7dac..c5d16d6 100644 --- a/RedisMessagePipeline.Demo/Program.cs +++ b/RedisMessagePipeline.Demo/Program.cs @@ -7,7 +7,7 @@ using RedLockNet.SERedis.Configuration; using StackExchange.Redis; -ConnectionMultiplexer redis = ConnectionMultiplexer.Connect("localhost:6379"); +ConnectionMultiplexer redis = ConnectionMultiplexer.Connect(""); RedLockMultiplexer lockMultiplexer = new RedLockMultiplexer(redis); IDatabase db = redis.GetDatabase(); @@ -15,8 +15,8 @@ RedLockFactory lockFactory = RedLockFactory.Create(new List { lockMultiplexer }); RedisPipelineFactory factory = new RedisPipelineFactory(loggerFactory, lockFactory, db); -var consumer = factory.CreateConsumer(new MyMessageHandler(), new RedisPipelineConsumerSettings("my-messages")); -var admin = factory.CreateAdmin(new RedisPipelineAdminSettings("my-messages")); +var consumer = factory.CreateConsumer(new MyMessageHandler(), new RedisPipelineConsumerSettings("my-messages") { Type = EnPipelineType.QUEUE_SCHEDULE}); +var admin = factory.CreateAdmin(new RedisPipelineAdminSettings("my-messages") { Type = EnPipelineType.QUEUE_SCHEDULE }); // ---- Administrate the pipeline ----- @@ -26,11 +26,11 @@ // Push messages for (int i = 0; i < 10; i++) { - await admin.PushAsync($"message:{i}"); + await admin.AddScheduleAsync($"{i}", DateTime.Now.AddSeconds(i * 10), $"message:{i}"); } // Resume the pipeline, skipping problematic messages if necessary -await admin.ResumeAsync(1, CancellationToken.None); +await admin.ResumeAsync(0, CancellationToken.None); // ------ Start the consumer to process messages ------ @@ -42,7 +42,7 @@ class MyMessageHandler : IRedisPipelineHandler public async Task HandleAsync(RedisValue redisValue, CancellationToken cancellationToken) { bool success = Random.Shared.Next(0, 100) > 50; - Console.WriteLine($"Processing task: {redisValue}, success: {success}"); + Console.WriteLine($"{DateTime.Now} Processing task: {redisValue}, success: {success}"); await Task.Delay(300, cancellationToken); return success; } diff --git a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs index 7c911ec..9e2f0e5 100644 --- a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs +++ b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs @@ -1,4 +1,5 @@ using StackExchange.Redis; +using System; using System.Threading; using System.Threading.Tasks; @@ -6,7 +7,8 @@ namespace RedisMessagePipeline.Admin { public interface IRedisPipelineAdmin { - Task PushAsync(RedisValue redisValue); + Task PushQueueAsync(RedisValue redisValue); + Task AddScheduleAsync(RedisValue keyValue, DateTimeOffset schedule, RedisValue redisValue); Task StopAsync(); Task CleanAsync(CancellationToken cancellationToken); Task ResumeAsync(int skip, CancellationToken cancellationToken); diff --git a/RedisMessagePipeline/Admin/RedisPipelineAdminSettings.cs b/RedisMessagePipeline/Admin/RedisPipelineAdminSettings.cs index cb89909..c3cbd63 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineAdminSettings.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineAdminSettings.cs @@ -1,4 +1,6 @@ -namespace RedisMessagePipeline.Admin +using RedisMessagePipeline.Factory; + +namespace RedisMessagePipeline.Admin { /// /// Settings for RedisPipelineAdmin to manage Redis pipelines. @@ -10,6 +12,7 @@ public RedisPipelineAdminSettings(string resource) Resource = resource; } public string Resource { get; set; } + public EnPipelineType Type { get; set; } = EnPipelineType.QUEUE; public RedisPipelineLockSettings LockSettings { get; set; } = new RedisPipelineLockSettings(); } } diff --git a/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs new file mode 100644 index 0000000..0803eb5 --- /dev/null +++ b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs @@ -0,0 +1,114 @@ +using Microsoft.Extensions.Logging; +using RedLockNet; +using StackExchange.Redis; +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; + +namespace RedisMessagePipeline.Admin +{ + /// + /// Admin functionality to manage operations on the Redis pipeline, such as starting, stopping, and cleaning. + /// + public abstract class RedisPipelineBaseAdmin : IRedisPipelineAdmin + { + protected readonly ILogger logger; + protected readonly RedisPipelineAdminSettings settings; + protected readonly IDistributedLockFactory lockFactory; + protected readonly IDatabase database; + + internal RedisPipelineBaseAdmin( + ILogger logger, + RedisPipelineAdminSettings settings, + IDistributedLockFactory lockFactory, + IDatabase database) + { + this.logger = logger; + this.settings = settings; + this.lockFactory = lockFactory; + this.database = database; + } + + /// + /// Pushes a new message to the Redis pipeline. + /// + public virtual Task PushQueueAsync(RedisValue redisValue) + { + if (this.settings.Type != Factory.EnPipelineType.QUEUE) + { + throw new Exception("This queue is configured for scheduling; use AddScheduleAsync."); + } + + logger.LogDebug("Push a new message '{message}' to '{resource}' redis pipeline", redisValue, settings.Resource); + + return Task.FromResult(-1L); + + } + + /// + /// Stops the Redis pipeline. + /// + public Task StopAsync() + { + logger.LogDebug("Redis pipeline '{resource}' has been stopped", settings.Resource); + + RedisKey key = RedisPipelineExtensions.StateKey(settings.Resource); + return database.StringSetAsync(key, RedisPipelineExtensions.STATE_STOPPED); + } + + /// + /// Cleans up resources used by the Redis pipeline. + /// + public abstract Task CleanAsync(CancellationToken cancellationToken); + + /// + /// Resumes operations of the Redis pipeline after a stop. + /// + public abstract Task ResumeAsync(int skip, CancellationToken cancellationToken); + + public virtual Task AddScheduleAsync(RedisValue keyValue, DateTimeOffset schedule, RedisValue redisValue) + { + if (this.settings.Type != Factory.EnPipelineType.QUEUE_SCHEDULE) + { + throw new Exception("The queue is not configured for scheduling; use PushQueueAsync."); + } + + logger.LogDebug("Redis pipeline '{resource}' has been stopped", settings.Resource); + + return Task.CompletedTask; + } + + protected async Task RemoveByPatternInBatchesAsync(string pattern, int batchSize = 500) + { + var endpoints = database.Multiplexer.GetEndPoints(); + long totalRemovidas = 0; + + foreach (var endpoint in endpoints) + { + var server = database.Multiplexer.GetServer(endpoint); + + if (!server.IsConnected || server.IsReplica) + continue; + + var batch = new List(batchSize); + + foreach (var key in server.Keys(pattern: pattern, pageSize: batchSize)) + { + batch.Add(key); + if (batch.Count >= batchSize) + { + totalRemovidas += await database.KeyDeleteAsync(batch.ToArray()); + batch.Clear(); + } + } + if (batch.Count > 0) + { + totalRemovidas += await database.KeyDeleteAsync(batch.ToArray()); + } + } + + return totalRemovidas; + } + } +} diff --git a/RedisMessagePipeline/Admin/RedisPipelineAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs similarity index 68% rename from RedisMessagePipeline/Admin/RedisPipelineAdmin.cs rename to RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs index 3aa3913..91f758a 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs @@ -10,51 +10,32 @@ namespace RedisMessagePipeline.Admin /// /// Admin functionality to manage operations on the Redis pipeline, such as starting, stopping, and cleaning. /// - public class RedisPipelineAdmin : IRedisPipelineAdmin + public class RedisPipelineQueueAdmin : RedisPipelineBaseAdmin { - private readonly ILogger logger; - private readonly RedisPipelineAdminSettings settings; - private readonly IDistributedLockFactory lockFactory; - private readonly IDatabase database; - - internal RedisPipelineAdmin( - ILogger logger, + internal RedisPipelineQueueAdmin( + ILogger logger, RedisPipelineAdminSettings settings, IDistributedLockFactory lockFactory, IDatabase database) + : base(logger, settings, lockFactory, database) { - this.logger = logger; - this.settings = settings; - this.lockFactory = lockFactory; - this.database = database; - } - /// - /// Pushes a new message to the Redis pipeline. - /// - public Task PushAsync(RedisValue redisValue) - { - logger.LogDebug("Push a new message '{message}' to '{resource}' redis pipeline", redisValue, settings.Resource); - - RedisKey key = RedisPipelineExtensions.MessagesKey(settings.Resource); - return database.ListRightPushAsync(key, redisValue); } /// - /// Stops the Redis pipeline. + /// Pushes a new message to the Redis pipeline. /// - public Task StopAsync() + public override async Task PushQueueAsync(RedisValue redisValue) { - logger.LogDebug("Redis pipeline '{resource}' has been stopped", settings.Resource); - - RedisKey key = RedisPipelineExtensions.StateKey(settings.Resource); - return database.StringSetAsync(key, RedisPipelineExtensions.STATE_STOPPED); + await base.PushQueueAsync(redisValue); + RedisKey key = RedisPipelineExtensions.MessagesListKey(settings.Resource); + return await database.ListRightPushAsync(key, redisValue); } /// /// Cleans up resources used by the Redis pipeline. /// - public async Task CleanAsync(CancellationToken cancellationToken) + public override async Task CleanAsync(CancellationToken cancellationToken) { using (IRedLock locker = await lockFactory.CreateLockAsync( resource: settings.Resource, @@ -73,7 +54,7 @@ await database.KeyDeleteAsync(new RedisKey[] { RedisPipelineExtensions.FailureKey(settings.Resource), RedisPipelineExtensions.StateKey(settings.Resource), - RedisPipelineExtensions.MessagesKey(settings.Resource), + RedisPipelineExtensions.MessagesListKey(settings.Resource), }); logger.LogDebug("Redis pipeline '{resource}' has been cleaned up", settings.Resource); @@ -83,7 +64,7 @@ await database.KeyDeleteAsync(new RedisKey[] /// /// Resumes operations of the Redis pipeline after a stop. /// - public async Task ResumeAsync(int skip, CancellationToken cancellationToken) + public override async Task ResumeAsync(int skip, CancellationToken cancellationToken) { using (IRedLock locker = await lockFactory.CreateLockAsync( resource: settings.Resource, @@ -106,10 +87,10 @@ public async Task ResumeAsync(int skip, CancellationToken cancellationToken) } ITransaction transaction = database.CreateTransaction(); - RedisKey messagesKey = RedisPipelineExtensions.MessagesKey(settings.Resource); + RedisKey messagesKey = RedisPipelineExtensions.MessagesListKey(settings.Resource); RedisKey stateKey = RedisPipelineExtensions.StateKey(settings.Resource); Task[] transactionTasks = new Task[] { - transaction.ListLeftPopAsync(messagesKey, count: skip), + skip > 0 ? transaction.ListLeftPopAsync(messagesKey, count: skip) : Task.CompletedTask, transaction.StringSetAsync(stateKey, 0) }; await transaction.ExecuteAsync(); diff --git a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs new file mode 100644 index 0000000..b1fb88f --- /dev/null +++ b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs @@ -0,0 +1,124 @@ +using Microsoft.Extensions.Logging; +using RedLockNet; +using StackExchange.Redis; +using System; +using System.Security.Cryptography; +using System.Threading; +using System.Threading.Tasks; + +namespace RedisMessagePipeline.Admin +{ + /// + /// Admin functionality to manage operations on the Redis pipeline, such as starting, stopping, and cleaning. + /// + public class RedisPipelineScheduleAdmin : RedisPipelineBaseAdmin + { + internal RedisPipelineScheduleAdmin( + ILogger logger, + RedisPipelineAdminSettings settings, + IDistributedLockFactory lockFactory, + IDatabase database) + : base(logger, settings, lockFactory, database) + { + + } + + public override async Task AddScheduleAsync(RedisValue keyValue, DateTimeOffset schedule, RedisValue redisValue) + { + await base.AddScheduleAsync(keyValue, schedule, redisValue); + double score = schedule.ToUnixTimeMilliseconds(); + RedisKey key = RedisPipelineExtensions.MessagesSortKey(settings.Resource); + await database.SortedSetAddAsync(key, keyValue, score); + key = RedisPipelineExtensions.MessageKey(settings.Resource, keyValue); + await database.StringSetAsync(key, redisValue); + } + + + /// + /// Cleans up resources used by the Redis pipeline. + /// + public override async Task CleanAsync(CancellationToken cancellationToken) + { + using (IRedLock locker = await lockFactory.CreateLockAsync( + resource: settings.Resource, + expiryTime: settings.LockSettings.ExpiryTime, + waitTime: settings.LockSettings.WaitTime, + retryTime: settings.LockSettings.RetryTime, + cancellationToken)) + { + if (!locker.IsAcquired) + { + logger.LogError("Cannot acquire redlock for Redis pipeline '{resource}'", settings.Resource); + throw new InvalidOperationException("Cannot acquire redlock"); + } + + await database.KeyDeleteAsync(new RedisKey[] + { + RedisPipelineExtensions.FailureKey(settings.Resource), + RedisPipelineExtensions.StateKey(settings.Resource), + RedisPipelineExtensions.MessagesSortKey(settings.Resource), + }); + + await RemoveByPatternInBatchesAsync($"{RedisPipelineExtensions.MessageKey(settings.Resource)}*"); + + logger.LogDebug("Redis pipeline '{resource}' has been cleaned up", settings.Resource); + } + } + + /// + /// Resumes operations of the Redis pipeline after a stop. + /// + public override async Task ResumeAsync(int skip, CancellationToken cancellationToken) + { + using (IRedLock locker = await lockFactory.CreateLockAsync( + resource: settings.Resource, + expiryTime: settings.LockSettings.ExpiryTime, + waitTime: settings.LockSettings.WaitTime, + retryTime: settings.LockSettings.RetryTime, + cancellationToken)) + { + if (!locker.IsAcquired) + { + logger.LogError("Cannot acquire redlock for Redis pipeline '{resource}'", settings.Resource); + throw new InvalidOperationException("Unable to acquire redlock"); + } + + RedisValue state = await database.StringGetAsync(RedisPipelineExtensions.StateKey(settings.Resource)); + if (!RedisPipelineExtensions.IsStopped(state)) + { + logger.LogError("Cannot resume '{resource}' redis pipeline that has not stopped", settings.Resource); + throw new InvalidOperationException("Cannot resume a pipeline that has not stopped"); + } + + ITransaction transaction = database.CreateTransaction(); + RedisKey stateKey = RedisPipelineExtensions.StateKey(settings.Resource); + Task[] transactionTasks = new Task[] { + skip > 0 ? RemoveSkip(transaction, skip) : Task.CompletedTask, + transaction.StringSetAsync(stateKey, 0) + }; + await transaction.ExecuteAsync(); + await Task.WhenAll(transactionTasks); + + logger.LogDebug("Redis pipeline '{resource}' has been resumed", settings.Resource); + } + } + + private async Task RemoveSkip(ITransaction transaction, int skip) + { + RedisKey messagesKey = RedisPipelineExtensions.MessagesSortKey(settings.Resource); + RedisValue[] values = await transaction.SortedSetRangeByScoreAsync(messagesKey, stop: 0, order: Order.Ascending, take: skip); + if (values == null || values.Length <= 0) + { + return; + } + + foreach (RedisValue value in values) + { + RedisValue message = await database.SortedSetRemoveAsync(RedisPipelineExtensions.MessagesSortKey(settings.Resource), value.ToString()); + await transaction.KeyDeleteAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value.ToString())); + } + + } + + } +} diff --git a/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs b/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs new file mode 100644 index 0000000..1402fdd --- /dev/null +++ b/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs @@ -0,0 +1,92 @@ +using Microsoft.Extensions.Logging; +using RedLockNet; +using StackExchange.Redis; +using System; +using System.Text.Json; +using System.Threading; +using System.Threading.Tasks; + +namespace RedisMessagePipeline.Consumer +{ + /// + /// Consumes messages from a Redis pipeline and processes them according to the specified handler logic. + /// + public abstract class RedisBasePipelineConsumer : IRedisPipelineConsumer + { + protected readonly ILogger logger; + private readonly IRedisPipelineHandler handler; + protected readonly IDistributedLockFactory lockFactory; + protected readonly IDatabase database; + protected readonly RedisPipelineConsumerSettings settings; + + internal RedisBasePipelineConsumer( + ILogger logger, + IRedisPipelineHandler handler, + RedisPipelineConsumerSettings settings, + IDistributedLockFactory lockFactory, + IDatabase database) + { + this.logger = logger; + this.handler = handler; + this.lockFactory = lockFactory; + this.database = database; + this.settings = settings; + } + + /// + /// Executes the consumer processing, continually polling for and handling new messages. + /// + public async Task ExecuteAsync(CancellationToken cancellationToken) + { + logger.LogDebug("RedisPipelineConsumer '{resource}' has been executed.", settings.Resource); + + while (!cancellationToken.IsCancellationRequested) + { + bool success = await PollAsync(cancellationToken); + if (!success) + { + await Task.Delay(settings.PullInterval, cancellationToken); + } + } + } + + /// + /// Polls for new messages, processes them, and handles any resulting state changes. + /// + protected abstract Task PollAsync(CancellationToken cancellationToken); + + + /// + /// Attempts to process a single message and handle its result. + /// + protected async Task HandleMessageAsync(RedisValue message, CancellationToken cancellationToken) + { + bool success = false; + try + { + success = await handler.HandleAsync(message, cancellationToken); + } + catch (Exception ex) + { + logger.LogError(ex, "Handle message '{message}' from redis pipeline '{resource}' has been failed.", message, settings.Resource); + await StoreFailureAsync(message, ex); + } + + return success; + } + + /// + /// Records a failure in processing to the Redis failure log. + /// + private async Task StoreFailureAsync(RedisValue message, Exception ex) + { + RedisPipelineFailure failure = new RedisPipelineFailure + { + Exception = ex.Message, + Message = message, + Timestamp = DateTime.UtcNow.Ticks + }; + await database.StringSetAsync(RedisPipelineExtensions.FailureKey(settings.Resource), JsonSerializer.Serialize(failure)); + } + } +} diff --git a/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs b/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs index b048c6c..6da5499 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs @@ -1,7 +1,11 @@ -using System; +using RedisMessagePipeline.Factory; +using System; namespace RedisMessagePipeline.Consumer { + + + /// /// Configuration settings for RedisPipelineConsumer, including resource identifiers and retry logic. /// @@ -12,7 +16,9 @@ public RedisPipelineConsumerSettings(string resource) Resource = resource; } public string Resource { get; set; } + public string Reserved { get; set; } public int MaxRetries { get; set; } = int.MaxValue; + public EnPipelineType Type { get; set; } = EnPipelineType.QUEUE; /// /// Interval between unsuccessful handling or empty fetching attempts. diff --git a/RedisMessagePipeline/Consumer/RedisPipelineConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs similarity index 56% rename from RedisMessagePipeline/Consumer/RedisPipelineConsumer.cs rename to RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs index 3812414..0ab0110 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs @@ -2,7 +2,6 @@ using RedLockNet; using StackExchange.Redis; using System; -using System.Text.Json; using System.Threading; using System.Threading.Tasks; @@ -11,49 +10,50 @@ namespace RedisMessagePipeline.Consumer /// /// Consumes messages from a Redis pipeline and processes them according to the specified handler logic. /// - public class RedisPipelineConsumer : IRedisPipelineConsumer + public class RedisPipelineQueueConsumer : RedisBasePipelineConsumer { - private readonly ILogger logger; - private readonly IRedisPipelineHandler handler; - private readonly IDistributedLockFactory lockFactory; - private readonly IDatabase database; - private readonly RedisPipelineConsumerSettings settings; - - internal RedisPipelineConsumer( - ILogger logger, + internal RedisPipelineQueueConsumer( + ILogger logger, IRedisPipelineHandler handler, RedisPipelineConsumerSettings settings, IDistributedLockFactory lockFactory, IDatabase database) + : base(logger, handler, settings, lockFactory, database) { - this.logger = logger; - this.handler = handler; - this.lockFactory = lockFactory; - this.database = database; - this.settings = settings; + } /// - /// Executes the consumer processing, continually polling for and handling new messages. + /// Polls for new messages, processes them, and handles any resulting state changes. /// - public async Task ExecuteAsync(CancellationToken cancellationToken) + protected override async Task PollAsync(CancellationToken cancellationToken) { - logger.LogDebug("RedisPipelineConsumer '{resource}' has been executed.", settings.Resource); - - while (!cancellationToken.IsCancellationRequested) + try { - bool success = await PollAsync(cancellationToken); - if (!success) + //realiza a reserva do item para a fila exclusiva + await TryDequeueAndReserveAsync(cancellationToken); + + RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved)); + if (message.IsNull) { - await Task.Delay(settings.PullInterval, cancellationToken); + return false; } + + bool success = await HandleMessageAsync(message, cancellationToken); + if (success) + { + await HandleSuccessAsync(); + return true; + } + return false; + } + catch (Exception) + { + return false; } } - /// - /// Polls for new messages, processes them, and handles any resulting state changes. - /// - private async Task PollAsync(CancellationToken cancellationToken) + private async Task TryDequeueAndReserveAsync(CancellationToken cancellationToken) { using (IRedLock locker = await lockFactory.CreateLockAsync( resource: settings.Resource, @@ -73,62 +73,32 @@ private async Task PollAsync(CancellationToken cancellationToken) return false; } - RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesKey(settings.Resource)); - if (message.IsNull) - { - return false; - } - bool success = await HandleMessageAsync(message, cancellationToken); - - if (success) + //verifica se este algum item reservado na fila de processamento. + long count = await database.ListLengthAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved)); + if (count > 0) { - await HandleSuccessAsync(); + //existe item na fila de reserva, ativa processamento return true; } - await HandleFailureAsync(message, state); - return false; - } - } + //envia objeto para a fila de reserva, para processamento + RedisValue item = await database.ListMoveAsync(RedisPipelineExtensions.MessagesListKey(settings.Resource), RedisPipelineExtensions.MessagesListKey(settings.Reserved), ListSide.Left, ListSide.Right); + if (item.IsNull) + { + // fila principal vazia + return false; + } - /// - /// Attempts to process a single message and handle its result. - /// - private async Task HandleMessageAsync(RedisValue message, CancellationToken cancellationToken) - { - bool success = false; - try - { - success = await handler.HandleAsync(message, cancellationToken); + return true; } - catch (Exception ex) - { - logger.LogError(ex, "Handle message '{message}' from redis pipeline '{resource}' has been failed.", message, settings.Resource); - await StoreFailureAsync(message, ex); - } - - return success; } - /// - /// Records a failure in processing to the Redis failure log. - /// - private async Task StoreFailureAsync(RedisValue message, Exception ex) - { - RedisPipelineFailure failure = new RedisPipelineFailure - { - Exception = ex.Message, - Message = message, - Timestamp = DateTime.UtcNow.Ticks - }; - await database.StringSetAsync(RedisPipelineExtensions.FailureKey(settings.Resource), JsonSerializer.Serialize(failure)); - } /// /// Handles successful message processing by resetting the pipeline state and clearing failures. /// - private async Task HandleSuccessAsync() + protected async Task HandleSuccessAsync() { ITransaction transaction = database.CreateTransaction(); Task[] transactionTasks = new Task[] { @@ -142,11 +112,11 @@ private async Task HandleSuccessAsync() /// /// Handles message processing failures by retrying or stopping the pipeline based on the retry policy. /// - private async Task HandleFailureAsync(RedisValue message, RedisValue state) + protected async Task HandleFailureAsync(RedisValue message, RedisValue state) { logger.LogWarning("Handle message '{message}' from redis pipeline '{resource}' has been failed.", message, settings.Resource); - await database.ListLeftPushAsync(RedisPipelineExtensions.MessagesKey(settings.Resource), message); + await database.ListLeftPushAsync(RedisPipelineExtensions.MessagesListKey(settings.Resource), message); RedisKey stateKey = RedisPipelineExtensions.StateKey(settings.Resource); if (state.IsNull) { diff --git a/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs new file mode 100644 index 0000000..0f78d8d --- /dev/null +++ b/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs @@ -0,0 +1,171 @@ +using Microsoft.Extensions.Logging; +using RedLockNet; +using StackExchange.Redis; +using System; +using System.Threading; +using System.Threading.Tasks; + +namespace RedisMessagePipeline.Consumer +{ + /// + /// Consumes messages from a Redis pipeline and processes them according to the specified handler logic. + /// + public class RedisPipelineScheduleConsumer : RedisBasePipelineConsumer + { + internal RedisPipelineScheduleConsumer( + ILogger logger, + IRedisPipelineHandler handler, + RedisPipelineConsumerSettings settings, + IDistributedLockFactory lockFactory, + IDatabase database) + : base(logger, handler, settings, lockFactory, database) + { + + } + + /// + /// Polls for new messages, processes them, and handles any resulting state changes. + /// + protected override async Task PollAsync(CancellationToken cancellationToken) + { + try + { + //realiza a reserva do item para a fila exclusiva + await TryDequeueAndReserveAsync(cancellationToken); + + RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved)); + if (message.IsNull) + { + return false; + } + + bool success = await HandleMessageAsync(message, cancellationToken); + if (success) + { + await HandleSuccessAsync(); + return true; + } + return false; + } + catch (Exception) + { + return false; + } + } + + + private async Task TryDequeueAndReserveAsync(CancellationToken cancellationToken) + { + using (IRedLock locker = await lockFactory.CreateLockAsync( + resource: settings.Resource, + expiryTime: settings.LockSettings.ExpiryTime, + waitTime: settings.LockSettings.WaitTime, + retryTime: settings.LockSettings.RetryTime, + cancellationToken)) + { + + if (!locker.IsAcquired) + { + return false; + } + + RedisValue state = await database.StringGetAsync( + RedisPipelineExtensions.StateKey(settings.Resource)); + + if (RedisPipelineExtensions.IsStopped(state)) + { + return false; + } + + + //verifica se este algum item reservado na fila de processamento. + long count = await database.ListLengthAsync( + RedisPipelineExtensions.MessagesListKey(settings.Reserved)); + if (count > 0) + { + //existe item na fila de reserva, ativa processamento + return true; + } + + double now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + + RedisValue[] values = await database.SortedSetRangeByScoreAsync( + RedisPipelineExtensions.MessagesSortKey(settings.Resource), + stop: now, + order: Order.Ascending, + take: 1); + + if (values == null || values.Length == 0 || values[0].IsNull) + { + return false; + } + + string id = values[0].ToString(); + bool removed = await database.SortedSetRemoveAsync( + RedisPipelineExtensions.MessagesSortKey(settings.Resource), + id); + + if (!removed) + { + return false; + } + + RedisValue message = await database.StringGetAsync( + RedisPipelineExtensions.MessageKey(settings.Resource, id)); + + if (message.IsNull) + { + return false; + } + + //envia para a file de processamento + long result = await database.ListRightPushAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved), message); + if (result <= 0) + { + return false; + } + + await database.KeyDeleteAsync(RedisPipelineExtensions.MessageKey(settings.Resource, id)); + + return true; + } + } + + + /// + /// Handles successful message processing by resetting the pipeline state and clearing failures. + /// + private async Task HandleSuccessAsync() + { + ITransaction transaction = database.CreateTransaction(); + Task[] transactionTasks = new Task[] { + transaction.StringSetAsync(RedisPipelineExtensions.StateKey(settings.Resource), 0), + transaction.KeyDeleteAsync(RedisPipelineExtensions.FailureKey(settings.Resource)) + }; + await transaction.ExecuteAsync(); + await Task.WhenAll(transactionTasks); + } + + /* + /// + /// Handles message processing failures by retrying or stopping the pipeline based on the retry policy. + /// + private async Task HandleFailureAsync(RedisValue message, RedisValue state, RedisValue value) + { + logger.LogWarning("Handle message '{message}' from redis pipeline '{resource}' has been failed.", message, settings.Resource); + RedisKey stateKey = RedisPipelineExtensions.StateKey(settings.Resource); + int count = 1; + if (!state.IsNull && int.TryParse(state, out count)) + { + count++; + } + count = (count <= 0 ? 1 : count); + await database.StringSetAsync(stateKey, $"{count}"); + + var novaData = DateTime.UtcNow.AddMinutes(count * 60); + var score = new DateTimeOffset(novaData).ToUnixTimeMilliseconds(); + await database.SortedSetAddAsync(RedisPipelineExtensions.MessagesSortKey(settings.Resource), value, score); + } + */ + } +} diff --git a/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs b/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs index 500ee9a..6c7fb91 100644 --- a/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs +++ b/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs @@ -3,6 +3,12 @@ namespace RedisMessagePipeline.Factory { + public enum EnPipelineType + { + QUEUE = 0, + QUEUE_SCHEDULE = 1 + } + public interface IRedisPipelineFactory { IRedisPipelineConsumer CreateConsumer(IRedisPipelineHandler handler, RedisPipelineConsumerSettings settings); diff --git a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs index d747995..c0bf45d 100644 --- a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs +++ b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs @@ -27,7 +27,13 @@ public RedisPipelineFactory(ILoggerFactory loggerFactory, IDistributedLockFactor /// public IRedisPipelineConsumer CreateConsumer(IRedisPipelineHandler handler, RedisPipelineConsumerSettings settings) { - return new RedisPipelineConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); + switch (settings.Type) + { + case EnPipelineType.QUEUE_SCHEDULE: + return new RedisPipelineScheduleConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); + default: + return new RedisPipelineQueueConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); + } } /// @@ -35,7 +41,19 @@ public IRedisPipelineConsumer CreateConsumer(IRedisPipelineHandler handler, Redi /// public IRedisPipelineAdmin CreateAdmin(RedisPipelineAdminSettings settings) { - return new RedisPipelineAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); + switch (settings.Type) + { + case EnPipelineType.QUEUE_SCHEDULE: + return new RedisPipelineScheduleAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); + default: + return new RedisPipelineQueueAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); + } + + } + + + + } } diff --git a/RedisMessagePipeline/RedisMessagePipeline.csproj b/RedisMessagePipeline/RedisMessagePipeline.csproj index 51a0c68..439c9f2 100644 --- a/RedisMessagePipeline/RedisMessagePipeline.csproj +++ b/RedisMessagePipeline/RedisMessagePipeline.csproj @@ -20,8 +20,8 @@ - - + + diff --git a/RedisMessagePipeline/RedisPipelineExtensions.cs b/RedisMessagePipeline/RedisPipelineExtensions.cs index ae745e0..f6e2987 100644 --- a/RedisMessagePipeline/RedisPipelineExtensions.cs +++ b/RedisMessagePipeline/RedisPipelineExtensions.cs @@ -11,8 +11,11 @@ public static class RedisPipelineExtensions public const string STATE_STOPPED = "STOPPED"; public static bool IsStopped(RedisValue redisValue) => redisValue.HasValue && redisValue == STATE_STOPPED; - public static RedisKey MessagesKey(string resource) => $"{resource}:messages"; + public static RedisKey MessagesListKey(string resource) => $"{resource}:messages"; public static RedisKey StateKey(string resource) => $"{resource}:state"; public static RedisKey FailureKey(string resource) => $"{resource}:failure"; + public static RedisKey MessagesSortKey(string resource) => $"{resource}:sortkeys"; + public static RedisKey MessageKey(string resource) => $"{resource}:message:"; + public static RedisKey MessageKey(string resource, string id) => $"{MessageKey(resource)}{id}"; } }