From a14981e0322c5f074e49876c544b9bc1f62ff2e1 Mon Sep 17 00:00:00 2001 From: Haiping Chen Date: Mon, 18 Nov 2024 22:31:18 +0000 Subject: [PATCH] Add Redis timestamp --- .../Infrastructures/Events/RedisPublisher.cs | 8 ++++++-- .../Infrastructures/Events/RedisSubscriber.cs | 7 ++----- 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs index f6daa59d..08ce275c 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs @@ -24,7 +24,11 @@ public class RedisPublisher : IEventPublisher { var db = _redis.GetDatabase(); // Add a message to the stream, keeping only the latest 1 million messages - await db.StreamAddAsync(channel, "message", message, + await db.StreamAddAsync(channel, + [ + new NameValueEntry("message", message), + new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o")) + ], maxLength: 1000 * 10000); _logger.LogInformation($"Published message {channel} {message}"); @@ -50,7 +54,7 @@ public class RedisPublisher : IEventPublisher } catch (Exception ex) { - _logger.LogError($"Error processing message: {ex.Message}, event id: {channel} {entry.Id}"); + _logger.LogError($"Error processing message: {ex.Message}, event id: {channel} {entry.Id}\r\n{ex}"); } } } diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs index 726b9358..bf6652b1 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs @@ -50,6 +50,7 @@ public class RedisSubscriber : IEventSubscriber foreach (var entry in entries) { _logger.LogInformation($"Consumer {Environment.MachineName} received: {channel} {entry.Values[0].Value}"); + await db.StreamAcknowledgeAsync(channel, group, entry.Id); try { @@ -60,11 +61,7 @@ public class RedisSubscriber : IEventSubscriber } catch (Exception ex) { - _logger.LogError($"Error processing message: {ex.Message}, event id: {channel} {entry.Id}"); - } - finally - { - await db.StreamAcknowledgeAsync(channel, group, entry.Id); + _logger.LogError($"Error processing message: {ex.Message}, event id: {channel} {entry.Id}\r\n{ex}"); } }