From f4bc6dc080646e35c33789149c7f64fca4639056 Mon Sep 17 00:00:00 2001 From: xRain Date: Fri, 26 Jul 2024 05:03:28 +0800 Subject: [PATCH] fix synapse use mqtt --- Configurations/SimApiSynapseOptions.cs | 18 +-- SimApi.csproj | 2 +- SimApiExtensions.cs | 30 +++++ Synapse/EventClient.cs | 32 +++--- Synapse/EventServer.cs | 45 ++++---- Synapse/RpcClient.cs | 71 ++++++------ Synapse/RpcServer.cs | 122 ++++++++++---------- Synapse/Synapse.cs | 150 +++++++++---------------- 8 files changed, 221 insertions(+), 249 deletions(-) diff --git a/Configurations/SimApiSynapseOptions.cs b/Configurations/SimApiSynapseOptions.cs index 1bee12c..bd3dc45 100644 --- a/Configurations/SimApiSynapseOptions.cs +++ b/Configurations/SimApiSynapseOptions.cs @@ -2,28 +2,22 @@ namespace SimApi.Configurations; public class SimApiSynapseOptions { - public string MqHost { get; set; } + /// + /// Mqtt服务器的Websocket地址 + /// + public string Websocket { get; set; } - public int MqPort { get; set; } + public string Username { get; set; } - public string MqUser { get; set; } - - public string MqPass { get; set; } - - public string MqVHost { get; set; } = "/"; + public string Password { get; set; } public string SysName { get; set; } - public string AppName { get; set; } public string AppId { get; set; } public int RpcTimeout { get; set; } = 3; - public ushort EventProcessorNum { get; set; } = 20; - - public ushort RpcProcessorNum { get; set; } = 20; - public bool DisableEventClient { get; set; } = false; public bool DisableRpcClient { get; set; } = false; diff --git a/SimApi.csproj b/SimApi.csproj index 55eed49..bd0efdc 100644 --- a/SimApi.csproj +++ b/SimApi.csproj @@ -26,7 +26,7 @@ - + diff --git a/SimApiExtensions.cs b/SimApiExtensions.cs index 4e6dbf8..ee1ac44 100644 --- a/SimApiExtensions.cs +++ b/SimApiExtensions.cs @@ -5,6 +5,7 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.OpenApi.Models; using SimApi.Middlewares; using Microsoft.AspNetCore.HttpOverrides; +using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using SimApi.Configurations; using SimApi.Logger; @@ -189,6 +190,35 @@ public static class SimApiExtensions return builder; } + public static IHost UseSimApi(this IHost builder) + { + var options = builder.Services.GetRequiredService(); + + var logger = builder.Services.GetRequiredService>(); + + logger.LogInformation("当前时区: {LocalId}", TimeZoneInfo.Local.Id); + + //请求一下检测存储错误 + if (options.EnableSimApiStorage) + { + logger.LogInformation("开始配置SimApiStorage..."); + builder.Services.GetService(); + } + + if (options.EnableLowerUrl) + { + logger.LogInformation("开始配置使用URL小写..."); + } + + if (options.EnableSynapse) + { + var synapse = builder.Services.GetRequiredService(); + synapse.Init(); + } + + return builder; + } + /// /// 使用所有SimApi自定义中间件 /// diff --git a/Synapse/EventClient.cs b/Synapse/EventClient.cs index 2e3f87e..1b6f89c 100644 --- a/Synapse/EventClient.cs +++ b/Synapse/EventClient.cs @@ -1,31 +1,25 @@ -using System; -using System.Text; -using System.Text.Encodings.Web; using System.Text.Json; -using System.Text.Unicode; +using System.Threading; using Microsoft.Extensions.Logging; +using MQTTnet; using SimApi.Helpers; namespace SimApi; public partial class Synapse { - private void RunEventClient() - { - EventClientChannel = CreateChannel(0, "EventClient"); - } - - private void FireEvent(string eventName, object param) + private bool FireEvent(string eventName, object param, bool retain = false) { var paramJson = JsonSerializer.Serialize(param, SimApiUtil.JsonOption); - var router = $"event.{Options.AppName}.{eventName}"; - var props = EventClientChannel.CreateBasicProperties(); - props.AppId = Options.AppId; - props.MessageId = Guid.NewGuid().ToString(); - props.ReplyTo = Options.AppName; - props.Type = eventName; - EventClientChannel.BasicPublish(Options.SysName, router, false, props, Encoding.UTF8.GetBytes(paramJson)); - Logger.LogDebug("Event Publish: {OptionsAppName}.{EventName}\n{ParamJson}", Options.AppName, eventName, - paramJson); + var topic = $"{Options.SysName}/{Options.AppName}/event/{eventName}"; + var message = new MqttApplicationMessageBuilder() + .WithTopic(topic) + .WithPayload(paramJson) + .WithRetainFlag(retain) + .Build(); + if (!Client!.IsConnected) return false; + Client!.PublishAsync(message, CancellationToken.None).Wait(); + logger.LogDebug("Event Publish: {Event}@{App} {Json}", eventName, Options.AppName, paramJson); + return true; } } \ No newline at end of file diff --git a/Synapse/EventServer.cs b/Synapse/EventServer.cs index 5b2b843..4dfb91c 100644 --- a/Synapse/EventServer.cs +++ b/Synapse/EventServer.cs @@ -1,10 +1,11 @@ using System; using System.Linq; -using System.Text; using System.Text.Json; +using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; -using RabbitMQ.Client.Events; +using MQTTnet; +using MQTTnet.Protocol; using SimApi.Helpers; namespace SimApi; @@ -13,38 +14,36 @@ public partial class Synapse { private void RunEventServer() { - EventServerChannel = CreateChannel(Options.EventProcessorNum, "EventServer"); - var queue = $"{Options.SysName}_{Options.AppName}_event"; - EventServerChannel.QueueDeclare(queue, true, false, true, null); - foreach (var ev in EventRegistry.Where(ev => !ev.Key.Contains('*') && !ev.Key.Contains('#'))) + var esTopicPrefix = $"{Options.SysName}/{Options.AppName}/event/"; + var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() + .WithTopicFilter(o => + o.WithTopic($"$queue/{esTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe)) + .Build(); + Client.ApplicationMessageReceivedAsync += e => { - EventServerChannel.QueueBind(queue, Options.SysName, $"event.{ev.Key}", null); - } - var consumer = new EventingBasicConsumer(EventServerChannel); - consumer.Received += (ch, ea) => - { - var reqBody = Encoding.UTF8.GetString(ea.Body.ToArray()); - Logger.LogDebug("Event Receive: {BasicPropertiesReplyTo}.{BasicPropertiesType}\n{S}", - ea.BasicProperties.ReplyTo, ea.BasicProperties.Type, reqBody); + if (!e.ApplicationMessage.Topic.StartsWith(esTopicPrefix)) return Task.CompletedTask; + var reqBody = e.ApplicationMessage.ConvertPayloadToString(); + var eventName = e.ApplicationMessage.Topic.Replace(esTopicPrefix, string.Empty); + logger.LogDebug("Synapse Event Receive: {AppName}.{EventName}\n{Body}", + Options.AppName, eventName, reqBody); - var key = ea.RoutingKey.Replace("event.", string.Empty); - var method = EventRegistry.FirstOrDefault(x => x.Key == key); - var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class); + var method = EventRegistry.FirstOrDefault(x => x.Key == eventName); + if (method == null) return Task.CompletedTask; + var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method!.Class); var mt = callClass.GetType().GetMethod(method.Method); - var pt = mt.GetParameters()[0].ParameterType; + var pt = mt!.GetParameters()[0].ParameterType; try { mt.Invoke(callClass, pt == typeof(string) ? new object[] { reqBody } : new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }); - EventServerChannel.BasicAck(ea.DeliveryTag, false); } - catch (Exception e) + catch (Exception ex) { - Logger.LogError("Event Processor Error: {Err}", e.InnerException); - EventServerChannel.BasicNack(ea.DeliveryTag, false, false); + logger.LogError("SynapseEvent Processor Error: {Err}", ex.InnerException); } + return Task.CompletedTask; }; - EventServerChannel.BasicConsume(queue, false, "", false, false, null, consumer); + Client.SubscribeAsync(eventSubOpts).Wait(); } } \ No newline at end of file diff --git a/Synapse/RpcClient.cs b/Synapse/RpcClient.cs index 767db34..968acf4 100644 --- a/Synapse/RpcClient.cs +++ b/Synapse/RpcClient.cs @@ -1,9 +1,10 @@ using System; using System.Collections.Generic; -using System.Text; using System.Text.Json; +using System.Threading; +using System.Threading.Tasks; using Microsoft.Extensions.Logging; -using RabbitMQ.Client.Events; +using MQTTnet; using SimApi.Communications; using SimApi.Helpers; @@ -11,42 +12,43 @@ namespace SimApi; public partial class Synapse { - private Dictionary ResponseCache { get; } = new(); + private Dictionary ResponseCache { get; } = new(); private void RunRpcClient() { - RpcClientChannel = CreateChannel(0, "RpcClient"); - var queue = $"{Options.SysName}_{Options.AppName}_client_{Options.AppId}"; - var router = $"client.{Options.AppName}.{Options.AppId}"; - RpcClientChannel.QueueDeclare(queue, true, false, true, null); - RpcClientChannel.QueueBind(queue, Options.SysName, router, null); - var consumer = new EventingBasicConsumer(RpcClientChannel); - consumer.Received += (ch, ea) => + var rcTopic = $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/"; + var rcSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() + .WithTopicFilter(o => o.WithTopic($"{rcTopic}+")).Build(); + Client.ApplicationMessageReceivedAsync += e => { - ResponseCache.Add(ea.BasicProperties.CorrelationId, ea.Body.ToArray()); - RpcClientChannel.BasicAck(ea.DeliveryTag, false); - Logger.LogDebug( - "RPC Response: ({BasicPropertiesCorrelationId}) {BasicPropertiesType}@{BasicPropertiesReplyTo} -> {OptionsAppName}\n{S}", - ea.BasicProperties.CorrelationId, ea.BasicProperties.Type, ea.BasicProperties.ReplyTo, Options.AppName, - Encoding.UTF8.GetString(ea.Body.ToArray())); + if (!e.ApplicationMessage.Topic.StartsWith(rcTopic)) return Task.CompletedTask; + var reqBody = e.ApplicationMessage.ConvertPayloadToString(); + var messageId = e.ApplicationMessage.Topic.Replace(rcTopic, string.Empty); + ResponseCache.Add(messageId, reqBody); + logger.LogDebug("Synapse RPC Client Message: ({BasicPropertiesCorrelationId}) => {S}", messageId, reqBody); + return Task.CompletedTask; }; - RpcClientChannel.BasicConsume(queue, false, "", false, false, null, consumer); + Client.SubscribeAsync(rcSubOpts).Wait(); } private string FireRpc(string app, string action, object param) { var paramJson = JsonSerializer.Serialize(param, SimApiUtil.JsonOption); - var router = $"server.{app}"; string response; - var props = RpcClientChannel.CreateBasicProperties(); - props.AppId = Options.AppId; - props.MessageId = Guid.NewGuid().ToString(); - props.Type = action; - props.ReplyTo = Options.AppName; - RpcClientChannel.BasicPublish(Options.SysName, router, false, props, Encoding.UTF8.GetBytes(paramJson)); - Logger.LogDebug("RPC Request: ({PropsMessageId}) {OptionsAppName} -> {Action}@{App}\n{ParamJson}", - props.MessageId, - Options.AppName, action, app, paramJson); + var topic = $"{Options.SysName}/{app}/rpc/server/{action}"; + var messageId = Guid.NewGuid().ToString(); + var message = new MqttApplicationMessageBuilder() + .WithTopic(topic) + .WithPayload(paramJson) + .WithResponseTopic($"{Options.AppName},{Options.AppId}") + .WithContentType(messageId) + .WithRetainFlag(false) + .Build(); + if (!Client!.IsConnected) return null; + Client!.PublishAsync(message, CancellationToken.None).Wait(); + logger.LogDebug( + "Synapse RPC Client Request: ({PropsMessageId}) {OptionsAppName} -> {Action}@{App}\n{ParamJson}", + messageId, Options.AppName, action, app, paramJson); var ts = SimApiUtil.TimestampNow; while (true) { @@ -55,13 +57,16 @@ public partial class Synapse response = JsonSerializer.Serialize(new SimApiBaseResponse(502, "timeout"), SimApiUtil.JsonOption); break; } - if (ResponseCache.TryGetValue(props.MessageId, out var value)) - { - response = Encoding.UTF8.GetString(value); - ResponseCache.Remove(props.MessageId); - break; - } + + if (!ResponseCache.TryGetValue(messageId, out var value)) continue; + response = value; + ResponseCache.Remove(messageId); + logger.LogDebug( + "Synapse RPC Client Response: ({BasicPropertiesCorrelationId}) {BasicPropertiesType}@{BasicPropertiesReplyTo} -> {OptionsAppName}\n{S}", + messageId, action, app, Options.AppName, response); + break; } + return response; } } \ No newline at end of file diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs index 4981640..e983729 100644 --- a/Synapse/RpcServer.cs +++ b/Synapse/RpcServer.cs @@ -1,10 +1,12 @@ using System; using System.Linq; using System.Reflection; -using System.Text; +using System.Threading; +using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; -using RabbitMQ.Client.Events; +using MQTTnet; +using MQTTnet.Protocol; using SimApi.Communications; using SimApi.Exceptions; using SimApi.Helpers; @@ -16,78 +18,70 @@ public partial class Synapse { private void RunRpcServer() { - RpcServerChannel = CreateChannel(Options.RpcProcessorNum, "RpcServer"); - var queue = $"{Options.SysName}_{Options.AppName}_server"; - var router = $"server.{Options.AppName}"; - RpcServerChannel.QueueDeclare(queue, true, false, true, null); - RpcServerChannel.QueueBind(queue, Options.SysName, router, null); - var consumer = new EventingBasicConsumer(RpcServerChannel); - consumer.Received += (ch, ea) => + 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 => { - var reqBody = Encoding.UTF8.GetString(ea.Body.ToArray()); - Logger.LogDebug( - "RPC Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}", - ea.BasicProperties.MessageId, ea.BasicProperties.ReplyTo, ea.BasicProperties.Type, Options.AppName, + if (!e.ApplicationMessage.Topic.StartsWith(rsTopicPrefix)) return Task.CompletedTask; + var reqBody = e.ApplicationMessage.ConvertPayloadToString(); + var action = e.ApplicationMessage.Topic.Replace(rsTopicPrefix, string.Empty); + var appInfo = e.ApplicationMessage.ResponseTopic.Split(","); + logger.LogDebug( + "Synapse RPC Server Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}", + e.ApplicationMessage.ContentType, appInfo[0], action, Options.AppName, reqBody); var res = new SimApiBaseResponse(404, "method not found"); - var method = RpcRegistry.FirstOrDefault(x => x.Key == ea.BasicProperties.Type); - if (method != null) + var method = RpcRegistry.FirstOrDefault(x => x.Key == action); + if (method == null) return Task.CompletedTask; + var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class); + var mt = callClass.GetType().GetMethod(method.Method); + try { - var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class); - var mt = callClass.GetType().GetMethod(method.Method); - var param = Array.Empty(); - try + var pt = mt!.GetParameters()[0].ParameterType; + var param = pt == typeof(string) + ? [reqBody] + : new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }; + var ret = mt.Invoke(callClass, param); + res = new SimApiBaseResponse { - var pt = mt.GetParameters()[0].ParameterType; - if (pt == typeof(string)) - { - param = new[] { reqBody }; - } - else - { - var paramObj = JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption); - param = new[] { paramObj }; - } - - var ret = mt.Invoke(callClass, param); - res = new SimApiBaseResponse - { - Data = ret - }; + Data = ret + }; + } + catch (TargetInvocationException ex) + { + if (ex.InnerException is SimApiException ie) + { + logger.LogDebug("Synapse RPC调用错误: {Err}", ie.Message); + res = new SimApiBaseResponse(ie.Code, ie.Message); } - catch (TargetInvocationException e) + else { - if (e.InnerException is SimApiException ie) - { - Logger.LogDebug("RPC调用错误: {Err}", ie.Message); - res = new SimApiBaseResponse(ie.Code, ie.Message); - } - else - { - res = new SimApiBaseResponse(500, e.Message); - } - } - catch (Exception e) - { - Logger.LogDebug("RPC调用失败: {Err}", e.Message); - res = new SimApiBaseResponse(500, e.Message); + res = new SimApiBaseResponse(500, ex.Message); } } - + catch (Exception ex) + { + logger.LogDebug("Synapse RPC调用失败: {Err}", ex.Message); + res = new SimApiBaseResponse(500, ex.Message); + } var returnJson = JsonSerializer.Serialize((object)res, SimApiUtil.JsonOption); - var reply = $"client.{ea.BasicProperties.ReplyTo}.{ea.BasicProperties.AppId}"; - var props = RpcServerChannel.CreateBasicProperties(); - props.AppId = Options.AppId; - props.CorrelationId = ea.BasicProperties.MessageId; - props.MessageId = Guid.NewGuid().ToString(); - props.ReplyTo = Options.AppName; - props.Type = ea.BasicProperties.Type; - RpcServerChannel.BasicPublish(Options.SysName, reply, false, props, Encoding.UTF8.GetBytes(returnJson)); - Logger.LogDebug( - "Rpc Return: ({BasicPropertiesMessageId}) {BasicPropertiesType}@{OptionsAppName} -> {BasicPropertiesReplyTo}\n{ReturnJson}", - ea.BasicProperties.MessageId, ea.BasicProperties.Type, Options.AppName, ea.BasicProperties.ReplyTo, - returnJson); + var reply = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{e.ApplicationMessage.ContentType}"; + var message = new MqttApplicationMessageBuilder() + .WithTopic(reply) + .WithPayload(returnJson) + .WithRetainFlag(false) + .Build(); + if (!Client!.IsConnected) return Task.CompletedTask; + Client!.PublishAsync(message, CancellationToken.None).Wait(); + logger.LogDebug( + "Synapse Rpc Server Return: ({BasicPropertiesMessageId}) {BasicPropertiesType}@{OptionsAppName} -> {BasicPropertiesReplyTo}\n{ReturnJson}", + e.ApplicationMessage.ContentType, action, Options.AppName, appInfo[0], returnJson); + + return Task.CompletedTask; }; - RpcServerChannel.BasicConsume(queue, true, "", false, false, null, consumer); + Client.SubscribeAsync(rsSubOpts).Wait(); } } \ No newline at end of file diff --git a/Synapse/Synapse.cs b/Synapse/Synapse.cs index 6fe9b30..8382c36 100644 --- a/Synapse/Synapse.cs +++ b/Synapse/Synapse.cs @@ -3,9 +3,12 @@ using System.Collections.Generic; using System.Diagnostics; using System.Linq; using System.Text.Json; +using System.Threading; +using System.Threading.Tasks; using Microsoft.Extensions.Logging; -using RabbitMQ.Client; -using RabbitMQ.Client.Exceptions; +using MQTTnet; +using MQTTnet.Client; +using MQTTnet.Formatter; using SimApi.Attributes; using SimApi.Communications; using SimApi.Configurations; @@ -13,69 +16,51 @@ using SimApi.Helpers; namespace SimApi; -public partial class Synapse +public partial class Synapse(SimApiOptions simApiOptions, ILogger logger, IServiceProvider sp) { - private IServiceProvider Sp { get; } + private IServiceProvider Sp { get; } = sp; - private SimApiSynapseOptions Options { get; } + private SimApiSynapseOptions Options { get; } = simApiOptions.SimApiSynapseOptions; - private ILogger Logger { get; } - - private IConnection Connection { get; set; } - - private IModel EventClientChannel { get; set; } - - private IModel EventServerChannel { get; set; } - - private IModel RpcClientChannel { get; set; } - - private IModel RpcServerChannel { get; set; } + private MqttFactory MqttFactory { get; } = new(); + public IMqttClient Client { get; set; } private List EventRegistry { get; set; } private List RpcRegistry { get; set; } - public Synapse(SimApiOptions simApiOptions, ILogger logger, IServiceProvider sp) - { - Logger = logger; - Sp = sp; - Options = simApiOptions.SimApiSynapseOptions; - } - public void Init() { ProcessAttribute(); - Logger.LogDebug("Synapse初始化配置信息: {Json}", SimApiUtil.Json(Options)); + logger.LogDebug("Synapse初始化配置信息: {Json}", SimApiUtil.Json(Options)); if (string.IsNullOrEmpty(Options.AppName) || string.IsNullOrEmpty(Options.SysName)) { - Logger.LogCritical("Synapse初始化失败: AppName or SysName 错误"); + logger.LogCritical("Synapse初始化失败: AppName 和 SysName 不能为空"); } Options.AppId ??= Guid.NewGuid().ToString(); - Logger.LogInformation("System Name: {SysName}\nApp Name: {AppName}\nAppId: {AppId}", Options.SysName, - Options.AppName, Options.AppId); + logger.LogInformation("Synapse Sys Name: {SysName}\nSynapse App Name: {AppName}\nSynapse App Id: {AppId}", + Options.SysName, Options.AppName, Options.AppId); CreateConnection(); - CheckAndCreateExchange(); //事件客户端 if (Options.DisableEventClient) { - Logger.LogWarning("Event Client Disabled: DisableEventClient set true"); + logger.LogWarning("Synapse Event Client Disabled: DisableEventClient set true"); } else { - RunEventClient(); - Logger.LogInformation("Event Client Ready"); + logger.LogInformation("Synapse Event Client Ready"); } //RPC客户端 if (Options.DisableRpcClient) { - Logger.LogWarning("Rpc Client Disabled: DisableEventClient set true"); + logger.LogWarning("Synapse Rpc Client Disabled: DisableEventClient set true"); } else { RunRpcClient(); - Logger.LogInformation("Rpc Client Ready, Client Timeout: {OptionsRpcTimeout}s", Options.RpcTimeout); + logger.LogInformation("Synapse Rpc Client Ready, Client Timeout: {OptionsRpcTimeout}s", Options.RpcTimeout); } if (RpcRegistry.Count > 0) @@ -92,10 +77,10 @@ public partial class Synapse public SimApiBaseResponse Rpc(string appName, string method, dynamic param) { - var res = new SimApiBaseResponse(500, "Rpc Client Disabled!"); + var res = new SimApiBaseResponse(500, "Synapse Rpc Client Disabled!"); if (Options.DisableRpcClient) { - Logger.LogError("Rpc Client Disabled!"); + logger.LogError("Synapse Rpc Client Disabled!"); } else { @@ -115,7 +100,7 @@ public partial class Synapse { if (Options.DisableEventClient) { - Logger.LogError("Event Client Disabled!"); + logger.LogError("Synapse Event Client Disabled!"); } else { @@ -125,73 +110,43 @@ public partial class Synapse private void CreateConnection() { - var factory = new ConnectionFactory + Client = MqttFactory.CreateMqttClient(); + var clientOpts = new MqttClientOptionsBuilder().WithProtocolVersion(MqttProtocolVersion.V500) + .WithWebSocketServer(o => o.WithUri(Options.Websocket)) + .WithCredentials(Options.Username, Options.Password) + .WithClientId($"{Options.AppName}:{Options.AppId}") + .Build(); + Client!.ConnectAsync(clientOpts, CancellationToken.None).Wait(); + Client.ConnectedAsync += _ => { - HostName = Options.MqHost, - Port = Options.MqPort, - VirtualHost = Options.MqVHost, - UserName = Options.MqUser, - Password = Options.MqPass + logger.LogInformation("Synapse MQTT[{AppName}:{AppId}] 连接成功...", Options.AppName, Options.AppId); + return Task.CompletedTask; }; - try + //重连 + Client.DisconnectedAsync += async _ => { - Connection = factory.CreateConnection(); - Logger.LogInformation("连接RabbitMQ服务器成功"); - } - catch (BrokerUnreachableException e) - { - Logger.LogError("连接RabbitMQ失败: \n{Err}", e); - } - } - - private IModel CreateChannel(ushort processNum = 0, string desc = "unknow") - { - IModel channel = null; - try - { - var log = $"Channel [{desc}] 创建成功..."; - channel = Connection.CreateModel(); - if (processNum != 0) + logger.LogError("Synapse MQTT[{AppName}:{AppId}] 断开连接,开始重连...", Options.AppName, Options.AppId); + await Task.Delay(TimeSpan.FromSeconds(5)); + try { - channel.BasicQos(0, processNum, false); - log += $"最大处理器数量: {processNum}"; + logger.LogInformation("Synapse MQTT[{AppName}:{AppId}] 开始连接MQTT服务器...", Options.AppName, Options.AppId); + Client.ConnectAsync(clientOpts).Wait(); } - - Logger.LogInformation(log); - } - catch (ConnectFailureException e) - { - Logger.LogError("Channel [{{Desc}}] 创建失败...\n {0}", e); - } - - return channel; - } - - private void CheckAndCreateExchange() - { - var channel = CreateChannel(0, "Exchange"); - try - { - channel.ExchangeDeclare(Options.SysName, ExchangeType.Topic, true, true, null); - Logger.LogDebug("Register Exchange Success"); - } - catch (ConnectFailureException e) - { - Logger.LogError("Failed to declare Exchange.\n {Err}", e); - } - - channel.Close(); - Logger.LogDebug("Exchange Channel Closed"); + catch + { + logger.LogError("Synapse MQTT[{AppName}:{AppId}] 重连失败...", Options.AppName, Options.AppId); + } + }; } private void ProcessAttribute() { - EventRegistry = new List(); - RpcRegistry = new List(); + EventRegistry = []; + RpcRegistry = []; var stackTrace = new StackTrace(); - var callingMethod = stackTrace.GetFrame(stackTrace.FrameCount - 1).GetMethod(); - var assembly = callingMethod.DeclaringType.Assembly; - var types = assembly.GetTypes(); // 获取程序集中的所有类型 + var callingMethod = stackTrace.GetFrame(stackTrace.FrameCount - 1)?.GetMethod(); + var assembly = callingMethod?.DeclaringType?.Assembly; + var types = assembly!.GetTypes(); // 获取程序集中的所有类型 foreach (var type in types) { var methods = type.GetMethods(); // 获取类型中的所有方法 @@ -233,15 +188,16 @@ public partial class Synapse (current, ev) => current + $"\n |- {ev.Key} -> {ev.Method}@{ev.Class.Name}"); var rpcList = RpcRegistry.Aggregate(string.Empty, (current, ev) => current + $"\n |- {ev.Key} -> {ev.Method}@{ev.Class.Name}"); - Logger.LogInformation(">>> Synapse System 读取Event方法:{Event}\n>>>Synapse System 读取RPC方法:{Rpc}", events, rpcList); + if (EventRegistry.Count > 0) logger.LogInformation(">>> Synapse System 读取Event方法:{Event}", events); + if (RpcRegistry.Count > 0) logger.LogInformation(">>> Synapse System 读取RPC方法:{Rpc}", rpcList); } } public class RegisterItem { - public string Key { get; set; } + public string Key { get; init; } - public Type Class { get; set; } + public Type Class { get; init; } - public string Method { get; set; } + public string Method { get; init; } } \ No newline at end of file