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()