From b4e25f309670313b5de32fd6478bb1bb0926f6e9 Mon Sep 17 00:00:00 2001 From: Haiping Chen Date: Mon, 11 Nov 2024 04:13:50 +0000 Subject: [PATCH] ReDispatchAsync --- .../Infrastructures/Events/IEventPublisher.cs | 2 ++ .../Infrastructures/Events/RedisPublisher.cs | 25 +++++++++++++++++++ 2 files changed, 27 insertions(+) diff --git a/src/Infrastructure/BotSharp.Abstraction/Infrastructures/Events/IEventPublisher.cs b/src/Infrastructure/BotSharp.Abstraction/Infrastructures/Events/IEventPublisher.cs index 471d810e..e242937e 100644 --- a/src/Infrastructure/BotSharp.Abstraction/Infrastructures/Events/IEventPublisher.cs +++ b/src/Infrastructure/BotSharp.Abstraction/Infrastructures/Events/IEventPublisher.cs @@ -11,4 +11,6 @@ public interface IEventPublisher Task BroadcastAsync(string channel, string message); Task PublishAsync(string channel, string message); + + Task ReDispatchAsync(string channel, int count = 10, string order = "asc"); } diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs index 4f3ad405..f6daa59d 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs @@ -29,4 +29,29 @@ public class RedisPublisher : IEventPublisher _logger.LogInformation($"Published message {channel} {message}"); } + + public async Task ReDispatchAsync(string channel, int count = 10, string order = "asc") + { + var db = _redis.GetDatabase(); + + var entries = await db.StreamRangeAsync(channel, "-", "+", count: count, messageOrder: order == "asc" ? Order.Ascending : Order.Descending); + foreach (var entry in entries) + { + _logger.LogInformation($"Fetched message: {channel} {entry.Values[0].Value} ({entry.Id})"); + + try + { + var messageId = await db.StreamAddAsync(channel, "message", entry.Values[0].Value); + + _logger.LogWarning($"ReDispatched message: {channel} {entry.Values[0].Value} ({messageId})"); + + // Optionally delete the message to save space + await db.StreamDeleteAsync(channel, [entry.Id]); + } + catch (Exception ex) + { + _logger.LogError($"Error processing message: {ex.Message}, event id: {channel} {entry.Id}"); + } + } + } }