ReDispatchAsync

This commit is contained in:
Haiping Chen 2024-11-11 04:13:50 +00:00
parent d2636a6eaf
commit b4e25f3096
2 changed files with 27 additions and 0 deletions

View file

@ -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");
}

View file

@ -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}");
}
}
}
}