Added small example to prove functionality for receiving masstransit messages

This commit is contained in:
Sören Uhrbach 2021-04-13 22:55:11 +02:00 committed by Sipke Schoorstra
parent e51a4e9834
commit 2bc60d9344
7 changed files with 176 additions and 0 deletions

View file

@ -0,0 +1,37 @@
using MassTransit;
using System.Threading.Tasks;
using Microsoft.AspNetCore.Mvc;
using Elsa.Samples.MassTransitRabbitMq.Messages;
namespace Elsa.Samples.MassTransitRabbitMq.Controllers
{
[ApiController]
[Route("trigger-message")]
public class TriggerMessageController : Controller
{
private readonly IPublishEndpoint _publishEndpoint;
public TriggerMessageController(IPublishEndpoint publishEndpoint)
{
_publishEndpoint = publishEndpoint;
}
[Route("first")]
[HttpPost]
public async Task<IActionResult> TriggerFirstMessage()
{
var message = new FirstMessage();
await _publishEndpoint.Publish(message);
return Ok();
}
[Route("second")]
[HttpPost]
public async Task<IActionResult> TriggerSecondMessage()
{
var message = new SecondMessage();
await _publishEndpoint.Publish(message);
return Ok();
}
}
}

View file

@ -0,0 +1,14 @@
<Project Sdk="Microsoft.NET.Sdk.Web">
<PropertyGroup>
<TargetFramework>net5.0</TargetFramework>
<IsPackable>false</IsPackable>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj" />
<ProjectReference Include="..\..\..\activities\Elsa.Activities.Console\Elsa.Activities.Console.csproj" />
<ProjectReference Include="..\..\..\activities\Elsa.Activities.MassTransit\Elsa.Activities.MassTransit.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,14 @@
using System;
namespace Elsa.Samples.MassTransitRabbitMq.Messages
{
public class SecondMessage
{
public Guid CorrelationId { get; private set; }
public SecondMessage()
{
CorrelationId = Guid.Parse("e9ca46dd-36b9-4fc4-b7db-3bb7190e4488");
}
}
}

View file

@ -0,0 +1,14 @@
using System;
namespace Elsa.Samples.MassTransitRabbitMq.Messages
{
public class FirstMessage
{
public Guid CorrelationId { get; private set; }
public FirstMessage()
{
CorrelationId = Guid.Parse("e9ca46dd-36b9-4fc4-b7db-3bb7190e4488");
}
}
}

View file

@ -0,0 +1,14 @@
using Microsoft.AspNetCore.Hosting;
using Microsoft.Extensions.Hosting;
namespace Elsa.Samples.MassTransitRabbitMq
{
public class Program
{
public static void Main(string[] args) => CreateHostBuilder(args).Build().Run();
public static IHostBuilder CreateHostBuilder(string[] args) =>
Host.CreateDefaultBuilder(args)
.ConfigureWebHostDefaults(webBuilder => { webBuilder.UseStartup<Startup>(); });
}
}

View file

@ -0,0 +1,54 @@
using System;
using MassTransit;
using Elsa.Samples.MassTransitRabbitMq.Messages;
using Elsa.Samples.MassTransitRabbitMq.Workflows;
using Microsoft.AspNetCore.Builder;
using Microsoft.Extensions.DependencyInjection;
using Elsa.Activities.MassTransit.Extensions;
using Elsa.Activities.MassTransit.Consumers;
namespace Elsa.Samples.MassTransitRabbitMq
{
public class Startup
{
private Type CreateWorkflowConsumer(Type messageType) => typeof(WorkflowConsumer<>).MakeGenericType(messageType);
public void ConfigureServices(IServiceCollection services)
{
services.AddControllers();
services
.AddMassTransit(x =>
{
// Add workflow consumer for message
x.AddConsumer(CreateWorkflowConsumer(typeof(FirstMessage)));
x.AddConsumer(CreateWorkflowConsumer(typeof(SecondMessage)));
// Configure rabbitmq
x.UsingRabbitMq((ctx, cfg) =>
{
cfg.ConfigureEndpoints(ctx);
cfg.Host("rabbitmq://guest:guest@localhost");
});
})
.AddMassTransitHostedService();
services
.AddElsa(options => options
.AddConsoleActivities()
.AddMassTransitActivities()
.AddWorkflow<TestWorkflow>()
);
}
public void Configure(IApplicationBuilder app)
{
app.UseRouting();
app.UseEndpoints(endpoints =>
{
endpoints.MapControllers();
});
}
}
}

View file

@ -0,0 +1,29 @@
using Elsa.Activities.Console;
using Elsa.Activities.MassTransit;
using Elsa.Activities.Primitives;
using Elsa.Builders;
using Elsa.Samples.MassTransitRabbitMq.Messages;
using Elsa.Activities.ControlFlow;
using System;
namespace Elsa.Samples.MassTransitRabbitMq.Workflows
{
public class TestWorkflow : IWorkflow
{
public void Build(IWorkflowBuilder builder)
{
builder
.ReceiveMassTransitMessage(
activity => activity.Set(x => x.MessageType, x => typeof(FirstMessage))
)
.Correlate(activity => activity.Set(x => x.Value, context => context.GetInput<FirstMessage>().CorrelationId.ToString()))
.WriteLine(context => $"Received first message")
// Wait until second message received with the same correlation id
.ReceiveMassTransitMessage(
activity => activity.Set(x => x.MessageType, x => typeof(SecondMessage))
)
.WriteLine(context => $"Received second message");
}
}
}