From 7dc2ac6a867aa7e372553e033a0f2f3605933130 Mon Sep 17 00:00:00 2001 From: xRain Date: Thu, 5 Sep 2024 12:59:57 +0800 Subject: [PATCH] add headers --- Synapse/RpcClient.cs | 23 +++++++++++++++-------- Synapse/RpcServer.cs | 39 ++++++++++++++++++++++++--------------- Synapse/Synapse.cs | 27 ++++++++++++++++++++------- 3 files changed, 59 insertions(+), 30 deletions(-) diff --git a/Synapse/RpcClient.cs b/Synapse/RpcClient.cs index 67bb82c..a7d5e6d 100644 --- a/Synapse/RpcClient.cs +++ b/Synapse/RpcClient.cs @@ -1,10 +1,12 @@ using System; using System.Collections.Generic; +using System.Linq; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using Microsoft.Extensions.Logging; using MQTTnet; +using MQTTnet.Packets; using SimApi.Communications; using SimApi.Helpers; @@ -38,7 +40,7 @@ public partial class Synapse Client!.SubscribeAsync(rcSubOpts).Wait(); } - private string? FireRpc(string app, string action, object? param) + private string? FireRpc(string app, string action, object? param, Dictionary? headers = null) { string paramJson; if (param is string strParam) @@ -54,18 +56,23 @@ public partial class Synapse var messageId = Guid.NewGuid().ToString(); var tcs = new TaskCompletionSource(); ResponseCompletionSources.Add(messageId, tcs); - var message = new MqttApplicationMessageBuilder() + var messageBuilder = new MqttApplicationMessageBuilder() .WithTopic(topic) .WithPayload(paramJson) - .WithResponseTopic($"{Options.AppName},{Options.AppId}") - .WithContentType(messageId) - .WithRetainFlag(false) - .Build(); + .WithResponseTopic($"{Options.AppName},{Options.AppId},{messageId}") + .WithContentType("application/json") + .WithRetainFlag(false); + foreach (var h in headers ?? new Dictionary()) + { + messageBuilder.WithUserProperty(h.Key, h.Value); + } + + var message = messageBuilder.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); + "Synapse RPC Client Request: ({PropsMessageId}) {OptionsAppName} -> {Action}@{App}\n{ParamJson}\nHeaders: {Headers}", + messageId, Options.AppName, action, app, paramJson, JsonSerializer.Serialize(headers)); string response; try diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs index 47ae074..2199ebd 100644 --- a/Synapse/RpcServer.cs +++ b/Synapse/RpcServer.cs @@ -28,8 +28,7 @@ public partial class Synapse 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); + appInfo[2], appInfo[0], action, Options.AppName, reqBody); SimApiBaseResponse res; var method = RpcRegistry.FirstOrDefault(x => x.Key == action); if (method == null) @@ -44,17 +43,27 @@ public partial class Synapse { var methodParams = mt!.GetParameters(); object? ret; - if (methodParams.Length == 0) + switch (methodParams.Length) { - 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); + 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 @@ -71,7 +80,7 @@ public partial class Synapse } else { - logger.LogError("Synapse RPC 方法异常: {Err}\n{Stack}", ex.Message,ex.StackTrace); + logger.LogError("Synapse RPC 方法异常: {Err}\n{Stack}", ex.Message, ex.StackTrace); res = new SimApiBaseResponse(500, ex.Message); } } @@ -83,7 +92,7 @@ public partial class Synapse } var returnJson = JsonSerializer.Serialize((object)res, SimApiUtil.JsonOption); - var reply = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{e.ApplicationMessage.ContentType}"; + var reply = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{appInfo[2]}"; var message = new MqttApplicationMessageBuilder() .WithTopic(reply) .WithPayload(returnJson) @@ -93,7 +102,7 @@ public partial class Synapse 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); + appInfo[2], action, Options.AppName, appInfo[0], returnJson); return Task.CompletedTask; }; diff --git a/Synapse/Synapse.cs b/Synapse/Synapse.cs index 9addd43..287aefb 100644 --- a/Synapse/Synapse.cs +++ b/Synapse/Synapse.cs @@ -87,9 +87,11 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge /// /// /// + /// /// /// - public SimApiBaseResponse Rpc(string appName, string method, dynamic? param = null) + public SimApiBaseResponse Rpc(string appName, string method, dynamic? param = null, + Dictionary? headers = null) { var res = new SimApiBaseResponse(500, "Synapse Rpc Client Disabled!"); if (Options.DisableRpcClient) @@ -98,7 +100,7 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge } else { - var data = FireRpc(appName, method, param); + var data = FireRpc(appName, method, param, headers); res = JsonSerializer.Deserialize>(data, SimApiUtil.JsonOption); } @@ -111,10 +113,12 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge /// /// /// + /// /// - public SimApiBaseResponse Rpc(string appName, string method, dynamic? param = null) + public SimApiBaseResponse Rpc(string appName, string method, dynamic? param = null, + Dictionary? headers = null) { - return Rpc(appName, method, param); + return Rpc(appName, method, param, headers); } @@ -141,7 +145,7 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge } /// - /// 发送一个事件 + /// 发送一个事件x /// /// /// @@ -271,10 +275,19 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger logge var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(tmp.Class!); var mt = callClass.GetType().GetMethod(tmp.Method); - if (mt!.GetParameters().Length > 1) + if (mt!.GetParameters().Length > 2) { logger.LogError( - "Synapse Rpc Register Error: Only one or none parameter supported. {Key} -> {Method}@{Class}", + "Synapse Rpc Register Error: Only 1,2 or none parameter supported. {Key} -> {Method}@{Class}", + tmp.Key, tmp.Method, tmp.Class.Name); + continue; + } + + if (mt.GetParameters().Length == 2 && + mt.GetParameters()[1].ParameterType != typeof(Dictionary)) + { + logger.LogError( + "Synapse Rpc Register Error: RpcMethod Parameter 2 must be Dictionary. {Key} -> {Method}@{Class}", tmp.Key, tmp.Method, tmp.Class.Name); continue; }