add config store

This commit is contained in:
2024-07-26 08:11:21 +08:00
parent 82f7e0fc00
commit 929ddaced3
21 changed files with 302 additions and 153 deletions
+39 -25
View File
@@ -1,11 +1,11 @@
using System;
using System.Linq;
using System.Text.Json;
using System.Text.RegularExpressions;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using MQTTnet;
using MQTTnet.Protocol;
using SimApi.Helpers;
namespace SimApi;
@@ -14,36 +14,50 @@ public partial class Synapse
{
private void RunEventServer()
{
var esTopicPrefix = $"{Options.SysName}/{Options.AppName}/event/";
var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
.WithTopicFilter(o =>
o.WithTopic($"$queue/{esTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
.Build();
Client.ApplicationMessageReceivedAsync += e =>
var esTopicPrefix = $"{Options.SysName}/event/";
Client!.ApplicationMessageReceivedAsync += e =>
{
if (!e.ApplicationMessage.Topic.StartsWith(esTopicPrefix)) return Task.CompletedTask;
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
var eventName = e.ApplicationMessage.Topic.Replace(esTopicPrefix, string.Empty);
logger.LogDebug("Synapse Event Receive: {AppName}.{EventName}\n{Body}",
Options.AppName, eventName, reqBody);
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();
foreach (var method in methods)
{
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!);
var mt = callClass.GetType().GetMethod(method.Method!);
try
{
if (mt!.GetParameters().Length == 2)
{
var pt = mt!.GetParameters()[0].ParameterType;
mt.Invoke(callClass, pt == typeof(string)
? [eventName, reqBody]
: [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]);
}
else
{
mt.Invoke(callClass, [eventName]);
}
}
catch (Exception ex)
{
logger.LogError("Synapse Event Processor Error: {Err}", ex.InnerException);
}
}
var method = EventRegistry.FirstOrDefault(x => x.Key == eventName);
if (method == null) return Task.CompletedTask;
var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method!.Class);
var mt = callClass.GetType().GetMethod(method.Method);
var pt = mt!.GetParameters()[0].ParameterType;
try
{
mt.Invoke(callClass, pt == typeof(string)
? new object[] { reqBody }
: new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) });
}
catch (Exception ex)
{
logger.LogError("SynapseEvent Processor Error: {Err}", ex.InnerException);
}
return Task.CompletedTask;
};
Client.SubscribeAsync(eventSubOpts).Wait();
foreach (var ev in EventRegistry)
{
var topic = $"$queue/{esTopicPrefix}{ev.Key}";
var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
.WithTopicFilter(o => o.WithTopic(topic)).Build();
Client.SubscribeAsync(evSubOpts).Wait();
logger.LogDebug("Synapse Event Register Event Success: {EventName}\nFull Topic: {Topic}", ev.Key, topic);
}
}
}