Compare commits

..
2 Commits
Author SHA1 Message Date
xrain 141846d79d 1. RpcServer can return not found , 2. Event server process all parllels 2024-08-05 19:55:14 +08:00
xrain a4bae0226e fix bug 2024-08-03 19:37:40 +08:00
2 changed files with 65 additions and 57 deletions
+25 -23
View File
@@ -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,32 +27,33 @@ 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);
};
SubRpcServerTopic();
SubEventServerTopic();
}
private void 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();
+40 -34
View File
@@ -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<object?>
{
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<object?>
{
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}";
@@ -93,7 +99,7 @@ public partial class Synapse
SubRpcServerTopic();
}
public void SubRpcServerTopic()
private void SubRpcServerTopic()
{
var rsSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
.WithTopicFilter(o =>