Add Azure Service Bus for MassTransit (#4624)

* Split RabbitMQ from MassTransit and add Azure Service Bus package

* Use Send topology for commands
This commit is contained in:
Sipke Schoorstra 2023-11-15 20:47:11 +01:00 committed by GitHub
parent 3a4b15be13
commit 21e438bdd1
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
22 changed files with 248 additions and 57 deletions

View file

@ -285,6 +285,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.CSharp", "src\modules\
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Python", "src\modules\Elsa.Python\Elsa.Python.csproj", "{790E94F2-5393-47DF-AC52-D9247F5B243A}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.MassTransit.RabbitMq", "src\modules\Elsa.MassTransit.RabbitMq\Elsa.MassTransit.RabbitMq.csproj", "{169BEA3D-2A81-47EE-A6C1-3F8719EEC1F6}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.MassTransit.AzureServiceBus", "src\modules\Elsa.MassTransit.AzureServiceBus\Elsa.MassTransit.AzureServiceBus.csproj", "{5B1AF00E-1030-4BBA-9C1A-1DC5AEFDD44A}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -715,6 +719,14 @@ Global
{790E94F2-5393-47DF-AC52-D9247F5B243A}.Debug|Any CPU.Build.0 = Debug|Any CPU
{790E94F2-5393-47DF-AC52-D9247F5B243A}.Release|Any CPU.ActiveCfg = Release|Any CPU
{790E94F2-5393-47DF-AC52-D9247F5B243A}.Release|Any CPU.Build.0 = Release|Any CPU
{169BEA3D-2A81-47EE-A6C1-3F8719EEC1F6}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{169BEA3D-2A81-47EE-A6C1-3F8719EEC1F6}.Debug|Any CPU.Build.0 = Debug|Any CPU
{169BEA3D-2A81-47EE-A6C1-3F8719EEC1F6}.Release|Any CPU.ActiveCfg = Release|Any CPU
{169BEA3D-2A81-47EE-A6C1-3F8719EEC1F6}.Release|Any CPU.Build.0 = Release|Any CPU
{5B1AF00E-1030-4BBA-9C1A-1DC5AEFDD44A}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{5B1AF00E-1030-4BBA-9C1A-1DC5AEFDD44A}.Debug|Any CPU.Build.0 = Debug|Any CPU
{5B1AF00E-1030-4BBA-9C1A-1DC5AEFDD44A}.Release|Any CPU.ActiveCfg = Release|Any CPU
{5B1AF00E-1030-4BBA-9C1A-1DC5AEFDD44A}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -846,6 +858,8 @@ Global
{732BF088-6AD7-4C4D-9A48-8074253596D4} = {B818988E-639C-4E6E-85C1-B231BCAD9DAB}
{24331E82-D7AF-45B1-ACF0-CA6C3B0B77DC} = {6EF07978-A6D2-40EB-891D-7D70C5F37E76}
{790E94F2-5393-47DF-AC52-D9247F5B243A} = {6EF07978-A6D2-40EB-891D-7D70C5F37E76}
{169BEA3D-2A81-47EE-A6C1-3F8719EEC1F6} = {DD089B8B-DA73-492A-9010-F772D1C178DA}
{5B1AF00E-1030-4BBA-9C1A-1DC5AEFDD44A} = {DD089B8B-DA73-492A-9010-F772D1C178DA}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {D4B5CEAA-7D70-4FCB-A68E-B03FBE5E0E5E}

View file

@ -15,6 +15,8 @@
<ProjectReference Include="..\..\modules\Elsa.CSharp\Elsa.CSharp.csproj" />
<ProjectReference Include="..\..\modules\Elsa.EntityFrameworkCore.Sqlite\Elsa.EntityFrameworkCore.Sqlite.csproj" />
<ProjectReference Include="..\..\modules\Elsa.EntityFrameworkCore.SqlServer\Elsa.EntityFrameworkCore.SqlServer.csproj" />
<ProjectReference Include="..\..\modules\Elsa.MassTransit.AzureServiceBus\Elsa.MassTransit.AzureServiceBus.csproj" />
<ProjectReference Include="..\..\modules\Elsa.MassTransit.RabbitMq\Elsa.MassTransit.RabbitMq.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Python\Elsa.Python.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Quartz.EntityFrameworkCore.Sqlite\Elsa.Quartz.EntityFrameworkCore.Sqlite.csproj" />
<ProjectReference Include="..\..\modules\Elsa.FileStorage\Elsa.FileStorage.csproj" />
@ -41,7 +43,7 @@
</ItemGroup>
<ItemGroup>
<PackageReference Include="Azure.Identity" Version="1.8.2"/>
<PackageReference Include="Azure.Identity" Version="1.10.4"/>
<PackageReference Include="FluentStorage.Azure.Blobs" Version="5.2.2" />
<PackageReference Include="Proto.Persistence.Sqlite" Version="1.4.0"/>
<PackageReference Include="Proto.Persistence.SqlServer" Version="1.4.0"/>

View file

@ -28,6 +28,8 @@ const bool useDapper = false;
const bool useProtoActor = false;
const bool useHangfire = false;
const bool useQuartz = true;
const bool useMassTransitAzureServiceBus = true;
const bool useMassTransitRabbitMq = false;
var builder = WebApplication.CreateBuilder(args);
var services = builder.Services;
@ -37,6 +39,8 @@ var identityTokenSection = identitySection.GetSection("Tokens");
var sqliteConnectionString = configuration.GetConnectionString("Sqlite")!;
var sqlServerConnectionString = configuration.GetConnectionString("SqlServer")!;
var mongoDbConnectionString = configuration.GetConnectionString("MongoDb")!;
var azureServiceBusConnectionString = configuration.GetConnectionString("AzureServiceBus")!;
var rabbitMqConnectionString = configuration.GetConnectionString("RabbitMq")!;
// Add Elsa services.
services
@ -186,14 +190,18 @@ services
if (useQuartz)
{
elsa.UseQuartz(quartz =>
{
quartz.UseSqlite(sqliteConnectionString);
});
elsa.UseQuartz(quartz => { quartz.UseSqlite(sqliteConnectionString); });
}
elsa.InstallDropIns(options => options.DropInRootDirectory = Path.Combine(Directory.GetCurrentDirectory(), "App_Data", "DropIns"));
elsa.UseMassTransit(massTransit =>
{
if (useMassTransitAzureServiceBus)
massTransit.UseAzureServiceBus(azureServiceBusConnectionString);
else if (useMassTransitRabbitMq)
massTransit.UseRabbitMq(rabbitMqConnectionString);
});
elsa.InstallDropIns(options => options.DropInRootDirectory = Path.Combine(Directory.GetCurrentDirectory(), "App_Data", "DropIns"));
elsa.AddSwagger();
});

