fix synapse use mqtt
This commit is contained in:
+22
-23
@@ -1,10 +1,11 @@
|
||||
using System;
|
||||
using System.Linq;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
using System.Threading.Tasks;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using RabbitMQ.Client.Events;
|
||||
using MQTTnet;
|
||||
using MQTTnet.Protocol;
|
||||
using SimApi.Helpers;
|
||||
|
||||
namespace SimApi;
|
||||
@@ -13,38 +14,36 @@ public partial class Synapse
|
||||
{
|
||||
private void RunEventServer()
|
||||
{
|
||||
EventServerChannel = CreateChannel(Options.EventProcessorNum, "EventServer");
|
||||
var queue = $"{Options.SysName}_{Options.AppName}_event";
|
||||
EventServerChannel.QueueDeclare(queue, true, false, true, null);
|
||||
foreach (var ev in EventRegistry.Where(ev => !ev.Key.Contains('*') && !ev.Key.Contains('#')))
|
||||
var esTopicPrefix = $"{Options.SysName}/{Options.AppName}/event/";
|
||||
var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
||||
.WithTopicFilter(o =>
|
||||
o.WithTopic($"$queue/{esTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
|
||||
.Build();
|
||||
Client.ApplicationMessageReceivedAsync += e =>
|
||||
{
|
||||
EventServerChannel.QueueBind(queue, Options.SysName, $"event.{ev.Key}", null);
|
||||
}
|
||||
var consumer = new EventingBasicConsumer(EventServerChannel);
|
||||
consumer.Received += (ch, ea) =>
|
||||
{
|
||||
var reqBody = Encoding.UTF8.GetString(ea.Body.ToArray());
|
||||
Logger.LogDebug("Event Receive: {BasicPropertiesReplyTo}.{BasicPropertiesType}\n{S}",
|
||||
ea.BasicProperties.ReplyTo, ea.BasicProperties.Type, reqBody);
|
||||
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);
|
||||
|
||||
var key = ea.RoutingKey.Replace("event.", string.Empty);
|
||||
var method = EventRegistry.FirstOrDefault(x => x.Key == key);
|
||||
var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class);
|
||||
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;
|
||||
var pt = mt!.GetParameters()[0].ParameterType;
|
||||
try
|
||||
{
|
||||
mt.Invoke(callClass, pt == typeof(string)
|
||||
? new object[] { reqBody }
|
||||
: new[] { JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption) });
|
||||
EventServerChannel.BasicAck(ea.DeliveryTag, false);
|
||||
}
|
||||
catch (Exception e)
|
||||
catch (Exception ex)
|
||||
{
|
||||
Logger.LogError("Event Processor Error: {Err}", e.InnerException);
|
||||
EventServerChannel.BasicNack(ea.DeliveryTag, false, false);
|
||||
logger.LogError("SynapseEvent Processor Error: {Err}", ex.InnerException);
|
||||
}
|
||||
return Task.CompletedTask;
|
||||
};
|
||||
EventServerChannel.BasicConsume(queue, false, "", false, false, null, consumer);
|
||||
Client.SubscribeAsync(eventSubOpts).Wait();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user