From 7d39ab5e0b001b8f728252fdc66f66801eb34ce8 Mon Sep 17 00:00:00 2001 From: xRain Date: Sat, 17 May 2025 20:38:39 +0800 Subject: [PATCH] add task.run --- Synapse/EventServer.cs | 59 +++++++++-------- Synapse/RpcServer.cs | 147 +++++++++++++++++++++-------------------- 2 files changed, 107 insertions(+), 99 deletions(-) diff --git a/Synapse/EventServer.cs b/Synapse/EventServer.cs index 6db0129..eccd900 100644 --- a/Synapse/EventServer.cs +++ b/Synapse/EventServer.cs @@ -19,39 +19,42 @@ public partial class Synapse { Client!.ApplicationMessageReceivedAsync += async e => { - 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); - var methods = EventRegistry - .Where(x => Regex.IsMatch(eventName, - "^" + Regex.Escape(x.Key!).Replace("\\+", "[^/]+").Replace("\\#", ".*") + "$")) - .ToArray(); - var tasks = methods.Select(method => Task.Run(() => - { - var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!); - var mt = callClass.GetType().GetMethod(method.Method!); - try + await Task.Run(async () => + { + 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); + var methods = EventRegistry + .Where(x => Regex.IsMatch(eventName, + "^" + Regex.Escape(x.Key!).Replace("\\+", "[^/]+").Replace("\\#", ".*") + "$")) + .ToArray(); + 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) - ? new object?[] { eventName, reqBody } - : new object?[] { eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }); + if (mt!.GetParameters().Length == 2) + { + var pt = mt.GetParameters()[1].ParameterType; + mt.Invoke(callClass, pt == typeof(string) + ? [eventName, reqBody] + : [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]); + } + else + { + mt.Invoke(callClass, new object[] { eventName }); + } } - else + catch (Exception ex) { - mt.Invoke(callClass, new object[] { 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); - } - })) - .ToList(); - await Task.WhenAll(tasks); + })) + .ToList(); + await Task.WhenAll(tasks); + }); }; SubEventServerTopic(); } diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs index 2199ebd..519da0e 100644 --- a/Synapse/RpcServer.cs +++ b/Synapse/RpcServer.cs @@ -22,88 +22,93 @@ public partial class Synapse { Client!.ApplicationMessageReceivedAsync += e => { - if (!e.ApplicationMessage.Topic.StartsWith(RpcServerTopicPrefix)) return Task.CompletedTask; - var reqBody = e.ApplicationMessage.ConvertPayloadToString(); - 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}", - appInfo[2], appInfo[0], action, Options.AppName, reqBody); - SimApiBaseResponse res; - var method = RpcRegistry.FirstOrDefault(x => x.Key == action); - if (method == null) + Task.Run(() => { - res = new SimApiBaseResponse(404, "method not found"); - } - else - { - var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!); - var mt = callClass.GetType().GetMethod(method.Method!); - try + if (!e.ApplicationMessage.Topic.StartsWith(RpcServerTopicPrefix)) return; + var reqBody = e.ApplicationMessage.ConvertPayloadToString(); + 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}", + appInfo[2], appInfo[0], action, Options.AppName, reqBody); + SimApiBaseResponse res; + var method = RpcRegistry.FirstOrDefault(x => x.Key == action); + if (method == null) { - var methodParams = mt!.GetParameters(); - object? ret; - switch (methodParams.Length) - { - case 1: - var pt = mt.GetParameters()[0].ParameterType; - var param = pt == typeof(string) - ? [reqBody] - : new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }; - ret = mt.Invoke(callClass, param); - break; - case 2: - var headerData = - e.ApplicationMessage.UserProperties.ToDictionary(x => x.Name, x => x.Value); - var pt2 = mt.GetParameters()[0].ParameterType; - var param2 = pt2 == typeof(string) - ? [reqBody] - : new[] { JsonSerializer.Deserialize(reqBody, pt2, SimApiUtil.JsonOption), headerData }; - ret = mt.Invoke(callClass, param2); - break; - default: - ret = mt.Invoke(callClass, []); - break; - } - - 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; + switch (methodParams.Length) + { + case 1: + var pt = mt.GetParameters()[0].ParameterType; + var param = pt == typeof(string) + ? [reqBody] + : new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) }; + ret = mt.Invoke(callClass, param); + break; + case 2: + var headerData = + e.ApplicationMessage.UserProperties.ToDictionary(x => x.Name, x => x.Value); + var pt2 = mt.GetParameters()[0].ParameterType; + var param2 = pt2 == typeof(string) + ? [reqBody] + : new[] + { + JsonSerializer.Deserialize(reqBody, pt2, SimApiUtil.JsonOption), headerData + }; + ret = mt.Invoke(callClass, param2); + break; + default: + ret = mt.Invoke(callClass, []); + break; + } + + res = new SimApiBaseResponse + { + Data = ret + }; } - else + catch (TargetInvocationException ex) { - logger.LogError("Synapse RPC 方法异常: {Err}\n{Stack}", ex.Message, ex.StackTrace); + if (ex.InnerException is SimApiException ie) + { + logger.LogDebug("Synapse RPC调用错误: {Err}", ie.Message); + res = new SimApiBaseResponse(ie.Code, ie.Message); + } + else + { + logger.LogError("Synapse RPC 方法异常: {Err}\n{Stack}", ex.Message, ex.StackTrace); + 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]}/{appInfo[2]}"; - 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}", - appInfo[2], action, Options.AppName, appInfo[0], returnJson); + var returnJson = JsonSerializer.Serialize((object)res, SimApiUtil.JsonOption); + var reply = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{appInfo[2]}"; + var message = new MqttApplicationMessageBuilder() + .WithTopic(reply) + .WithPayload(returnJson) + .WithRetainFlag(false) + .Build(); + if (!Client.IsConnected) return; + Client.PublishAsync(message, CancellationToken.None).Wait(); + logger.LogDebug( + "Synapse Rpc Server Return: ({BasicPropertiesMessageId}) {BasicPropertiesType}@{OptionsAppName} -> {BasicPropertiesReplyTo}\n{ReturnJson}", + appInfo[2], action, Options.AppName, appInfo[0], returnJson); + }); return Task.CompletedTask; }; SubRpcServerTopic();