From 617a0119f461a5ec4e1eda61cf2d8d187abb9c21 Mon Sep 17 00:00:00 2001 From: YouWeiDH Date: Sat, 14 Dec 2024 18:36:51 +0800 Subject: [PATCH 1/2] hdong: create token by current user. --- .../BotSharp.Abstraction/Users/IUserService.cs | 1 + .../BotSharp.Core/Users/Services/UserService.cs | 14 ++++++++++++++ 2 files changed, 15 insertions(+) diff --git a/src/Infrastructure/BotSharp.Abstraction/Users/IUserService.cs b/src/Infrastructure/BotSharp.Abstraction/Users/IUserService.cs index 134ca25b..5e56b431 100644 --- a/src/Infrastructure/BotSharp.Abstraction/Users/IUserService.cs +++ b/src/Infrastructure/BotSharp.Abstraction/Users/IUserService.cs @@ -17,6 +17,7 @@ public interface IUserService Task GetAffiliateToken(string authorization); Task GetAdminToken(string authorization); Task GetToken(string authorization); + Task CreateTokenByUser(User user); Task GetMyProfile(); Task VerifyUserNameExisting(string userName); Task VerifyEmailExisting(string email); diff --git a/src/Infrastructure/BotSharp.Core/Users/Services/UserService.cs b/src/Infrastructure/BotSharp.Core/Users/Services/UserService.cs index a87c55fa..18deb7d8 100644 --- a/src/Infrastructure/BotSharp.Core/Users/Services/UserService.cs +++ b/src/Infrastructure/BotSharp.Core/Users/Services/UserService.cs @@ -506,6 +506,20 @@ public class UserService : IUserService return token; } + public async Task CreateTokenByUser(User user) + { + var accessToken = GenerateJwtToken(user); + var jwt = new JwtSecurityTokenHandler().ReadJwtToken(accessToken); + var token = new Token + { + AccessToken = accessToken, + ExpireTime = jwt.Payload.Exp.Value, + TokenType = "Bearer", + Scope = "api" + }; + return token; + } + public async Task VerifyUserNameExisting(string userName) { if (string.IsNullOrEmpty(userName)) From b3bdc01fc8b23940001d3cbfc542331e3f6a4b8b Mon Sep 17 00:00:00 2001 From: Haiping Chen Date: Sun, 15 Dec 2024 16:27:37 +0000 Subject: [PATCH 2/2] Add machine and retry to Redis Event message --- .../Infrastructures/Events/RedisPublisher.cs | 26 +++++++++++------ .../Infrastructures/Events/RedisSubscriber.cs | 28 ++++++++++++------- 2 files changed, 36 insertions(+), 18 deletions(-) diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs index 89de2453..602fde53 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisPublisher.cs @@ -39,10 +39,7 @@ public class RedisPublisher : IEventPublisher // Add a message to the stream, keeping only the latest 1 million messages var messageId = await db.StreamAddAsync(channel, - [ - new NameValueEntry("message", message), - new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o")) - ], + AssembleMessage(message), maxLength: 1000 * 10000); _logger.LogInformation($"Published message {channel} {message} ({messageId})"); @@ -82,6 +79,17 @@ public class RedisPublisher : IEventPublisher return exists; } + private NameValueEntry[] AssembleMessage(RedisValue message, int retry = 0) + { + return + [ + new NameValueEntry("message", message), + new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o")), + new NameValueEntry("machine", Environment.MachineName), + new NameValueEntry("retry", retry), + ]; + } + public async Task ReDispatchAsync(string channel, int count = 10, string order = "asc") { var db = _redis.GetDatabase(); @@ -93,10 +101,12 @@ public class RedisPublisher : IEventPublisher try { - var messageId = await db.StreamAddAsync(channel, [ - new NameValueEntry("message", entry.Values[0].Value), - new NameValueEntry("timestamp", DateTime.UtcNow.ToString("o")) - ]); + var message = entry.Values.First(x => x.Name == "message").Value; + var retryKv = entry.Values.FirstOrDefault(x => x.Name == "retry"); + int.TryParse(retryKv.Value, out int retry); + var messageId = await db.StreamAddAsync(channel, + AssembleMessage(message, retry: retry + 1), + maxLength: 1000 * 10000); _logger.LogWarning($"ReDispatched message: {channel} {entry.Values[0].Value} ({messageId})"); diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs index 9a3c9982..fb035aeb 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/Events/RedisSubscriber.cs @@ -58,23 +58,31 @@ public class RedisSubscriber : IEventSubscriber break; } - if (priorityEnabled) + try { - if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received) > 0) + if (priorityEnabled) { - continue; - } + if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received) > 0) + { + continue; + } - if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received) > 0) + if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received) > 0) + { + continue; + } + + await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received); + } + else { - continue; + await HandleGroupMessage(db, channel, group, consumer, received); } - - await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received); } - else + catch (Exception ex) { - await HandleGroupMessage(db, channel, group, consumer, received); + _logger.LogError($"Error processing message: {ex.Message}\r\n{ex}"); + await Task.Delay(1000 * 60); } } }