View file

@ -1,21 +1,24 @@
{
"Logging": {
"LogLevel": {
"Default": "Warning",
"Default": "Debug",
"Elsa.Mediator": "Warning",
"Elsa.Workflows.Runtime.HostedServices": "Information",
"MassTransit": "Warning",
"MassTransit": "Debug",
"Microsoft.Extensions.Http": "Warning",
"Microsoft.Hosting.Lifetime": "Information",
"Microsoft.EntityFrameworkCore": "Warning",
"Microsoft.AspNetCore": "Warning",
"Quartz": "Warning",
"System.Net.Http": "Warning"
}
},
"AllowedHosts": "*",
"ConnectionStrings": {
"Sqlite": "Data Source=App_Data/elsa.sqlite.db;Cache=Shared;",
"MongoDb": "mongodb://localhost:27017/elsa-workflows"
"MongoDb": "mongodb://localhost:27017/elsa-workflows",
"AzureServiceBus": "",
"RabbitMq": ""
},
"Smtp": {
"Host": "localhost",

View file

@ -16,7 +16,4 @@
<ProjectReference Include="..\Elsa.MassTransit\Elsa.MassTransit.csproj" />
</ItemGroup>
</Project>

View file

@ -1,11 +1,13 @@
using Elsa.Alterations.Features;
using Elsa.Alterations.MassTransit.Consumers;
using Elsa.Alterations.MassTransit.Messages;
using Elsa.Alterations.MassTransit.Services;
using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.MassTransit.Features;
using MassTransit;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Alterations.MassTransit.Features;
@ -32,6 +34,9 @@ public class MassTransitAlterationsFeature : FeatureBase
/// <inheritdoc />
public override void Apply()
{
var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer<RunAlterationJobConsumer>();
var queueAddress = new Uri($"queue:elsa-{queueName}");
EndpointConvention.Map<RunAlterationJob>(queueAddress);
Services.AddSingleton<MassTransitAlterationJobDispatcher>();
}
}

View file

@ -23,6 +23,6 @@ public class MassTransitAlterationJobDispatcher : IAlterationJobDispatcher
public async ValueTask DispatchAsync(string jobId, CancellationToken cancellationToken = default)
{
var message = new RunAlterationJob(jobId);
await _bus.Publish(message, cancellationToken);
await _bus.Send(message, cancellationToken);
}
}

