diff --git a/Elsa.sln b/Elsa.sln
index 97074488c..20db3e94a 100644
--- a/Elsa.sln
+++ b/Elsa.sln
@@ -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}
diff --git a/src/bundles/Elsa.WorkflowServer.Web/Elsa.WorkflowServer.Web.csproj b/src/bundles/Elsa.WorkflowServer.Web/Elsa.WorkflowServer.Web.csproj
index a1d644e92..254fe3739 100644
--- a/src/bundles/Elsa.WorkflowServer.Web/Elsa.WorkflowServer.Web.csproj
+++ b/src/bundles/Elsa.WorkflowServer.Web/Elsa.WorkflowServer.Web.csproj
@@ -15,6 +15,8 @@
+
+
@@ -41,7 +43,7 @@
-
+
diff --git a/src/bundles/Elsa.WorkflowServer.Web/Program.cs b/src/bundles/Elsa.WorkflowServer.Web/Program.cs
index 35fea0cb0..917ceb429 100644
--- a/src/bundles/Elsa.WorkflowServer.Web/Program.cs
+++ b/src/bundles/Elsa.WorkflowServer.Web/Program.cs
@@ -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();
});
diff --git a/src/bundles/Elsa.WorkflowServer.Web/appsettings.json b/src/bundles/Elsa.WorkflowServer.Web/appsettings.json
index b4853ecc6..edcfc36f2 100644
--- a/src/bundles/Elsa.WorkflowServer.Web/appsettings.json
+++ b/src/bundles/Elsa.WorkflowServer.Web/appsettings.json
@@ -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",
diff --git a/src/modules/Elsa.Alterations.MassTransit/Elsa.Alterations.MassTransit.csproj b/src/modules/Elsa.Alterations.MassTransit/Elsa.Alterations.MassTransit.csproj
index e228ff551..51354f7d6 100644
--- a/src/modules/Elsa.Alterations.MassTransit/Elsa.Alterations.MassTransit.csproj
+++ b/src/modules/Elsa.Alterations.MassTransit/Elsa.Alterations.MassTransit.csproj
@@ -16,7 +16,4 @@
-
-
-
diff --git a/src/modules/Elsa.Alterations.MassTransit/Features/AlterationsMassTransitFeature.cs b/src/modules/Elsa.Alterations.MassTransit/Features/AlterationsMassTransitFeature.cs
index c2a80a0c3..1cba0ae11 100644
--- a/src/modules/Elsa.Alterations.MassTransit/Features/AlterationsMassTransitFeature.cs
+++ b/src/modules/Elsa.Alterations.MassTransit/Features/AlterationsMassTransitFeature.cs
@@ -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
///
public override void Apply()
{
+ var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer();
+ var queueAddress = new Uri($"queue:elsa-{queueName}");
+ EndpointConvention.Map(queueAddress);
Services.AddSingleton();
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.Alterations.MassTransit/Services/MassTransitAlterationJobDispatcher.cs b/src/modules/Elsa.Alterations.MassTransit/Services/MassTransitAlterationJobDispatcher.cs
index edf79941b..a2e232758 100644
--- a/src/modules/Elsa.Alterations.MassTransit/Services/MassTransitAlterationJobDispatcher.cs
+++ b/src/modules/Elsa.Alterations.MassTransit/Services/MassTransitAlterationJobDispatcher.cs
@@ -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);
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj b/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj
new file mode 100644
index 000000000..9b1428656
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Elsa.MassTransit.AzureServiceBus.csproj
@@ -0,0 +1,22 @@
+
+
+
+
+
+
+ net6.0;net7.0
+
+ Provides integration of Azure Service Bus with MassTransit's integration of Elsa.
+
+ elsa module service-bus masstransit azure-service-bus
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs
new file mode 100644
index 000000000..88413926b
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Extensions/ModuleExtensions.cs
@@ -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;
+
+///
+/// Provides extensions to that enables and configures MassTransit and the Azure Service Bus transport.
+///
+[PublicAPI]
+public static class ModuleExtensions
+{
+ ///
+ /// Enable and configure the Azure Service Bus transport for MassTransit.
+ ///
+ public static MassTransitFeature UseAzureServiceBus(this MassTransitFeature feature, string? connectionString)
+ {
+ feature.Module.Configure((Action)Configure);
+ return feature;
+
+ void Configure(AzureServiceBusFeature bus)
+ {
+ bus.ConnectionString = connectionString;
+ }
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs
new file mode 100644
index 000000000..b10a3bb54
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Features/AzureServiceBusFeature.cs
@@ -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;
+
+///
+/// Configures MassTransit to use the Azure Service Bus transport.
+///
+[DependsOn(typeof(MassTransitFeature))]
+public class AzureServiceBusFeature : FeatureBase
+{
+ ///
+ public AzureServiceBusFeature(IModule module) : base(module)
+ {
+ }
+
+ ///
+ /// An Azure Service Bus connection string.
+ ///
+ public string? ConnectionString { get; set; }
+
+ ///
+ public override void Configure()
+ {
+ Module.Configure(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));
+ });
+ };
+ });
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/FodyWeavers.xml b/src/modules/Elsa.MassTransit.AzureServiceBus/FodyWeavers.xml
new file mode 100644
index 000000000..00e1d9a1c
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/FodyWeavers.xml
@@ -0,0 +1,3 @@
+
+
+
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs
new file mode 100644
index 000000000..26df271c5
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.AzureServiceBus/Options/AzureServiceBusOptions.cs
@@ -0,0 +1,8 @@
+namespace Elsa.MassTransit.AzureServiceBus.Options;
+
+///
+/// Provides settings to the RabbitMQ broker for MassTransit.
+///
+public class AzureServiceBusOptions
+{
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Elsa.MassTransit.RabbitMq.csproj b/src/modules/Elsa.MassTransit.RabbitMq/Elsa.MassTransit.RabbitMq.csproj
new file mode 100644
index 000000000..f23d1507e
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.RabbitMq/Elsa.MassTransit.RabbitMq.csproj
@@ -0,0 +1,22 @@
+
+
+
+
+
+
+ net6.0;net7.0
+
+ Provides integration of RabbitMQ with MassTransit's integration of Elsa.
+
+ elsa module service-bus masstransit rabbitmq
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs
new file mode 100644
index 000000000..432ecce57
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.RabbitMq/Extensions/ModuleExtensions.cs
@@ -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;
+
+///
+/// Provides extensions to that enables and configures MassTransit and the RabbitMQ transport.
+///
+public static class ModuleExtensions
+{
+ ///
+ /// Enable and configure the RabbitMQ transport for MassTransit.
+ ///
+ public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) => feature.UseRabbitMq(new Uri(connectionString), null);
+
+ ///
+ /// Enable and configure the RabbitMQ transport for MassTransit.
+ ///
+ public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri connectionString) => feature.UseRabbitMq(connectionString, null);
+
+ ///
+ /// Enable and configure the RabbitMQ transport for MassTransit.
+ ///
+ public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, RabbitMqOptions options) => feature.UseRabbitMq(null, options);
+
+ ///
+ /// Enable and configure the RabbitMQ transport for MassTransit.
+ ///
+ private static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri? connectionString, RabbitMqOptions? options)
+ {
+ feature.Module.Configure((Action) Configure);
+ return feature;
+
+ void Configure(RabbitMqServiceBusFeature bus)
+ {
+ bus.ConnectionString = connectionString;
+ bus.Options = options;
+ }
+ }
+}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit/Features/RabbitMqServiceBusFeature.cs b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs
similarity index 88%
rename from src/modules/Elsa.MassTransit/Features/RabbitMqServiceBusFeature.cs
rename to src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs
index f1966a8c0..c488f2994 100644
--- a/src/modules/Elsa.MassTransit/Features/RabbitMqServiceBusFeature.cs
+++ b/src/modules/Elsa.MassTransit.RabbitMq/Features/RabbitMqServiceBusFeature.cs
@@ -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;
///
-/// Configures MassTransit to use the RabbitMQ broker.
+/// Configures MassTransit to use the RabbitMQ transport.
///
[DependsOn(typeof(MassTransitFeature))]
public class RabbitMqServiceBusFeature : FeatureBase
diff --git a/src/modules/Elsa.MassTransit.RabbitMq/FodyWeavers.xml b/src/modules/Elsa.MassTransit.RabbitMq/FodyWeavers.xml
new file mode 100644
index 000000000..00e1d9a1c
--- /dev/null
+++ b/src/modules/Elsa.MassTransit.RabbitMq/FodyWeavers.xml
@@ -0,0 +1,3 @@
+
+
+
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit/Options/RabbitMqOptions.cs b/src/modules/Elsa.MassTransit.RabbitMq/Options/RabbitMqOptions.cs
similarity index 86%
rename from src/modules/Elsa.MassTransit/Options/RabbitMqOptions.cs
rename to src/modules/Elsa.MassTransit.RabbitMq/Options/RabbitMqOptions.cs
index 2f8c3ca6c..17b26851e 100644
--- a/src/modules/Elsa.MassTransit/Options/RabbitMqOptions.cs
+++ b/src/modules/Elsa.MassTransit.RabbitMq/Options/RabbitMqOptions.cs
@@ -1,4 +1,4 @@
-namespace Elsa.MassTransit.Options;
+namespace Elsa.MassTransit.RabbitMq.Options;
///
/// Provides settings to the RabbitMQ broker for MassTransit.
diff --git a/src/modules/Elsa.MassTransit/Elsa.MassTransit.csproj b/src/modules/Elsa.MassTransit/Elsa.MassTransit.csproj
index 0513cfa7d..58e1bf9b7 100644
--- a/src/modules/Elsa.MassTransit/Elsa.MassTransit.csproj
+++ b/src/modules/Elsa.MassTransit/Elsa.MassTransit.csproj
@@ -1,7 +1,7 @@
-
-
+
+
net6.0;net7.0
@@ -12,14 +12,12 @@
-
-
+
-
-
-
+
+
diff --git a/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs b/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs
index 49db276a5..1b696d661 100644
--- a/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs
+++ b/src/modules/Elsa.MassTransit/Extensions/ModuleExtensions.cs
@@ -24,34 +24,4 @@ public static class ModuleExtensions
module.Configure(massTransit => massTransit.AddConsumer());
return module;
}
-
- ///
- /// Enable and configure the RabbitMQ broker for MassTransit.
- ///
- public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, string connectionString) => feature.UseRabbitMq(new Uri(connectionString), null);
-
- ///
- /// Enable and configure the RabbitMQ broker for MassTransit.
- ///
- public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, Uri connectionString) => feature.UseRabbitMq(connectionString, null);
-
- ///
- /// Enable and configure the RabbitMQ broker for MassTransit.
- ///
- public static MassTransitFeature UseRabbitMq(this MassTransitFeature feature, RabbitMqOptions options) => feature.UseRabbitMq(null, options);
-
- ///
- /// Enable and configure the RabbitMQ broker for MassTransit.
- ///
- 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) Configure);
- return feature;
- }
}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs
index 6a1b73375..4e45ab77a 100644
--- a/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs
+++ b/src/modules/Elsa.MassTransit/Features/MassTransitWorkflowDispatcherFeature.cs
@@ -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
///
public override void Apply()
{
+ var queueName = KebabCaseEndpointNameFormatter.Instance.Consumer();
+ var queueAddress = new Uri($"queue:elsa-{queueName}");
+ EndpointConvention.Map(queueAddress);
+ EndpointConvention.Map(queueAddress);
+ EndpointConvention.Map(queueAddress);
+ EndpointConvention.Map(queueAddress);
+
Services.AddSingleton();
}
}
\ No newline at end of file
diff --git a/src/modules/Elsa.MassTransit/Implementations/MassTransitWorkflowDispatcher.cs b/src/modules/Elsa.MassTransit/Implementations/MassTransitWorkflowDispatcher.cs
index 9112a4866..9f6d02f15 100644
--- a/src/modules/Elsa.MassTransit/Implementations/MassTransitWorkflowDispatcher.cs
+++ b/src/modules/Elsa.MassTransit/Implementations/MassTransitWorkflowDispatcher.cs
@@ -23,7 +23,7 @@ public class MassTransitWorkflowDispatcher : IWorkflowDispatcher
///
public async Task 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
///
public async Task 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
///
public async Task 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
///
public async Task DispatchAsync(DispatchResumeWorkflowsRequest request, CancellationToken cancellationToken = default)
{
- await _bus.Publish(new DispatchResumeWorkflows(
+ await _bus.Send(new DispatchResumeWorkflows(
request.ActivityTypeName,
request.BookmarkPayload,
request.CorrelationId,
diff --git a/src/samples/aspnet/Elsa.Samples.AspNet.MassTransitActivities/Elsa.Samples.AspNet.MassTransitActivities.csproj b/src/samples/aspnet/Elsa.Samples.AspNet.MassTransitActivities/Elsa.Samples.AspNet.MassTransitActivities.csproj
index 427b622f3..d655de1c0 100644
--- a/src/samples/aspnet/Elsa.Samples.AspNet.MassTransitActivities/Elsa.Samples.AspNet.MassTransitActivities.csproj
+++ b/src/samples/aspnet/Elsa.Samples.AspNet.MassTransitActivities/Elsa.Samples.AspNet.MassTransitActivities.csproj
@@ -12,6 +12,7 @@
+