Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a4bae0226e | ||
|
|
15f74b6570 | ||
|
|
ac4619a789 | ||
|
|
4d0f801f21 | ||
|
|
683cf9176a | ||
|
|
1ee8d941e9 | ||
|
|
1d474f596a |
@@ -54,7 +54,7 @@ public class SimApiBaseResponse(int code = 200, string message = "成功")
|
||||
/// <typeparam name="T"></typeparam>
|
||||
public class SimApiBasePageResponse<T>() : SimApiBaseResponse
|
||||
{
|
||||
public T List { get; set; }
|
||||
public T? List { get; set; }
|
||||
public int Page { get; set; } = 1;
|
||||
public int Count { get; set; } = 20;
|
||||
public int Total { get; set; }
|
||||
@@ -74,7 +74,7 @@ public class SimApiBasePageResponse<T>() : SimApiBaseResponse
|
||||
/// <typeparam name="T"></typeparam>
|
||||
public class SimApiBaseResponse<T>() : SimApiBaseResponse
|
||||
{
|
||||
public T Data { get; set; }
|
||||
public T? Data { get; set; }
|
||||
|
||||
public SimApiBaseResponse(T data) : this()
|
||||
{
|
||||
|
||||
@@ -11,17 +11,17 @@ public class SimApiDocGroupOption
|
||||
/// <summary>
|
||||
/// 文档标识
|
||||
/// </summary>
|
||||
public string Id { get; set; }
|
||||
public string Id { get; set; } = null!;
|
||||
|
||||
/// <summary>
|
||||
/// 文档名称
|
||||
/// </summary>
|
||||
public string Name { get; set; }
|
||||
public string Name { get; set; } = null!;
|
||||
|
||||
/// <summary>
|
||||
/// 文档描述
|
||||
/// </summary>
|
||||
public string Description { get; set; }
|
||||
public string Description { get; set; } = null!;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -33,11 +33,11 @@ public class SimApiAuthOption
|
||||
|
||||
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>
|
||||
|
||||
@@ -14,13 +14,13 @@ public class SimApiOptions
|
||||
/// 启用SimApiAuth,一个简单的基于Header Token的认证方式。
|
||||
/// default: false
|
||||
/// </summary>
|
||||
public bool EnableSimApiAuth { get; set; } = false;
|
||||
public bool EnableSimApiAuth { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// 启用在线文档,启用后 访问 /swagger 可以查看对应的api文档。
|
||||
/// default: false
|
||||
/// </summary>
|
||||
public bool EnableSimApiDoc { get; set; } = false;
|
||||
public bool EnableSimApiDoc { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// 启用异常拦截,启用后,所有的异常将被通过json反馈。
|
||||
@@ -32,7 +32,7 @@ public class SimApiOptions
|
||||
/// 开启S3兼容的存储系统。
|
||||
/// default: false
|
||||
/// </summary>
|
||||
public bool EnableSimApiStorage { get; set; } = false;
|
||||
public bool EnableSimApiStorage { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// 开启ForwardHeaders,开启后可以透传负载均衡的Headers
|
||||
@@ -50,12 +50,12 @@ public class SimApiOptions
|
||||
/// 启用格式化的 Console Logger
|
||||
/// default: false
|
||||
/// </summary>
|
||||
public bool EnableLogger { get; set; } = false;
|
||||
public bool EnableLogger { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// 是否启用Synapse
|
||||
/// </summary>
|
||||
public bool EnableSynapse { get; set; } = false;
|
||||
public bool EnableSynapse { get; set; }
|
||||
|
||||
|
||||
/// <summary>
|
||||
@@ -70,7 +70,7 @@ public class SimApiOptions
|
||||
|
||||
public SimApiSynapseOptions SimApiSynapseOptions { get; set; } = new();
|
||||
|
||||
public void ConfigureSimApiSynapse(Action<SimApiSynapseOptions> options = null)
|
||||
public void ConfigureSimApiSynapse(Action<SimApiSynapseOptions>? options = null)
|
||||
{
|
||||
options?.Invoke(SimApiSynapseOptions);
|
||||
}
|
||||
@@ -80,12 +80,12 @@ public class SimApiOptions
|
||||
SimApiSynapseOptions = options;
|
||||
}
|
||||
|
||||
public void ConfigureSimApiDoc(Action<SimApiDocOptions> options = null)
|
||||
public void ConfigureSimApiDoc(Action<SimApiDocOptions>? options = null)
|
||||
{
|
||||
options?.Invoke(SimApiDocOptions);
|
||||
}
|
||||
|
||||
public void ConfigureSimApiStorage(Action<SimApiStorageOptions> options = null)
|
||||
public void ConfigureSimApiStorage(Action<SimApiStorageOptions>? options = null)
|
||||
{
|
||||
options?.Invoke(SimApiStorageOptions);
|
||||
}
|
||||
|
||||
@@ -13,6 +13,12 @@ public class SimApiSynapseOptions
|
||||
public string? AppId { get; set; }
|
||||
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 DisableEventClient { get; set; } = false;
|
||||
public bool DisableRpcClient { get; set; } = false;
|
||||
|
||||
@@ -26,14 +26,12 @@ public class SimApiBaseController : Controller
|
||||
/// <param name="context"></param>
|
||||
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 =>
|
||||
item.ValidationState == ModelValidationState.Invalid))
|
||||
{
|
||||
Error(400, item.Errors.First().ErrorMessage);
|
||||
break;
|
||||
}
|
||||
Error(400, item.Errors.First().ErrorMessage);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -112,7 +112,7 @@ public class SimApiStorage
|
||||
/// <param name="path"></param>
|
||||
/// <returns></returns>
|
||||
/// <exception cref="Exception"></exception>
|
||||
public string FullUrl(string path)
|
||||
public string? FullUrl(string? path)
|
||||
{
|
||||
if (string.IsNullOrEmpty(path)) return path;
|
||||
if(path.StartsWith("http://") || path.StartsWith("https://")) return path;
|
||||
|
||||
+10
-3
@@ -1,13 +1,11 @@
|
||||
using System;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using SimApi.Helpers;
|
||||
|
||||
namespace SimApi.Logger;
|
||||
|
||||
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 void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception,
|
||||
@@ -33,3 +31,12 @@ public class SimApiLogger(string name) : ILogger
|
||||
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,
|
||||
Action<SimApiOptions> options = null)
|
||||
Action<SimApiOptions>? options = null)
|
||||
{
|
||||
var simApiOptions = new SimApiOptions();
|
||||
options?.Invoke(simApiOptions);
|
||||
@@ -224,11 +224,11 @@ public static class SimApiExtensions
|
||||
/// </summary>
|
||||
/// <param name="builder"></param>
|
||||
/// <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);
|
||||
if (options.EnableForwardHeaders)
|
||||
@@ -277,7 +277,7 @@ public static class SimApiExtensions
|
||||
if (options.EnableSimApiStorage)
|
||||
{
|
||||
logger.LogInformation("开始配置SimApiStorage...");
|
||||
builder.ApplicationServices.GetService<SimApiStorage>();
|
||||
builder.Services.GetService<SimApiStorage>();
|
||||
}
|
||||
|
||||
if (options.EnableLowerUrl)
|
||||
@@ -287,7 +287,7 @@ public static class SimApiExtensions
|
||||
|
||||
if (options.EnableSynapse)
|
||||
{
|
||||
var synapse = builder.ApplicationServices.GetRequiredService<Synapse>();
|
||||
var synapse = builder.Services.GetRequiredService<Synapse>();
|
||||
synapse.Init();
|
||||
}
|
||||
|
||||
|
||||
+15
-8
@@ -13,8 +13,11 @@ public partial class Synapse
|
||||
public event EventHandler<ConfigStoreItem>? OnConfigChanged;
|
||||
private Dictionary<string, string> CurrentConfig { get; } = new();
|
||||
|
||||
private string ConfigStoreTopicPrefix => $"{Options.SysName}/synapse-config-store/";
|
||||
|
||||
private bool FireSetConfig(string key, string value)
|
||||
{
|
||||
if (key.Contains('#') || key.Contains('+')) return false;
|
||||
var topic = $"{Options.SysName}/synapse-config-store/{key}";
|
||||
var message = new MqttApplicationMessageBuilder()
|
||||
.WithTopic(topic)
|
||||
@@ -34,22 +37,26 @@ public partial class Synapse
|
||||
|
||||
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 =>
|
||||
{
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(csTopicPrefix)) return Task.CompletedTask;
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(ConfigStoreTopicPrefix)) return Task.CompletedTask;
|
||||
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;
|
||||
OnConfigChanged?.Invoke(this, new ConfigStoreItem(eventName, reqBody));
|
||||
logger.LogDebug("Synapse Config Changed: {Config} => {Data}", eventName, reqBody);
|
||||
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);
|
||||
|
||||
+18
-8
@@ -12,14 +12,16 @@ namespace SimApi;
|
||||
|
||||
public partial class Synapse
|
||||
{
|
||||
|
||||
private string EventServerTopicPrefix => $"{Options.SysName}/event/";
|
||||
|
||||
private void RunEventServer()
|
||||
{
|
||||
var esTopicPrefix = $"{Options.SysName}/event/";
|
||||
Client!.ApplicationMessageReceivedAsync += e =>
|
||||
{
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(esTopicPrefix)) return Task.CompletedTask;
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(EventServerTopicPrefix)) return Task.CompletedTask;
|
||||
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);
|
||||
var methods = EventRegistry
|
||||
.Where(x => Regex.IsMatch(eventName,
|
||||
@@ -33,7 +35,7 @@ public partial class Synapse
|
||||
{
|
||||
if (mt!.GetParameters().Length == 2)
|
||||
{
|
||||
var pt = mt!.GetParameters()[0].ParameterType;
|
||||
var pt = mt.GetParameters()[1].ParameterType;
|
||||
mt.Invoke(callClass, pt == typeof(string)
|
||||
? [eventName, reqBody]
|
||||
: [eventName, JsonSerializer.Deserialize(reqBody, pt, SimApiUtil.JsonOption)]);
|
||||
@@ -45,18 +47,26 @@ public partial class Synapse
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
logger.LogError("Synapse Event Processor Error: {Err}", ex.InnerException);
|
||||
logger.LogError("Synapse Event Processor Error: {Err}\n{Stack}", ex.Message,ex.StackTrace);
|
||||
}
|
||||
}
|
||||
|
||||
return Task.CompletedTask;
|
||||
};
|
||||
SubEventServerTopic();
|
||||
}
|
||||
|
||||
private void SubEventServerTopic()
|
||||
{
|
||||
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()
|
||||
.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);
|
||||
}
|
||||
}
|
||||
|
||||
+12
-6
@@ -14,22 +14,28 @@ public partial class Synapse
|
||||
{
|
||||
private Dictionary<string, TaskCompletionSource<string>> ResponseCompletionSources { get; } = new();
|
||||
|
||||
private string EventClientTopicPrefix => $"{Options.SysName}/{Options.AppName}/rpc/client/{Options.AppId}/";
|
||||
|
||||
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 =>
|
||||
{
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(rcTopic)) return Task.CompletedTask;
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(EventClientTopicPrefix)) return Task.CompletedTask;
|
||||
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;
|
||||
tcs.SetResult(reqBody);
|
||||
ResponseCompletionSources.Remove(messageId);
|
||||
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)
|
||||
|
||||
+14
-8
@@ -16,18 +16,15 @@ namespace SimApi;
|
||||
|
||||
public partial class Synapse
|
||||
{
|
||||
private string RpcServerTopicPrefix => $"{Options.SysName}/{Options.AppName}/rpc/server/";
|
||||
|
||||
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 =>
|
||||
{
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(rsTopicPrefix)) return Task.CompletedTask;
|
||||
if (!e.ApplicationMessage.Topic.StartsWith(RpcServerTopicPrefix)) return Task.CompletedTask;
|
||||
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(",");
|
||||
logger.LogDebug(
|
||||
"Synapse RPC Server Receive: ({BasicPropertiesMessageId}) {BasicPropertiesReplyTo} -> {BasicPropertiesType}@{OptionsAppName}\n{S}",
|
||||
@@ -93,6 +90,15 @@ public partial class Synapse
|
||||
|
||||
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();
|
||||
}
|
||||
}
|
||||
+6
-1
@@ -165,6 +165,11 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger<Synapse> logge
|
||||
Client.ConnectedAsync += _ =>
|
||||
{
|
||||
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;
|
||||
};
|
||||
//重连
|
||||
@@ -213,7 +218,7 @@ public partial class Synapse(SimApiOptions simApiOptions, ILogger<Synapse> logge
|
||||
|
||||
var callClass = sp.CreateScope().ServiceProvider.GetRequiredService(tmp.Class!);
|
||||
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(
|
||||
"Synapse Event Register Error: Only one or two parameter supported. {Key} -> {Method}@{Class}",
|
||||
|
||||
Reference in New Issue
Block a user