View file

@ -0,0 +1,22 @@
<Project Sdk="Microsoft.NET.Sdk">
<Import Project="..\..\..\common.props" />
<Import Project="..\..\..\configureawait.props" />
<PropertyGroup>
<TargetFrameworks>net6.0;net7.0</TargetFrameworks>
<Description>
Provides integration of Azure Service Bus with MassTransit's integration of Elsa.
</Description>
<PackageTags>elsa module service-bus masstransit azure-service-bus</PackageTags>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.MassTransit\Elsa.MassTransit.csproj" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="MassTransit.Azure.ServiceBus.Core" Version="8.1.2" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,30 @@
using Elsa.Features.Services;
using Elsa.MassTransit.AzureServiceBus.Features;
using Elsa.MassTransit.AzureServiceBus.Options;
using Elsa.MassTransit.Features;
using Elsa.MassTransit.Options;
using JetBrains.Annotations;
// ReSharper disable once CheckNamespace
namespace Elsa.Extensions;
/// <summary>
/// Provides extensions to <see cref="IModule"/> that enables and configures MassTransit and the Azure Service Bus transport.
/// </summary>
[PublicAPI]
public static class ModuleExtensions
{
/// <summary>
/// Enable and configure the Azure Service Bus transport for MassTransit.
/// </summary>
public static MassTransitFeature UseAzureServiceBus(this MassTransitFeature feature, string? connectionString)
{
feature.Module.Configure((Action<AzureServiceBusFeature>)Configure);
return feature;
void Configure(AzureServiceBusFeature bus)
{
bus.ConnectionString = connectionString;
}
}
}

View file

@ -0,0 +1,50 @@
using Azure.Messaging.ServiceBus.Administration;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.MassTransit.Features;
using Elsa.MassTransit.Messages;
using Elsa.MassTransit.Options;
using Elsa.Workflows.Runtime.Activities;
using MassTransit;
using MassTransit.Configuration;
namespace Elsa.MassTransit.AzureServiceBus.Features;
/// <summary>
/// Configures MassTransit to use the Azure Service Bus transport.
/// </summary>
[DependsOn(typeof(MassTransitFeature))]
public class AzureServiceBusFeature : FeatureBase
{
/// <inheritdoc />
public AzureServiceBusFeature(IModule module) : base(module)
{
}
/// <summary>
/// An Azure Service Bus connection string.
/// </summary>
public string? ConnectionString { get; set; }
/// <inheritdoc />
public override void Configure()
{
Module.Configure<MassTransitFeature>(massTransitFeature =>
{
massTransitFeature.BusConfigurator = configure =>
{
configure.AddServiceBusMessageScheduler();
configure.UsingAzureServiceBus((context, serviceBus) =>
{
if (ConnectionString != null)
serviceBus.Host(ConnectionString);
serviceBus.UseServiceBusMessageScheduler();
serviceBus.ConfigureEndpoints(context, new KebabCaseEndpointNameFormatter("Elsa", false));
});
};
});
}
}

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,8 @@
namespace Elsa.MassTransit.AzureServiceBus.Options;
/// <summary>
/// Provides settings to the RabbitMQ broker for MassTransit.
/// </summary>
public class AzureServiceBusOptions
{
}

