Redis proto cluster member store (#4241)
* Implemented redis as a proto cluster member store * Added component tests for redis protocluster member store --------- Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com>
This commit is contained in:
parent
c11aa6cf5f
commit
90e560bc73
7
Elsa.sln
7
Elsa.sln
|
|
@ -230,6 +230,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.AspNet.MongoDb
|
|||
EndProject
|
||||
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Workflows.Api.Rpc", "src\modules\Elsa.Workflows.Api.Rpc\Elsa.Workflows.Api.Rpc.csproj", "{0A94D69A-61AC-4654-83AB-B24445628A99}"
|
||||
EndProject
|
||||
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.ProtoCluster.ComponentTests", "test\Elsa.ProtoCluster.ComponentTests\Elsa.ProtoCluster.ComponentTests.csproj", "{2430CB5F-7D07-4A9E-BD40-EC1B111B85FF}"
|
||||
EndProject
|
||||
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.AspNet.DynamicActivityProvider", "src\samples\aspnet\Elsa.Samples.AspNet.DynamicActivityProvider\Elsa.Samples.AspNet.DynamicActivityProvider.csproj", "{85E13383-7C39-4719-AAC0-0B357C3A97C7}"
|
||||
EndProject
|
||||
Global
|
||||
|
|
@ -578,6 +580,10 @@ Global
|
|||
{0A94D69A-61AC-4654-83AB-B24445628A99}.Debug|Any CPU.Build.0 = Debug|Any CPU
|
||||
{0A94D69A-61AC-4654-83AB-B24445628A99}.Release|Any CPU.ActiveCfg = Release|Any CPU
|
||||
{0A94D69A-61AC-4654-83AB-B24445628A99}.Release|Any CPU.Build.0 = Release|Any CPU
|
||||
{2430CB5F-7D07-4A9E-BD40-EC1B111B85FF}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
|
||||
{2430CB5F-7D07-4A9E-BD40-EC1B111B85FF}.Debug|Any CPU.Build.0 = Debug|Any CPU
|
||||
{2430CB5F-7D07-4A9E-BD40-EC1B111B85FF}.Release|Any CPU.ActiveCfg = Release|Any CPU
|
||||
{2430CB5F-7D07-4A9E-BD40-EC1B111B85FF}.Release|Any CPU.Build.0 = Release|Any CPU
|
||||
{85E13383-7C39-4719-AAC0-0B357C3A97C7}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
|
||||
{85E13383-7C39-4719-AAC0-0B357C3A97C7}.Debug|Any CPU.Build.0 = Debug|Any CPU
|
||||
{85E13383-7C39-4719-AAC0-0B357C3A97C7}.Release|Any CPU.ActiveCfg = Release|Any CPU
|
||||
|
|
@ -686,6 +692,7 @@ Global
|
|||
{3A377A6A-F735-4010-9E00-72E8BEB8F1F2} = {9B4F139F-7D26-435C-A561-89E65A67A8E5}
|
||||
{DAACE93A-866E-428E-B2A9-AFAC79F9C7A4} = {56C2FFB8-EA54-45B5-A095-4A78142EB4B5}
|
||||
{0A94D69A-61AC-4654-83AB-B24445628A99} = {B08B4E00-C2AB-48F3-8389-449F42AEF179}
|
||||
{2430CB5F-7D07-4A9E-BD40-EC1B111B85FF} = {08B41FFA-CEE3-46A7-B5C0-3EB65D37A16C}
|
||||
{85E13383-7C39-4719-AAC0-0B357C3A97C7} = {56C2FFB8-EA54-45B5-A095-4A78142EB4B5}
|
||||
EndGlobalSection
|
||||
EndGlobal
|
||||
|
|
|
|||
|
|
@ -17,5 +17,6 @@
|
|||
<PackageReference Include="Azure.ResourceManager.AppContainers" Version="1.0.3" />
|
||||
<PackageReference Include="Azure.ResourceManager.Resources" Version="1.4.0" />
|
||||
<PackageReference Include="Proto.Cluster" Version="1.1.0" />
|
||||
<PackageReference Include="StackExchange.Redis" Version="2.6.122" />
|
||||
</ItemGroup>
|
||||
</Project>
|
||||
|
|
|
|||
|
|
@ -19,9 +19,10 @@ public static class ServiceCollectionExtensions
|
|||
/// </summary>
|
||||
/// <param name="services">The service collection to add the provider to.</param>
|
||||
/// <param name="armClientProvider">An <see cref="IArmClientProvider"/> to create <see cref="ArmClient"/> instances.</param>
|
||||
/// <param name="configureMemberStore">An optional configuration for the member store.</param>
|
||||
/// <param name="configure">An optional action to configure the provider options.</param>
|
||||
/// <returns>The service collection.</returns>
|
||||
public static IServiceCollection AddAzureContainerAppsProvider(this IServiceCollection services, IArmClientProvider? armClientProvider = default, [AllowNull] Action<AzureContainerAppsProviderOptions> configure = null)
|
||||
public static IServiceCollection AddAzureContainerAppsProvider(this IServiceCollection services, IArmClientProvider? armClientProvider = default, Action<IServiceCollection>? configureMemberStore = null, Action<AzureContainerAppsProviderOptions>? configure = null)
|
||||
{
|
||||
var configureOptions = configure ?? (_ => { });
|
||||
services.Configure(configureOptions);
|
||||
|
|
@ -31,14 +32,10 @@ public static class ServiceCollectionExtensions
|
|||
if (armClientProvider != null)
|
||||
services.AddSingleton(armClientProvider);
|
||||
|
||||
// Register the default member store.
|
||||
services.AddSingleton<IClusterMemberStore, ResourceTagsClusterMemberStore>(sp =>
|
||||
{
|
||||
var clientProvider = sp.GetRequiredService<IArmClientProvider>();
|
||||
var logger = sp.GetRequiredService<ILogger<ResourceTagsClusterMemberStore>>();
|
||||
var options = sp.GetRequiredService<IOptions<AzureContainerAppsProviderOptions>>().Value;
|
||||
return new ResourceTagsClusterMemberStore(clientProvider, logger, options.ResourceGroupName, options.SubscriptionId);
|
||||
});
|
||||
if (configureMemberStore != null)
|
||||
configureMemberStore.Invoke(services);
|
||||
else
|
||||
services.AddResourceTagsMemberStore();
|
||||
|
||||
return services;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,6 @@
|
|||
namespace Proto.Cluster.AzureContainerApps.Stores.Redis;
|
||||
|
||||
/// <summary>
|
||||
/// A member with a cluster name and kind values.
|
||||
/// </summary>
|
||||
public record ClusterMember(string Host, int Port, string ClusterName, IEnumerable<string> Kinds);
|
||||
|
|
@ -0,0 +1,62 @@
|
|||
using System.Text.Json;
|
||||
using JetBrains.Annotations;
|
||||
using StackExchange.Redis;
|
||||
|
||||
namespace Proto.Cluster.AzureContainerApps.Stores.Redis;
|
||||
|
||||
/// <summary>
|
||||
/// Stores cluster member information in Redis database.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public class RedisClusterMemberStore : IClusterMemberStore
|
||||
{
|
||||
private const string ClusterKey = "proto:cluster:members";
|
||||
private readonly IDatabase _database;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="RedisClusterMemberStore"/> class.
|
||||
/// </summary>
|
||||
public RedisClusterMemberStore(ConnectionMultiplexer connectionMultiplexer) => _database = connectionMultiplexer.GetDatabase();
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask<ICollection<Member>> ListAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
var entries = await _database.HashGetAllAsync(ClusterKey);
|
||||
|
||||
return entries.Select(entry =>
|
||||
{
|
||||
var clusterMember = Deserialize(entry.Value!);
|
||||
return new Member
|
||||
{
|
||||
Id = entry.Name,
|
||||
Host = clusterMember.Host,
|
||||
Port = clusterMember.Port,
|
||||
Kinds = { clusterMember.Kinds }
|
||||
};
|
||||
}).ToList();
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask RegisterAsync(string clusterName, Member member, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var clusterMember = new ClusterMember(member.Host, member.Port, clusterName, member.Kinds);
|
||||
var serialized = Serialize(clusterMember);
|
||||
|
||||
await _database.HashSetAsync(ClusterKey, member.Id, serialized);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask UnregisterAsync(string memberId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await _database.HashDeleteAsync(ClusterKey, memberId);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async ValueTask ClearAsync(string clusterName, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await _database.KeyDeleteAsync(ClusterKey);
|
||||
}
|
||||
|
||||
private static string Serialize(ClusterMember member) => JsonSerializer.Serialize(member);
|
||||
private static ClusterMember Deserialize(string data) => JsonSerializer.Deserialize<ClusterMember>(data)!;
|
||||
}
|
||||
|
|
@ -0,0 +1,24 @@
|
|||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.DependencyInjection.Extensions;
|
||||
using StackExchange.Redis;
|
||||
|
||||
namespace Proto.Cluster.AzureContainerApps.Stores.Redis;
|
||||
|
||||
/// <summary>
|
||||
/// Adds extension methods to <see cref="IServiceCollection"/> for registering the Azure Container Apps provider
|
||||
/// </summary>
|
||||
public static class ServiceCollectionExtensions
|
||||
{
|
||||
/// <summary>
|
||||
/// Adds the <see cref="RedisClusterMemberStore"/> to the service collection.
|
||||
/// </summary>
|
||||
/// <param name="services">The service collection to add the provider to.</param>
|
||||
/// <param name="connectionString">Connection string for the Redis client.</param>
|
||||
/// <returns>The service collection.</returns>
|
||||
public static IServiceCollection AddRedisClusterMemberStore(this IServiceCollection services, string connectionString)
|
||||
{
|
||||
services.AddSingleton<ConnectionMultiplexer>(sp => ConnectionMultiplexer.Connect(connectionString));
|
||||
services.AddSingleton<IClusterMemberStore, RedisClusterMemberStore>();
|
||||
return services;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,6 +1,8 @@
|
|||
using System.Diagnostics.CodeAnalysis;
|
||||
using Microsoft.Extensions.DependencyInjection;
|
||||
using Microsoft.Extensions.DependencyInjection.Extensions;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Microsoft.Extensions.Options;
|
||||
|
||||
namespace Proto.Cluster.AzureContainerApps.Stores.ResourceTags;
|
||||
|
||||
|
|
@ -20,8 +22,16 @@ public static class ServiceCollectionExtensions
|
|||
var configureOptions = configure ?? (_ => { });
|
||||
services.Configure(configureOptions);
|
||||
services.ConfigureOptions<ResourceTagsMemberStoreOptionsValidator>();
|
||||
services.Replace(new ServiceDescriptor(typeof(IClusterMemberStore), typeof(ResourceTagsClusterMemberStore), ServiceLifetime.Singleton));
|
||||
|
||||
services.AddSingleton<IClusterMemberStore, ResourceTagsClusterMemberStore>(sp =>
|
||||
{
|
||||
var clientProvider = sp.GetRequiredService<IArmClientProvider>();
|
||||
var logger = sp.GetRequiredService<ILogger<ResourceTagsClusterMemberStore>>();
|
||||
var options = sp.GetRequiredService<IOptions<AzureContainerAppsProviderOptions>>().Value;
|
||||
return new ResourceTagsClusterMemberStore(clientProvider, logger, options.ResourceGroupName,
|
||||
options.SubscriptionId);
|
||||
});
|
||||
|
||||
return services;
|
||||
}
|
||||
}
|
||||
|
|
@ -7,6 +7,7 @@ using Elsa.ProtoActor.Protos;
|
|||
using Google.Protobuf.WellKnownTypes;
|
||||
using Microsoft.Data.Sqlite;
|
||||
using Proto.Cluster.AzureContainerApps;
|
||||
using Proto.Cluster.AzureContainerApps.Stores.Redis;
|
||||
using Proto.Persistence.Sqlite;
|
||||
using Proto.Remote;
|
||||
using Proto.Remote.GrpcNet;
|
||||
|
|
@ -15,13 +16,17 @@ var builder = WebApplication.CreateBuilder(args);
|
|||
var services = builder.Services;
|
||||
var configuration = builder.Configuration;
|
||||
var sqliteConnectionString = configuration.GetConnectionString("Sqlite")!;
|
||||
var redisConnectionString = configuration.GetConnectionString("Redis")!;
|
||||
var identitySection = configuration.GetSection("Identity");
|
||||
var identityTokenSection = identitySection.GetSection("Tokens");
|
||||
var protoActorSection = configuration.GetSection("ProtoActor");
|
||||
var protoActorClusterSection = protoActorSection.GetSection("Cluster");
|
||||
|
||||
// Configure Proto Actor cluster provider services.
|
||||
services.AddAzureContainerAppsProvider(ArmClientProviders.DefaultAzureCredential, options => protoActorClusterSection.GetSection("AzureContainerApps").Bind(options));
|
||||
services.AddAzureContainerAppsProvider(
|
||||
ArmClientProviders.DefaultAzureCredential,
|
||||
sc => sc.AddRedisClusterMemberStore(redisConnectionString),
|
||||
options => protoActorClusterSection.GetSection("AzureContainerApps").Bind(options));
|
||||
|
||||
// Add Elsa services.
|
||||
services
|
||||
|
|
|
|||
|
|
@ -7,7 +7,8 @@
|
|||
},
|
||||
"AllowedHosts": "*",
|
||||
"ConnectionStrings": {
|
||||
"Sqlite": "Data Source=elsa.sqlite.db;Cache=Shared;"
|
||||
"Sqlite": "Data Source=elsa.sqlite.db;Cache=Shared;",
|
||||
"Redis": "localhost:6379"
|
||||
},
|
||||
"Identity": {
|
||||
"Tokens": {
|
||||
|
|
|
|||
|
|
@ -0,0 +1,63 @@
|
|||
using Proto.Cluster;
|
||||
using Proto.Cluster.AzureContainerApps.Stores.Redis;
|
||||
using StackExchange.Redis;
|
||||
|
||||
namespace Elsa.ProtoCluster.ComponentTests.AzureContainerAppTests;
|
||||
|
||||
public class RedisClusterMemberStoreTests : IClassFixture<RedisFixture>
|
||||
{
|
||||
private readonly RedisClusterMemberStore _memberStore;
|
||||
private readonly Member _member;
|
||||
|
||||
public RedisClusterMemberStoreTests(RedisFixture redisFixture)
|
||||
{
|
||||
var multiplexer = ConnectionMultiplexer.Connect(redisFixture.GetConnectionString());
|
||||
_memberStore = new RedisClusterMemberStore(multiplexer);
|
||||
|
||||
_member = new Member
|
||||
{
|
||||
Id = "id1",
|
||||
Host = "localhost",
|
||||
Port = 8000,
|
||||
Kinds = { "kind1", "kind2" }
|
||||
};
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Invoking RegisterAsync should register member")]
|
||||
public async Task RegisterAsync_ShouldRegisterMember()
|
||||
{
|
||||
await _memberStore.RegisterAsync("cluster", _member);
|
||||
|
||||
var members = await _memberStore.ListAsync();
|
||||
Assert.Contains(members, m => m.Id == _member.Id);
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Invoking UnregisterAsync should unregister member")]
|
||||
public async Task UnregisterAsync_ShouldUnregisterMember()
|
||||
{
|
||||
await _memberStore.RegisterAsync("cluster", _member);
|
||||
await _memberStore.UnregisterAsync(_member.Id);
|
||||
|
||||
var members = await _memberStore.ListAsync();
|
||||
Assert.DoesNotContain(members, m => m.Id == _member.Id);
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Invoking ClearAsync should clear all members")]
|
||||
public async Task ClearAsync_ShouldClearAllMembers()
|
||||
{
|
||||
await _memberStore.RegisterAsync("cluster", _member);
|
||||
await _memberStore.ClearAsync("cluster");
|
||||
|
||||
var members = await _memberStore.ListAsync();
|
||||
Assert.Empty(members);
|
||||
}
|
||||
|
||||
[Fact(DisplayName = "Invoking ListAsync should return all registered members")]
|
||||
public async Task ListAsync_ShouldReturnAllRegisteredMembers()
|
||||
{
|
||||
await _memberStore.RegisterAsync("cluster", _member);
|
||||
|
||||
var members = await _memberStore.ListAsync();
|
||||
Assert.Contains(members, m => m.Id == _member.Id);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,31 @@
|
|||
using DotNet.Testcontainers.Builders;
|
||||
using Testcontainers.Redis;
|
||||
|
||||
namespace Elsa.ProtoCluster.ComponentTests.AzureContainerAppTests;
|
||||
|
||||
public class RedisFixture : IAsyncLifetime
|
||||
{
|
||||
private RedisContainer? _redisContainer;
|
||||
|
||||
public async Task InitializeAsync()
|
||||
{
|
||||
_redisContainer = new RedisBuilder()
|
||||
.WithPortBinding(6379, true)
|
||||
.Build();
|
||||
|
||||
await _redisContainer.StartAsync().ConfigureAwait(false);
|
||||
}
|
||||
|
||||
public string GetConnectionString()
|
||||
{
|
||||
return _redisContainer != null ?
|
||||
_redisContainer.GetConnectionString() :
|
||||
throw new InvalidOperationException("Redis container not initialized.");
|
||||
}
|
||||
|
||||
public async Task DisposeAsync()
|
||||
{
|
||||
if(_redisContainer != null)
|
||||
await _redisContainer.StopAsync();
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,31 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>net7.0</TargetFramework>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<Nullable>enable</Nullable>
|
||||
|
||||
<IsPackable>false</IsPackable>
|
||||
<IsTestProject>true</IsTestProject>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Microsoft.NET.Test.Sdk" Version="17.5.0" />
|
||||
<PackageReference Include="Testcontainers" Version="3.3.0" />
|
||||
<PackageReference Include="Testcontainers.Redis" Version="3.3.0" />
|
||||
<PackageReference Include="xunit" Version="2.4.2" />
|
||||
<PackageReference Include="xunit.runner.visualstudio" Version="2.4.5">
|
||||
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
|
||||
<PrivateAssets>all</PrivateAssets>
|
||||
</PackageReference>
|
||||
<PackageReference Include="coverlet.collector" Version="3.2.0">
|
||||
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
|
||||
<PrivateAssets>all</PrivateAssets>
|
||||
</PackageReference>
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\src\modules\Elsa.ProtoActor.Cluster.AzureContainerApps\Elsa.ProtoActor.Cluster.AzureContainerApps.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
1
test/Elsa.ProtoCluster.ComponentTests/Usings.cs
Normal file
1
test/Elsa.ProtoCluster.ComponentTests/Usings.cs
Normal file
|
|
@ -0,0 +1 @@
|
|||
global using Xunit;
|
||||
Loading…
Reference in a new issue