diff --git a/Attributes/SynapseEventAttribute.cs b/Attributes/SynapseEventAttribute.cs new file mode 100644 index 0000000..173c983 --- /dev/null +++ b/Attributes/SynapseEventAttribute.cs @@ -0,0 +1,14 @@ +using System; + +namespace SimApi.Attributes; + +[AttributeUsage(AttributeTargets.Method)] +public class SynapseEventAttribute : Attribute +{ + public string Name { get; } + + public SynapseEventAttribute(string name) + { + Name = name; + } +} \ No newline at end of file diff --git a/Attributes/SynapseRpcAttribute.cs b/Attributes/SynapseRpcAttribute.cs new file mode 100644 index 0000000..b82a5dd --- /dev/null +++ b/Attributes/SynapseRpcAttribute.cs @@ -0,0 +1,19 @@ +using System; + +namespace SimApi.Attributes; + +[AttributeUsage(AttributeTargets.Method)] +public class SynapseRpcAttribute : Attribute +{ + public string Name { get; } + + public SynapseRpcAttribute() + { + Name = null; + } + + public SynapseRpcAttribute(string name) + { + Name = name; + } +} \ No newline at end of file diff --git a/Configurations/SimApiOptions.cs b/Configurations/SimApiOptions.cs index 2ab28b6..51a69fc 100644 --- a/Configurations/SimApiOptions.cs +++ b/Configurations/SimApiOptions.cs @@ -52,6 +52,12 @@ namespace SimApi.Configs /// public bool EnableLogger { get; set; } = false; + /// + /// 是否启用Synapse + /// + public bool EnableSynapse { get; set; } = false; + + /// /// Swagger文档相关配置,需要启用 EnableSimApiDoc /// @@ -62,6 +68,17 @@ namespace SimApi.Configs /// public SimApiStorageOptions SimApiStorageOptions { get; set; } = new SimApiStorageOptions(); + public SimApiSynapseOptions SimApiSynapseOptions { get; set; } = new SimApiSynapseOptions(); + + public void ConfigureSimApiSynapse(Action options = null) + { + options?.Invoke(SimApiSynapseOptions); + } + + public void ConfigureSimApiSynapse(SimApiSynapseOptions options) + { + SimApiSynapseOptions = options; + } public void ConfigureSimApiDoc(Action options = null) { diff --git a/Configurations/SimApiSynapseOptions.cs b/Configurations/SimApiSynapseOptions.cs new file mode 100644 index 0000000..a4b6866 --- /dev/null +++ b/Configurations/SimApiSynapseOptions.cs @@ -0,0 +1,30 @@ +namespace SimApi.Configs; + +public class SimApiSynapseOptions +{ + public string MqHost { get; set; } + + public int MqPort { get; set; } + + public string MqUser { get; set; } + + public string MqPass { get; set; } + + public string MqVHost { 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; +} \ No newline at end of file diff --git a/SimApi.csproj b/SimApi.csproj index da9f3c1..2273f62 100644 --- a/SimApi.csproj +++ b/SimApi.csproj @@ -22,12 +22,11 @@ - - + diff --git a/SimApiExtensions.cs b/SimApiExtensions.cs index 960dfb2..ce2e3da 100644 --- a/SimApiExtensions.cs +++ b/SimApiExtensions.cs @@ -42,6 +42,11 @@ namespace SimApi policy => { policy.AllowAnyHeader().AllowAnyMethod().AllowAnyOrigin(); })); } + if (simApiOptions.EnableSynapse) + { + builder.AddSingleton(); + } + // 使用SimApiDoc if (simApiOptions.EnableSimApiDoc) { @@ -77,10 +82,7 @@ namespace SimApi Id = "HeaderToken" } }, - new[] - { - "readAccess", "writeAccess" - } + new[] { "readAccess", "writeAccess" } } }); haveSimApiAuth = true; @@ -143,10 +145,7 @@ namespace SimApi Id = "oauth2" } }, - new[] - { - "SimApiAuth" - } + new[] { "SimApiAuth" } } }); } @@ -254,6 +253,12 @@ namespace SimApi logger.LogInformation("开始配置使用URL小写..."); } + if (options.EnableSynapse) + { + var synapse = builder.ApplicationServices.GetRequiredService(); + synapse.Init(); + } + return builder; } } diff --git a/Synapse/EventClient.cs b/Synapse/EventClient.cs new file mode 100644 index 0000000..f2f5ef8 --- /dev/null +++ b/Synapse/EventClient.cs @@ -0,0 +1,34 @@ +using System; +using System.Text; +using System.Text.Encodings.Web; +using System.Text.Json; +using System.Text.Unicode; +using Microsoft.Extensions.Logging; + +namespace SimApi; + +public partial class Synapse +{ + private void RunEventClient() + { + EventClientChannel = CreateChannel(0, "EventClient"); + } + + private void FireEvent(string eventName, object param) + { + var paramJson = JsonSerializer.Serialize(param, new JsonSerializerOptions + { + Encoder = JavaScriptEncoder.Create(UnicodeRanges.All), + PropertyNamingPolicy = JsonNamingPolicy.CamelCase, + }); + 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); + } +} \ No newline at end of file diff --git a/Synapse/EventServer.cs b/Synapse/EventServer.cs new file mode 100644 index 0000000..926bbcf --- /dev/null +++ b/Synapse/EventServer.cs @@ -0,0 +1,59 @@ +using System; +using System.Linq; +using System.Text; +using System.Text.Json; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using RabbitMQ.Client.Events; + +namespace SimApi; + +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('#'))) + { + 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); + + var key = ea.RoutingKey.Replace("event.", string.Empty); + var method = EventRegistry.FirstOrDefault(x => x.Key == key); + if (method != null) + { + var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class); + var mt = callClass.GetType().GetMethod(method.Method); + if (mt != null) + { + var pt = mt.GetParameters()[0].ParameterType; + try + { + mt.Invoke(callClass, pt == typeof(string) + ? new object[] { reqBody } + : new[] { JsonSerializer.Deserialize(reqBody, pt) }); + EventServerChannel.BasicAck(ea.DeliveryTag, false); + } + catch (Exception) + { + EventServerChannel.BasicNack(ea.DeliveryTag, false, true); + } + } + else + { + Logger.LogError("Event Callback not available: {Ev}", method.Key); + EventServerChannel.BasicNack(ea.DeliveryTag, false, false); + } + } + }; + EventServerChannel.BasicConsume(queue, true, "", false, false, null, consumer); + } +} \ No newline at end of file diff --git a/Synapse/RpcClient.cs b/Synapse/RpcClient.cs new file mode 100644 index 0000000..ebd5716 --- /dev/null +++ b/Synapse/RpcClient.cs @@ -0,0 +1,74 @@ +using System; +using System.Collections.Generic; +using System.Text; +using System.Text.Encodings.Web; +using System.Text.Json; +using System.Text.Unicode; +using Microsoft.Extensions.Logging; +using RabbitMQ.Client.Events; +using SimApi.Communications; +using SimApi.Helpers; + +namespace SimApi; + +public partial class Synapse +{ + 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) => + { + 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())); + }; + RpcClientChannel.BasicConsume(queue, false, "", false, false, null, consumer); + } + + private object FireRpc(string app, string action, object param) + { + var paramJson = JsonSerializer.Serialize(param, new JsonSerializerOptions + { + Encoder = JavaScriptEncoder.Create(UnicodeRanges.All), + PropertyNamingPolicy = JsonNamingPolicy.CamelCase, + }); + var router = $"server.{app}"; + SimApiBaseResponse 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 ts = SimApiUtil.TimestampNow; + while (true) + { + if (SimApiUtil.TimestampNow - ts > Options.RpcTimeout) + { + response = new SimApiBaseResponse(502, "timeout"); + break; + } + if (ResponseCache.TryGetValue(props.MessageId, out var value)) + { + response = JsonSerializer.Deserialize>( + Encoding.UTF8.GetString(value)); + ResponseCache.Remove(props.MessageId); + break; + } + } + return response; + } +} \ No newline at end of file diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs new file mode 100644 index 0000000..81aec62 --- /dev/null +++ b/Synapse/RpcServer.cs @@ -0,0 +1,85 @@ +using System; +using System.Linq; +using System.Text; +using System.Text.Encodings.Web; +using System.Text.Json; +using System.Text.Unicode; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using RabbitMQ.Client.Events; +using SimApi.Communications; +using SimApi.Exceptions; +using JsonSerializer = System.Text.Json.JsonSerializer; + +namespace SimApi; + +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 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, + reqBody); + var res = new SimApiBaseResponse(404, "method not found"); + var method = RpcRegistry.FirstOrDefault(x => x.Key == ea.BasicProperties.Type); + if (method != null) + { + var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class); + var mt = callClass.GetType().GetMethod(method.Method); + if (mt != null) + { + try + { + var pt = mt.GetParameters()[0].ParameterType; + if (pt == typeof(string)) + { + var ret = mt.Invoke(callClass, new[] { reqBody }); + res = new SimApiBaseResponse(ret); + } + else + { + var ret = mt.Invoke(callClass, new[] { JsonSerializer.Deserialize(reqBody, pt) }); + res = new SimApiBaseResponse(ret); + } + } + catch (SimApiException e) + { + res = new SimApiBaseResponse(e.Code, e.Message); + } + catch (Exception e) + { + res = new SimApiBaseResponse(500, e.ToString()); + } + } + } + var returnJson = JsonSerializer.Serialize((SimApiBaseResponse)res, new JsonSerializerOptions + { + Encoder = JavaScriptEncoder.Create(UnicodeRanges.All), + PropertyNamingPolicy = JsonNamingPolicy.CamelCase, + }); + 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); + }; + RpcServerChannel.BasicConsume(queue, true, "", false, false, null, consumer); + } +} \ No newline at end of file diff --git a/Synapse/Synapse.cs b/Synapse/Synapse.cs new file mode 100644 index 0000000..384bed1 --- /dev/null +++ b/Synapse/Synapse.cs @@ -0,0 +1,230 @@ +using System; +using System.Collections.Generic; +using System.Diagnostics; +using System.Linq; +using Microsoft.Extensions.Logging; +using RabbitMQ.Client; +using RabbitMQ.Client.Exceptions; +using SimApi.Attributes; +using SimApi.Communications; +using SimApi.Configs; +using SimApi.Helpers; + +namespace SimApi; + +public partial class Synapse +{ + private IServiceProvider Sp { get; } + + private SimApiSynapseOptions Options { get; } + + 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 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)); + if (string.IsNullOrEmpty(Options.AppName) || string.IsNullOrEmpty(Options.SysName)) + { + Logger.LogCritical("Synapse初始化失败: AppName or SysName 错误"); + } + Options.AppId ??= Guid.NewGuid().ToString(); + Logger.LogInformation("System Name: {SysName}\nApp Name: {AppName}\nAppId: {AppId}", Options.SysName, + Options.AppName, Options.AppId); + CreateConnection(); + CheckAndCreateExchange(); + //事件客户端 + if (Options.DisableEventClient) + { + Logger.LogWarning("Event Client Disabled: DisableEventClient set true"); + } + else + { + RunEventClient(); + Logger.LogInformation("Event Client Ready"); + } + + //RPC客户端 + if (Options.DisableRpcClient) + { + Logger.LogWarning("Rpc Client Disabled: DisableEventClient set true"); + } + else + { + RunRpcClient(); + Logger.LogInformation("Rpc Client Ready, Client Timeout: {OptionsRpcTimeout}s", Options.RpcTimeout); + } + RunRpcServer(); + RunEventServer(); + } + + + public SimApiBaseResponse Rpc(string appName, string method, dynamic param) + { + var res = new SimApiBaseResponse(500, "Rpc Client Disabled!"); + if (Options.DisableRpcClient) + { + Logger.LogError("Rpc Client Disabled!"); + } + else + { + res = FireRpc(appName, method, param); + } + return res as SimApiBaseResponse; + } + + public SimApiBaseResponse Rpc(string appName, string method, dynamic param) + { + return Rpc(appName, method, param); + } + + public void Event(string eventName, dynamic param) + { + if (Options.DisableEventClient) + { + Logger.LogError("Event Client Disabled!"); + } + else + { + FireEvent(eventName, param); + } + } + + private void CreateConnection() + { + var factory = new ConnectionFactory + { + HostName = Options.MqHost, + Port = Options.MqPort, + VirtualHost = Options.MqVHost, + UserName = Options.MqUser, + Password = Options.MqPass + }; + try + { + 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) + { + channel.BasicQos(0, processNum, false); + log += $"最大处理器数量: {processNum}"; + } + 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"); + } + + private void ProcessAttribute() + { + EventRegistry = new List(); + RpcRegistry = new List(); + var stackTrace = new StackTrace(); + 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(); // 获取类型中的所有方法 + foreach (var method in methods) + { + if (method.IsDefined(typeof(SynapseEventAttribute), false)) + { + var attribute = + (SynapseEventAttribute)Attribute.GetCustomAttribute(method, typeof(SynapseEventAttribute)); + if (attribute != null) + { + EventRegistry.Add(new RegisterItem + { + Key = attribute.Name ?? method.Name, + Class = type, + Method = method.Name + }); + } + } + if (method.IsDefined(typeof(SynapseRpcAttribute), false)) + { + var attribute = + (SynapseRpcAttribute)Attribute.GetCustomAttribute(method, typeof(SynapseRpcAttribute)); + if (attribute != null) + { + RpcRegistry.Add(new RegisterItem + { + Key = attribute.Name ?? $"{type.Name}.{method.Name}", + Class = type, + Method = method.Name + }); + } + } + } + } + var events = EventRegistry.Aggregate(string.Empty, + (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); + } +} + +public class RegisterItem +{ + public string Key { get; set; } + + public Type Class { get; set; } + + public string Method { get; set; } +} \ No newline at end of file