Merge remote-tracking branch 'origin/master' into jason_dev
This commit is contained in:
commit
09b1f3e234
|
|
@ -6,7 +6,7 @@ public interface IEventSubscriber
|
|||
{
|
||||
Task SubscribeAsync(string channel, Func<string, string, Task> received);
|
||||
|
||||
Task SubscribeAsync(string channel, string group, bool priorityEnabled,
|
||||
Task SubscribeAsync(string channel, string group, int? port, bool priorityEnabled,
|
||||
Func<string, string, Task> received,
|
||||
CancellationToken? stoppingToken = null);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2,18 +2,19 @@ namespace BotSharp.Abstraction.Repositories;
|
|||
|
||||
public class BotSharpDatabaseSettings : DatabaseBasicSettings
|
||||
{
|
||||
public string[] Assemblies { get; set; }
|
||||
public string FileRepository { get; set; }
|
||||
public string BotSharpMongoDb { get; set; }
|
||||
public string TablePrefix { get; set; }
|
||||
public DbConnectionSetting BotSharp { get; set; }
|
||||
public string Redis { get; set; }
|
||||
public string[] Assemblies { get; set; } = [];
|
||||
public string FileRepository { get; set; } = string.Empty;
|
||||
public string BotSharpMongoDb { get; set; } = string.Empty;
|
||||
public string TablePrefix { get; set; } = string.Empty;
|
||||
public DbConnectionSetting BotSharp { get; set; } = new();
|
||||
public string Redis { get; set; } = string.Empty;
|
||||
public bool EnableReplica { get; set; } = true;
|
||||
}
|
||||
|
||||
public class DatabaseBasicSettings
|
||||
{
|
||||
public string Default { get; set; }
|
||||
public DbConnectionSetting DefaultConnection { get; set; }
|
||||
public string Default { get; set; } = string.Empty;
|
||||
public DbConnectionSetting DefaultConnection { get; set; } = new();
|
||||
public bool EnableSqlLog { get; set; }
|
||||
public bool EnableSensitiveDataLogging { get; set; }
|
||||
public bool EnableRetryOnFailure { get; set; }
|
||||
|
|
@ -23,9 +24,11 @@ public class DbConnectionSetting
|
|||
{
|
||||
public string Master { get; set; }
|
||||
public string[] Slavers { get; set; }
|
||||
public int ConnectionTimeout { get; set; } = 30;
|
||||
public int ExecutionTimeout { get; set; } = 30;
|
||||
|
||||
public DbConnectionSetting()
|
||||
{
|
||||
Slavers = new string[0];
|
||||
Slavers = [];
|
||||
}
|
||||
}
|
||||
14
src/Infrastructure/BotSharp.Abstraction/Utilities/MathExt.cs
Normal file
14
src/Infrastructure/BotSharp.Abstraction/Utilities/MathExt.cs
Normal file
|
|
@ -0,0 +1,14 @@
|
|||
namespace BotSharp.Abstraction.Utilities;
|
||||
|
||||
public static class MathExt
|
||||
{
|
||||
public static int Max(int a, int b, int c)
|
||||
{
|
||||
return Math.Max(Math.Max(a, b), c);
|
||||
}
|
||||
|
||||
public static long Max(long a, long b, long c)
|
||||
{
|
||||
return Math.Max(Math.Max(a, b), c);
|
||||
}
|
||||
}
|
||||
|
|
@ -190,7 +190,7 @@
|
|||
<ItemGroup>
|
||||
<PackageReference Include="Aspects.Cache" Version="2.0.4" />
|
||||
<PackageReference Include="DistributedLock.Redis" Version="1.0.3" />
|
||||
<PackageReference Include="EntityFrameworkCore.BootKit" Version="8.5.1" />
|
||||
<PackageReference Include="EntityFrameworkCore.BootKit" Version="8.6.0" />
|
||||
<PackageReference Include="Fluid.Core" Version="2.11.1" />
|
||||
<PackageReference Include="Microsoft.Extensions.Caching.Memory" Version="8.0.1" />
|
||||
<PackageReference Include="Microsoft.Extensions.Http" Version="8.0.0" />
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ public class RedisSubscriber : IEventSubscriber
|
|||
});
|
||||
}
|
||||
|
||||
public async Task SubscribeAsync(string channel, string group, bool priorityEnabled,
|
||||
public async Task SubscribeAsync(string channel, string group, int? port, bool priorityEnabled,
|
||||
Func<string, string, Task> received,
|
||||
CancellationToken? stoppingToken = null)
|
||||
{
|
||||
|
|
@ -42,6 +42,12 @@ public class RedisSubscriber : IEventSubscriber
|
|||
await CreateConsumerGroup(db, channel, group);
|
||||
}
|
||||
|
||||
var consumer = Environment.MachineName;
|
||||
if (port.HasValue)
|
||||
{
|
||||
consumer += $"-{port}";
|
||||
}
|
||||
|
||||
while (true)
|
||||
{
|
||||
await Task.Delay(100);
|
||||
|
|
@ -54,28 +60,28 @@ public class RedisSubscriber : IEventSubscriber
|
|||
|
||||
if (priorityEnabled)
|
||||
{
|
||||
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, received) > 0)
|
||||
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.High}", group, consumer, received) > 0)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, received) > 0)
|
||||
if (await HandleGroupMessage(db, $"{channel}-{EventPriority.Medium}", group, consumer, received) > 0)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, received);
|
||||
await HandleGroupMessage(db, $"{channel}-{EventPriority.Low}", group, consumer, received);
|
||||
}
|
||||
else
|
||||
{
|
||||
await HandleGroupMessage(db, channel, group, received);
|
||||
await HandleGroupMessage(db, channel, group, consumer, received);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<int> HandleGroupMessage(IDatabase db, string channel, string group, Func<string, string, Task> received)
|
||||
private async Task<int> HandleGroupMessage(IDatabase db, string channel, string group, string consumer, Func<string, string, Task> received)
|
||||
{
|
||||
var entries = await db.StreamReadGroupAsync(channel, group, Environment.MachineName, count: 1);
|
||||
var entries = await db.StreamReadGroupAsync(channel, group, consumer, count: 1);
|
||||
foreach (var entry in entries)
|
||||
{
|
||||
_logger.LogInformation($"Consumer {Environment.MachineName} received: {channel} {entry.Values[0].Value}");
|
||||
|
|
|
|||
|
|
@ -11,8 +11,7 @@
|
|||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Aspire.MongoDB.Driver" Version="8.0.1" />
|
||||
<PackageReference Include="MongoDB.Driver" Version="2.28.0" />
|
||||
<PackageReference Include="MongoDB.Driver" Version="3.0.0" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
|
|
|
|||
Loading…
Reference in a new issue