Compare commits

...
3 Commits
Author SHA1 Message Date
xrain 683cf9176a add event loadbalancing 2024-07-27 01:33:17 +08:00
xrain 1ee8d941e9 fix eventserver bug 2024-07-26 19:18:59 +08:00
xrain 1d474f596a config fix 2024-07-26 08:56:01 +08:00
4 changed files with 21 additions and 9 deletions
+6
View File
@@ -13,6 +13,12 @@ public class SimApiSynapseOptions
public string? AppId { get; set; } public string? AppId { get; set; }
public int RpcTimeout { get; set; } = 3; public int RpcTimeout { get; set; } = 3;
/// <summary>
/// Event是否使用负载均衡
/// 也就是订阅$queue主题,消息会分发给不同的AppId
/// 如果false,多个AppId都可以同时收到消息
/// </summary>
public bool EventLoadBalancing { get; set; } = false;
public bool EnableConfigStore { get; set; } = true; public bool EnableConfigStore { get; set; } = true;
public bool DisableEventClient { get; set; } = false; public bool DisableEventClient { get; set; } = false;
public bool DisableRpcClient { get; set; } = false; public bool DisableRpcClient { get; set; } = false;
+5 -5
View File
@@ -224,11 +224,11 @@ public static class SimApiExtensions
/// </summary> /// </summary>
/// <param name="builder"></param> /// <param name="builder"></param>
/// <returns></returns> /// <returns></returns>
public static IApplicationBuilder UseSimApi(this IApplicationBuilder builder) public static WebApplication UseSimApi(this WebApplication builder)
{ {
var options = builder.ApplicationServices.GetRequiredService<SimApiOptions>(); var options = builder.Services.GetRequiredService<SimApiOptions>();
var logger = builder.ApplicationServices.GetRequiredService<ILogger<SimApiOptions>>(); var logger = builder.Services.GetRequiredService<ILogger<SimApiOptions>>();
logger.LogInformation("当前时区: {LocalId}", TimeZoneInfo.Local.Id); logger.LogInformation("当前时区: {LocalId}", TimeZoneInfo.Local.Id);
if (options.EnableForwardHeaders) if (options.EnableForwardHeaders)
@@ -277,7 +277,7 @@ public static class SimApiExtensions
if (options.EnableSimApiStorage) if (options.EnableSimApiStorage)
{ {
logger.LogInformation("开始配置SimApiStorage..."); logger.LogInformation("开始配置SimApiStorage...");
builder.ApplicationServices.GetService<SimApiStorage>(); builder.Services.GetService<SimApiStorage>();
} }
if (options.EnableLowerUrl) if (options.EnableLowerUrl)
@@ -287,7 +287,7 @@ public static class SimApiExtensions
if (options.EnableSynapse) if (options.EnableSynapse)
{ {
var synapse = builder.ApplicationServices.GetRequiredService<Synapse>(); var synapse = builder.Services.GetRequiredService<Synapse>();
synapse.Init(); synapse.Init();
} }
+2 -1
View File
@@ -15,6 +15,7 @@ public partial class Synapse
private bool FireSetConfig(string key, string value) private bool FireSetConfig(string key, string value)
{ {
if (key.Contains('#') || key.Contains('+')) return false;
var topic = $"{Options.SysName}/synapse-config-store/{key}"; var topic = $"{Options.SysName}/synapse-config-store/{key}";
var message = new MqttApplicationMessageBuilder() var message = new MqttApplicationMessageBuilder()
.WithTopic(topic) .WithTopic(topic)
@@ -37,7 +38,7 @@ public partial class Synapse
var csTopicPrefix = $"{Options.SysName}/synapse-config-store/"; var csTopicPrefix = $"{Options.SysName}/synapse-config-store/";
var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
.WithTopicFilter(o => .WithTopicFilter(o =>
o.WithTopic($"{csTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe)) o.WithTopic($"{csTopicPrefix}#").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
.Build(); .Build();
Client!.ApplicationMessageReceivedAsync += e => Client!.ApplicationMessageReceivedAsync += e =>
{ {
+8 -3
View File
@@ -33,7 +33,7 @@ public partial class Synapse
{ {
if (mt!.GetParameters().Length == 2) if (mt!.GetParameters().Length == 2)
{ {
var pt = mt!.GetParameters()[0].ParameterType; var pt = mt.GetParameters()[1].ParameterType;
mt.Invoke(callClass, pt == typeof(string) mt.Invoke(callClass, pt == typeof(string)
? [eventName, reqBody] ? [eventName, reqBody]
: [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]); : [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]);
@@ -45,7 +45,7 @@ public partial class Synapse
} }
catch (Exception ex) catch (Exception ex)
{ {
logger.LogError("Synapse Event Processor Error: {Err}", ex.InnerException); logger.LogError("Synapse Event Processor Error: {Err}\n{Stack}", ex.Message,ex.StackTrace);
} }
} }
@@ -53,7 +53,12 @@ public partial class Synapse
}; };
foreach (var ev in EventRegistry) foreach (var ev in EventRegistry)
{ {
var topic = $"$queue/{esTopicPrefix}{ev.Key}";
var topic = $"{esTopicPrefix}{ev.Key}";
if (Options.EventLoadBalancing)
{
topic = "$queue/" + topic;
}
var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
.WithTopicFilter(o => o.WithTopic(topic)).Build(); .WithTopicFilter(o => o.WithTopic(topic)).Build();
Client.SubscribeAsync(evSubOpts).Wait(); Client.SubscribeAsync(evSubOpts).Wait();