From 4b20841ea726accaa80730217d698d0c396d10de Mon Sep 17 00:00:00 2001 From: Peter Davis Date: Fri, 15 Nov 2019 19:49:20 +1030 Subject: [PATCH] MassTransit Fixes (#148) * Fix issue #134 where consumer could not be created * Fix issue where message was null in ReceiveMassTransitMessage input variables are case sensitive * Add simple test script to execute the workflow * Add message type as input from consumer * Allow passwords to be specified in amqp URI --- .../Activities/ReceiveMassTransitMessage.cs | 9 ++--- .../Elsa.Activities.MassTransit/Constants.cs | 12 +++++++ .../Consumers/WorkflowConsumer.cs | 11 +++--- .../Extensions/ServiceCollectionExtensions.cs | 34 ++++++++++++------- .../Options/RabbitMqOptions.cs | 2 ++ src/samples/Sample08/Sample08.http | 16 +++++++++ .../Sample08/Workflows/CreateOrderWorkflow.cs | 2 +- src/samples/Sample08/appsettings.json | 4 ++- 8 files changed, 67 insertions(+), 23 deletions(-) create mode 100644 src/activities/Elsa.Activities.MassTransit/Constants.cs create mode 100644 src/samples/Sample08/Sample08.http diff --git a/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs b/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs index 70257a558..5e9d5671b 100644 --- a/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs +++ b/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs @@ -32,10 +32,11 @@ namespace Elsa.Activities.MassTransit.Activities protected override bool OnCanExecute(WorkflowExecutionContext context) { - var message = context.Workflow.Input["message"]; + var messageTypeName = context.Workflow.Input[Constants.MessageTypeNameInputKey]; + var messageInputType = System.Type.GetType(messageTypeName.ToString()); var messageType = MessageType; - - return message != null && messageType != null && message.GetType() == messageType; + + return messageInputType != null && messageType != null && messageInputType == messageType; } protected override ActivityExecutionResult OnExecute(WorkflowExecutionContext context) @@ -46,7 +47,7 @@ namespace Elsa.Activities.MassTransit.Activities protected override Task OnResumeAsync(WorkflowExecutionContext context, CancellationToken cancellationToken) { - var message = context.Workflow.Input["message"]; + var message = context.Workflow.Input[Constants.MessageInputKey]; context.SetLastResult(message); return Task.FromResult(Done()); diff --git a/src/activities/Elsa.Activities.MassTransit/Constants.cs b/src/activities/Elsa.Activities.MassTransit/Constants.cs new file mode 100644 index 000000000..f2613700f --- /dev/null +++ b/src/activities/Elsa.Activities.MassTransit/Constants.cs @@ -0,0 +1,12 @@ +using System; +using System.Collections.Generic; +using System.Text; + +namespace Elsa.Activities.MassTransit +{ + internal static class Constants + { + public const string MessageInputKey = "Message"; + public const string MessageTypeNameInputKey = "MessageType"; + } +} diff --git a/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs b/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs index 6097b7064..a6b6ea4bf 100644 --- a/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs +++ b/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs @@ -15,19 +15,20 @@ namespace Elsa.Activities.MassTransit.Consumers { this.workflowInvoker = workflowInvoker; } - + public async Task Consume(ConsumeContext context) { var message = context.Message; var activityType = nameof(ReceiveMassTransitMessage); var input = new Variables(); - input.SetVariable("Message", message); - + input.SetVariable(Constants.MessageInputKey, message); + input.SetVariable(Constants.MessageTypeNameInputKey, typeof(T).AssemblyQualifiedName); + var correlationId = context.CorrelationId?.ToString(); - + await workflowInvoker.TriggerAsync( - activityType, + activityType, input, correlationId, x => ReceiveMassTransitMessage.GetMessageType(x) == message.GetType(), diff --git a/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs index 63d0ce9c3..360ebf8d3 100644 --- a/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs @@ -1,4 +1,4 @@ -using System; +using System; using System.Collections.Generic; using Elsa.Activities.MassTransit.Activities; using Elsa.Activities.MassTransit.Consumers; @@ -19,38 +19,48 @@ namespace Elsa.Activities.MassTransit.Extensions { var optionsBuilder = services.AddOptions(); options?.Invoke(optionsBuilder); - + services .AddActivity() .AddActivity(); services.AddSingleton(); services.AddSingleton(); - + services.AddMassTransit( massTransit => { + foreach (var messageType in messageTypes) + { + massTransit.AddConsumer(CreateConsumerType(messageType)); + } + massTransit.AddBus(sp => CreateUsingRabbitMq(massTransit, sp, messageTypes)); }); services.AddSingleton(); - foreach (var messageType in messageTypes) - { - var consumerType = CreateConsumerType(messageType); - services.AddSingleton(consumerType); - } - return services; } private static IBusControl CreateUsingRabbitMq(IServiceCollectionConfigurator massTransit, IServiceProvider sp, IEnumerable messageTypes) { + return Bus.Factory.CreateUsingRabbitMq( bus => { - var options = sp.GetRequiredService>(); - var host = bus.Host(new Uri(options.Value.Host), _ => { }); + var options = sp.GetRequiredService>().Value; + var host = bus.Host(new Uri(options.Host), h => + { + if (!string.IsNullOrEmpty(options.Username)) + { + h.Username(options.Username); + + if (!string.IsNullOrEmpty(options.Password)) + h.Password(options.Password); + } + }); + foreach (var messageType in messageTypes) { @@ -63,7 +73,7 @@ namespace Elsa.Activities.MassTransit.Extensions endpoint => { endpoint.PrefetchCount = 16; - endpoint.Consumer(consumerType, sp.GetRequiredService); + endpoint.ConfigureConsumer(sp, consumerType); MapEndpointConvention(messageType, endpoint.InputAddress); }); } diff --git a/src/activities/Elsa.Activities.MassTransit/Options/RabbitMqOptions.cs b/src/activities/Elsa.Activities.MassTransit/Options/RabbitMqOptions.cs index c80c5519b..715f7fe9f 100644 --- a/src/activities/Elsa.Activities.MassTransit/Options/RabbitMqOptions.cs +++ b/src/activities/Elsa.Activities.MassTransit/Options/RabbitMqOptions.cs @@ -3,5 +3,7 @@ namespace Elsa.Activities.MassTransit.Options public class RabbitMqOptions { public string Host { get; set; } + public string Username { get; set; } + public string Password { get; set; } } } \ No newline at end of file diff --git a/src/samples/Sample08/Sample08.http b/src/samples/Sample08/Sample08.http new file mode 100644 index 000000000..864f8b8e7 --- /dev/null +++ b/src/samples/Sample08/Sample08.http @@ -0,0 +1,16 @@ +// Use VS Code Rest Client +// https://marketplace.visualstudio.com/items?itemName=humao.rest-client + +// Click Send Requst below: +POST http://localhost:59862/orders +Content-Type: application/json + +{ + "id": "order-1", + "customer": { + "name": "Jimmy Drop Tables", + "email": "drop-tables@where.com" + }, + "product": "MassTransit", + "amount": 99.95 +} diff --git a/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs b/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs index db9a282a4..06e117d0a 100644 --- a/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs +++ b/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs @@ -32,7 +32,7 @@ namespace Sample08.Workflows activity => { activity.VariableName = "order"; - activity.ValueExpression = new JavaScriptExpression("lastResult().Content"); + activity.ValueExpression = new JavaScriptExpression("lastResult().Body"); } ) .Then(activity => diff --git a/src/samples/Sample08/appsettings.json b/src/samples/Sample08/appsettings.json index d10c4eb0d..6435d5db6 100644 --- a/src/samples/Sample08/appsettings.json +++ b/src/samples/Sample08/appsettings.json @@ -11,7 +11,9 @@ }, "MassTransit": { "RabbitMq": { - "Host": "rabbitmq://localhost:5672" + "Host": "rabbitmq://localhost:5672", + "Username": "guest", + "Password": "guest" } } }