Trigger workflows from masstransit consumers

This commit is contained in:
Sören Uhrbach 2021-04-13 21:32:44 +02:00 committed by Sipke Schoorstra
parent 5825ce6bc1
commit e51a4e9834
4 changed files with 56 additions and 10 deletions

View file

@ -18,9 +18,12 @@ namespace Elsa.Activities.MassTransit
public Type? MessageType { get; set; }
protected override bool OnCanExecute(ActivityExecutionContext context) => context.Input?.GetType() == MessageType;
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => Suspend();
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context)
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => context.WorkflowExecutionContext.IsFirstPass ? ExecuteInternal(context) : Suspend();
protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) => ExecuteInternal(context);
private IActivityExecutionResult ExecuteInternal(ActivityExecutionContext context)
{
var message = context.Input;
return Done(message);

View file

@ -0,0 +1,26 @@
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Bookmarks;
namespace Elsa.Activities.MassTransit.Bookmarks
{
public class MessageReceivedBookmark : IBookmark
{
public string MessageType { get; set; } = default!;
public string? CorrelationId { get; set; }
}
public class MessageReceivedTriggerProvider : BookmarkProvider<MessageReceivedBookmark, ReceiveMassTransitMessage>
{
public override async ValueTask<IEnumerable<IBookmark>> GetBookmarksAsync(BookmarkProviderContext<ReceiveMassTransitMessage> context, CancellationToken cancellationToken) =>
new[]
{
new MessageReceivedBookmark
{
MessageType = (await context.Activity.GetPropertyValueAsync(x => x.MessageType, cancellationToken))!.Name,
CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId
}
};
}
}

View file

@ -1,9 +1,12 @@
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Elsa.Activities.MassTransit.Bookmarks;
using Elsa.Activities.MassTransit.Consumers.MessageCorrelation;
using Elsa.Dispatch;
using Elsa.Services;
using MassTransit;
using MediatR;
namespace Elsa.Activities.MassTransit.Consumers
{
@ -18,11 +21,11 @@ namespace Elsa.Activities.MassTransit.Consumers
new PropertyCorrelationIdSelector<T>("CommandId")
};
private readonly IWorkflowRunner _workflowRunner;
private readonly IMediator _mediator;
public WorkflowConsumer(IWorkflowRunner workflowRunner)
public WorkflowConsumer(IMediator mediator)
{
_workflowRunner = workflowRunner;
_mediator = mediator;
}
public async Task Consume(ConsumeContext<T> context)
@ -34,11 +37,21 @@ namespace Elsa.Activities.MassTransit.Consumers
if (item.TryGetCorrelationId(message, out correlationId))
break;
throw new NotImplementedException();
// await _workflowScheduler.TriggerWorkflowsAsync<ReceiveMassTransitMessage>(
// message,
// correlationId?.ToString(),
// cancellationToken: context.CancellationToken);
var bookmark = new MessageReceivedBookmark {
MessageType = message.GetType().Name,
CorrelationId = correlationId.ToString()
};
var trigger = new MessageReceivedBookmark {
MessageType = message.GetType().Name
};
await _mediator.Send(new TriggerWorkflowsRequest(
nameof(ReceiveMassTransitMessage),
bookmark,
trigger,
message,
correlationId.ToString()
));
}
}
}

View file

@ -1,7 +1,9 @@
using System;
using System.Collections.Generic;
using Elsa.Activities.MassTransit.Bookmarks;
using Elsa.Activities.MassTransit.Consumers;
using Elsa.Activities.MassTransit.Options;
using Elsa.Bookmarks;
using MassTransit;
using MassTransit.ConsumeConfigurators;
using Microsoft.Extensions.DependencyInjection;
@ -17,6 +19,8 @@ namespace Elsa.Activities.MassTransit.Extensions
throw new ArgumentNullException(nameof(options));
}
options.Services.AddBookmarkProvider<MessageReceivedTriggerProvider>();
return options
.AddActivity<PublishMassTransitMessage>()
.AddActivity<ReceiveMassTransitMessage>()