View file

@ -0,0 +1,22 @@
<Project Sdk="Microsoft.NET.Sdk">
<Import Project="..\..\..\common.props" />
<Import Project="..\..\..\configureawait.props" />
<PropertyGroup>
<TargetFrameworks>net6.0;net7.0</TargetFrameworks>
<Description>
Provides integration of RabbitMQ with MassTransit's integration of Elsa.
</Description>
<PackageTags>elsa module service-bus masstransit rabbitmq</PackageTags>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.MassTransit\Elsa.MassTransit.csproj" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="MassTransit.RabbitMQ" Version="8.1.2" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,44 @@
using Elsa.Features.Services;
using Elsa.MassTransit.Features;
using Elsa.MassTransit.Options;
using Elsa.MassTransit.RabbitMq.Features;
using Elsa.MassTransit.RabbitMq.Options;
// ReSharper disable once CheckNamespace
namespace Elsa.Extensions;
/// <summary>
/// Provides extensions to <see cref="IModule"/> that enables and configures MassTransit and the RabbitMQ transport.
/// </summary>
public static class ModuleExtensions
{
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) => feature.UseRabbitMq(new Uri(connectionString), null);
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri connectionString) => feature.UseRabbitMq(connectionString, null);
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, RabbitMqOptions options) => feature.UseRabbitMq(null, options);
/// <summary>
/// Enable and configure the RabbitMQ transport for MassTransit.
/// </summary>
private static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri? connectionString, RabbitMqOptions? options)
{
feature.Module.Configure((Action<RabbitMqServiceBusFeature>) Configure);
return feature;
void Configure(RabbitMqServiceBusFeature bus)
{
bus.ConnectionString = connectionString;
bus.Options = options;
}
}
}

View file

@ -1,13 +1,15 @@
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.MassTransit.Features;
using Elsa.MassTransit.Options;
using Elsa.MassTransit.RabbitMq.Options;
using MassTransit;
namespace Elsa.MassTransit.Features;
namespace Elsa.MassTransit.RabbitMq.Features;
/// <summary>
/// Configures MassTransit to use the RabbitMQ broker.
/// Configures MassTransit to use the RabbitMQ transport.
/// </summary>
[DependsOn(typeof(MassTransitFeature))]
public class RabbitMqServiceBusFeature : FeatureBase

View file

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

View file

@ -1,4 +1,4 @@
namespace Elsa.MassTransit.Options;
namespace Elsa.MassTransit.RabbitMq.Options;
/// <summary>
/// Provides settings to the RabbitMQ broker for MassTransit.

View file

@ -1,7 +1,7 @@
<Project Sdk="Microsoft.NET.Sdk">
<Import Project="..\..\..\common.props" />
<Import Project="..\..\..\configureawait.props" />
<Import Project="..\..\..\common.props"/>
<Import Project="..\..\..\configureawait.props"/>
<PropertyGroup>
<TargetFrameworks>net6.0;net7.0</TargetFrameworks>
@ -12,14 +12,12 @@
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\common\Elsa.Api.Common\Elsa.Api.Common.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj"/>
</ItemGroup>
<ItemGroup>
<PackageReference Include="MassTransit.Extensions.DependencyInjection" Version="7.3.1" />
<PackageReference Include="MassTransit" Version="8.0.13" />
<PackageReference Include="MassTransit.RabbitMQ" Version="8.0.13" />
<PackageReference Include="MassTransit.Extensions.DependencyInjection" Version="7.3.1"/>
<PackageReference Include="MassTransit" Version="8.1.2"/>
</ItemGroup>
</Project>

View file

