From 68df1f14c465535fba7db6fc460d170723c90f4a Mon Sep 17 00:00:00 2001 From: Haiping Chen Date: Sun, 29 Dec 2024 16:23:31 +0000 Subject: [PATCH] Fix ErrorChannel --- .../Infrastructures/Events/RedisSubscriber.cs | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs index b3d910d9..0dedb7c6 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs @@ -64,21 +64,21 @@ public class RedisSubscriber : IEventSubscriber { if (priorityEnabled) { - if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received) > 0) + if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received, $"{channel}-Error") > 0) { continue; } - if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received) > 0) + if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received, $"{channel}-Error") > 0) { continue; } - await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received); + await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received, $"{channel}-Error"); } else { - await HandleGroupMessage(db, channel, group, consumer, received); + await HandleGroupMessage(db, channel, group, consumer, received, $"{channel}-Error"); } } catch (Exception ex) @@ -89,7 +89,7 @@ public class RedisSubscriber : IEventSubscriber } } - private async Task HandleGroupMessage(IDatabase db, string channel, string group, string consumer, Func received) + private async Task HandleGroupMessage(IDatabase db, string channel, string group, string consumer, Func received, string errorChannel) { var entries = await db.StreamReadGroupAsync(channel, group, consumer, count: 1); foreach (var entry in entries) @@ -106,7 +106,7 @@ public class RedisSubscriber : IEventSubscriber _logger.LogError($"Error processing message: {ex.Message}, event id: {channel} {entry.Id} {entry.Values[0].Value}"); // Add a message to the Error stream, keeping only the latest 1 million messages - await db.StreamAddAsync($"{channel}-Error", + await db.StreamAddAsync(errorChannel, RedisPublisher.AssembleErrorMessage(entry.Values[0].Value, ex.Message), messageId: entry.Id, maxLength: 1000 * 10000);