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
This commit is contained in:
parent
cccf5c831d
commit
4b20841ea7
|
|
@ -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<ActivityExecutionResult> 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());
|
||||
|
|
|
|||
12
src/activities/Elsa.Activities.MassTransit/Constants.cs
Normal file
12
src/activities/Elsa.Activities.MassTransit/Constants.cs
Normal file
|
|
@ -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";
|
||||
}
|
||||
}
|
||||
|
|
@ -15,19 +15,20 @@ namespace Elsa.Activities.MassTransit.Consumers
|
|||
{
|
||||
this.workflowInvoker = workflowInvoker;
|
||||
}
|
||||
|
||||
|
||||
public async Task Consume(ConsumeContext<T> 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(),
|
||||
|
|
|
|||
|
|
@ -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<RabbitMqOptions>();
|
||||
options?.Invoke(optionsBuilder);
|
||||
|
||||
|
||||
services
|
||||
.AddActivity<SendMassTransitMessage>()
|
||||
.AddActivity<ReceiveMassTransitMessage>();
|
||||
|
||||
services.AddSingleton<SimplifiedBusHealthCheck>();
|
||||
services.AddSingleton<ReceiveEndpointHealthCheck>();
|
||||
|
||||
|
||||
services.AddMassTransit(
|
||||
massTransit =>
|
||||
{
|
||||
foreach (var messageType in messageTypes)
|
||||
{
|
||||
massTransit.AddConsumer(CreateConsumerType(messageType));
|
||||
}
|
||||
|
||||
massTransit.AddBus(sp => CreateUsingRabbitMq(massTransit, sp, messageTypes));
|
||||
});
|
||||
|
||||
services.AddSingleton<IHostedService, MassTransitHostedService>();
|
||||
|
||||
foreach (var messageType in messageTypes)
|
||||
{
|
||||
var consumerType = CreateConsumerType(messageType);
|
||||
services.AddSingleton(consumerType);
|
||||
}
|
||||
|
||||
return services;
|
||||
}
|
||||
|
||||
private static IBusControl CreateUsingRabbitMq(IServiceCollectionConfigurator massTransit, IServiceProvider sp, IEnumerable<Type> messageTypes)
|
||||
{
|
||||
|
||||
return Bus.Factory.CreateUsingRabbitMq(
|
||||
bus =>
|
||||
{
|
||||
var options = sp.GetRequiredService<IOptions<RabbitMqOptions>>();
|
||||
var host = bus.Host(new Uri(options.Value.Host), _ => { });
|
||||
var options = sp.GetRequiredService<IOptions<RabbitMqOptions>>().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);
|
||||
});
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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; }
|
||||
}
|
||||
}
|
||||
16
src/samples/Sample08/Sample08.http
Normal file
16
src/samples/Sample08/Sample08.http
Normal file
|
|
@ -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
|
||||
}
|
||||
|
|
@ -32,7 +32,7 @@ namespace Sample08.Workflows
|
|||
activity =>
|
||||
{
|
||||
activity.VariableName = "order";
|
||||
activity.ValueExpression = new JavaScriptExpression<object>("lastResult().Content");
|
||||
activity.ValueExpression = new JavaScriptExpression<object>("lastResult().Body");
|
||||
}
|
||||
)
|
||||
.Then<SendMassTransitMessage>(activity =>
|
||||
|
|
|
|||
|
|
@ -11,7 +11,9 @@
|
|||
},
|
||||
"MassTransit": {
|
||||
"RabbitMq": {
|
||||
"Host": "rabbitmq://localhost:5672"
|
||||
"Host": "rabbitmq://localhost:5672",
|
||||
"Username": "guest",
|
||||
"Password": "guest"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue