Implement Redis distributed cache signal

This commit is contained in:
Sipke Schoorstra 2021-05-22 13:58:21 +02:00
parent b5453b23a2
commit 2ec654d521
22 changed files with 256 additions and 59 deletions

View file

@ -241,8 +241,6 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "locking", "locking", "{DBBD
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.DistributedLocking.SqlServer", "src\locking\Elsa.DistributedLocking.SqlServer\Elsa.DistributedLocking.SqlServer.csproj", "{B09A6E42-EF42-4693-BA0E-5EE1D7F81FFD}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.DistributedLocking.Redis", "src\locking\Elsa.DistributedLocking.Redis\Elsa.DistributedLocking.Redis.csproj", "{7E5514AC-3633-4237-9490-0992417F2273}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.DistributedLocking.AzureBlob", "src\locking\Elsa.DistributedLocking.AzureBlob\Elsa.DistributedLocking.AzureBlob.csproj", "{495BE954-E6EA-41A9-8054-A5CA5DD0924E}"
EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "blazor", "blazor", "{D86B94DC-A53C-4A67-A820-828DD359C49B}"
@ -293,6 +291,10 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "caching", "caching", "{9937
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Caching.Rebus", "src\caching\Elsa.Caching.Rebus\Elsa.Caching.Rebus.csproj", "{61C16CA0-B190-4642-A81A-5C03705CA2C1}"
EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "providers", "providers", "{EC5D4AFD-3F7F-4B51-9C38-F18C60618731}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Providers.Redis", "src\providers\Elsa.Providers.Redis\Elsa.Providers.Redis.csproj", "{9ACA06DA-9AFE-41A2-8109-EFED8B9B14A1}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -631,10 +633,6 @@ Global
{B09A6E42-EF42-4693-BA0E-5EE1D7F81FFD}.Debug|Any CPU.Build.0 = Debug|Any CPU
{B09A6E42-EF42-4693-BA0E-5EE1D7F81FFD}.Release|Any CPU.ActiveCfg = Release|Any CPU
{B09A6E42-EF42-4693-BA0E-5EE1D7F81FFD}.Release|Any CPU.Build.0 = Release|Any CPU
{7E5514AC-3633-4237-9490-0992417F2273}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{7E5514AC-3633-4237-9490-0992417F2273}.Debug|Any CPU.Build.0 = Debug|Any CPU
{7E5514AC-3633-4237-9490-0992417F2273}.Release|Any CPU.ActiveCfg = Release|Any CPU
{7E5514AC-3633-4237-9490-0992417F2273}.Release|Any CPU.Build.0 = Release|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Debug|Any CPU.Build.0 = Debug|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Release|Any CPU.ActiveCfg = Release|Any CPU
@ -687,6 +685,10 @@ Global
{61C16CA0-B190-4642-A81A-5C03705CA2C1}.Debug|Any CPU.Build.0 = Debug|Any CPU
{61C16CA0-B190-4642-A81A-5C03705CA2C1}.Release|Any CPU.ActiveCfg = Release|Any CPU
{61C16CA0-B190-4642-A81A-5C03705CA2C1}.Release|Any CPU.Build.0 = Release|Any CPU
{9ACA06DA-9AFE-41A2-8109-EFED8B9B14A1}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{9ACA06DA-9AFE-41A2-8109-EFED8B9B14A1}.Debug|Any CPU.Build.0 = Debug|Any CPU
{9ACA06DA-9AFE-41A2-8109-EFED8B9B14A1}.Release|Any CPU.ActiveCfg = Release|Any CPU
{9ACA06DA-9AFE-41A2-8109-EFED8B9B14A1}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -802,7 +804,6 @@ Global
{A80FC28D-D865-428D-AF12-7A393D407803} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180}
{DBBD242E-4437-4BDB-919F-A70839BE75FA} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925}
{B09A6E42-EF42-4693-BA0E-5EE1D7F81FFD} = {DBBD242E-4437-4BDB-919F-A70839BE75FA}
{7E5514AC-3633-4237-9490-0992417F2273} = {DBBD242E-4437-4BDB-919F-A70839BE75FA}
{495BE954-E6EA-41A9-8054-A5CA5DD0924E} = {DBBD242E-4437-4BDB-919F-A70839BE75FA}
{D86B94DC-A53C-4A67-A820-828DD359C49B} = {4673732F-2853-47BD-91B8-C95C229D2C89}
{C869CC72-9A98-4246-9A76-6A50F48AFD85} = {4673732F-2853-47BD-91B8-C95C229D2C89}
@ -820,6 +821,8 @@ Global
{3E2423CF-50E6-4D2B-8749-17B1EF540FE4} = {22E75696-6FE9-436A-9097-EE21C603F818}
{9937FB02-72D9-4FCD-B31E-8E60CFF1B37F} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925}
{61C16CA0-B190-4642-A81A-5C03705CA2C1} = {9937FB02-72D9-4FCD-B31E-8E60CFF1B37F}
{EC5D4AFD-3F7F-4B51-9C38-F18C60618731} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925}
{9ACA06DA-9AFE-41A2-8109-EFED8B9B14A1} = {EC5D4AFD-3F7F-4B51-9C38-F18C60618731}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -1,7 +1,15 @@
<Project Sdk="Microsoft.NET.Sdk">
<Import Project="..\..\..\common.props" />
<Import Project="..\..\..\configureawait.props" />
<PropertyGroup>
<TargetFramework>net5.0</TargetFramework>
<TargetFramework>netstandard2.1</TargetFramework>
<Description>
Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application.
This package contains a distributed cache signal decorator using Rebus to trigger signals.
</Description>
<PackageTags>elsa, workflows, rebus, cache</PackageTags>
</PropertyGroup>
<ItemGroup>

View file

@ -7,7 +7,7 @@ namespace Elsa.Caching.Rebus.Extensions
{
public static class ServiceCollectionExtensions
{
public static ElsaOptionsBuilder AddRebusCacheSignal(this ElsaOptionsBuilder builder)
public static ElsaOptionsBuilder UseRebusCacheSignal(this ElsaOptionsBuilder builder)
{
var services = builder.Services;
services.Decorate<ICacheSignal, RebusCacheSignal>();

View file

@ -0,0 +1,3 @@
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd">
<ConfigureAwait />
</Weavers>

View file

@ -0,0 +1,17 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
using System.ComponentModel;
// ReSharper disable once CheckNamespace
namespace System.Runtime.CompilerServices
{
/// <summary>
/// Reserved to be used by the compiler for tracking metadata.
/// This class should not be used by developers in source code.
/// </summary>
[EditorBrowsable(EditorBrowsableState.Never)]
internal static class IsExternalInit
{
}
}

View file

@ -1,26 +0,0 @@
<Project Sdk="Microsoft.NET.Sdk">
<Import Project="..\..\..\common.props" />
<Import Project="..\..\..\configureawait.props" />
<PropertyGroup>
<TargetFramework>netstandard2.1</TargetFramework>
<RootNamespace>Elsa</RootNamespace>
<Description>
Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application.
This package provides a distributed locking provider using Redis.
</Description>
<PackageTags>elsa, workflows, distributed lock, redis</PackageTags>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="DistributedLock.Redis" Version="1.0.1" />
<PackageReference Include="RedLock.net" Version="2.2.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa.Abstractions\Elsa.Abstractions.csproj" />
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
</ItemGroup>
</Project>

View file

@ -1,3 +0,0 @@
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd">
<ConfigureAwait ContinueOnCapturedContext="false" />
</Weavers>

View file

@ -0,0 +1,25 @@
<Project Sdk="Microsoft.NET.Sdk">
<Import Project="..\..\..\common.props" />
<Import Project="..\..\..\configureawait.props" />
<PropertyGroup>
<TargetFramework>netstandard2.1</TargetFramework>
<RootNamespace>Elsa</RootNamespace>
<Description>
Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application.
This package provides a Redis implementation of a distributed lock provider and distributed cache signaler.
</Description>
<PackageTags>elsa, workflows, distributed lock, cache, redis</PackageTags>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="DistributedLock.Redis" Version="1.0.1" />
<PackageReference Include="RedLock.net" Version="2.2.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,23 @@
using Elsa.Caching;
using Elsa.Runtime;
using Elsa.Services;
using Elsa.StartupTasks;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Extensions
{
public static class RedisCacheSignalElsaOptionsBuilderExtensions
{
public static ElsaOptionsBuilder UseRedisCacheSignal(this ElsaOptionsBuilder builder)
{
var services = builder.Services;
services
.AddSingleton<RedisBus>()
.AddStartupTask<SubscribeToRedisCacheSignals>()
.Decorate<ICacheSignal, RedisCacheSignal>();
return builder;
}
}
}

View file

@ -7,17 +7,13 @@ using RedLockNet.SERedis;
using RedLockNet.SERedis.Configuration;
using StackExchange.Redis;
namespace Elsa
namespace Elsa.Extensions
{
public static class DistributedLockingOptionsExtensions
public static class RedisDistributedLockingOptionsExtensions
{
public static DistributedLockingOptionsBuilder UseRedisLockProvider(this DistributedLockingOptionsBuilder options, string connectionString)
public static DistributedLockingOptionsBuilder UseRedisLockProvider(this DistributedLockingOptionsBuilder options)
{
options
.Services
.UseStackExchangeConnectionMultiplexer(connectionString)
.UseRedLockFactory();
options.Services.AddRedLockFactory();
options.UseProviderFactory(CreateRedisDistributedLockFactory);
return options;
@ -29,9 +25,7 @@ namespace Elsa
return name => new RedisDistributedLock(name, multiplexer.GetDatabase());
}
private static IServiceCollection UseStackExchangeConnectionMultiplexer(this IServiceCollection services, string connectionString) => services.AddSingleton<IConnectionMultiplexer>(ConnectionMultiplexer.Connect(connectionString));
private static IServiceCollection UseRedLockFactory(this IServiceCollection services) =>
private static IServiceCollection AddRedLockFactory(this IServiceCollection services) =>
services.AddSingleton<IDistributedLockFactory, RedLockFactory>(sp => RedLockFactory.Create(
new[]
{

View file

@ -0,0 +1,13 @@
using Microsoft.Extensions.DependencyInjection;
using StackExchange.Redis;
namespace Elsa.Extensions
{
public static class RedisServiceCollectionExtensions
{
public static IServiceCollection AddRedis(this IServiceCollection services, string connectionString)
{
return services.AddSingleton<IConnectionMultiplexer>(ConnectionMultiplexer.Connect(connectionString));
}
}
}

View file

@ -0,0 +1,3 @@
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd">
<ConfigureAwait />
</Weavers>

View file

@ -0,0 +1,17 @@
// Licensed to the .NET Foundation under one or more agreements.
// The .NET Foundation licenses this file to you under the MIT license.
using System.ComponentModel;
// ReSharper disable once CheckNamespace
namespace System.Runtime.CompilerServices
{
/// <summary>
/// Reserved to be used by the compiler for tracking metadata.
/// This class should not be used by developers in source code.
/// </summary>
[EditorBrowsable(EditorBrowsableState.Never)]
internal static class IsExternalInit
{
}
}

View file

@ -0,0 +1,27 @@
using System.Threading.Tasks;
using Elsa.Caching;
using Microsoft.Extensions.Primitives;
namespace Elsa.Services
{
public class RedisCacheSignal : ICacheSignal
{
private readonly ICacheSignal _cacheSignal;
private readonly RedisBus _redisBus;
public RedisCacheSignal(ICacheSignal cacheSignal, RedisBus redisBus)
{
_cacheSignal = cacheSignal;
_redisBus = redisBus;
}
public IChangeToken GetToken(string key) => _cacheSignal.GetToken(key);
public void TriggerToken(string key) => _cacheSignal.TriggerToken(key);
public async ValueTask TriggerTokenAsync(string key)
{
await _cacheSignal.TriggerTokenAsync(key);
await _redisBus.PublishAsync(nameof(CacheSignal), key);
}
}
}

View file

@ -0,0 +1,67 @@
using System;
using System.Diagnostics;
using System.Linq;
using System.Net;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;
namespace Elsa.Services
{
public class RedisBus
{
/// <summary>
/// TODO: Design multi-tenancy.
/// </summary>
private const string? TenantId = default;
private readonly string _hostName;
private readonly string _channelPrefix;
private readonly string _messagePrefix;
private readonly IConnectionMultiplexer _multiplexer;
private readonly ILogger _logger;
public RedisBus(IConnectionMultiplexer multiplexer, IContainerNameAccessor containerNameAccessor, ILogger<RedisBus> logger)
{
_multiplexer = multiplexer;
_hostName = Dns.GetHostName() + ':' + Process.GetCurrentProcess().Id;
_channelPrefix = (TenantId ?? "Default") + ':';
_messagePrefix = _hostName + '/';
_logger = logger;
}
public async Task SubscribeAsync(string channel, Action<string, string> handler)
{
try
{
var subscriber = _multiplexer.GetSubscriber();
await subscriber.SubscribeAsync(_channelPrefix + channel, (redisChannel, redisValue) =>
{
var tokens = redisValue.ToString().Split('/').ToArray();
if (tokens.Length != 2 || tokens[0].Length == 0 || tokens[0].Equals(_hostName, StringComparison.OrdinalIgnoreCase))
return;
handler(channel, tokens[1]);
});
}
catch (Exception e)
{
_logger.LogError(e, "Unable to subscribe to channel {ChannelName}", _channelPrefix + channel);
}
}
public async Task PublishAsync(string channel, string message)
{
try
{
await _multiplexer.GetSubscriber().PublishAsync(_channelPrefix + channel, _messagePrefix + message);
}
catch (Exception e)
{
_logger.LogError(e, "Unable to publish to channel {ChannelName}", _channelPrefix + channel);
}
}
}
}

View file

@ -0,0 +1,22 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Caching;
using Elsa.Services;
namespace Elsa.StartupTasks
{
public class SubscribeToRedisCacheSignals : IStartupTask
{
private readonly RedisBus _redisBus;
private readonly ICacheSignal _cacheSignal;
public SubscribeToRedisCacheSignals(RedisBus redisBus, ICacheSignal cacheSignal)
{
_redisBus = redisBus;
_cacheSignal = cacheSignal;
}
public int Order => 0;
public async Task ExecuteAsync(CancellationToken cancellationToken = default) => await _redisBus.SubscribeAsync(nameof(CacheSignal), (channel, message) => _cacheSignal.TriggerToken(message));
}
}

View file

@ -18,6 +18,7 @@
<ProjectReference Include="..\..\..\persistence\Elsa.Persistence.EntityFramework\Elsa.Persistence.EntityFramework.Sqlite\Elsa.Persistence.EntityFramework.Sqlite.csproj" />
<ProjectReference Include="..\..\..\persistence\Elsa.Persistence.EntityFramework\Elsa.Persistence.EntityFramework.SqlServer\Elsa.Persistence.EntityFramework.SqlServer.csproj" />
<ProjectReference Include="..\..\..\persistence\Elsa.Persistence.YesSql\Elsa.Persistence.YesSql.csproj" />
<ProjectReference Include="..\..\..\providers\Elsa.Providers.Redis\Elsa.Providers.Redis.csproj" />
<ProjectReference Include="..\..\..\server\Elsa.Server.Api\Elsa.Server.Api.csproj" />
<ProjectReference Include="..\..\..\server\Elsa.Server.Hangfire\Elsa.Server.Hangfire.csproj" />
<ProjectReference Include="..\..\..\server\Elsa.Server.Orleans\Elsa.Server.Orleans.csproj" />

View file

@ -1,6 +1,7 @@
using System;
using Elsa.Activities.UserTask.Extensions;
using Elsa.Caching.Rebus.Extensions;
using Elsa.Extensions;
using Elsa.Persistence.EntityFramework.Core.Extensions;
using Elsa.Persistence.EntityFramework.PostgreSql;
using Elsa.Persistence.EntityFramework.Sqlite;
@ -37,11 +38,13 @@ namespace Elsa.Samples.Server.Host
services
.AddActivityPropertyOptionsProvider<VehicleActivity>()
.AddRuntimeSelectItemsProvider<VehicleActivity>()
.AddRedis(Configuration.GetConnectionString("Redis"))
.AddElsa(elsa => elsa
.WithContainerName(Configuration["ContainerName"] ?? System.Environment.MachineName)
//.WithContainerName(Configuration["ContainerName"] ?? System.Environment.MachineName)
.UseEntityFrameworkPersistence(ef => ef.UseSqlite())
.UseRabbitMq(Configuration.GetConnectionString("RabbitMq"))
.AddRebusCacheSignal()
//.UseRabbitMq(Configuration.GetConnectionString("RabbitMq"))
//.UseRebusCacheSignal()
//.UseRedisCacheSignal()
.AddConsoleActivities()
.AddHttpActivities(elsaSection.GetSection("Http").Bind)
.AddEmailActivities(elsaSection.GetSection("Smtp").Bind)

View file

@ -9,7 +9,8 @@
},
"AllowedHosts": "*",
"ConnectionStrings": {
"RabbitMq": "amqp://localhost:5672"
"RabbitMq": "amqp://localhost:5672",
"Redis": "localhost:6379,abortConnect=false"
},
"Elsa": {
"Http": {

View file

@ -16,13 +16,10 @@
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj" />
<ProjectReference Include="..\..\..\locking\Elsa.DistributedLocking.AzureBlob\Elsa.DistributedLocking.AzureBlob.csproj" />
<ProjectReference Include="..\..\..\activities\Elsa.Activities.Temporal.Common\Elsa.Activities.Temporal.Common.csproj" />
<ProjectReference Include="..\..\..\locking\Elsa.DistributedLocking.Redis\Elsa.DistributedLocking.Redis.csproj" />
<ProjectReference Include="..\..\..\locking\Elsa.DistributedLocking.SqlServer\Elsa.DistributedLocking.SqlServer.csproj" />
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj" />
<ProjectReference Include="..\..\..\activities\Elsa.Activities.Temporal.Common\Elsa.Activities.Temporal.Common.csproj" />
<ProjectReference Include="..\..\..\locking\Elsa.DistributedLocking.AzureBlob\Elsa.DistributedLocking.AzureBlob.csproj" />
<ProjectReference Include="..\..\..\locking\Elsa.DistributedLocking.Redis\Elsa.DistributedLocking.Redis.csproj" />
<ProjectReference Include="..\..\..\locking\Elsa.DistributedLocking.SqlServer\Elsa.DistributedLocking.SqlServer.csproj" />
<ProjectReference Include="..\..\..\providers\Elsa.Providers.Redis\Elsa.Providers.Redis.csproj" />
</ItemGroup>
</Project>

View file

@ -1,6 +1,7 @@
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using System.Threading.Tasks;
using Elsa.Extensions;
namespace Elsa.Samples.DistributedLock
{
@ -18,7 +19,8 @@ namespace Elsa.Samples.DistributedLock
(_, services) =>
{
services
.AddElsa(elsa => elsa.ConfigureDistributedLockProvider(options => options.UseRedisLockProvider("localhost:6379,abortConnect=false"))
.AddRedis("localhost:6379,abortConnect=false")
.AddElsa(elsa => elsa.ConfigureDistributedLockProvider(options => options.UseRedisLockProvider())
.AddConsoleActivities()
.AddQuartzTemporalActivities()
.AddWorkflow<RecurringWorkflow>());

View file

@ -3,7 +3,7 @@ using Rebus.Config;
namespace Elsa.Rebus.RabbitMq.Extensions
{
public static class ElsaOptionsExtensions
public static class ElsaOptionsBuilderExtensions
{
public static ElsaOptionsBuilder UseRabbitMq(this ElsaOptionsBuilder elsaOptions, string connectionString) => elsaOptions.UseServiceBus(context => ConfigureRabbitMqEndpoint(context, connectionString));
private static void ConfigureRabbitMqEndpoint(ServiceBusEndpointConfigurationContext context, string connectionString) => context.Configurer.Transport(t => t.UseRabbitMq(connectionString, context.QueueName));