From b58094ea5afcadfa3c656d0d182e786a6a3c719b Mon Sep 17 00:00:00 2001 From: Wanderson Alves dos Santos Date: Tue, 17 Mar 2026 21:33:17 -0300 Subject: [PATCH 01/10] =?UTF-8?q?Cria=C3=A7=C3=A3o=20filas=20de=20Agendame?= =?UTF-8?q?nto?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- RedisMessagePipeline.Demo/Program.cs | 12 +- .../Admin/IRedisPipelineAdmin.cs | 4 +- .../Admin/RedisPipelineAdminSettings.cs | 5 +- .../Admin/RedisPipelineBaseAdmin.cs | 114 ++++++++++++++++ ...ineAdmin.cs => RedisPipelineQueueAdmin.cs} | 46 ++----- .../Admin/RedisPipelineSheduleAdmin.cs | 124 +++++++++++++++++ .../Consumer/RedisBasePipelineConsumer.cs | 92 +++++++++++++ .../Consumer/RedisPipelineConsumerSettings.cs | 7 +- ...sumer.cs => RedisPipelineQueueConsumer.cs} | 81 ++--------- .../Consumer/RedisPipelineSheduleConsumer.cs | 127 ++++++++++++++++++ .../Factory/IRedisPipelineFactory.cs | 6 + .../Factory/RedisPipelineFactory.cs | 23 +++- .../RedisPipelineExtensions.cs | 5 +- 13 files changed, 532 insertions(+), 114 deletions(-) create mode 100644 RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs rename RedisMessagePipeline/Admin/{RedisPipelineAdmin.cs => RedisPipelineQueueAdmin.cs} (70%) create mode 100644 RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs create mode 100644 RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs rename RedisMessagePipeline/Consumer/{RedisPipelineConsumer.cs => RedisPipelineQueueConsumer.cs} (54%) create mode 100644 RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs diff --git a/RedisMessagePipeline.Demo/Program.cs b/RedisMessagePipeline.Demo/Program.cs index 74b7dac..c2e8efd 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("redis.hml.prospexti.com.br:6379,password=Zerbeto18, abortConnect=false"); 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.AddSheduleAsync($"{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..fb6d8d9 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 AddSheduleAsync(RedisValue keyValue, DateTime shedule, 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..19430e4 --- /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.CompletedTask; + + } + + /// + /// 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 AddSheduleAsync(RedisValue keyValue, DateTime shedule, 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 70% rename from RedisMessagePipeline/Admin/RedisPipelineAdmin.cs rename to RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs index 3aa3913..449e420 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs @@ -10,51 +10,33 @@ 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) + public override Task PushQueueAsync(RedisValue redisValue) { - logger.LogDebug("Push a new message '{message}' to '{resource}' redis pipeline", redisValue, settings.Resource); + base.PushQueueAsync(redisValue); - RedisKey key = RedisPipelineExtensions.MessagesKey(settings.Resource); + RedisKey key = RedisPipelineExtensions.MessagesListKey(settings.Resource); return database.ListRightPushAsync(key, redisValue); } - - /// - /// 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 async Task CleanAsync(CancellationToken cancellationToken) + public override async Task CleanAsync(CancellationToken cancellationToken) { using (IRedLock locker = await lockFactory.CreateLockAsync( resource: settings.Resource, @@ -73,7 +55,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 +65,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 +88,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/RedisPipelineSheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs new file mode 100644 index 0000000..1c373a6 --- /dev/null +++ b/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs @@ -0,0 +1,124 @@ +using Microsoft.Extensions.Logging; +using RedLockNet; +using StackExchange.Redis; +using System; +using System.Threading; +using System.Threading.Tasks; +using System.Transactions; + +namespace RedisMessagePipeline.Admin +{ + /// + /// Admin functionality to manage operations on the Redis pipeline, such as starting, stopping, and cleaning. + /// + public class RedisPipelineSheduleAdmin : RedisPipelineBaseAdmin + { + internal RedisPipelineSheduleAdmin( + ILogger logger, + RedisPipelineAdminSettings settings, + IDistributedLockFactory lockFactory, + IDatabase database) + : base(logger, settings, lockFactory, database) + { + + } + + public override async Task AddSheduleAsync(RedisValue keyValue, DateTime shedule, RedisValue redisValue) + { + await base.AddSheduleAsync(keyValue, shedule, redisValue); + var score = new DateTimeOffset(shedule).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..9d34d48 --- /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..ad349a6 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. /// @@ -13,6 +17,7 @@ public RedisPipelineConsumerSettings(string resource) } public string Resource { 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 54% rename from RedisMessagePipeline/Consumer/RedisPipelineConsumer.cs rename to RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs index 3812414..bbdec15 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs @@ -11,49 +11,23 @@ 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. - /// - 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. /// - private async Task PollAsync(CancellationToken cancellationToken) + protected override async Task PollAsync(CancellationToken cancellationToken) { using (IRedLock locker = await lockFactory.CreateLockAsync( resource: settings.Resource, @@ -73,7 +47,7 @@ private async Task PollAsync(CancellationToken cancellationToken) return false; } - RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesKey(settings.Resource)); + RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesListKey(settings.Resource)); if (message.IsNull) { return false; @@ -87,48 +61,15 @@ private async Task PollAsync(CancellationToken cancellationToken) return true; } - await HandleFailureAsync(message, state); + //await HandleFailureAsync(message, state); 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); - } - 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 +83,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/RedisPipelineSheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs new file mode 100644 index 0000000..f6a9e10 --- /dev/null +++ b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs @@ -0,0 +1,127 @@ +using Microsoft.Extensions.Logging; +using RedLockNet; +using StackExchange.Redis; +using System; +using System.Drawing; +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 class RedisPipelineSheduleConsumer : RedisBasePipelineConsumer + { + internal RedisPipelineSheduleConsumer( + 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) + { + 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; + } + + var stop = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + RedisValue[] values = await database.SortedSetRangeByScoreAsync( + RedisPipelineExtensions.MessagesSortKey(settings.Resource), + stop: stop, + order: Order.Ascending, + take: 1 + ); + + if (values == null || values.Length <= 0 || values[0].IsNull) + { + return false; + } + + string value = values[0].ToString(); + + // Reserva simples: remove da fila + RedisValue message = await database.SortedSetRemoveAsync(RedisPipelineExtensions.MessagesSortKey(settings.Resource), value); + if (message.IsNull) + { + return false; //item já foi removido + } + + message = await database.StringGetAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value)); + if (message.IsNull) + { + return false; + } + + bool success = await HandleMessageAsync(message, cancellationToken); + if (success) + { + await HandleSuccessAsync(value); + return true; + } + //await HandleFailureAsync(message, state, value); + return false; + + } + } + + + /// + /// Handles successful message processing by resetting the pipeline state and clearing failures. + /// + private async Task HandleSuccessAsync(RedisValue value) + { + ITransaction transaction = database.CreateTransaction(); + Task[] transactionTasks = new Task[] { + transaction.StringSetAsync(RedisPipelineExtensions.StateKey(settings.Resource), 0), + transaction.KeyDeleteAsync(RedisPipelineExtensions.FailureKey(settings.Resource)), + transaction.KeyDeleteAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value)) + }; + 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..4748477 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..30b458b 100644 --- a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs +++ b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs @@ -3,6 +3,7 @@ using RedisMessagePipeline.Consumer; using RedLockNet; using StackExchange.Redis; +using System.Runtime.InteropServices; namespace RedisMessagePipeline.Factory { @@ -27,7 +28,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 RedisPipelineSheduleConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); + default: + return new RedisPipelineQueueConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); + } } /// @@ -35,7 +42,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 RedisPipelineSheduleAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); + default: + return new RedisPipelineQueueAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); + } + + } + + + + } } diff --git a/RedisMessagePipeline/RedisPipelineExtensions.cs b/RedisMessagePipeline/RedisPipelineExtensions.cs index ae745e0..ccaad2b 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}"; } } From b2cf91d69b65d424b4a74cb17460cea7f16edc92 Mon Sep 17 00:00:00 2001 From: Wanderson Alves dos Santos Date: Tue, 17 Mar 2026 21:34:05 -0300 Subject: [PATCH 02/10] Ok. --- RedisMessagePipeline.Demo/Program.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/RedisMessagePipeline.Demo/Program.cs b/RedisMessagePipeline.Demo/Program.cs index c2e8efd..4872467 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("redis.hml.prospexti.com.br:6379,password=Zerbeto18, abortConnect=false"); +ConnectionMultiplexer redis = ConnectionMultiplexer.Connect(""); RedLockMultiplexer lockMultiplexer = new RedLockMultiplexer(redis); IDatabase db = redis.GetDatabase(); From e0bb0ef3f282694643b5f852a75da6f68b5eb6b3 Mon Sep 17 00:00:00 2001 From: Wanderson Alves dos Santos Date: Fri, 27 Mar 2026 21:01:44 -0300 Subject: [PATCH 03/10] =?UTF-8?q?Alteral=C3=A7oes?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Admin/RedisPipelineQueueAdmin.cs | 13 ++++++------- .../Admin/RedisPipelineSheduleAdmin.cs | 1 - .../Consumer/RedisBasePipelineConsumer.cs | 2 +- .../Consumer/RedisPipelineQueueConsumer.cs | 4 +--- .../Consumer/RedisPipelineSheduleConsumer.cs | 16 +++++++--------- .../Factory/IRedisPipelineFactory.cs | 2 +- .../Factory/RedisPipelineFactory.cs | 5 ++--- RedisMessagePipeline/RedisMessagePipeline.csproj | 4 ++-- RedisMessagePipeline/RedisPipelineExtensions.cs | 2 +- 9 files changed, 21 insertions(+), 28 deletions(-) diff --git a/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs index 449e420..4a94895 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs @@ -17,22 +17,21 @@ internal RedisPipelineQueueAdmin( RedisPipelineAdminSettings settings, IDistributedLockFactory lockFactory, IDatabase database) - :base(logger, settings, lockFactory, database) + : base(logger, settings, lockFactory, database) { - + } /// /// Pushes a new message to the Redis pipeline. /// - public override Task PushQueueAsync(RedisValue redisValue) + public override async Task PushQueueAsync(RedisValue redisValue) { - base.PushQueueAsync(redisValue); - + await base.PushQueueAsync(redisValue); RedisKey key = RedisPipelineExtensions.MessagesListKey(settings.Resource); - return database.ListRightPushAsync(key, redisValue); + await database.ListRightPushAsync(key, redisValue); } - + /// /// Cleans up resources used by the Redis pipeline. /// diff --git a/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs index 1c373a6..2df8a55 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs @@ -4,7 +4,6 @@ using System; using System.Threading; using System.Threading.Tasks; -using System.Transactions; namespace RedisMessagePipeline.Admin { diff --git a/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs b/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs index 9d34d48..1402fdd 100644 --- a/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisBasePipelineConsumer.cs @@ -54,7 +54,7 @@ public async Task ExecuteAsync(CancellationToken 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. diff --git a/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs index bbdec15..6ea1a93 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs @@ -1,8 +1,6 @@ using Microsoft.Extensions.Logging; using RedLockNet; using StackExchange.Redis; -using System; -using System.Text.Json; using System.Threading; using System.Threading.Tasks; @@ -21,7 +19,7 @@ internal RedisPipelineQueueConsumer( IDatabase database) : base(logger, handler, settings, lockFactory, database) { - + } /// diff --git a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs index f6a9e10..511dde8 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs @@ -2,8 +2,6 @@ using RedLockNet; using StackExchange.Redis; using System; -using System.Drawing; -using System.Text.Json; using System.Threading; using System.Threading.Tasks; @@ -20,9 +18,9 @@ internal RedisPipelineSheduleConsumer( RedisPipelineConsumerSettings settings, IDistributedLockFactory lockFactory, IDatabase database) - :base(logger, handler, settings, lockFactory, database) + : base(logger, handler, settings, lockFactory, database) { - + } /// @@ -61,7 +59,7 @@ protected override async Task PollAsync(CancellationToken cancellationToke return false; } - string value = values[0].ToString(); + string value = values[0].ToString(); // Reserva simples: remove da fila RedisValue message = await database.SortedSetRemoveAsync(RedisPipelineExtensions.MessagesSortKey(settings.Resource), value); @@ -73,7 +71,7 @@ protected override async Task PollAsync(CancellationToken cancellationToke message = await database.StringGetAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value)); if (message.IsNull) { - return false; + return false; } bool success = await HandleMessageAsync(message, cancellationToken); @@ -84,11 +82,11 @@ protected override async Task PollAsync(CancellationToken cancellationToke } //await HandleFailureAsync(message, state, value); return false; - + } } - + /// /// Handles successful message processing by resetting the pipeline state and clearing failures. /// @@ -97,7 +95,7 @@ private async Task HandleSuccessAsync(RedisValue value) ITransaction transaction = database.CreateTransaction(); Task[] transactionTasks = new Task[] { transaction.StringSetAsync(RedisPipelineExtensions.StateKey(settings.Resource), 0), - transaction.KeyDeleteAsync(RedisPipelineExtensions.FailureKey(settings.Resource)), + transaction.KeyDeleteAsync(RedisPipelineExtensions.FailureKey(settings.Resource)), transaction.KeyDeleteAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value)) }; await transaction.ExecuteAsync(); diff --git a/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs b/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs index 4748477..6c7fb91 100644 --- a/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs +++ b/RedisMessagePipeline/Factory/IRedisPipelineFactory.cs @@ -7,7 +7,7 @@ public enum EnPipelineType { QUEUE = 0, QUEUE_SCHEDULE = 1 - } + } public interface IRedisPipelineFactory { diff --git a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs index 30b458b..1fc5488 100644 --- a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs +++ b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs @@ -3,7 +3,6 @@ using RedisMessagePipeline.Consumer; using RedLockNet; using StackExchange.Redis; -using System.Runtime.InteropServices; namespace RedisMessagePipeline.Factory { @@ -50,11 +49,11 @@ public IRedisPipelineAdmin CreateAdmin(RedisPipelineAdminSettings settings) 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 ccaad2b..f6e2987 100644 --- a/RedisMessagePipeline/RedisPipelineExtensions.cs +++ b/RedisMessagePipeline/RedisPipelineExtensions.cs @@ -11,7 +11,7 @@ public static class RedisPipelineExtensions public const string STATE_STOPPED = "STOPPED"; public static bool IsStopped(RedisValue redisValue) => redisValue.HasValue && redisValue == STATE_STOPPED; - public static RedisKey MessagesListKey(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"; From bcf4b2e71d215dd002f1c1d21dac44d6ee580cf1 Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Sun, 29 Mar 2026 15:52:18 -0300 Subject: [PATCH 04/10] =?UTF-8?q?Corre=C3=A7=C3=B5es?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../Consumer/RedisPipelineConsumerSettings.cs | 1 + .../Consumer/RedisPipelineQueueConsumer.cs | 51 ++++++++--- .../Consumer/RedisPipelineSheduleConsumer.cs | 88 ++++++++++++++----- 3 files changed, 109 insertions(+), 31 deletions(-) diff --git a/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs b/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs index ad349a6..6da5499 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineConsumerSettings.cs @@ -16,6 +16,7 @@ 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; diff --git a/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs index 6ea1a93..0ab0110 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineQueueConsumer.cs @@ -1,6 +1,7 @@ using Microsoft.Extensions.Logging; using RedLockNet; using StackExchange.Redis; +using System; using System.Threading; using System.Threading.Tasks; @@ -26,6 +27,33 @@ internal RedisPipelineQueueConsumer( /// 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, @@ -45,25 +73,28 @@ protected override async Task PollAsync(CancellationToken cancellationToke return false; } - RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesListKey(settings.Resource)); - if (message.IsNull) + + //verifica se este algum item reservado na fila de processamento. + long count = await database.ListLengthAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved)); + if (count > 0) { - return false; + //existe item na fila de reserva, ativa processamento + return true; } - bool success = await HandleMessageAsync(message, cancellationToken); - - if (success) + //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) { - await HandleSuccessAsync(); - return true; + // fila principal vazia + return false; } - //await HandleFailureAsync(message, state); - return false; + return true; } } + /// /// Handles successful message processing by resetting the pipeline state and clearing failures. /// diff --git a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs index 511dde8..1f13eba 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs @@ -27,6 +27,34 @@ internal RedisPipelineSheduleConsumer( /// 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.MessageKey(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, @@ -35,54 +63,71 @@ protected override async Task PollAsync(CancellationToken cancellationToke retryTime: settings.LockSettings.RetryTime, cancellationToken)) { + if (!locker.IsAcquired) { return false; } - RedisValue state = await database.StringGetAsync(RedisPipelineExtensions.StateKey(settings.Resource)); + RedisValue state = await database.StringGetAsync( + RedisPipelineExtensions.StateKey(settings.Resource)); + if (RedisPipelineExtensions.IsStopped(state)) { return false; } - var stop = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + + //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; + } + + var now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + RedisValue[] values = await database.SortedSetRangeByScoreAsync( RedisPipelineExtensions.MessagesSortKey(settings.Resource), - stop: stop, + stop: now, order: Order.Ascending, - take: 1 - ); + take: 1); - if (values == null || values.Length <= 0 || values[0].IsNull) + if (values == null || values.Length == 0 || values[0].IsNull) { return false; } - string value = values[0].ToString(); + string id = values[0].ToString(); + bool removed = await database.SortedSetRemoveAsync( + RedisPipelineExtensions.MessagesSortKey(settings.Resource), + id); - // Reserva simples: remove da fila - RedisValue message = await database.SortedSetRemoveAsync(RedisPipelineExtensions.MessagesSortKey(settings.Resource), value); - if (message.IsNull) + if (!removed) { - return false; //item já foi removido + return false; } - message = await database.StringGetAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value)); + RedisValue message = await database.StringGetAsync( + RedisPipelineExtensions.MessageKey(settings.Resource, id)); + if (message.IsNull) { return false; } - bool success = await HandleMessageAsync(message, cancellationToken); - if (success) + //envia para a file de processamento + long result = await database.ListRightPushAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved), message); + if (result <= 0) { - await HandleSuccessAsync(value); - return true; + return false; } - //await HandleFailureAsync(message, state, value); - return false; + + await database.KeyDeleteAsync(RedisPipelineExtensions.MessageKey(settings.Resource, id)); + return true; } } @@ -90,18 +135,18 @@ protected override async Task PollAsync(CancellationToken cancellationToke /// /// Handles successful message processing by resetting the pipeline state and clearing failures. /// - private async Task HandleSuccessAsync(RedisValue value) + 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)), - transaction.KeyDeleteAsync(RedisPipelineExtensions.MessageKey(settings.Resource, value)) + 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. /// @@ -121,5 +166,6 @@ private async Task HandleFailureAsync(RedisValue message, RedisValue state, Redi var score = new DateTimeOffset(novaData).ToUnixTimeMilliseconds(); await database.SortedSetAddAsync(RedisPipelineExtensions.MessagesSortKey(settings.Resource), value, score); } + */ } } From 7ad93b3cc5a33f720e258f78c9f55637ae4146aa Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Sun, 29 Mar 2026 16:55:26 -0300 Subject: [PATCH 05/10] Ok. --- RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs index 1f13eba..6b73e0b 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs @@ -33,7 +33,7 @@ protected override async Task PollAsync(CancellationToken cancellationToke //realiza a reserva do item para a fila exclusiva await TryDequeueAndReserveAsync(cancellationToken); - RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessageKey(settings.Reserved)); + RedisValue message = await database.ListLeftPopAsync(RedisPipelineExtensions.MessagesListKey(settings.Reserved)); if (message.IsNull) { return false; From e4aa8d19b4cadb91fd6a4a0aaffe21a7c7170ea2 Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Mon, 6 Apr 2026 14:51:34 -0300 Subject: [PATCH 06/10] oK. --- RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs | 2 +- RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs | 4 ++-- RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs | 4 ++-- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs index fb6d8d9..cd362e8 100644 --- a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs +++ b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs @@ -7,7 +7,7 @@ namespace RedisMessagePipeline.Admin { public interface IRedisPipelineAdmin { - Task PushQueueAsync(RedisValue redisValue); + Task PushQueueAsync(RedisValue redisValue); Task AddSheduleAsync(RedisValue keyValue, DateTime shedule, RedisValue redisValue); Task StopAsync(); Task CleanAsync(CancellationToken cancellationToken); diff --git a/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs index 19430e4..8f7c8e4 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs @@ -33,7 +33,7 @@ internal RedisPipelineBaseAdmin( /// /// Pushes a new message to the Redis pipeline. /// - public virtual Task PushQueueAsync(RedisValue redisValue) + public virtual Task PushQueueAsync(RedisValue redisValue) { if (this.settings.Type != Factory.EnPipelineType.QUEUE) { @@ -42,7 +42,7 @@ public virtual Task PushQueueAsync(RedisValue redisValue) logger.LogDebug("Push a new message '{message}' to '{resource}' redis pipeline", redisValue, settings.Resource); - return Task.CompletedTask; + return Task.FromResult(-1L); } diff --git a/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs index 4a94895..91f758a 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineQueueAdmin.cs @@ -25,11 +25,11 @@ internal RedisPipelineQueueAdmin( /// /// Pushes a new message to the Redis pipeline. /// - public override async Task PushQueueAsync(RedisValue redisValue) + public override async Task PushQueueAsync(RedisValue redisValue) { await base.PushQueueAsync(redisValue); RedisKey key = RedisPipelineExtensions.MessagesListKey(settings.Resource); - await database.ListRightPushAsync(key, redisValue); + return await database.ListRightPushAsync(key, redisValue); } /// From ac53902d71dc884a61d6144de5430422758f50b5 Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Mon, 6 Apr 2026 15:45:36 -0300 Subject: [PATCH 07/10] Ok. --- RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs | 2 +- RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs | 2 +- ...SheduleAdmin.cs => RedisPipelineScheduleAdmin.cs} | 12 ++++++------ ...eConsumer.cs => RedisPipelineScheduleConsumer.cs} | 4 ++-- RedisMessagePipeline/Factory/RedisPipelineFactory.cs | 4 ++-- 5 files changed, 12 insertions(+), 12 deletions(-) rename RedisMessagePipeline/Admin/{RedisPipelineSheduleAdmin.cs => RedisPipelineScheduleAdmin.cs} (91%) rename RedisMessagePipeline/Consumer/{RedisPipelineSheduleConsumer.cs => RedisPipelineScheduleConsumer.cs} (98%) diff --git a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs index cd362e8..008f20d 100644 --- a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs +++ b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs @@ -8,7 +8,7 @@ namespace RedisMessagePipeline.Admin public interface IRedisPipelineAdmin { Task PushQueueAsync(RedisValue redisValue); - Task AddSheduleAsync(RedisValue keyValue, DateTime shedule, RedisValue redisValue); + Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, RedisValue redisValue); Task StopAsync(); Task CleanAsync(CancellationToken cancellationToken); Task ResumeAsync(int skip, CancellationToken cancellationToken); diff --git a/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs index 8f7c8e4..1daab66 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs @@ -67,7 +67,7 @@ public Task StopAsync() /// public abstract Task ResumeAsync(int skip, CancellationToken cancellationToken); - public virtual Task AddSheduleAsync(RedisValue keyValue, DateTime shedule, RedisValue redisValue) + public virtual Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, RedisValue redisValue) { if (this.settings.Type != Factory.EnPipelineType.QUEUE_SCHEDULE) { diff --git a/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs similarity index 91% rename from RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs rename to RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs index 2df8a55..de9f913 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineSheduleAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs @@ -10,10 +10,10 @@ namespace RedisMessagePipeline.Admin /// /// Admin functionality to manage operations on the Redis pipeline, such as starting, stopping, and cleaning. /// - public class RedisPipelineSheduleAdmin : RedisPipelineBaseAdmin + public class RedisPipelineScheduleAdmin : RedisPipelineBaseAdmin { - internal RedisPipelineSheduleAdmin( - ILogger logger, + internal RedisPipelineScheduleAdmin( + ILogger logger, RedisPipelineAdminSettings settings, IDistributedLockFactory lockFactory, IDatabase database) @@ -22,10 +22,10 @@ internal RedisPipelineSheduleAdmin( } - public override async Task AddSheduleAsync(RedisValue keyValue, DateTime shedule, RedisValue redisValue) + public override async Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, RedisValue redisValue) { - await base.AddSheduleAsync(keyValue, shedule, redisValue); - var score = new DateTimeOffset(shedule).ToUnixTimeMilliseconds(); + await base.AddScheduleAsync(keyValue, schedule, redisValue); + var score = new DateTimeOffset(schedule).ToUnixTimeMilliseconds(); RedisKey key = RedisPipelineExtensions.MessagesSortKey(settings.Resource); await database.SortedSetAddAsync(key, keyValue, score); key = RedisPipelineExtensions.MessageKey(settings.Resource, keyValue); diff --git a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs similarity index 98% rename from RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs rename to RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs index 6b73e0b..c00ac7c 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineSheduleConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs @@ -10,9 +10,9 @@ namespace RedisMessagePipeline.Consumer /// /// Consumes messages from a Redis pipeline and processes them according to the specified handler logic. /// - public class RedisPipelineSheduleConsumer : RedisBasePipelineConsumer + public class RedisPipelineScheduleConsumer : RedisBasePipelineConsumer { - internal RedisPipelineSheduleConsumer( + internal RedisPipelineScheduleConsumer( ILogger logger, IRedisPipelineHandler handler, RedisPipelineConsumerSettings settings, diff --git a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs index 1fc5488..c0bf45d 100644 --- a/RedisMessagePipeline/Factory/RedisPipelineFactory.cs +++ b/RedisMessagePipeline/Factory/RedisPipelineFactory.cs @@ -30,7 +30,7 @@ public IRedisPipelineConsumer CreateConsumer(IRedisPipelineHandler handler, Redi switch (settings.Type) { case EnPipelineType.QUEUE_SCHEDULE: - return new RedisPipelineSheduleConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); + return new RedisPipelineScheduleConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); default: return new RedisPipelineQueueConsumer(loggerFactory.CreateLogger(), handler, settings, lockFactory, database); } @@ -44,7 +44,7 @@ public IRedisPipelineAdmin CreateAdmin(RedisPipelineAdminSettings settings) switch (settings.Type) { case EnPipelineType.QUEUE_SCHEDULE: - return new RedisPipelineSheduleAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); + return new RedisPipelineScheduleAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); default: return new RedisPipelineQueueAdmin(loggerFactory.CreateLogger(), settings, lockFactory, database); } From 730ec47035d912d3950aac567d57e91032cae0bc Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Mon, 6 Apr 2026 17:16:08 -0300 Subject: [PATCH 08/10] ok. --- RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs | 5 ++++- .../Consumer/RedisPipelineScheduleConsumer.cs | 2 +- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs index de9f913..ac53b02 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs @@ -25,7 +25,10 @@ internal RedisPipelineScheduleAdmin( public override async Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, RedisValue redisValue) { await base.AddScheduleAsync(keyValue, schedule, redisValue); - var score = new DateTimeOffset(schedule).ToUnixTimeMilliseconds(); + + DateTime scheduleUtc = schedule.ToUniversalTime(); + double score = new DateTimeOffset(scheduleUtc).ToUnixTimeMilliseconds(); + RedisKey key = RedisPipelineExtensions.MessagesSortKey(settings.Resource); await database.SortedSetAddAsync(key, keyValue, score); key = RedisPipelineExtensions.MessageKey(settings.Resource, keyValue); diff --git a/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs b/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs index c00ac7c..0f78d8d 100644 --- a/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs +++ b/RedisMessagePipeline/Consumer/RedisPipelineScheduleConsumer.cs @@ -87,7 +87,7 @@ private async Task TryDequeueAndReserveAsync(CancellationToken cancellatio return true; } - var now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + double now = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); RedisValue[] values = await database.SortedSetRangeByScoreAsync( RedisPipelineExtensions.MessagesSortKey(settings.Resource), From 1e2448d32f1fc54f22c5cb9e07b97f3b73662a97 Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Mon, 6 Apr 2026 18:03:13 -0300 Subject: [PATCH 09/10] ok. --- RedisMessagePipeline.Demo/Program.cs | 2 +- RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs | 9 +++++++++ 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/RedisMessagePipeline.Demo/Program.cs b/RedisMessagePipeline.Demo/Program.cs index 4872467..c5d16d6 100644 --- a/RedisMessagePipeline.Demo/Program.cs +++ b/RedisMessagePipeline.Demo/Program.cs @@ -26,7 +26,7 @@ // Push messages for (int i = 0; i < 10; i++) { - await admin.AddSheduleAsync($"{i}", DateTime.Now.AddSeconds(i * 10), $"message:{i}"); + await admin.AddScheduleAsync($"{i}", DateTime.Now.AddSeconds(i * 10), $"message:{i}"); } // Resume the pipeline, skipping problematic messages if necessary diff --git a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs index ac53b02..a9bb951 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs @@ -2,6 +2,7 @@ using RedLockNet; using StackExchange.Redis; using System; +using System.Security.Cryptography; using System.Threading; using System.Threading.Tasks; @@ -26,6 +27,14 @@ public override async Task AddScheduleAsync(RedisValue keyValue, DateTime schedu { await base.AddScheduleAsync(keyValue, schedule, redisValue); + logger.LogError($"schedule O: {schedule:O}"); + logger.LogError($"schedule: {schedule:yyyy-MM-dd HH:mm:ss}"); + logger.LogError($"schedule kind: {schedule.Kind}"); + logger.LogError($"schedule utc: {schedule.ToUniversalTime():yyyy-MM-dd HH:mm:ss}"); + logger.LogError($"now local: {DateTime.Now:yyyy-MM-dd HH:mm:ss}"); + logger.LogError($"now utc: {DateTime.UtcNow:yyyy-MM-dd HH:mm:ss}"); + + DateTime scheduleUtc = schedule.ToUniversalTime(); double score = new DateTimeOffset(scheduleUtc).ToUnixTimeMilliseconds(); From cef30ce132c3c7b1d0f9c187b0d5aeae7f232a03 Mon Sep 17 00:00:00 2001 From: prospextecnologia Date: Mon, 6 Apr 2026 18:43:10 -0300 Subject: [PATCH 10/10] Ok. --- RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs | 2 +- .../Admin/RedisPipelineBaseAdmin.cs | 2 +- .../Admin/RedisPipelineScheduleAdmin.cs | 15 ++------------- 3 files changed, 4 insertions(+), 15 deletions(-) diff --git a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs index 008f20d..9e2f0e5 100644 --- a/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs +++ b/RedisMessagePipeline/Admin/IRedisPipelineAdmin.cs @@ -8,7 +8,7 @@ namespace RedisMessagePipeline.Admin public interface IRedisPipelineAdmin { Task PushQueueAsync(RedisValue redisValue); - Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, 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/RedisPipelineBaseAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs index 1daab66..0803eb5 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineBaseAdmin.cs @@ -67,7 +67,7 @@ public Task StopAsync() /// public abstract Task ResumeAsync(int skip, CancellationToken cancellationToken); - public virtual Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, RedisValue redisValue) + public virtual Task AddScheduleAsync(RedisValue keyValue, DateTimeOffset schedule, RedisValue redisValue) { if (this.settings.Type != Factory.EnPipelineType.QUEUE_SCHEDULE) { diff --git a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs index a9bb951..b1fb88f 100644 --- a/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs +++ b/RedisMessagePipeline/Admin/RedisPipelineScheduleAdmin.cs @@ -23,21 +23,10 @@ internal RedisPipelineScheduleAdmin( } - public override async Task AddScheduleAsync(RedisValue keyValue, DateTime schedule, RedisValue redisValue) + public override async Task AddScheduleAsync(RedisValue keyValue, DateTimeOffset schedule, RedisValue redisValue) { await base.AddScheduleAsync(keyValue, schedule, redisValue); - - logger.LogError($"schedule O: {schedule:O}"); - logger.LogError($"schedule: {schedule:yyyy-MM-dd HH:mm:ss}"); - logger.LogError($"schedule kind: {schedule.Kind}"); - logger.LogError($"schedule utc: {schedule.ToUniversalTime():yyyy-MM-dd HH:mm:ss}"); - logger.LogError($"now local: {DateTime.Now:yyyy-MM-dd HH:mm:ss}"); - logger.LogError($"now utc: {DateTime.UtcNow:yyyy-MM-dd HH:mm:ss}"); - - - DateTime scheduleUtc = schedule.ToUniversalTime(); - double score = new DateTimeOffset(scheduleUtc).ToUnixTimeMilliseconds(); - + double score = schedule.ToUnixTimeMilliseconds(); RedisKey key = RedisPipelineExtensions.MessagesSortKey(settings.Resource); await database.SortedSetAddAsync(key, keyValue, score); key = RedisPipelineExtensions.MessageKey(settings.Resource, keyValue);