diff --git a/Communications/SimApiBaseResponse.cs b/Communications/SimApiBaseResponse.cs index e36186c..2c7a01a 100644 --- a/Communications/SimApiBaseResponse.cs +++ b/Communications/SimApiBaseResponse.cs @@ -54,7 +54,7 @@ public class SimApiBaseResponse(int code = 200, string message = "成功") /// public class SimApiBasePageResponse() : SimApiBaseResponse { - public T List { get; set; } + public T? List { get; set; } public int Page { get; set; } = 1; public int Count { get; set; } = 20; public int Total { get; set; } @@ -74,7 +74,7 @@ public class SimApiBasePageResponse() : SimApiBaseResponse /// public class SimApiBaseResponse() : SimApiBaseResponse { - public T Data { get; set; } + public T? Data { get; set; } public SimApiBaseResponse(T data) : this() { diff --git a/Configurations/SimApiDocOptions.cs b/Configurations/SimApiDocOptions.cs index 7112892..d921caa 100644 --- a/Configurations/SimApiDocOptions.cs +++ b/Configurations/SimApiDocOptions.cs @@ -11,17 +11,17 @@ public class SimApiDocGroupOption /// /// 文档标识 /// - public string Id { get; set; } + public string Id { get; set; } = null!; /// /// 文档名称 /// - public string Name { get; set; } + public string Name { get; set; } = null!; /// /// 文档描述 /// - public string Description { get; set; } + public string Description { get; set; } = null!; } /// @@ -33,11 +33,11 @@ public class SimApiAuthOption public string Description { get; set; } = "认证服务器颁发的AccessToken"; - public string AuthorizationUrl { get; set; } + public string AuthorizationUrl { get; set; } = null!; - public string TokenUrl { get; set; } + public string TokenUrl { get; set; } = null!; - public Dictionary Scopes { get; set; } + public Dictionary Scopes { get; set; } = null!; } /// diff --git a/Configurations/SimApiOptions.cs b/Configurations/SimApiOptions.cs index 92d3e32..6141870 100644 --- a/Configurations/SimApiOptions.cs +++ b/Configurations/SimApiOptions.cs @@ -14,13 +14,13 @@ public class SimApiOptions /// 启用SimApiAuth,一个简单的基于Header Token的认证方式。 /// default: false /// - public bool EnableSimApiAuth { get; set; } = false; + public bool EnableSimApiAuth { get; set; } /// /// 启用在线文档,启用后 访问 /swagger 可以查看对应的api文档。 /// default: false /// - public bool EnableSimApiDoc { get; set; } = false; + public bool EnableSimApiDoc { get; set; } /// /// 启用异常拦截,启用后,所有的异常将被通过json反馈。 @@ -32,7 +32,7 @@ public class SimApiOptions /// 开启S3兼容的存储系统。 /// default: false /// - public bool EnableSimApiStorage { get; set; } = false; + public bool EnableSimApiStorage { get; set; } /// /// 开启ForwardHeaders,开启后可以透传负载均衡的Headers @@ -50,12 +50,12 @@ public class SimApiOptions /// 启用格式化的 Console Logger /// default: false /// - public bool EnableLogger { get; set; } = false; + public bool EnableLogger { get; set; } /// /// 是否启用Synapse /// - public bool EnableSynapse { get; set; } = false; + public bool EnableSynapse { get; set; } /// @@ -70,7 +70,7 @@ public class SimApiOptions public SimApiSynapseOptions SimApiSynapseOptions { get; set; } = new(); - public void ConfigureSimApiSynapse(Action options = null) + public void ConfigureSimApiSynapse(Action? options = null) { options?.Invoke(SimApiSynapseOptions); } @@ -80,12 +80,12 @@ public class SimApiOptions SimApiSynapseOptions = options; } - public void ConfigureSimApiDoc(Action options = null) + public void ConfigureSimApiDoc(Action? options = null) { options?.Invoke(SimApiDocOptions); } - public void ConfigureSimApiStorage(Action options = null) + public void ConfigureSimApiStorage(Action? options = null) { options?.Invoke(SimApiStorageOptions); } diff --git a/Logger/SimApiLogger.cs b/Logger/SimApiLogger.cs index 21e8730..8cbd919 100644 --- a/Logger/SimApiLogger.cs +++ b/Logger/SimApiLogger.cs @@ -1,13 +1,11 @@ using System; using Microsoft.Extensions.Logging; -using SimApi.Helpers; namespace SimApi.Logger; public class SimApiLogger(string name) : ILogger { - public IDisposable BeginScope(TState state) => default!; - + public IDisposable BeginScope(TState state) where TState : notnull => default!; public bool IsEnabled(LogLevel logLevel) => true; public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, @@ -32,4 +30,13 @@ public class SimApiLogger(string name) : ILogger Console.WriteLine(message); Console.ResetColor(); } +} + +public class EmptyDisposable : IDisposable +{ + public static EmptyDisposable Instance { get; } = new EmptyDisposable(); + + public void Dispose() + { + } } \ No newline at end of file diff --git a/SimApiExtensions.cs b/SimApiExtensions.cs index a2aad95..6e04b25 100644 --- a/SimApiExtensions.cs +++ b/SimApiExtensions.cs @@ -19,7 +19,7 @@ public static class SimApiExtensions { //**********快捷添加************** public static IServiceCollection AddSimApi(this IServiceCollection builder, - Action options = null) + Action? options = null) { var simApiOptions = new SimApiOptions(); options?.Invoke(simApiOptions); diff --git a/Synapse/ConfigStore.cs b/Synapse/ConfigStore.cs index 8b20997..644c838 100644 --- a/Synapse/ConfigStore.cs +++ b/Synapse/ConfigStore.cs @@ -13,6 +13,8 @@ public partial class Synapse public event EventHandler? OnConfigChanged; private Dictionary CurrentConfig { get; } = new(); + private string ConfigStoreTopicPrefix => $"{Options.SysName}/synapse-config-store/"; + private bool FireSetConfig(string key, string value) { if (key.Contains('#') || key.Contains('+')) return false; @@ -35,22 +37,26 @@ public partial class Synapse private void RunConfigStoreServer() { - var csTopicPrefix = $"{Options.SysName}/synapse-config-store/"; - var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() - .WithTopicFilter(o => - o.WithTopic($"{csTopicPrefix}#").WithRetainHandling(MqttRetainHandling.SendAtSubscribe)) - .Build(); Client!.ApplicationMessageReceivedAsync += e => { - if (!e.ApplicationMessage.Topic.StartsWith(csTopicPrefix)) return Task.CompletedTask; + if (!e.ApplicationMessage.Topic.StartsWith(ConfigStoreTopicPrefix)) return Task.CompletedTask; var reqBody = e.ApplicationMessage.ConvertPayloadToString(); - var eventName = e.ApplicationMessage.Topic.Replace(csTopicPrefix, string.Empty); + var eventName = e.ApplicationMessage.Topic.Replace(ConfigStoreTopicPrefix, string.Empty); CurrentConfig[eventName] = reqBody; OnConfigChanged?.Invoke(this, new ConfigStoreItem(eventName, reqBody)); logger.LogDebug("Synapse Config Changed: {Config} => {Data}", eventName, reqBody); return Task.CompletedTask; }; - Client.SubscribeAsync(eventSubOpts).Wait(); + SubConfigStoreServerTopic(); + } + + private void SubConfigStoreServerTopic() + { + var csSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() + .WithTopicFilter(o => + o.WithTopic($"{ConfigStoreTopicPrefix}#").WithRetainHandling(MqttRetainHandling.SendAtSubscribe)) + .Build(); + Client!.SubscribeAsync(csSubOpts).Wait(); } public record ConfigStoreItem(string Key, string Value); diff --git a/Synapse/EventServer.cs b/Synapse/EventServer.cs index d196cff..ecd9f2e 100644 --- a/Synapse/EventServer.cs +++ b/Synapse/EventServer.cs @@ -12,14 +12,16 @@ namespace SimApi; public partial class Synapse { + + private string EventServerTopicPrefix => $"{Options.SysName}/event/"; + private void RunEventServer() { - var esTopicPrefix = $"{Options.SysName}/event/"; Client!.ApplicationMessageReceivedAsync += e => { - if (!e.ApplicationMessage.Topic.StartsWith(esTopicPrefix)) return Task.CompletedTask; + if (!e.ApplicationMessage.Topic.StartsWith(EventServerTopicPrefix)) return Task.CompletedTask; var reqBody = e.ApplicationMessage.ConvertPayloadToString(); - var eventName = e.ApplicationMessage.Topic.Replace(esTopicPrefix, string.Empty); + var eventName = e.ApplicationMessage.Topic.Replace(EventServerTopicPrefix, string.Empty); logger.LogDebug("Synapse Event Receive: {EventName}\n{Body}", eventName, reqBody); var methods = EventRegistry .Where(x => Regex.IsMatch(eventName, @@ -48,20 +50,23 @@ public partial class Synapse logger.LogError("Synapse Event Processor Error: {Err}\n{Stack}", ex.Message,ex.StackTrace); } } - return Task.CompletedTask; }; + SubRpcServerTopic(); + } + + private void SubEventServerTopic() + { foreach (var ev in EventRegistry) { - - var topic = $"{esTopicPrefix}{ev.Key}"; + var topic = $"{EventServerTopicPrefix}{ev.Key}"; if (Options.EventLoadBalancing) { topic = "$queue/" + topic; } var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() .WithTopicFilter(o => o.WithTopic(topic)).Build(); - Client.SubscribeAsync(evSubOpts).Wait(); + Client!.SubscribeAsync(evSubOpts).Wait(); logger.LogDebug("Synapse Event Register Event Success: {EventName}\nFull Topic: {Topic}", ev.Key, topic); } } diff --git a/Synapse/RpcClient.cs b/Synapse/RpcClient.cs index 83836a0..67bb82c 100644 --- a/Synapse/RpcClient.cs +++ b/Synapse/RpcClient.cs @@ -14,22 +14,28 @@ public partial class Synapse { private Dictionary> ResponseCompletionSources { get; } = new(); + private string EventClientTopicPrefix => $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/"; + private void RunRpcClient() { - var rcTopic = $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/"; - var rcSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() - .WithTopicFilter(o => o.WithTopic($"{rcTopic}+")).Build(); Client!.ApplicationMessageReceivedAsync += e => { - if (!e.ApplicationMessage.Topic.StartsWith(rcTopic)) return Task.CompletedTask; + if (!e.ApplicationMessage.Topic.StartsWith(EventClientTopicPrefix)) return Task.CompletedTask; var reqBody = e.ApplicationMessage.ConvertPayloadToString(); - var messageId = e.ApplicationMessage.Topic.Replace(rcTopic, string.Empty); + var messageId = e.ApplicationMessage.Topic.Replace(EventClientTopicPrefix, string.Empty); if (!ResponseCompletionSources.TryGetValue(messageId, out var tcs)) return Task.CompletedTask; tcs.SetResult(reqBody); ResponseCompletionSources.Remove(messageId); return Task.CompletedTask; }; - Client.SubscribeAsync(rcSubOpts).Wait(); + SubRpcClientTopic(); + } + + private void SubRpcClientTopic() + { + var rcSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() + .WithTopicFilter(o => o.WithTopic($"{EventClientTopicPrefix}+")).Build(); + Client!.SubscribeAsync(rcSubOpts).Wait(); } private string? FireRpc(string app, string action, object? param) diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs index 22a36f7..f4957b2 100644 --- a/Synapse/RpcServer.cs +++ b/Synapse/RpcServer.cs @@ -16,18 +16,15 @@ namespace SimApi; public partial class Synapse { + private string RpcServerTopicPrefix => $"{Options.SysName}/{Options.AppName}/rpc/server/"; + private void RunRpcServer() { - var rsTopicPrefix = $"{Options.SysName}/{Options.AppName}/rpc/server/"; - var rsSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() - .WithTopicFilter(o => - o.WithTopic($"$queue/{rsTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe)) - .Build(); Client!.ApplicationMessageReceivedAsync += e => { - if (!e.ApplicationMessage.Topic.StartsWith(rsTopicPrefix)) return Task.CompletedTask; + if (!e.ApplicationMessage.Topic.StartsWith(RpcServerTopicPrefix)) return Task.CompletedTask; var reqBody = e.ApplicationMessage.ConvertPayloadToString(); - var action = e.ApplicationMessage.Topic.Replace(rsTopicPrefix, string.Empty); + var action = e.ApplicationMessage.Topic.Replace(RpcServerTopicPrefix, string.Empty); var appInfo = e.ApplicationMessage.ResponseTopic.Split(","); logger.LogDebug( "Synapse RPC Server Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}", @@ -93,6 +90,15 @@ public partial class Synapse return Task.CompletedTask; }; - Client.SubscribeAsync(rsSubOpts).Wait(); + SubRpcServerTopic(); + } + + public void SubRpcServerTopic() + { + var rsSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() + .WithTopicFilter(o => + o.WithTopic($"$queue/{RpcServerTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe)) + .Build(); + Client!.SubscribeAsync(rsSubOpts).Wait(); } } \ No newline at end of file diff --git a/Synapse/Synapse.cs b/Synapse/Synapse.cs index c2b4bfb..06f7c1b 100644 --- a/Synapse/Synapse.cs +++ b/Synapse/Synapse.cs @@ -41,8 +41,45 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge Options.SysName, Options.AppName, Options.AppId); CreateConnection(); ProcessAttribute(); + //事件客户端 + if (Options.DisableEventClient) + { + logger.LogWarning("Synapse Event Client Disabled: DisableEventClient set true"); + } + else + { + logger.LogInformation("Synapse Event Client Ready"); + } + + //RPC客户端 + if (Options.DisableRpcClient) + { + logger.LogWarning("Synapse Rpc Client Disabled: DisableEventClient set true"); + } + else + { + RunRpcClient(); + logger.LogInformation("Synapse Rpc Client Ready, Client Timeout: {OptionsRpcTimeout}s", Options.RpcTimeout); + } + + if (RpcRegistry.Count > 0) + { + RunRpcServer(); + } + + if (EventRegistry.Count > 0) + { + RunEventServer(); + } + + if (Options.EnableConfigStore) + { + RunConfigStoreServer(); + logger.LogInformation("Synapse Config Store Ready [{SysName}] ...", Options.SysName); + } } - + + /// /// 调用RPC使用明确的返回值类型 /// @@ -128,49 +165,18 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge Client.ConnectedAsync += _ => { logger.LogInformation("Synapse MQTT[{AppName}:{AppId}] 连接成功...", Options.AppName, Options.AppId); - //事件客户端 - if (Options.DisableEventClient) - { - logger.LogWarning("Synapse Event Client Disabled: DisableEventClient set true"); - } - else - { - logger.LogInformation("Synapse Event Client Ready"); - } - //RPC客户端 - if (Options.DisableRpcClient) - { - logger.LogWarning("Synapse Rpc Client Disabled: DisableEventClient set true"); - } - else - { - RunRpcClient(); - logger.LogInformation("Synapse Rpc Client Ready, Client Timeout: {OptionsRpcTimeout}s", Options.RpcTimeout); - } - - if (RpcRegistry.Count > 0) - { - RunRpcServer(); - } - - if (EventRegistry.Count > 0) - { - RunEventServer(); - } - - if (Options.EnableConfigStore) - { - RunConfigStoreServer(); - logger.LogInformation("Synapse Config Store Ready [{SysName}] ...", Options.SysName); - } + if (!Options.DisableRpcClient) SubRpcClientTopic(); + if (RpcRegistry.Count > 0) SubRpcServerTopic(); + if (EventRegistry.Count > 0) SubEventServerTopic(); + if (Options.EnableConfigStore) SubConfigStoreServerTopic(); return Task.CompletedTask; }; //重连 Client.DisconnectedAsync += async _ => { logger.LogError("Synapse MQTT[{AppName}:{AppId}] 断开连接,开始重连...", Options.AppName, Options.AppId); - await Task.Delay(TimeSpan.FromSeconds(3)); + await Task.Delay(TimeSpan.FromSeconds(5)); try { logger.LogInformation("Synapse MQTT[{AppName}:{AppId}] 开始连接MQTT服务器...", Options.AppName, Options.AppId);