From 141846d79d0e6438e7704fc4e5264d1ad414873a Mon Sep 17 00:00:00 2001 From: xRain Date: Mon, 5 Aug 2024 19:55:14 +0800 Subject: [PATCH] 1. RpcServer can return not found , 2. Event server process all parllels --- Synapse/EventServer.cs | 46 ++++++++++++++------------- Synapse/RpcServer.cs | 72 +++++++++++++++++++++++------------------- 2 files changed, 63 insertions(+), 55 deletions(-) diff --git a/Synapse/EventServer.cs b/Synapse/EventServer.cs index 0b92e07..6db0129 100644 --- a/Synapse/EventServer.cs +++ b/Synapse/EventServer.cs @@ -1,4 +1,5 @@ using System; +using System.Collections.Generic; using System.Linq; using System.Text.Json; using System.Text.RegularExpressions; @@ -12,14 +13,13 @@ namespace SimApi; public partial class Synapse { - private string EventServerTopicPrefix => $"{Options.SysName}/event/"; private void RunEventServer() { - Client!.ApplicationMessageReceivedAsync += e => + Client!.ApplicationMessageReceivedAsync += async e => { - if (!e.ApplicationMessage.Topic.StartsWith(EventServerTopicPrefix)) return Task.CompletedTask; + if (!e.ApplicationMessage.Topic.StartsWith(EventServerTopicPrefix)) return; var reqBody = e.ApplicationMessage.ConvertPayloadToString(); var eventName = e.ApplicationMessage.Topic.Replace(EventServerTopicPrefix, string.Empty); logger.LogDebug("Synapse Event Receive: {EventName}\n{Body}", eventName, reqBody); @@ -27,30 +27,31 @@ public partial class Synapse .Where(x => Regex.IsMatch(eventName, "^" + Regex.Escape(x.Key!).Replace("\\+", "[^/]+").Replace("\\#", ".*") + "$")) .ToArray(); - foreach (var method in methods) - { - var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!); - var mt = callClass.GetType().GetMethod(method.Method!); - try + var tasks = methods.Select(method => Task.Run(() => { - if (mt!.GetParameters().Length == 2) + var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!); + var mt = callClass.GetType().GetMethod(method.Method!); + try { - var pt = mt.GetParameters()[1].ParameterType; - mt.Invoke(callClass, pt == typeof(string) - ? [eventName, reqBody] - : [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]); + if (mt!.GetParameters().Length == 2) + { + var pt = mt.GetParameters()[1].ParameterType; + mt.Invoke(callClass, pt == typeof(string) + ? new object?[] { eventName, reqBody } + : new object?[] { eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }); + } + else + { + mt.Invoke(callClass, new object[] { eventName }); + } } - else + catch (Exception ex) { - mt.Invoke(callClass, [eventName]); + logger.LogError("Synapse Event Processor Error: {Err}\n{Stack}", ex.Message, ex.StackTrace); } - } - catch (Exception ex) - { - logger.LogError("Synapse Event Processor Error: {Err}\n{Stack}", ex.Message,ex.StackTrace); - } - } - return Task.CompletedTask; + })) + .ToList(); + await Task.WhenAll(tasks); }; SubEventServerTopic(); } @@ -64,6 +65,7 @@ public partial class Synapse { topic = "$queue/" + topic; } + var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder() .WithTopicFilter(o => o.WithTopic(topic)).Build(); Client!.SubscribeAsync(evSubOpts).Wait(); diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs index f365fee..17b1664 100644 --- a/Synapse/RpcServer.cs +++ b/Synapse/RpcServer.cs @@ -32,48 +32,54 @@ public partial class Synapse reqBody); SimApiBaseResponse res; 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 + if (method == null) { - var methodParams = mt!.GetParameters(); - object? ret; - if (methodParams.Length == 0) - { - ret = mt.Invoke(callClass, []); - } - else - { - var pt = mt.GetParameters()[0].ParameterType; - var param = pt == typeof(string) - ? [reqBody] - : new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }; - ret = mt.Invoke(callClass, param); - } - - res = new SimApiBaseResponse - { - Data = ret - }; + res = new SimApiBaseResponse(404, "method not found"); } - catch (TargetInvocationException ex) + else { - if (ex.InnerException is SimApiException ie) + var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!); + var mt = callClass.GetType().GetMethod(method.Method!); + try { - logger.LogDebug("Synapse RPC调用错误: {Err}", ie.Message); - res = new SimApiBaseResponse(ie.Code, ie.Message); + var methodParams = mt!.GetParameters(); + object? ret; + if (methodParams.Length == 0) + { + ret = mt.Invoke(callClass, []); + } + else + { + var pt = mt.GetParameters()[0].ParameterType; + var param = pt == typeof(string) + ? [reqBody] + : new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }; + ret = mt.Invoke(callClass, param); + } + + res = new SimApiBaseResponse + { + Data = ret + }; } - else + catch (TargetInvocationException ex) { + if (ex.InnerException is SimApiException ie) + { + logger.LogDebug("Synapse RPC调用错误: {Err}", ie.Message); + res = new SimApiBaseResponse(ie.Code, ie.Message); + } + else + { + res = new SimApiBaseResponse(500, ex.Message); + } + } + catch (Exception ex) + { + logger.LogDebug("Synapse RPC调用失败: {Err}", ex.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 = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{e.ApplicationMessage.ContentType}";