Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
141846d79d | ||
|
|
a4bae0226e | ||
|
|
15f74b6570 | ||
|
|
ac4619a789 | ||
|
|
4d0f801f21 | ||
|
|
683cf9176a | ||
|
|
1ee8d941e9 | ||
|
|
1d474f596a |
@@ -54,7 +54,7 @@ public class SimApiBaseResponse(int code = 200, string message = "成功")
|
|||||||
/// <typeparam name="T"></typeparam>
|
/// <typeparam name="T"></typeparam>
|
||||||
public class SimApiBasePageResponse<T>() : SimApiBaseResponse
|
public class SimApiBasePageResponse<T>() : SimApiBaseResponse
|
||||||
{
|
{
|
||||||
public T List { get; set; }
|
public T? List { get; set; }
|
||||||
public int Page { get; set; } = 1;
|
public int Page { get; set; } = 1;
|
||||||
public int Count { get; set; } = 20;
|
public int Count { get; set; } = 20;
|
||||||
public int Total { get; set; }
|
public int Total { get; set; }
|
||||||
@@ -74,7 +74,7 @@ public class SimApiBasePageResponse<T>() : SimApiBaseResponse
|
|||||||
/// <typeparam name="T"></typeparam>
|
/// <typeparam name="T"></typeparam>
|
||||||
public class SimApiBaseResponse<T>() : SimApiBaseResponse
|
public class SimApiBaseResponse<T>() : SimApiBaseResponse
|
||||||
{
|
{
|
||||||
public T Data { get; set; }
|
public T? Data { get; set; }
|
||||||
|
|
||||||
public SimApiBaseResponse(T data) : this()
|
public SimApiBaseResponse(T data) : this()
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -11,17 +11,17 @@ public class SimApiDocGroupOption
|
|||||||
/// <summary>
|
/// <summary>
|
||||||
/// 文档标识
|
/// 文档标识
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public string Id { get; set; }
|
public string Id { get; set; } = null!;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// 文档名称
|
/// 文档名称
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public string Name { get; set; }
|
public string Name { get; set; } = null!;
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// 文档描述
|
/// 文档描述
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public string Description { get; set; }
|
public string Description { get; set; } = null!;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -33,11 +33,11 @@ public class SimApiAuthOption
|
|||||||
|
|
||||||
public string Description { get; set; } = "认证服务器颁发的AccessToken";
|
public string Description { get; set; } = "认证服务器颁发的AccessToken";
|
||||||
|
|
||||||
public string AuthorizationUrl { get; set; }
|
public string AuthorizationUrl { get; set; } = null!;
|
||||||
|
|
||||||
public string TokenUrl { get; set; }
|
public string TokenUrl { get; set; } = null!;
|
||||||
|
|
||||||
public Dictionary<string, string> Scopes { get; set; }
|
public Dictionary<string, string> Scopes { get; set; } = null!;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
|
|||||||
@@ -14,13 +14,13 @@ public class SimApiOptions
|
|||||||
/// 启用SimApiAuth,一个简单的基于Header Token的认证方式。
|
/// 启用SimApiAuth,一个简单的基于Header Token的认证方式。
|
||||||
/// default: false
|
/// default: false
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public bool EnableSimApiAuth { get; set; } = false;
|
public bool EnableSimApiAuth { get; set; }
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// 启用在线文档,启用后 访问 /swagger 可以查看对应的api文档。
|
/// 启用在线文档,启用后 访问 /swagger 可以查看对应的api文档。
|
||||||
/// default: false
|
/// default: false
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public bool EnableSimApiDoc { get; set; } = false;
|
public bool EnableSimApiDoc { get; set; }
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// 启用异常拦截,启用后,所有的异常将被通过json反馈。
|
/// 启用异常拦截,启用后,所有的异常将被通过json反馈。
|
||||||
@@ -32,7 +32,7 @@ public class SimApiOptions
|
|||||||
/// 开启S3兼容的存储系统。
|
/// 开启S3兼容的存储系统。
|
||||||
/// default: false
|
/// default: false
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public bool EnableSimApiStorage { get; set; } = false;
|
public bool EnableSimApiStorage { get; set; }
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// 开启ForwardHeaders,开启后可以透传负载均衡的Headers
|
/// 开启ForwardHeaders,开启后可以透传负载均衡的Headers
|
||||||
@@ -50,12 +50,12 @@ public class SimApiOptions
|
|||||||
/// 启用格式化的 Console Logger
|
/// 启用格式化的 Console Logger
|
||||||
/// default: false
|
/// default: false
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public bool EnableLogger { get; set; } = false;
|
public bool EnableLogger { get; set; }
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
/// 是否启用Synapse
|
/// 是否启用Synapse
|
||||||
/// </summary>
|
/// </summary>
|
||||||
public bool EnableSynapse { get; set; } = false;
|
public bool EnableSynapse { get; set; }
|
||||||
|
|
||||||
|
|
||||||
/// <summary>
|
/// <summary>
|
||||||
@@ -70,7 +70,7 @@ public class SimApiOptions
|
|||||||
|
|
||||||
public SimApiSynapseOptions SimApiSynapseOptions { get; set; } = new();
|
public SimApiSynapseOptions SimApiSynapseOptions { get; set; } = new();
|
||||||
|
|
||||||
public void ConfigureSimApiSynapse(Action<SimApiSynapseOptions> options = null)
|
public void ConfigureSimApiSynapse(Action<SimApiSynapseOptions>? options = null)
|
||||||
{
|
{
|
||||||
options?.Invoke(SimApiSynapseOptions);
|
options?.Invoke(SimApiSynapseOptions);
|
||||||
}
|
}
|
||||||
@@ -80,12 +80,12 @@ public class SimApiOptions
|
|||||||
SimApiSynapseOptions = options;
|
SimApiSynapseOptions = options;
|
||||||
}
|
}
|
||||||
|
|
||||||
public void ConfigureSimApiDoc(Action<SimApiDocOptions> options = null)
|
public void ConfigureSimApiDoc(Action<SimApiDocOptions>? options = null)
|
||||||
{
|
{
|
||||||
options?.Invoke(SimApiDocOptions);
|
options?.Invoke(SimApiDocOptions);
|
||||||
}
|
}
|
||||||
|
|
||||||
public void ConfigureSimApiStorage(Action<SimApiStorageOptions> options = null)
|
public void ConfigureSimApiStorage(Action<SimApiStorageOptions>? options = null)
|
||||||
{
|
{
|
||||||
options?.Invoke(SimApiStorageOptions);
|
options?.Invoke(SimApiStorageOptions);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,6 +13,12 @@ public class SimApiSynapseOptions
|
|||||||
public string? AppId { get; set; }
|
public string? AppId { get; set; }
|
||||||
public int RpcTimeout { get; set; } = 3;
|
public int RpcTimeout { get; set; } = 3;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Event是否使用负载均衡
|
||||||
|
/// 也就是订阅$queue主题,消息会分发给不同的AppId
|
||||||
|
/// 如果false,多个AppId都可以同时收到消息
|
||||||
|
/// </summary>
|
||||||
|
public bool EventLoadBalancing { get; set; } = false;
|
||||||
public bool EnableConfigStore { get; set; } = true;
|
public bool EnableConfigStore { get; set; } = true;
|
||||||
public bool DisableEventClient { get; set; } = false;
|
public bool DisableEventClient { get; set; } = false;
|
||||||
public bool DisableRpcClient { get; set; } = false;
|
public bool DisableRpcClient { get; set; } = false;
|
||||||
|
|||||||
@@ -26,14 +26,12 @@ public class SimApiBaseController : Controller
|
|||||||
/// <param name="context"></param>
|
/// <param name="context"></param>
|
||||||
public override void OnActionExecuting(ActionExecutingContext context)
|
public override void OnActionExecuting(ActionExecutingContext context)
|
||||||
{
|
{
|
||||||
if (!context.ModelState.IsValid)
|
if (context.ModelState.IsValid) return;
|
||||||
|
foreach (var item in context.ModelState.Values.Where(item =>
|
||||||
|
item.ValidationState == ModelValidationState.Invalid))
|
||||||
{
|
{
|
||||||
foreach (var item in context.ModelState.Values.Where(item =>
|
Error(400, item.Errors.First().ErrorMessage);
|
||||||
item.ValidationState == ModelValidationState.Invalid))
|
break;
|
||||||
{
|
|
||||||
Error(400, item.Errors.First().ErrorMessage);
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -112,7 +112,7 @@ public class SimApiStorage
|
|||||||
/// <param name="path"></param>
|
/// <param name="path"></param>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
/// <exception cref="Exception"></exception>
|
/// <exception cref="Exception"></exception>
|
||||||
public string FullUrl(string path)
|
public string? FullUrl(string? path)
|
||||||
{
|
{
|
||||||
if (string.IsNullOrEmpty(path)) return path;
|
if (string.IsNullOrEmpty(path)) return path;
|
||||||
if(path.StartsWith("http://") || path.StartsWith("https://")) return path;
|
if(path.StartsWith("http://") || path.StartsWith("https://")) return path;
|
||||||
|
|||||||
+10
-3
@@ -1,13 +1,11 @@
|
|||||||
using System;
|
using System;
|
||||||
using Microsoft.Extensions.Logging;
|
using Microsoft.Extensions.Logging;
|
||||||
using SimApi.Helpers;
|
|
||||||
|
|
||||||
namespace SimApi.Logger;
|
namespace SimApi.Logger;
|
||||||
|
|
||||||
public class SimApiLogger(string name) : ILogger
|
public class SimApiLogger(string name) : ILogger
|
||||||
{
|
{
|
||||||
public IDisposable BeginScope<TState>(TState state) => default!;
|
public IDisposable BeginScope<TState>(TState state) where TState : notnull => default!;
|
||||||
|
|
||||||
public bool IsEnabled(LogLevel logLevel) => true;
|
public bool IsEnabled(LogLevel logLevel) => true;
|
||||||
|
|
||||||
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception,
|
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception,
|
||||||
@@ -32,4 +30,13 @@ public class SimApiLogger(string name) : ILogger
|
|||||||
Console.WriteLine(message);
|
Console.WriteLine(message);
|
||||||
Console.ResetColor();
|
Console.ResetColor();
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public class EmptyDisposable : IDisposable
|
||||||
|
{
|
||||||
|
public static EmptyDisposable Instance { get; } = new EmptyDisposable();
|
||||||
|
|
||||||
|
public void Dispose()
|
||||||
|
{
|
||||||
|
}
|
||||||
}
|
}
|
||||||
+6
-6
@@ -19,7 +19,7 @@ public static class SimApiExtensions
|
|||||||
{
|
{
|
||||||
//**********快捷添加**************
|
//**********快捷添加**************
|
||||||
public static IServiceCollection AddSimApi(this IServiceCollection builder,
|
public static IServiceCollection AddSimApi(this IServiceCollection builder,
|
||||||
Action<SimApiOptions> options = null)
|
Action<SimApiOptions>? options = null)
|
||||||
{
|
{
|
||||||
var simApiOptions = new SimApiOptions();
|
var simApiOptions = new SimApiOptions();
|
||||||
options?.Invoke(simApiOptions);
|
options?.Invoke(simApiOptions);
|
||||||
@@ -224,11 +224,11 @@ public static class SimApiExtensions
|
|||||||
/// </summary>
|
/// </summary>
|
||||||
/// <param name="builder"></param>
|
/// <param name="builder"></param>
|
||||||
/// <returns></returns>
|
/// <returns></returns>
|
||||||
public static IApplicationBuilder UseSimApi(this IApplicationBuilder builder)
|
public static WebApplication UseSimApi(this WebApplication builder)
|
||||||
{
|
{
|
||||||
var options = builder.ApplicationServices.GetRequiredService<SimApiOptions>();
|
var options = builder.Services.GetRequiredService<SimApiOptions>();
|
||||||
|
|
||||||
var logger = builder.ApplicationServices.GetRequiredService<ILogger<SimApiOptions>>();
|
var logger = builder.Services.GetRequiredService<ILogger<SimApiOptions>>();
|
||||||
|
|
||||||
logger.LogInformation("当前时区: {LocalId}", TimeZoneInfo.Local.Id);
|
logger.LogInformation("当前时区: {LocalId}", TimeZoneInfo.Local.Id);
|
||||||
if (options.EnableForwardHeaders)
|
if (options.EnableForwardHeaders)
|
||||||
@@ -277,7 +277,7 @@ public static class SimApiExtensions
|
|||||||
if (options.EnableSimApiStorage)
|
if (options.EnableSimApiStorage)
|
||||||
{
|
{
|
||||||
logger.LogInformation("开始配置SimApiStorage...");
|
logger.LogInformation("开始配置SimApiStorage...");
|
||||||
builder.ApplicationServices.GetService<SimApiStorage>();
|
builder.Services.GetService<SimApiStorage>();
|
||||||
}
|
}
|
||||||
|
|
||||||
if (options.EnableLowerUrl)
|
if (options.EnableLowerUrl)
|
||||||
@@ -287,7 +287,7 @@ public static class SimApiExtensions
|
|||||||
|
|
||||||
if (options.EnableSynapse)
|
if (options.EnableSynapse)
|
||||||
{
|
{
|
||||||
var synapse = builder.ApplicationServices.GetRequiredService<Synapse>();
|
var synapse = builder.Services.GetRequiredService<Synapse>();
|
||||||
synapse.Init();
|
synapse.Init();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+15
-8
@@ -13,8 +13,11 @@ public partial class Synapse
|
|||||||
public event EventHandler<ConfigStoreItem>? OnConfigChanged;
|
public event EventHandler<ConfigStoreItem>? OnConfigChanged;
|
||||||
private Dictionary<string, string> CurrentConfig { get; } = new();
|
private Dictionary<string, string> CurrentConfig { get; } = new();
|
||||||
|
|
||||||
|
private string ConfigStoreTopicPrefix => $"{Options.SysName}/synapse-config-store/";
|
||||||
|
|
||||||
private bool FireSetConfig(string key, string value)
|
private bool FireSetConfig(string key, string value)
|
||||||
{
|
{
|
||||||
|
if (key.Contains('#') || key.Contains('+')) return false;
|
||||||
var topic = $"{Options.SysName}/synapse-config-store/{key}";
|
var topic = $"{Options.SysName}/synapse-config-store/{key}";
|
||||||
var message = new MqttApplicationMessageBuilder()
|
var message = new MqttApplicationMessageBuilder()
|
||||||
.WithTopic(topic)
|
.WithTopic(topic)
|
||||||
@@ -34,22 +37,26 @@ public partial class Synapse
|
|||||||
|
|
||||||
private void RunConfigStoreServer()
|
private void RunConfigStoreServer()
|
||||||
{
|
{
|
||||||
var csTopicPrefix = $"{Options.SysName}/synapse-config-store/";
|
|
||||||
var eventSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
|
||||||
.WithTopicFilter(o =>
|
|
||||||
o.WithTopic($"{csTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
|
|
||||||
.Build();
|
|
||||||
Client!.ApplicationMessageReceivedAsync += e =>
|
Client!.ApplicationMessageReceivedAsync += e =>
|
||||||
{
|
{
|
||||||
if (!e.ApplicationMessage.Topic.StartsWith(csTopicPrefix)) return Task.CompletedTask;
|
if (!e.ApplicationMessage.Topic.StartsWith(ConfigStoreTopicPrefix)) return Task.CompletedTask;
|
||||||
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
||||||
var eventName = e.ApplicationMessage.Topic.Replace(csTopicPrefix, string.Empty);
|
var eventName = e.ApplicationMessage.Topic.Replace(ConfigStoreTopicPrefix, string.Empty);
|
||||||
CurrentConfig[eventName] = reqBody;
|
CurrentConfig[eventName] = reqBody;
|
||||||
OnConfigChanged?.Invoke(this, new ConfigStoreItem(eventName, reqBody));
|
OnConfigChanged?.Invoke(this, new ConfigStoreItem(eventName, reqBody));
|
||||||
logger.LogDebug("Synapse Config Changed: {Config} => {Data}", eventName, reqBody);
|
logger.LogDebug("Synapse Config Changed: {Config} => {Data}", eventName, reqBody);
|
||||||
return Task.CompletedTask;
|
return Task.CompletedTask;
|
||||||
};
|
};
|
||||||
Client.SubscribeAsync(eventSubOpts).Wait();
|
SubConfigStoreServerTopic();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void SubConfigStoreServerTopic()
|
||||||
|
{
|
||||||
|
var csSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
||||||
|
.WithTopicFilter(o =>
|
||||||
|
o.WithTopic($"{ConfigStoreTopicPrefix}#").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
|
||||||
|
.Build();
|
||||||
|
Client!.SubscribeAsync(csSubOpts).Wait();
|
||||||
}
|
}
|
||||||
|
|
||||||
public record ConfigStoreItem(string Key, string Value);
|
public record ConfigStoreItem(string Key, string Value);
|
||||||
|
|||||||
+38
-26
@@ -1,4 +1,5 @@
|
|||||||
using System;
|
using System;
|
||||||
|
using System.Collections.Generic;
|
||||||
using System.Linq;
|
using System.Linq;
|
||||||
using System.Text.Json;
|
using System.Text.Json;
|
||||||
using System.Text.RegularExpressions;
|
using System.Text.RegularExpressions;
|
||||||
@@ -12,51 +13,62 @@ namespace SimApi;
|
|||||||
|
|
||||||
public partial class Synapse
|
public partial class Synapse
|
||||||
{
|
{
|
||||||
|
private string EventServerTopicPrefix => $"{Options.SysName}/event/";
|
||||||
|
|
||||||
private void RunEventServer()
|
private void RunEventServer()
|
||||||
{
|
{
|
||||||
var esTopicPrefix = $"{Options.SysName}/event/";
|
Client!.ApplicationMessageReceivedAsync += async e =>
|
||||||
Client!.ApplicationMessageReceivedAsync += e =>
|
|
||||||
{
|
{
|
||||||
if (!e.ApplicationMessage.Topic.StartsWith(esTopicPrefix)) return Task.CompletedTask;
|
if (!e.ApplicationMessage.Topic.StartsWith(EventServerTopicPrefix)) return;
|
||||||
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
||||||
var eventName = e.ApplicationMessage.Topic.Replace(esTopicPrefix, string.Empty);
|
var eventName = e.ApplicationMessage.Topic.Replace(EventServerTopicPrefix, string.Empty);
|
||||||
logger.LogDebug("Synapse Event Receive: {EventName}\n{Body}", eventName, reqBody);
|
logger.LogDebug("Synapse Event Receive: {EventName}\n{Body}", eventName, reqBody);
|
||||||
var methods = EventRegistry
|
var methods = EventRegistry
|
||||||
.Where(x => Regex.IsMatch(eventName,
|
.Where(x => Regex.IsMatch(eventName,
|
||||||
"^" + Regex.Escape(x.Key!).Replace("\\+", "[^/]+").Replace("\\#", ".*") + "$"))
|
"^" + Regex.Escape(x.Key!).Replace("\\+", "[^/]+").Replace("\\#", ".*") + "$"))
|
||||||
.ToArray();
|
.ToArray();
|
||||||
foreach (var method in methods)
|
var tasks = methods.Select(method => Task.Run(() =>
|
||||||
{
|
|
||||||
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!);
|
|
||||||
var mt = callClass.GetType().GetMethod(method.Method!);
|
|
||||||
try
|
|
||||||
{
|
{
|
||||||
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()[0].ParameterType;
|
if (mt!.GetParameters().Length == 2)
|
||||||
mt.Invoke(callClass, pt == typeof(string)
|
{
|
||||||
? [eventName, reqBody]
|
var pt = mt.GetParameters()[1].ParameterType;
|
||||||
: [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]);
|
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)
|
.ToList();
|
||||||
{
|
await Task.WhenAll(tasks);
|
||||||
logger.LogError("Synapse Event Processor Error: {Err}", ex.InnerException);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return Task.CompletedTask;
|
|
||||||
};
|
};
|
||||||
|
SubEventServerTopic();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void SubEventServerTopic()
|
||||||
|
{
|
||||||
foreach (var ev in EventRegistry)
|
foreach (var ev in EventRegistry)
|
||||||
{
|
{
|
||||||
var topic = $"$queue/{esTopicPrefix}{ev.Key}";
|
var topic = $"{EventServerTopicPrefix}{ev.Key}";
|
||||||
|
if (Options.EventLoadBalancing)
|
||||||
|
{
|
||||||
|
topic = "$queue/" + topic;
|
||||||
|
}
|
||||||
|
|
||||||
var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
var evSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
||||||
.WithTopicFilter(o => o.WithTopic(topic)).Build();
|
.WithTopicFilter(o => o.WithTopic(topic)).Build();
|
||||||
Client.SubscribeAsync(evSubOpts).Wait();
|
Client!.SubscribeAsync(evSubOpts).Wait();
|
||||||
logger.LogDebug("Synapse Event Register Event Success: {EventName}\nFull Topic: {Topic}", ev.Key, topic);
|
logger.LogDebug("Synapse Event Register Event Success: {EventName}\nFull Topic: {Topic}", ev.Key, topic);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+12
-6
@@ -14,22 +14,28 @@ public partial class Synapse
|
|||||||
{
|
{
|
||||||
private Dictionary<string, TaskCompletionSource<string>> ResponseCompletionSources { get; } = new();
|
private Dictionary<string, TaskCompletionSource<string>> ResponseCompletionSources { get; } = new();
|
||||||
|
|
||||||
|
private string EventClientTopicPrefix => $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/";
|
||||||
|
|
||||||
private void RunRpcClient()
|
private void RunRpcClient()
|
||||||
{
|
{
|
||||||
var rcTopic = $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/";
|
|
||||||
var rcSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
|
||||||
.WithTopicFilter(o => o.WithTopic($"{rcTopic}+")).Build();
|
|
||||||
Client!.ApplicationMessageReceivedAsync += e =>
|
Client!.ApplicationMessageReceivedAsync += e =>
|
||||||
{
|
{
|
||||||
if (!e.ApplicationMessage.Topic.StartsWith(rcTopic)) return Task.CompletedTask;
|
if (!e.ApplicationMessage.Topic.StartsWith(EventClientTopicPrefix)) return Task.CompletedTask;
|
||||||
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
||||||
var messageId = e.ApplicationMessage.Topic.Replace(rcTopic, string.Empty);
|
var messageId = e.ApplicationMessage.Topic.Replace(EventClientTopicPrefix, string.Empty);
|
||||||
if (!ResponseCompletionSources.TryGetValue(messageId, out var tcs)) return Task.CompletedTask;
|
if (!ResponseCompletionSources.TryGetValue(messageId, out var tcs)) return Task.CompletedTask;
|
||||||
tcs.SetResult(reqBody);
|
tcs.SetResult(reqBody);
|
||||||
ResponseCompletionSources.Remove(messageId);
|
ResponseCompletionSources.Remove(messageId);
|
||||||
return Task.CompletedTask;
|
return Task.CompletedTask;
|
||||||
};
|
};
|
||||||
Client.SubscribeAsync(rcSubOpts).Wait();
|
SubRpcClientTopic();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void SubRpcClientTopic()
|
||||||
|
{
|
||||||
|
var rcSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
||||||
|
.WithTopicFilter(o => o.WithTopic($"{EventClientTopicPrefix}+")).Build();
|
||||||
|
Client!.SubscribeAsync(rcSubOpts).Wait();
|
||||||
}
|
}
|
||||||
|
|
||||||
private string? FireRpc(string app, string action, object? param)
|
private string? FireRpc(string app, string action, object? param)
|
||||||
|
|||||||
+53
-41
@@ -16,18 +16,15 @@ namespace SimApi;
|
|||||||
|
|
||||||
public partial class Synapse
|
public partial class Synapse
|
||||||
{
|
{
|
||||||
|
private string RpcServerTopicPrefix => $"{Options.SysName}/{Options.AppName}/rpc/server/";
|
||||||
|
|
||||||
private void RunRpcServer()
|
private void RunRpcServer()
|
||||||
{
|
{
|
||||||
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 =>
|
Client!.ApplicationMessageReceivedAsync += e =>
|
||||||
{
|
{
|
||||||
if (!e.ApplicationMessage.Topic.StartsWith(rsTopicPrefix)) return Task.CompletedTask;
|
if (!e.ApplicationMessage.Topic.StartsWith(RpcServerTopicPrefix)) return Task.CompletedTask;
|
||||||
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
var reqBody = e.ApplicationMessage.ConvertPayloadToString();
|
||||||
var action = e.ApplicationMessage.Topic.Replace(rsTopicPrefix, string.Empty);
|
var action = e.ApplicationMessage.Topic.Replace(RpcServerTopicPrefix, string.Empty);
|
||||||
var appInfo = e.ApplicationMessage.ResponseTopic.Split(",");
|
var appInfo = e.ApplicationMessage.ResponseTopic.Split(",");
|
||||||
logger.LogDebug(
|
logger.LogDebug(
|
||||||
"Synapse RPC Server Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}",
|
"Synapse RPC Server Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}",
|
||||||
@@ -35,48 +32,54 @@ public partial class Synapse
|
|||||||
reqBody);
|
reqBody);
|
||||||
SimApiBaseResponse res;
|
SimApiBaseResponse res;
|
||||||
var method = RpcRegistry.FirstOrDefault(x => x.Key == action);
|
var method = RpcRegistry.FirstOrDefault(x => x.Key == action);
|
||||||
if (method == null) return Task.CompletedTask;
|
if (method == null)
|
||||||
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(method.Class!);
|
|
||||||
var mt = callClass.GetType().GetMethod(method.Method!);
|
|
||||||
try
|
|
||||||
{
|
{
|
||||||
var methodParams = mt!.GetParameters();
|
res = new SimApiBaseResponse(404, "method not found");
|
||||||
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
|
|
||||||
};
|
|
||||||
}
|
}
|
||||||
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);
|
var methodParams = mt!.GetParameters();
|
||||||
res = new SimApiBaseResponse(ie.Code, ie.Message);
|
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);
|
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 returnJson = JsonSerializer.Serialize((object)res, SimApiUtil.JsonOption);
|
||||||
var reply = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{e.ApplicationMessage.ContentType}";
|
var reply = $"{Options.SysName}/{appInfo[0]}/rpc/client/{appInfo[1]}/{e.ApplicationMessage.ContentType}";
|
||||||
@@ -93,6 +96,15 @@ public partial class Synapse
|
|||||||
|
|
||||||
return Task.CompletedTask;
|
return Task.CompletedTask;
|
||||||
};
|
};
|
||||||
Client.SubscribeAsync(rsSubOpts).Wait();
|
SubRpcServerTopic();
|
||||||
|
}
|
||||||
|
|
||||||
|
private void SubRpcServerTopic()
|
||||||
|
{
|
||||||
|
var rsSubOpts = MqttFactory.CreateSubscribeOptionsBuilder()
|
||||||
|
.WithTopicFilter(o =>
|
||||||
|
o.WithTopic($"$queue/{RpcServerTopicPrefix}+").WithRetainHandling(MqttRetainHandling.SendAtSubscribe))
|
||||||
|
.Build();
|
||||||
|
Client!.SubscribeAsync(rsSubOpts).Wait();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
+7
-2
@@ -165,6 +165,11 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger<Synapse> logge
|
|||||||
Client.ConnectedAsync += _ =>
|
Client.ConnectedAsync += _ =>
|
||||||
{
|
{
|
||||||
logger.LogInformation("Synapse MQTT[{AppName}:{AppId}] 连接成功...", Options.AppName, Options.AppId);
|
logger.LogInformation("Synapse MQTT[{AppName}:{AppId}] 连接成功...", Options.AppName, Options.AppId);
|
||||||
|
//RPC客户端
|
||||||
|
if (!Options.DisableRpcClient) SubRpcClientTopic();
|
||||||
|
if (RpcRegistry.Count > 0) SubRpcServerTopic();
|
||||||
|
if (EventRegistry.Count > 0) SubEventServerTopic();
|
||||||
|
if (Options.EnableConfigStore) SubConfigStoreServerTopic();
|
||||||
return Task.CompletedTask;
|
return Task.CompletedTask;
|
||||||
};
|
};
|
||||||
//重连
|
//重连
|
||||||
@@ -210,10 +215,10 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger<Synapse> logge
|
|||||||
logger.LogError("Synapse Event Register Error: {Key} Can't start or end of '/'", tmp.Key);
|
logger.LogError("Synapse Event Register Error: {Key} Can't start or end of '/'", tmp.Key);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(tmp.Class!);
|
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(tmp.Class!);
|
||||||
var mt = callClass.GetType().GetMethod(tmp.Method);
|
var mt = callClass.GetType().GetMethod(tmp.Method);
|
||||||
if (mt!.GetParameters().Length > 2 || mt.GetParameters().Length<1)
|
if (mt!.GetParameters().Length > 2 || mt.GetParameters().Length < 1)
|
||||||
{
|
{
|
||||||
logger.LogError(
|
logger.LogError(
|
||||||
"Synapse Event Register Error: Only one or two parameter supported. {Key} -> {Method}@{Class}",
|
"Synapse Event Register Error: Only one or two parameter supported. {Key} -> {Method}@{Class}",
|
||||||
|
|||||||
Reference in New Issue
Block a user