From 2a9c8d0b7e6e439231a39f2c3be95b1edd83f01e Mon Sep 17 00:00:00 2001 From: Konstantins Vedenins <77638996+constantine-v@users.noreply.github.com> Date: Wed, 22 Dec 2021 14:52:56 +0200 Subject: [PATCH] Added exchange name property, removed auto delete for bindings (#2622) Added Exchange Name property to give ability to subscribe/publish to specific exchange. Removed binding auto-delete to preserve binding in case of disconnect. Added RabbitMQ connection string format to connection string property tooltip. --- .../Activities/IRabbitMqActivity.cs | 1 + .../RabbitMqMessageReceived.cs | 14 +++++-- ...abbitMqMessageReceivedBuilderExtensions.cs | 9 +++++ .../RabbitMqMessageReceivedExtensions.cs | 6 +++ .../SendRabbitMqMessage.cs | 16 +++++--- .../SendRabbitMqMessageBuilderExtensions.cs | 38 +++++++++++++++++++ .../SendRabbitMqMessageExtensions.cs | 6 +++ .../Bookmarks/MessageReceivedBookmark.cs | 5 ++- .../Configuration/RabbitMqBusConfiguration.cs | 8 +++- .../Services/Client.cs | 6 +-- .../Services/RabbitMqQueueStarter.cs | 3 +- .../Services/Worker.cs | 2 +- .../Workflows/ConsumerWorkflow.cs | 2 +- .../Workflows/ProducerWorkflow.cs | 2 +- .../Workflows/SendAndReceiveWorkflow.cs | 4 +- 15 files changed, 101 insertions(+), 21 deletions(-) diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/IRabbitMqActivity.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/IRabbitMqActivity.cs index 545a92c9c..6c169dac1 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/IRabbitMqActivity.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/IRabbitMqActivity.cs @@ -6,6 +6,7 @@ namespace Elsa.Activities.RabbitMq public interface IRabbitMqActivity : IActivity { string ConnectionString { get; set; } + string ExchangeName { get; set; } string RoutingKey { get; set; } Dictionary Headers { get; set; } } diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceived.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceived.cs index a9f07c051..9577296e7 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceived.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceived.cs @@ -28,21 +28,27 @@ namespace Elsa.Activities.RabbitMq } [ActivityInput( - Hint = "Routing Key", + Hint = "Exchange to listen to", Order = 1, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public string ExchangeName { get; set; } = default!; + + [ActivityInput( + Hint = "Routing Key", + Order = 2, + SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string RoutingKey { get; set; } = default!; [ActivityInput( Hint = "List of headers that should be present in the message", - Order = 2, + Order = 3, UIHint = ActivityInputUIHints.Dictionary, DefaultSyntax = SyntaxNames.Json, SupportedSyntaxes = new[] { SyntaxNames.Json })] public Dictionary Headers { get; set; } = new Dictionary(); [ActivityInput( - Hint = "RabbitMQ connection string", + Hint = "RabbitMQ connection string [amqp://user:pass@host:10000/vhost] - https://www.rabbitmq.com/uri-spec.html", SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid }, Order = 2, Category = PropertyCategories.Configuration)] @@ -78,7 +84,7 @@ namespace Elsa.Activities.RabbitMq } private async Task StartClient() { - var config = new RabbitMqBusConfiguration(ConnectionString, RoutingKey, Headers); + var config = new RabbitMqBusConfiguration(ConnectionString, ExchangeName, RoutingKey, Headers); var client = await _messageReceiverClientFactory.GetReceiverAsync(config); client.StartClient(); diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedBuilderExtensions.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedBuilderExtensions.cs index fdf03eb03..20bf9041b 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedBuilderExtensions.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedBuilderExtensions.cs @@ -10,7 +10,16 @@ namespace Elsa.Activities.RabbitMq public static IActivityBuilder MessageReceived(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => builder.Then(setup, null, lineNumber, sourceFile); + public static IActivityBuilder MessageReceived(this IBuilder builder, string connectionString, string exchangeName, string routingKey, Dictionary headers, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.MessageReceived(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithRoutingKey(routingKey).WithHeaders(headers), lineNumber, sourceFile); + public static IActivityBuilder MessageReceived(this IBuilder builder, string connectionString, string routingKey, Dictionary headers, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => builder.MessageReceived(setup => setup.WithConnectionString(connectionString).WithRoutingKey(routingKey).WithHeaders(headers), lineNumber, sourceFile); + + public static IActivityBuilder MessageReceived(this IBuilder builder, string connectionString, string exchangeName, string routingKey, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.MessageReceived(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithRoutingKey(routingKey), lineNumber, sourceFile); + + public static IActivityBuilder MessageReceived(this IBuilder builder, string connectionString, string routingKey, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.MessageReceived(setup => setup.WithConnectionString(connectionString).WithRoutingKey(routingKey), lineNumber, sourceFile); } } diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedExtensions.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedExtensions.cs index ea8a74756..4dd3b7ef8 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedExtensions.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/RabbitMqMessageReceived/RabbitMqMessageReceivedExtensions.cs @@ -14,6 +14,12 @@ namespace Elsa.Activities.RabbitMq public static ISetupActivity WithConnectionString(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.ConnectionString, value!); public static ISetupActivity WithConnectionString(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.ConnectionString, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity messageReceived, string value) => messageReceived.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithRoutingKey(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.RoutingKey, value!); public static ISetupActivity WithRoutingKey(this ISetupActivity messageReceived, Func> value) => messageReceived.Set(x => x.RoutingKey, value!); public static ISetupActivity WithRoutingKey(this ISetupActivity messageReceived, Func value) => messageReceived.Set(x => x.RoutingKey, value!); diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessage.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessage.cs index fa7c0e070..389f3dc4d 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessage.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessage.cs @@ -27,14 +27,20 @@ namespace Elsa.Activities.RabbitMq } [ActivityInput( - Hint = "Topic", + Hint = "Exchange where message will be published", Order = 1, SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] + public string ExchangeName { get; set; } = default!; + + [ActivityInput( + Hint = "Topic", + Order = 2, + SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid })] public string RoutingKey { get; set; } = default!; [ActivityInput( Hint = "List of headers that should be present in the message", - Order = 2, + Order = 3, UIHint = ActivityInputUIHints.Dictionary, DefaultSyntax = SyntaxNames.Json, SupportedSyntaxes = new[] { SyntaxNames.Json })] @@ -42,13 +48,13 @@ namespace Elsa.Activities.RabbitMq [ActivityInput( Hint = "Message body", - Order = 3, + Order = 4, UIHint = ActivityInputUIHints.MultiLine, SupportedSyntaxes = new[] { SyntaxNames.Json })] public string Message { get; set; } = default!; [ActivityInput( - Hint = "RabbitMQ connection string", + Hint = "RabbitMQ connection string [amqp://user:pass@host:10000/vhost] - https://www.rabbitmq.com/uri-spec.html", SupportedSyntaxes = new[] { SyntaxNames.JavaScript, SyntaxNames.Liquid }, Order = 1, Category = PropertyCategories.Configuration)] @@ -56,7 +62,7 @@ namespace Elsa.Activities.RabbitMq protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context) { - var config = new RabbitMqBusConfiguration(ConnectionString, RoutingKey, Headers); + var config = new RabbitMqBusConfiguration(ConnectionString, ExchangeName, RoutingKey, Headers); var client = await _messageSenderClientFactory.GetSenderAsync(config); diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageBuilderExtensions.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageBuilderExtensions.cs index b4ba340e7..8efd151d2 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageBuilderExtensions.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageBuilderExtensions.cs @@ -12,6 +12,18 @@ namespace Elsa.Activities.RabbitMq public static IActivityBuilder SendTopicMessage(this IBuilder builder, Action> setup, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => builder.Then(setup, null, lineNumber, sourceFile); + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Dictionary headers, Func> message, [CallerLineNumber] int lineNumber = default, + [CallerFilePath] string? sourceFile = default) => builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithHeaders(headers).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Dictionary headers, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithHeaders(headers).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Dictionary headers, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithHeaders(headers).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Dictionary headers, string message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithHeaders(headers).WithMessage(message), lineNumber, sourceFile); + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string topic, Dictionary headers, Func> message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithTopic(topic).WithHeaders(headers).WithMessage(message), lineNumber, sourceFile); @@ -23,5 +35,31 @@ namespace Elsa.Activities.RabbitMq public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string topic, Dictionary headers, string message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithTopic(topic).WithHeaders(headers).WithMessage(message), lineNumber, sourceFile); + + + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Func> message, [CallerLineNumber] int lineNumber = default, + [CallerFilePath] string? sourceFile = default) => builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string exchangeName, string topic, string message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithExchangeName(exchangeName).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string topic, Func> message, [CallerLineNumber] int lineNumber = default, + [CallerFilePath] string? sourceFile = default) => builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string topic, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string topic, Func message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); + + public static IActivityBuilder SendTopicMessage(this IBuilder builder, string connectionString, string topic, string message, [CallerLineNumber] int lineNumber = default, [CallerFilePath] string? sourceFile = default) => + builder.SendTopicMessage(setup => setup.WithConnectionString(connectionString).WithTopic(topic).WithMessage(message), lineNumber, sourceFile); } } diff --git a/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageExtensions.cs b/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageExtensions.cs index e83694d6f..6ea9160c0 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageExtensions.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Activities/SendRabbitMqMessage/SendRabbitMqMessageExtensions.cs @@ -14,6 +14,12 @@ namespace Elsa.Activities.RabbitMq public static ISetupActivity WithConnectionString(this ISetupActivity sendMessage, Func value) => sendMessage.Set(x => x.ConnectionString, value!); public static ISetupActivity WithConnectionString(this ISetupActivity sendMessage, string value) => sendMessage.Set(x => x.ConnectionString, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity sendMessage, Func> value) => sendMessage.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity sendMessage, Func> value) => sendMessage.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity sendMessage, Func value) => sendMessage.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity sendMessage, Func value) => sendMessage.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithExchangeName(this ISetupActivity sendMessage, string value) => sendMessage.Set(x => x.ExchangeName, value!); + public static ISetupActivity WithTopic(this ISetupActivity sendMessage, Func> value) => sendMessage.Set(x => x.RoutingKey, value!); public static ISetupActivity WithTopic(this ISetupActivity sendMessage, Func> value) => sendMessage.Set(x => x.RoutingKey, value!); public static ISetupActivity WithTopic(this ISetupActivity sendMessage, Func value) => sendMessage.Set(x => x.RoutingKey, value!); diff --git a/src/activities/Elsa.Activities.RabbitMq/Bookmarks/MessageReceivedBookmark.cs b/src/activities/Elsa.Activities.RabbitMq/Bookmarks/MessageReceivedBookmark.cs index 0d45f5329..4f3888a4e 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Bookmarks/MessageReceivedBookmark.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Bookmarks/MessageReceivedBookmark.cs @@ -11,14 +11,16 @@ namespace Elsa.Activities.RabbitMq.Bookmarks { } - public MessageReceivedBookmark(string routingKey, string connectionString, Dictionary headers) + public MessageReceivedBookmark(string exchangeName, string routingKey, string connectionString, Dictionary headers) { + ExchangeName = exchangeName; RoutingKey = routingKey; ConnectionString = connectionString; Headers = headers ?? new Dictionary(); } + public string ExchangeName { get; set; } = default!; public string RoutingKey { get; set; } = default!; public string ConnectionString { get; set; } = default!; public Dictionary Headers { get; set; } = default!; @@ -31,6 +33,7 @@ namespace Elsa.Activities.RabbitMq.Bookmarks { Result(new MessageReceivedBookmark { + ExchangeName = (await context.ReadActivityPropertyAsync(x => x.ExchangeName, cancellationToken))!, RoutingKey = (await context.ReadActivityPropertyAsync(x => x.RoutingKey, cancellationToken))!, ConnectionString = (await context.ReadActivityPropertyAsync(x => x.ConnectionString, cancellationToken))!, Headers = (await context.ReadActivityPropertyAsync(x => x.Headers, cancellationToken))! diff --git a/src/activities/Elsa.Activities.RabbitMq/Configuration/RabbitMqBusConfiguration.cs b/src/activities/Elsa.Activities.RabbitMq/Configuration/RabbitMqBusConfiguration.cs index 95f56db9a..0d76a79b3 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Configuration/RabbitMqBusConfiguration.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Configuration/RabbitMqBusConfiguration.cs @@ -6,12 +6,14 @@ namespace Elsa.Activities.RabbitMq.Configuration public class RabbitMqBusConfiguration { public string ConnectionString { get; } + public string ExchangeName { get; } public string RoutingKey { get; } public Dictionary Headers { get; } - public RabbitMqBusConfiguration(string connectionString, string routingKey, Dictionary headers) + public RabbitMqBusConfiguration(string connectionString, string exchangeName, string routingKey, Dictionary headers) { ConnectionString = connectionString; + ExchangeName = exchangeName; RoutingKey = routingKey; Headers = headers ?? new Dictionary(); } @@ -20,7 +22,9 @@ namespace Elsa.Activities.RabbitMq.Configuration { var headersString = string.Concat(Headers.Select((x, y) => string.Concat(x, y))); - return System.HashCode.Combine(ConnectionString, RoutingKey, headersString); + return System.HashCode.Combine(ConnectionString, ExchangeName, RoutingKey, headersString); } + + public string TopicFullName => string.IsNullOrEmpty(ExchangeName) ? RoutingKey : $"{RoutingKey}@{ExchangeName}"; } } diff --git a/src/activities/Elsa.Activities.RabbitMq/Services/Client.cs b/src/activities/Elsa.Activities.RabbitMq/Services/Client.cs index 2c7de12ea..fd5c30a3c 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Services/Client.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Services/Client.cs @@ -42,18 +42,18 @@ namespace Elsa.Activities.RabbitMq.Services })) .Transport(t => { - t.UseRabbitMq(Configuration.ConnectionString, $"Elsa{Guid.NewGuid().ToString("n").ToUpper()}").InputQueueOptions(o => o.SetAutoDelete(autoDelete: true)); + t.UseRabbitMq(Configuration.ConnectionString, $"Elsa{Guid.NewGuid().ToString("n").ToUpper()}"); }) .Start(); - _bus.Advanced.Topics.Subscribe(Configuration.RoutingKey); + _bus.Advanced.Topics.Subscribe(Configuration.TopicFullName); } public async Task PublishMessage(string message) { if (_bus == null) ConfigureAsOneWayClient(); - await _bus!.Advanced.Topics.Publish(Configuration.RoutingKey, message, Configuration.Headers); + await _bus!.Advanced.Topics.Publish(Configuration.TopicFullName, message, Configuration.Headers); } public void Dispose() diff --git a/src/activities/Elsa.Activities.RabbitMq/Services/RabbitMqQueueStarter.cs b/src/activities/Elsa.Activities.RabbitMq/Services/RabbitMqQueueStarter.cs index 1cdd402bf..744432015 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Services/RabbitMqQueueStarter.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Services/RabbitMqQueueStarter.cs @@ -95,9 +95,10 @@ namespace Elsa.Activities.RabbitMq.Services { var connectionString = await activity.EvaluatePropertyValueAsync(x => x.ConnectionString, cancellationToken); var routingKey = await activity.EvaluatePropertyValueAsync(x => x.RoutingKey, cancellationToken); + var exchangeName = await activity.EvaluatePropertyValueAsync(x => x.ExchangeName, cancellationToken); var headers = await activity.EvaluatePropertyValueAsync(x => x.Headers, cancellationToken); - var config = new RabbitMqBusConfiguration(connectionString!, routingKey!, headers!); + var config = new RabbitMqBusConfiguration(connectionString!, exchangeName, routingKey!, headers!); yield return config!; } diff --git a/src/activities/Elsa.Activities.RabbitMq/Services/Worker.cs b/src/activities/Elsa.Activities.RabbitMq/Services/Worker.cs index 1452d0501..053e6bc2c 100644 --- a/src/activities/Elsa.Activities.RabbitMq/Services/Worker.cs +++ b/src/activities/Elsa.Activities.RabbitMq/Services/Worker.cs @@ -51,7 +51,7 @@ namespace Elsa.Activities.RabbitMq.Services var config = _client.Configuration; - var bookmark = new MessageReceivedBookmark(config.RoutingKey, config.ConnectionString, config.Headers); + var bookmark = new MessageReceivedBookmark(config.ExchangeName, config.RoutingKey, config.ConnectionString, config.Headers); var launchContext = new WorkflowsQuery(ActivityType, bookmark); await _workflowLaunchpad.UseServiceAsync(service => service.CollectAndDispatchWorkflowsAsync(launchContext, new WorkflowInput(message), cancellationToken)); diff --git a/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ConsumerWorkflow.cs b/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ConsumerWorkflow.cs index 758d34baf..cd26e04ed 100644 --- a/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ConsumerWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ConsumerWorkflow.cs @@ -19,7 +19,7 @@ namespace Elsa.Samples.RabbitMqWorker.Workflows { builder .Timer(Duration.FromSeconds(5)) - .MessageReceived(_connectionString, "Podcasts.Weather", default!) + .MessageReceived(_connectionString, "Podcasts.Weather") .WriteLine(context => { var message = context.GetInput(); diff --git a/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ProducerWorkflow.cs b/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ProducerWorkflow.cs index a4734eefd..9a6891ade 100644 --- a/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ProducerWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/ProducerWorkflow.cs @@ -20,7 +20,7 @@ namespace Elsa.Samples.RabbitMqWorker.Workflows builder .Timer(Duration.FromSeconds(5)) .WriteLine("Sending a weather update with the \"Podcasts.Weather\" topic.") - .SendTopicMessage(_connectionString, "Podcasts.Weather", default!, "Cloudy with a chance of meatballs"); + .SendTopicMessage(_connectionString, "Podcasts.Weather", "Cloudy with a chance of meatballs"); } } } \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/SendAndReceiveWorkflow.cs b/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/SendAndReceiveWorkflow.cs index 08f4a0f80..83bda8954 100644 --- a/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/SendAndReceiveWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.RabbitMqWorker/Workflows/SendAndReceiveWorkflow.cs @@ -22,8 +22,8 @@ namespace Elsa.Samples.RabbitMqWorker.Workflows return $"Start! - correlationId: {correlationId}"; }) - .SendTopicMessage(_connectionString, "Greetings", default!, "Greetings from RabbitMQ") - .MessageReceived(_connectionString, "Greetings", default!) + .SendTopicMessage(_connectionString, "Greetings", "Greetings from RabbitMQ") + .MessageReceived(_connectionString, "Greetings") .WriteLine(ctx => "End: " + (string?)ctx.Input); }