@ -24,34 +24,4 @@ public static class ModuleExtensions
module.Configure<MassTransitFeature>(massTransit => massTransit.AddConsumer<T>());
return module;
}
/// <summary>
/// Enable and configure the RabbitMQ broker for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) => feature.UseRabbitMq(new Uri(connectionString), null);
/// <summary>
/// Enable and configure the RabbitMQ broker for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri connectionString) => feature.UseRabbitMq(connectionString, null);
/// <summary>
/// Enable and configure the RabbitMQ broker for MassTransit.
/// </summary>
public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, RabbitMqOptions options) => feature.UseRabbitMq(null, options);
/// <summary>
/// Enable and configure the RabbitMQ broker for MassTransit.
/// </summary>
private static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri? connectionString, RabbitMqOptions? options)
{
void Configure(RabbitMqServiceBusFeature bus)
{
bus.ConnectionString = connectionString;
bus.Options = options;
}
feature.Module.Configure((Action<RabbitMqServiceBusFeature>) Configure);
return feature;
}
}

View file

@ -4,8 +4,10 @@ using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.MassTransit.Consumers;
using Elsa.MassTransit.Implementations;
using Elsa.MassTransit.Messages;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Features;
using MassTransit;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.MassTransit.Features;
@ -32,6 +34,13 @@ public class MassTransitWorkflowDispatcherFeature : FeatureBase
/// <inheritdoc />
public override void Apply()
{
var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer<DispatchWorkflowRequestConsumer>();
var queueAddress = new Uri($"queue:elsa-{queueName}");
EndpointConvention.Map<DispatchWorkflowDefinition>(queueAddress);
EndpointConvention.Map<DispatchWorkflowInstance>(queueAddress);
EndpointConvention.Map<DispatchTriggerWorkflows>(queueAddress);
EndpointConvention.Map<DispatchResumeWorkflows>(queueAddress);
Services.AddSingleton<MassTransitWorkflowDispatcher>();
}
}

View file

@ -23,7 +23,7 @@ public class MassTransitWorkflowDispatcher : IWorkflowDispatcher
/// <inheritdoc />
public async Task<DispatchWorkflowDefinitionResponse> DispatchAsync(DispatchWorkflowDefinitionRequest request, CancellationToken cancellationToken = default)
{
await _bus.Publish(new DispatchWorkflowDefinition(
await _bus.Send(new DispatchWorkflowDefinition(
request.DefinitionId,
request.VersionOptions,
request.Input,
@ -37,7 +37,7 @@ public class MassTransitWorkflowDispatcher : IWorkflowDispatcher
/// <inheritdoc />
public async Task<DispatchWorkflowInstanceResponse> DispatchAsync(DispatchWorkflowInstanceRequest request, CancellationToken cancellationToken = default)
{
await _bus.Publish(new DispatchWorkflowInstance(
await _bus.Send(new DispatchWorkflowInstance(
request.InstanceId,
request.BookmarkId,
request.ActivityId,
@ -53,7 +53,7 @@ public class MassTransitWorkflowDispatcher : IWorkflowDispatcher
/// <inheritdoc />
public async Task<DispatchTriggerWorkflowsResponse> DispatchAsync(DispatchTriggerWorkflowsRequest request, CancellationToken cancellationToken = default)
{
await _bus.Publish(new DispatchTriggerWorkflows(
await _bus.Send(new DispatchTriggerWorkflows(
request.ActivityTypeName,
request.BookmarkPayload,
request.CorrelationId,
@ -67,7 +67,7 @@ public class MassTransitWorkflowDispatcher : IWorkflowDispatcher
/// <inheritdoc />
public async Task<DispatchResumeWorkflowsResponse> DispatchAsync(DispatchResumeWorkflowsRequest request, CancellationToken cancellationToken = default)
{
await _bus.Publish(new DispatchResumeWorkflows(
await _bus.Send(new DispatchResumeWorkflows(
request.ActivityTypeName,
request.BookmarkPayload,
request.CorrelationId,

View file

@ -12,6 +12,7 @@
<ProjectReference Include="..\..\..\modules\Elsa.Http\Elsa.Http.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.Identity\Elsa.Identity.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.Liquid\Elsa.Liquid.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.MassTransit.RabbitMq\Elsa.MassTransit.RabbitMq.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.MassTransit\Elsa.MassTransit.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.Workflows.Api\Elsa.Workflows.Api.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.Workflows.Designer\Elsa.Workflows.Designer.csproj" />