Files
simapi-net/Synapse/EventServer.cs
T

78 lines
3.1 KiB
C#
Raw Normal View History

2023-10-20 12:33:18 +08:00
using System;
using System.Collections.Generic;
2023-10-20 12:33:18 +08:00
using System.Linq;
using System.Text.Json;
2024-07-26 08:11:21 +08:00
using System.Text.RegularExpressions;
2024-07-26 05:03:28 +08:00
using System.Threading.Tasks;
2023-10-20 12:33:18 +08:00
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
2024-07-26 05:03:28 +08:00
using MQTTnet;
2023-10-23 17:39:27 +08:00
using SimApi.Helpers;
2023-10-20 12:33:18 +08:00
namespace SimApi;
public partial class Synapse
{
2024-08-03 19:19:42 +08:00
private string EventServerTopicPrefix => $"{Options.SysName}/event/";
2023-10-20 12:33:18 +08:00
private void RunEventServer()
{
Client!.ApplicationMessageReceivedAsync += async e =>
2023-10-20 12:33:18 +08:00
{
2025-05-17 20:38:39 +08:00
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(() =>
2024-07-26 08:11:21 +08:00
{
2025-05-17 20:38:39 +08:00
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!);
var mt = callClass.GetType().GetMethod(method.Method!);
try
{
2025-05-17 20:38:39 +08:00
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 });
}
}
2025-05-17 20:38:39 +08:00
catch (Exception ex)
{
2025-05-17 20:38:39 +08:00
logger.LogError("Synapse Event Processor Error: {Err}\n{Stack}", ex.Message, ex.StackTrace);
}
2025-05-17 20:38:39 +08:00
}))
.ToList();
await Task.WhenAll(tasks);
});
2023-10-20 12:33:18 +08:00
};
2024-08-03 19:37:40 +08:00
SubEventServerTopic();
2024-08-03 19:19:42 +08:00
}
private void SubEventServerTopic()
{
2024-07-26 08:11:21 +08:00
foreach (var ev in EventRegistry)
{
2024-08-03 19:19:42 +08:00
var topic = $"{EventServerTopicPrefix}{ev.Key}";
2024-07-27 01:33:17 +08:00
if (Options.EventLoadBalancing)
{
topic = "$queue/" + topic;
}
2024-07-26 08:11:21 +08:00
var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
.WithTopicFilter(o => o.WithTopic(topic)).Build();
2024-08-03 19:19:42 +08:00
Client!.SubscribeAsync(evSubOpts).Wait();
2024-07-26 08:11:21 +08:00
logger.LogDebug("Synapse Event Register Event Success: {EventName}\nFull Topic: {Topic}", ev.Key, topic);
}
2023-10-20 12:33:18 +08:00
}
}