diff --git a/Configurations/SimApiSynapseOptions.cs b/Configurations/SimApiSynapseOptions.cs
index 1bee12c..bd3dc45 100644
--- a/Configurations/SimApiSynapseOptions.cs
+++ b/Configurations/SimApiSynapseOptions.cs
@@ -2,28 +2,22 @@ namespace SimApi.Configurations;
public class SimApiSynapseOptions
{
- public string MqHost { get; set; }
+ ///
+ /// Mqtt服务器的Websocket地址
+ ///
+ public string Websocket { get; set; }
- public int MqPort { get; set; }
+ public string Username { get; set; }
- public string MqUser { get; set; }
-
- public string MqPass { get; set; }
-
- public string MqVHost { get; set; } = "/";
+ public string Password { get; set; }
public string SysName { get; set; }
-
public string AppName { get; set; }
public string AppId { get; set; }
public int RpcTimeout { get; set; } = 3;
- public ushort EventProcessorNum { get; set; } = 20;
-
- public ushort RpcProcessorNum { get; set; } = 20;
-
public bool DisableEventClient { get; set; } = false;
public bool DisableRpcClient { get; set; } = false;
diff --git a/SimApi.csproj b/SimApi.csproj
index 55eed49..bd0efdc 100644
--- a/SimApi.csproj
+++ b/SimApi.csproj
@@ -26,7 +26,7 @@
-
+
diff --git a/SimApiExtensions.cs b/SimApiExtensions.cs
index 4e6dbf8..ee1ac44 100644
--- a/SimApiExtensions.cs
+++ b/SimApiExtensions.cs
@@ -5,6 +5,7 @@ using Microsoft.Extensions.DependencyInjection;
using Microsoft.OpenApi.Models;
using SimApi.Middlewares;
using Microsoft.AspNetCore.HttpOverrides;
+using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using SimApi.Configurations;
using SimApi.Logger;
@@ -189,6 +190,35 @@ public static class SimApiExtensions
return builder;
}
+ public static IHost UseSimApi(this IHost builder)
+ {
+ var options = builder.Services.GetRequiredService();
+
+ var logger = builder.Services.GetRequiredService>();
+
+ logger.LogInformation("当前时区: {LocalId}", TimeZoneInfo.Local.Id);
+
+ //请求一下检测存储错误
+ if (options.EnableSimApiStorage)
+ {
+ logger.LogInformation("开始配置SimApiStorage...");
+ builder.Services.GetService();
+ }
+
+ if (options.EnableLowerUrl)
+ {
+ logger.LogInformation("开始配置使用URL小写...");
+ }
+
+ if (options.EnableSynapse)
+ {
+ var synapse = builder.Services.GetRequiredService();
+ synapse.Init();
+ }
+
+ return builder;
+ }
+
///
/// 使用所有SimApi自定义中间件
///
diff --git a/Synapse/EventClient.cs b/Synapse/EventClient.cs
index 2e3f87e..1b6f89c 100644
--- a/Synapse/EventClient.cs
+++ b/Synapse/EventClient.cs
@@ -1,31 +1,25 @@
-using System;
-using System.Text;
-using System.Text.Encodings.Web;
using System.Text.Json;
-using System.Text.Unicode;
+using System.Threading;
using Microsoft.Extensions.Logging;
+using MQTTnet;
using SimApi.Helpers;
namespace SimApi;
public partial class Synapse
{
- private void RunEventClient()
- {
- EventClientChannel = CreateChannel(0, "EventClient");
- }
-
- private void FireEvent(string eventName, object param)
+ private bool FireEvent(string eventName, object param, bool retain = false)
{
var paramJson = JsonSerializer.Serialize(param, SimApiUtil.JsonOption);
- var router = $"event.{Options.AppName}.{eventName}";
- var props = EventClientChannel.CreateBasicProperties();
- props.AppId = Options.AppId;
- props.MessageId = Guid.NewGuid().ToString();
- props.ReplyTo = Options.AppName;
- props.Type = eventName;
- EventClientChannel.BasicPublish(Options.SysName, router, false, props, Encoding.UTF8.GetBytes(paramJson));
- Logger.LogDebug("Event Publish: {OptionsAppName}.{EventName}\n{ParamJson}", Options.AppName, eventName,
- paramJson);
+ var topic = $"{Options.SysName}/{Options.AppName}/event/{eventName}";
+ var message = new MqttApplicationMessageBuilder()
+ .WithTopic(topic)
+ .WithPayload(paramJson)
+ .WithRetainFlag(retain)
+ .Build();
+ if (!Client!.IsConnected) return false;
+ Client!.PublishAsync(message, CancellationToken.None).Wait();
+ logger.LogDebug("Event Publish: {Event}@{App} {Json}", eventName, Options.AppName, paramJson);
+ return true;
}
}
\ No newline at end of file
diff --git a/Synapse/EventServer.cs b/Synapse/EventServer.cs
index 5b2b843..4dfb91c 100644
--- a/Synapse/EventServer.cs
+++ b/Synapse/EventServer.cs
@@ -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();
}
}
\ No newline at end of file
diff --git a/Synapse/RpcClient.cs b/Synapse/RpcClient.cs
index 767db34..968acf4 100644
--- a/Synapse/RpcClient.cs
+++ b/Synapse/RpcClient.cs
@@ -1,9 +1,10 @@
using System;
using System.Collections.Generic;
-using System.Text;
using System.Text.Json;
+using System.Threading;
+using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
-using RabbitMQ.Client.Events;
+using MQTTnet;
using SimApi.Communications;
using SimApi.Helpers;
@@ -11,42 +12,43 @@ namespace SimApi;
public partial class Synapse
{
- private Dictionary ResponseCache { get; } = new();
+ private Dictionary ResponseCache { get; } = new();
private void RunRpcClient()
{
- RpcClientChannel = CreateChannel(0, "RpcClient");
- var queue = $"{Options.SysName}_{Options.AppName}_client_{Options.AppId}";
- var router = $"client.{Options.AppName}.{Options.AppId}";
- RpcClientChannel.QueueDeclare(queue, true, false, true, null);
- RpcClientChannel.QueueBind(queue, Options.SysName, router, null);
- var consumer = new EventingBasicConsumer(RpcClientChannel);
- consumer.Received += (ch, ea) =>
+ var rcTopic = $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/";
+ var rcSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
+ .WithTopicFilter(o => o.WithTopic($"{rcTopic}+")).Build();
+ Client.ApplicationMessageReceivedAsync += e =>
{
- ResponseCache.Add(ea.BasicProperties.CorrelationId, ea.Body.ToArray());
- RpcClientChannel.BasicAck(ea.DeliveryTag, false);
- Logger.LogDebug(
- "RPC Response: ({BasicPropertiesCorrelationId}) {BasicPropertiesType}@{BasicPropertiesReplyTo} -> {OptionsAppName}\n{S}",
- ea.BasicProperties.CorrelationId, ea.BasicProperties.Type, ea.BasicProperties.ReplyTo, Options.AppName,
- Encoding.UTF8.GetString(ea.Body.ToArray()));
+ if (!e.ApplicationMessage.Topic.StartsWith(rcTopic)) return Task.CompletedTask;
+ var reqBody = e.ApplicationMessage.ConvertPayloadToString();
+ var messageId = e.ApplicationMessage.Topic.Replace(rcTopic, string.Empty);
+ ResponseCache.Add(messageId, reqBody);
+ logger.LogDebug("Synapse RPC Client Message: ({BasicPropertiesCorrelationId}) => {S}", messageId, reqBody);
+ return Task.CompletedTask;
};
- RpcClientChannel.BasicConsume(queue, false, "", false, false, null, consumer);
+ Client.SubscribeAsync(rcSubOpts).Wait();
}
private string FireRpc(string app, string action, object param)
{
var paramJson = JsonSerializer.Serialize(param, SimApiUtil.JsonOption);
- var router = $"server.{app}";
string response;
- var props = RpcClientChannel.CreateBasicProperties();
- props.AppId = Options.AppId;
- props.MessageId = Guid.NewGuid().ToString();
- props.Type = action;
- props.ReplyTo = Options.AppName;
- RpcClientChannel.BasicPublish(Options.SysName, router, false, props, Encoding.UTF8.GetBytes(paramJson));
- Logger.LogDebug("RPC Request: ({PropsMessageId}) {OptionsAppName} -> {Action}@{App}\n{ParamJson}",
- props.MessageId,
- Options.AppName, action, app, paramJson);
+ var topic = $"{Options.SysName}/{app}/rpc/server/{action}";
+ var messageId = Guid.NewGuid().ToString();
+ var message = new MqttApplicationMessageBuilder()
+ .WithTopic(topic)
+ .WithPayload(paramJson)
+ .WithResponseTopic($"{Options.AppName},{Options.AppId}")
+ .WithContentType(messageId)
+ .WithRetainFlag(false)
+ .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);
var ts = SimApiUtil.TimestampNow;
while (true)
{
@@ -55,13 +57,16 @@ public partial class Synapse
response = JsonSerializer.Serialize(new SimApiBaseResponse(502, "timeout"), SimApiUtil.JsonOption);
break;
}
- if (ResponseCache.TryGetValue(props.MessageId, out var value))
- {
- response = Encoding.UTF8.GetString(value);
- ResponseCache.Remove(props.MessageId);
- break;
- }
+
+ if (!ResponseCache.TryGetValue(messageId, out var value)) continue;
+ response = value;
+ ResponseCache.Remove(messageId);
+ logger.LogDebug(
+ "Synapse RPC Client Response: ({BasicPropertiesCorrelationId}) {BasicPropertiesType}@{BasicPropertiesReplyTo} -> {OptionsAppName}\n{S}",
+ messageId, action, app, Options.AppName, response);
+ break;
}
+
return response;
}
}
\ No newline at end of file
diff --git a/Synapse/RpcServer.cs b/Synapse/RpcServer.cs
index 4981640..e983729 100644
--- a/Synapse/RpcServer.cs
+++ b/Synapse/RpcServer.cs
@@ -1,10 +1,12 @@
using System;
using System.Linq;
using System.Reflection;
-using System.Text;
+using System.Threading;
+using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
-using RabbitMQ.Client.Events;
+using MQTTnet;
+using MQTTnet.Protocol;
using SimApi.Communications;
using SimApi.Exceptions;
using SimApi.Helpers;
@@ -16,78 +18,70 @@ public partial class Synapse
{
private void RunRpcServer()
{
- RpcServerChannel = CreateChannel(Options.RpcProcessorNum, "RpcServer");
- var queue = $"{Options.SysName}_{Options.AppName}_server";
- var router = $"server.{Options.AppName}";
- RpcServerChannel.QueueDeclare(queue, true, false, true, null);
- RpcServerChannel.QueueBind(queue, Options.SysName, router, null);
- var consumer = new EventingBasicConsumer(RpcServerChannel);
- consumer.Received += (ch, ea) =>
+ var rsTopicPrefix = $"{Options.SysName}/{Options.AppName}/rpc/server/";
+ var rsSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
+ .WithTopicFilter(o =>
+ o.WithTopic($"$queue/{rsTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
+ .Build();
+ Client.ApplicationMessageReceivedAsync += e =>
{
- var reqBody = Encoding.UTF8.GetString(ea.Body.ToArray());
- Logger.LogDebug(
- "RPC Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}",
- ea.BasicProperties.MessageId, ea.BasicProperties.ReplyTo, ea.BasicProperties.Type, Options.AppName,
+ if (!e.ApplicationMessage.Topic.StartsWith(rsTopicPrefix)) return Task.CompletedTask;
+ var reqBody = e.ApplicationMessage.ConvertPayloadToString();
+ var action = e.ApplicationMessage.Topic.Replace(rsTopicPrefix, string.Empty);
+ 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);
var res = new SimApiBaseResponse(404, "method not found");
- var method = RpcRegistry.FirstOrDefault(x => x.Key == ea.BasicProperties.Type);
- if (method != null)
+ 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
{
- var callClass = Sp.CreateScope().ServiceProvider.GetRequiredService(method.Class);
- var mt = callClass.GetType().GetMethod(method.Method);
- var param = Array.Empty