From e51a4e98344a4c23d211ccc55a6c2340ee24a604 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?S=C3=B6ren=20Uhrbach?= Date: Tue, 13 Apr 2021 21:32:44 +0200 Subject: [PATCH] Trigger workflows from masstransit consumers --- .../ReceiveMassTransitMessage.cs | 7 +++-- .../Bookmarks/MessageReceivedBookmark.cs | 26 +++++++++++++++++ .../Consumers/WorkflowConsumer.cs | 29 ++++++++++++++----- .../Extensions/ServiceCollectionExtensions.cs | 4 +++ 4 files changed, 56 insertions(+), 10 deletions(-) create mode 100644 src/activities/Elsa.Activities.MassTransit/Bookmarks/MessageReceivedBookmark.cs diff --git a/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage/ReceiveMassTransitMessage.cs b/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage/ReceiveMassTransitMessage.cs index 6a2836403..55aca73b9 100644 --- a/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage/ReceiveMassTransitMessage.cs +++ b/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage/ReceiveMassTransitMessage.cs @@ -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); diff --git a/src/activities/Elsa.Activities.MassTransit/Bookmarks/MessageReceivedBookmark.cs b/src/activities/Elsa.Activities.MassTransit/Bookmarks/MessageReceivedBookmark.cs new file mode 100644 index 000000000..3d44d2fa5 --- /dev/null +++ b/src/activities/Elsa.Activities.MassTransit/Bookmarks/MessageReceivedBookmark.cs @@ -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 + { + public override async ValueTask> GetBookmarksAsync(BookmarkProviderContext context, CancellationToken cancellationToken) => + new[] + { + new MessageReceivedBookmark + { + MessageType = (await context.Activity.GetPropertyValueAsync(x => x.MessageType, cancellationToken))!.Name, + CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId + } + }; + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs b/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs index b521877cd..d99845508 100644 --- a/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs +++ b/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs @@ -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("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 context) @@ -34,11 +37,21 @@ namespace Elsa.Activities.MassTransit.Consumers if (item.TryGetCorrelationId(message, out correlationId)) break; - throw new NotImplementedException(); - // await _workflowScheduler.TriggerWorkflowsAsync( - // 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() + )); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs index 63e03c183..b76d6d439 100644 --- a/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.MassTransit/Extensions/ServiceCollectionExtensions.cs @@ -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(); + return options .AddActivity() .AddActivity()