From 2d54b7e33795873eeb797d679929622e00760b53 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 19 Nov 2020 22:21:17 +0100 Subject: [PATCH] Add Rebus activities --- Samples.sln | 14 +++++ .../Activities/MessageReceived.cs | 29 +++++++++ .../Activities/PublishMessage.cs | 34 +++++++++++ .../Activities/SendMessage.cs | 37 ++++++++++++ .../Consumers/MessageConsumer.cs | 29 +++++++++ .../Elsa.Activities.Rebus.csproj | 13 ++++ .../Elsa.Activities.Rebus.csproj.DotSettings | 2 + .../Extensions/ServiceCollectionExtensions.cs | 26 ++++++++ .../Triggers/MessageReceivedTrigger.cs | 22 +++++++ .../Activities/CronEvent/CronEvent.cs | 3 - .../Activities/TimerEvent/TimerEvent.cs | 7 +-- .../Triggers/CronEventTrigger.cs | 1 - .../Triggers/TimerEventTrigger.cs | 24 ++++++-- .../Services/IWorkflowRunner.cs | 7 ++- .../ElsaServiceCollectionExtensions.cs | 13 +++- src/core/Elsa.Core/Services/WorkflowRunner.cs | 44 +++++++++++--- .../StartupTasks/StartServiceBusTask.cs | 17 ++++-- .../Models/WorkflowModel.cs | 1 - .../Services/ButtonDescriptor.cs | 1 - .../Surrogates/VersionOptionsSurrogate.cs | 4 -- .../Workflows/DemoWorkflow.cs | 1 - .../Elsa.Samples.HelloWorldConsole/Program.cs | 1 - .../Elsa.Samples.RebusWorker.csproj | 17 ++++++ .../Messages/Greeting.cs | 9 +++ .../Elsa.Samples.RebusWorker/Program.cs | 50 ++++++++++++++++ .../Properties/launchSettings.json | 11 ++++ .../Workflows/ConsumerWorkflow.cs | 21 +++++++ .../Workflows/ProducerWorkflow.cs | 60 +++++++++++++++++++ .../appsettings.Development.json | 9 +++ .../Elsa.Samples.RebusWorker/appsettings.json | 9 +++ .../WorkflowDefinitions/PostTests.cs | 1 - 31 files changed, 479 insertions(+), 38 deletions(-) create mode 100644 src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs create mode 100644 src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs create mode 100644 src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs create mode 100644 src/activities/Elsa.Activities.Rebus/Consumers/MessageConsumer.cs create mode 100644 src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj create mode 100644 src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj.DotSettings create mode 100644 src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs create mode 100644 src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs create mode 100644 src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj create mode 100644 src/samples/Elsa.Samples.RebusWorker/Messages/Greeting.cs create mode 100644 src/samples/Elsa.Samples.RebusWorker/Program.cs create mode 100644 src/samples/Elsa.Samples.RebusWorker/Properties/launchSettings.json create mode 100644 src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs create mode 100644 src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs create mode 100644 src/samples/Elsa.Samples.RebusWorker/appsettings.Development.json create mode 100644 src/samples/Elsa.Samples.RebusWorker/appsettings.json diff --git a/Samples.sln b/Samples.sln index 448581277..2d8f16dcd 100644 --- a/Samples.sln +++ b/Samples.sln @@ -138,6 +138,10 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "ElsaDashboard.Application.W EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.ForkJoinTimerAndSignalHttp", "src\samples\Elsa.Samples.ForkJoinTimerAndSignalHttp\Elsa.Samples.ForkJoinTimerAndSignalHttp.csproj", "{922F1EB6-5C8F-45DC-82A8-C651E4554542}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Activities.Rebus", "src\activities\Elsa.Activities.Rebus\Elsa.Activities.Rebus.csproj", "{115DEB38-679F-467E-8F87-9E669DFD2EB8}" +EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.RebusWorker", "src\samples\Elsa.Samples.RebusWorker\Elsa.Samples.RebusWorker.csproj", "{A855C6B9-1548-4183-926C-75D80CEEF510}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -337,6 +341,14 @@ Global {922F1EB6-5C8F-45DC-82A8-C651E4554542}.Debug|Any CPU.Build.0 = Debug|Any CPU {922F1EB6-5C8F-45DC-82A8-C651E4554542}.Release|Any CPU.ActiveCfg = Release|Any CPU {922F1EB6-5C8F-45DC-82A8-C651E4554542}.Release|Any CPU.Build.0 = Release|Any CPU + {115DEB38-679F-467E-8F87-9E669DFD2EB8}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {115DEB38-679F-467E-8F87-9E669DFD2EB8}.Debug|Any CPU.Build.0 = Debug|Any CPU + {115DEB38-679F-467E-8F87-9E669DFD2EB8}.Release|Any CPU.ActiveCfg = Release|Any CPU + {115DEB38-679F-467E-8F87-9E669DFD2EB8}.Release|Any CPU.Build.0 = Release|Any CPU + {A855C6B9-1548-4183-926C-75D80CEEF510}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {A855C6B9-1548-4183-926C-75D80CEEF510}.Debug|Any CPU.Build.0 = Debug|Any CPU + {A855C6B9-1548-4183-926C-75D80CEEF510}.Release|Any CPU.ActiveCfg = Release|Any CPU + {A855C6B9-1548-4183-926C-75D80CEEF510}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -402,6 +414,8 @@ Global {8DD4F1E8-8AC8-4AB4-AC08-9456AA6EAF10} = {5837821B-CA71-40B6-A9F1-C25D318B4691} {F7181887-E6B5-4DA9-9598-8E9E806B5A20} = {5837821B-CA71-40B6-A9F1-C25D318B4691} {922F1EB6-5C8F-45DC-82A8-C651E4554542} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} + {115DEB38-679F-467E-8F87-9E669DFD2EB8} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} + {A855C6B9-1548-4183-926C-75D80CEEF510} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158} diff --git a/src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs b/src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs new file mode 100644 index 000000000..b5b735957 --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Activities/MessageReceived.cs @@ -0,0 +1,29 @@ +using System; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Services; +using Elsa.Services.Models; +using Rebus.Bus; + +namespace Elsa.Activities.Rebus +{ + [Trigger(Category = "Rebus", Description = "Triggered when a message is received.", Outcomes = new[] { OutcomeNames.Done })] + public class MessageReceived : Activity + { + private readonly IBus _bus; + + public MessageReceived(IBus bus) + { + _bus = bus; + } + + [ActivityProperty(Hint = "The type of message to receive.")] + public Type MessageType { get; set; } = default!; + + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) + { + var message = context.Input; + return Done(message); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs b/src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs new file mode 100644 index 000000000..d7ce72a13 --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Activities/PublishMessage.cs @@ -0,0 +1,34 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Services; +using Elsa.Services.Models; +using Rebus.Bus; + +namespace Elsa.Activities.Rebus +{ + [Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })] + public class PublishMessage : Activity + { + private readonly IBus _bus; + + public PublishMessage(IBus bus) + { + _bus = bus; + } + + [ActivityProperty(Hint = "The message to publish.")] + public object Message { get; set; } = default!; + + [ActivityProperty(Hint = "Optional headers to send along with the message.")] + public IDictionary? Headers { get; set; } + + protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) + { + await _bus.Publish(Message, Headers); + return Done(); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs b/src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs new file mode 100644 index 000000000..4d7d61f04 --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Activities/SendMessage.cs @@ -0,0 +1,37 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Services; +using Elsa.Services.Models; +using Rebus.Bus; + +namespace Elsa.Activities.Rebus +{ + [Action(Category = "Rebus", Description = "Publishes a message.", Outcomes = new[] { OutcomeNames.Done })] + public class SendMessage : Activity + { + private readonly IBus _bus; + + public SendMessage(IBus bus) + { + _bus = bus; + } + + [ActivityProperty(Hint = "The message to send.")] + public object Message { get; set; } = default!; + + [ActivityProperty(Hint = "Optional headers to send along with the message.")] + public IDictionary? Headers { get; set; } + + protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) + { + await _bus.Advanced.Routing.Send("greeting", Message, Headers); + + await _bus.Advanced.Topics.Subscribe("greeting"); + //await _bus.Send(Message, Headers); + return Done(); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Consumers/MessageConsumer.cs b/src/activities/Elsa.Activities.Rebus/Consumers/MessageConsumer.cs new file mode 100644 index 000000000..d3049f4dc --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Consumers/MessageConsumer.cs @@ -0,0 +1,29 @@ +using System.Threading.Tasks; +using Elsa.Activities.Rebus.Triggers; +using Elsa.Services; +using Rebus.Extensions; +using Rebus.Handlers; +using Rebus.Messages; +using Rebus.Pipeline; + +namespace Elsa.Activities.Rebus.Consumers +{ + public class MessageConsumer : IHandleMessages + { + private readonly IWorkflowScheduler _workflowScheduler; + + public MessageConsumer(IWorkflowScheduler workflowScheduler) + { + _workflowScheduler = workflowScheduler; + } + + public async Task Handle(T message) + { + var correlationId = MessageContext.Current.TransportMessage.Headers.GetValueOrNull(Headers.CorrelationId); + await _workflowScheduler.TriggerWorkflowsAsync( + x => x.MessageType == typeof(T).Name && (x.CorrelationId == null || x.CorrelationId == correlationId), + message, + correlationId); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj b/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj new file mode 100644 index 000000000..8a08de5b4 --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj @@ -0,0 +1,13 @@ + + + + netstandard2.0 + latest + enable + + + + + + + diff --git a/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj.DotSettings b/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj.DotSettings new file mode 100644 index 000000000..7230eb354 --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Elsa.Activities.Rebus.csproj.DotSettings @@ -0,0 +1,2 @@ + + True \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 000000000..9e48e40ff --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,26 @@ +using Elsa.Activities.Rebus.Consumers; +using Elsa.Activities.Rebus.Triggers; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Activities.Rebus.Extensions +{ + public static class ServiceCollectionExtensions + { + public static IServiceCollection AddRebusActivities(this IServiceCollection services) => + services + .AddTriggerProvider() + .AddActivity() + .AddActivity() + .AddActivity(); + + public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType(); + public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType(); + public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType(); + public static IServiceCollection AddRebusActivities(this IServiceCollection services) => services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType().AddMessageType(); + + public static IServiceCollection AddRebusActivities(this IServiceCollection services) => + services.AddRebusActivities().AddMessageType().AddMessageType().AddMessageType().AddMessageType().AddMessageType(); + + public static IServiceCollection AddMessageType(this IServiceCollection services) => services.AddConsumer>(); + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs b/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs new file mode 100644 index 000000000..8fde08bfd --- /dev/null +++ b/src/activities/Elsa.Activities.Rebus/Triggers/MessageReceivedTrigger.cs @@ -0,0 +1,22 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Triggers; + +namespace Elsa.Activities.Rebus.Triggers +{ + public class MessageReceivedTrigger : Trigger + { + public string MessageType { get; set; } = default!; + public string? CorrelationId { get; set; } + } + + public class MessageReceivedTriggerProvider : TriggerProvider + { + public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => + new MessageReceivedTrigger + { + 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.Timers/Activities/CronEvent/CronEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/CronEvent/CronEvent.cs index 6ab4dce92..36cbb9cfa 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/CronEvent/CronEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/CronEvent/CronEvent.cs @@ -34,9 +34,6 @@ namespace Elsa.Activities.Timers protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - if (context.WorkflowExecutionContext.IsFirstPass) - return Done(); - ExecuteAt = GetNextOccurrence(CronExpression); return Suspend(); } diff --git a/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs index c664b406b..d17e97976 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/TimerEvent/TimerEvent.cs @@ -20,17 +20,14 @@ namespace Elsa.Activities.Timers [ActivityProperty(Hint = "An expression that evaluates to a Duration value.")] public Duration Timeout { get; set; } = default!; - public Instant ExecuteAt + public Instant? ExecuteAt { - get => GetState(); + get => GetState(); set => SetState(value); } protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - if (context.WorkflowExecutionContext.IsFirstPass) - return Done(); - ExecuteAt = _clock.GetCurrentInstant().Plus(Timeout); return Suspend(); } diff --git a/src/activities/Elsa.Activities.Timers/Triggers/CronEventTrigger.cs b/src/activities/Elsa.Activities.Timers/Triggers/CronEventTrigger.cs index 9928bbebf..fd8017d31 100644 --- a/src/activities/Elsa.Activities.Timers/Triggers/CronEventTrigger.cs +++ b/src/activities/Elsa.Activities.Timers/Triggers/CronEventTrigger.cs @@ -1,7 +1,6 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Triggers; -using NCrontab; using NodaTime; namespace Elsa.Activities.Timers.Triggers diff --git a/src/activities/Elsa.Activities.Timers/Triggers/TimerEventTrigger.cs b/src/activities/Elsa.Activities.Timers/Triggers/TimerEventTrigger.cs index 37b1ed73e..dbcfe5492 100644 --- a/src/activities/Elsa.Activities.Timers/Triggers/TimerEventTrigger.cs +++ b/src/activities/Elsa.Activities.Timers/Triggers/TimerEventTrigger.cs @@ -1,4 +1,6 @@ -using Elsa.Triggers; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Triggers; using NodaTime; namespace Elsa.Activities.Timers.Triggers @@ -10,13 +12,27 @@ namespace Elsa.Activities.Timers.Triggers public class TimerEventTriggerProvider : TriggerProvider { - public override ITrigger GetTrigger(TriggerProviderContext context) + private readonly IClock _clock; + + public TimerEventTriggerProvider(IClock clock) { - var executeAt = context.GetActivity().GetState(x => x.ExecuteAt); + _clock = clock; + } + + public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) + { + var activity = context.GetActivity(); + var executeAt = activity.GetState(x => x.ExecuteAt); + + if (executeAt == null) + { + var timeout = await activity.GetPropertyValueAsync(x => x.Timeout, cancellationToken); + executeAt = _clock.GetCurrentInstant().Plus(timeout); + } return new TimerEventTrigger { - ExecuteAt = executeAt + ExecuteAt = executeAt.Value }; } } diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs b/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs index 61602bebc..57cdeb4a8 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowRunner.cs @@ -1,13 +1,18 @@ +using System; using System.Threading; using System.Threading.Tasks; using Elsa.Builders; using Elsa.Models; using Elsa.Services.Models; +using Elsa.Triggers; namespace Elsa.Services { public interface IWorkflowRunner { + Task TriggerWorkflowsAsync(Func predicate, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default) + where TTrigger : ITrigger; + ValueTask RunWorkflowAsync( WorkflowInstance workflowInstance, string? activityId = default, @@ -49,7 +54,7 @@ namespace Elsa.Services string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default); - + ValueTask RunWorkflowAsync( IWorkflow workflow, WorkflowInstance workflowInstance, diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index bb0efe59f..f1c361319 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -47,7 +47,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton(options.SignalFactory) .AddSingleton(options.StorageFactory) .AddPersistence(options.ConfigurePersistence); - + options.AddWorkflowsCore(); options.AddMediatR(); options.AddServiceBus(); @@ -70,7 +70,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddTransient() .AddTransient(sp => sp.GetRequiredService()); } - + public static IServiceCollection AddWorkflow(this IServiceCollection services, IWorkflow workflow) { return services @@ -78,7 +78,14 @@ namespace Microsoft.Extensions.DependencyInjection .AddTransient(sp => workflow); } - public static IServiceCollection AddConsumer(this IServiceCollection services) where TConsumer : class, IHandleMessages => services.AddTransient, TConsumer>(); + public static IServiceCollection AddConsumer(this IServiceCollection services) where TConsumer : class, IHandleMessages + { + return services + .AddTransient() + .AddTransient(sp => sp.GetRequiredService()) + .AddTransient, TConsumer>(sp => sp.GetRequiredService()); + } + private static IServiceCollection AddMediatR(this ElsaOptions options) => options.Services.AddMediatR(mediatr => mediatr.AsScoped(), typeof(IActivity)); private static ElsaOptions AddWorkflowsCore(this ElsaOptions configuration) diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index 4fa33430e..bb7d10ce3 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -6,11 +6,14 @@ using Elsa.ActivityResults; using Elsa.Builders; using Elsa.Events; using Elsa.Exceptions; +using Elsa.Extensions; using Elsa.Models; using Elsa.Services.Models; +using Elsa.Triggers; using MediatR; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; +using Open.Linq.AsyncExtensions; namespace Elsa.Services { @@ -23,6 +26,8 @@ namespace Elsa.Services private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowFactory _workflowFactory; + private readonly IWorkflowSelector _workflowSelector; + private readonly IWorkflowInstanceManager _workflowInstanceManager; private readonly Func _workflowBuilderFactory; private readonly IWorkflowContextManager _workflowContextManager; private readonly IMediator _mediator; @@ -32,11 +37,12 @@ namespace Elsa.Services public WorkflowRunner( IWorkflowRegistry workflowRegistry, IWorkflowFactory workflowFactory, + IWorkflowSelector workflowSelector, Func workflowBuilderFactory, IWorkflowContextManager workflowContextManager, IMediator mediator, IServiceProvider serviceProvider, - ILogger logger) + ILogger logger, IWorkflowInstanceManager workflowInstanceManager) { _workflowRegistry = workflowRegistry; _workflowFactory = workflowFactory; @@ -45,6 +51,28 @@ namespace Elsa.Services _mediator = mediator; _serviceProvider = serviceProvider; _logger = logger; + _workflowInstanceManager = workflowInstanceManager; + _workflowSelector = workflowSelector; + } + + public async Task TriggerWorkflowsAsync(Func predicate, object? input = default, string? correlationId = default, string? contextId = default, CancellationToken cancellationToken = default) + where TTrigger : ITrigger + { + var results = await _workflowSelector.SelectWorkflowsAsync(predicate, cancellationToken).ToList(); + + foreach (var result in results) + { + if (result.WorkflowInstanceId != null) + { + var workflowInstance = await _workflowInstanceManager.GetByIdAsync(result.WorkflowInstanceId, cancellationToken); + await RunWorkflowAsync(result.WorkflowBlueprint, workflowInstance, result.ActivityId, input, cancellationToken); + } + else + await RunWorkflowAsync(result.WorkflowBlueprint, result.ActivityId, input, correlationId, contextId, cancellationToken); + + if (result.Trigger.IsOneOff) + await _workflowSelector.RemoveTriggerAsync(result.Trigger, cancellationToken); + } } public async ValueTask RunWorkflowAsync( @@ -146,7 +174,7 @@ namespace Elsa.Services await ResumeWorkflowAsync(workflowExecutionContext, activity!, input, cancellationToken); break; } - + workflowInstance.ContextId = await SaveWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Burst, false, cancellationToken); await _mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken); @@ -173,11 +201,11 @@ namespace Elsa.Services var context = new LoadWorkflowContext(workflowBlueprint, workflowInstance); return await _workflowContextManager.LoadContext(context, cancellationToken); } - + private async ValueTask SaveWorkflowContextAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowContextFidelity fidelity, bool always, CancellationToken cancellationToken) { var workflowContext = workflowExecutionContext.WorkflowContext; - + if (!always && (workflowContext == null || workflowExecutionContext.WorkflowBlueprint.ContextOptions?.ContextFidelity != fidelity)) return workflowExecutionContext.WorkflowInstance.ContextId; @@ -230,16 +258,16 @@ namespace Elsa.Services var serviceProvider = scope.ServiceProvider; var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint; var workflowInstance = workflowExecutionContext.WorkflowInstance; - + while (workflowExecutionContext.HasScheduledActivities) { var scheduledActivity = workflowExecutionContext.PopScheduledActivity(); var currentActivityId = scheduledActivity.ActivityId; var activityBlueprint = workflowBlueprint.GetActivity(currentActivityId)!; - - if(workflowBlueprint.ContextOptions?.ContextFidelity == WorkflowContextFidelity.Activity || activityBlueprint.LoadWorkflowContext) + + if (workflowBlueprint.ContextOptions?.ContextFidelity == WorkflowContextFidelity.Activity || activityBlueprint.LoadWorkflowContext) workflowExecutionContext.WorkflowContext = await LoadWorkflowContextAsync(workflowBlueprint, workflowInstance, WorkflowContextFidelity.Activity, activityBlueprint.LoadWorkflowContext, cancellationToken); - + var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, serviceProvider, activityBlueprint, scheduledActivity.Input); var activity = await activityBlueprint.CreateActivityAsync(activityExecutionContext, cancellationToken); var result = await activityOperation(activityExecutionContext, activity, cancellationToken); diff --git a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs index af04e2c6b..d6d67b115 100644 --- a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs +++ b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs @@ -1,8 +1,10 @@ using System; +using System.Linq; using System.Threading; using System.Threading.Tasks; -using Elsa.Messages; using Elsa.Services; +using Microsoft.Extensions.DependencyInjection; +using Rebus.Handlers; using Rebus.ServiceProvider; namespace Elsa.StartupTasks @@ -11,12 +13,19 @@ namespace Elsa.StartupTasks { private readonly IServiceProvider _serviceProvider; public StartServiceBusTask(IServiceProvider serviceProvider) => _serviceProvider = serviceProvider; - + public Task ExecuteAsync(CancellationToken cancellationToken = default) { - _serviceProvider.UseRebus(x => x.Subscribe()); + var consumers = _serviceProvider.GetServices(); + var messageTypes = consumers.Select(consumer => consumer.GetType().GetInterfaces().First(x => x.GenericTypeArguments.Any()).GenericTypeArguments.First()); + + _serviceProvider.UseRebus(async bus => + { + foreach (var messageType in messageTypes) + await bus.Subscribe(messageType); + }); + return Task.CompletedTask; } } - } \ No newline at end of file diff --git a/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModel.cs b/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModel.cs index eda945e87..cacb2a868 100644 --- a/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModel.cs +++ b/src/dashboards/blazor/ElsaDashboard.Application/Models/WorkflowModel.cs @@ -2,7 +2,6 @@ using System.Collections.Generic; using System.Collections.Immutable; using System.Linq; -using Elsa.Client.Models; namespace ElsaDashboard.Application.Models { diff --git a/src/dashboards/blazor/ElsaDashboard.Application/Services/ButtonDescriptor.cs b/src/dashboards/blazor/ElsaDashboard.Application/Services/ButtonDescriptor.cs index 66bebbe56..0555e7bf6 100644 --- a/src/dashboards/blazor/ElsaDashboard.Application/Services/ButtonDescriptor.cs +++ b/src/dashboards/blazor/ElsaDashboard.Application/Services/ButtonDescriptor.cs @@ -1,6 +1,5 @@ using System; using System.Threading.Tasks; -using ElsaDashboard.Application.Extensions; using ElsaDashboard.Application.Models; using Microsoft.AspNetCore.Components; diff --git a/src/dashboards/blazor/ElsaDashboard.Shared/Surrogates/VersionOptionsSurrogate.cs b/src/dashboards/blazor/ElsaDashboard.Shared/Surrogates/VersionOptionsSurrogate.cs index db7814377..587f7a500 100644 --- a/src/dashboards/blazor/ElsaDashboard.Shared/Surrogates/VersionOptionsSurrogate.cs +++ b/src/dashboards/blazor/ElsaDashboard.Shared/Surrogates/VersionOptionsSurrogate.cs @@ -1,8 +1,4 @@ using Elsa.Client.Models; -using Newtonsoft.Json; -using Newtonsoft.Json.Linq; -using NodaTime; -using NodaTime.Serialization.JsonNet; using ProtoBuf; namespace ElsaDashboard.Shared.Surrogates diff --git a/src/samples/Elsa.Samples.ForkJoinTimerAndSignalHttp/Workflows/DemoWorkflow.cs b/src/samples/Elsa.Samples.ForkJoinTimerAndSignalHttp/Workflows/DemoWorkflow.cs index 3f55e6498..099c9ec20 100644 --- a/src/samples/Elsa.Samples.ForkJoinTimerAndSignalHttp/Workflows/DemoWorkflow.cs +++ b/src/samples/Elsa.Samples.ForkJoinTimerAndSignalHttp/Workflows/DemoWorkflow.cs @@ -3,7 +3,6 @@ using Elsa.Activities.ControlFlow; using Elsa.Activities.Timers; using Elsa.Builders; using Elsa.Services.Models; -using Microsoft.AspNetCore.Mvc.Filters; using NodaTime; namespace Elsa.Samples.ForkJoinTimerAndSignalHttp.Workflows diff --git a/src/samples/Elsa.Samples.HelloWorldConsole/Program.cs b/src/samples/Elsa.Samples.HelloWorldConsole/Program.cs index e459b3847..f4620c64e 100644 --- a/src/samples/Elsa.Samples.HelloWorldConsole/Program.cs +++ b/src/samples/Elsa.Samples.HelloWorldConsole/Program.cs @@ -1,5 +1,4 @@ using System.Threading.Tasks; -using AutoMapper; using Elsa.Extensions; using Elsa.Services; using Microsoft.Extensions.DependencyInjection; diff --git a/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj b/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj new file mode 100644 index 000000000..735d93108 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/Elsa.Samples.RebusWorker.csproj @@ -0,0 +1,17 @@ + + + + net5.0 + dotnet-Elsa.Samples.Rebus.AzureServiceBusWorker-34DD18B8-38B1-4F6F-ABE0-BCE3BEEB5FA6 + + + + + + + + + + + + diff --git a/src/samples/Elsa.Samples.RebusWorker/Messages/Greeting.cs b/src/samples/Elsa.Samples.RebusWorker/Messages/Greeting.cs new file mode 100644 index 000000000..81946a526 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/Messages/Greeting.cs @@ -0,0 +1,9 @@ +namespace Elsa.Samples.RebusWorker.Messages +{ + public class Greeting + { + public string From { get; set; } + public string To { get; set; } + public string Message { get; set; } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.RebusWorker/Program.cs b/src/samples/Elsa.Samples.RebusWorker/Program.cs new file mode 100644 index 000000000..0af51fc43 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/Program.cs @@ -0,0 +1,50 @@ +using System; +using Elsa.Activities.Rebus.Extensions; +using Elsa.Samples.RebusWorker.Messages; +using Elsa.Samples.RebusWorker.Workflows; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using NodaTime; +using Rebus.Config; +using Rebus.DataBus.InMem; +using Rebus.Logging; +using Rebus.Persistence.InMem; +using Rebus.Routing.TypeBased; +using Rebus.Transport.InMem; +using YesSql.Provider.Sqlite; + +namespace Elsa.Samples.RebusWorker +{ + public class Program + { + public static void Main(string[] args) + { + CreateHostBuilder(args).Build().Run(); + } + + public static IHostBuilder CreateHostBuilder(string[] args) => + Host.CreateDefaultBuilder(args) + .ConfigureServices((hostContext, services) => + { + services + .AddElsa(option => option + .UsePersistence(db => db.UseSqLite("Data Source=elsa.db;Cache=Shared")) + .ConfigureServiceBus(ConfigureRebus)) + .AddConsoleActivities() + .AddTimerActivities(options => options.SweepInterval = Duration.FromSeconds(1)) + .AddRebusActivities() + .AddWorkflow() + .AddWorkflow(); + }); + + private static RebusConfigurer ConfigureRebus(RebusConfigurer rebus, IServiceProvider serviceProvider) + { + return rebus + .Logging(logging => logging.ColoredConsole(LogLevel.Info)) + .Subscriptions(s => s.StoreInMemory(new InMemorySubscriberStore())) + .DataBus(s => s.StoreInMemory(new InMemDataStore())) + .Routing(r => r.TypeBased().Map("greeting")) + .Transport(t => t.UseInMemoryTransport(new InMemNetwork(), "inbox")); + } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.RebusWorker/Properties/launchSettings.json b/src/samples/Elsa.Samples.RebusWorker/Properties/launchSettings.json new file mode 100644 index 000000000..a535fa630 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/Properties/launchSettings.json @@ -0,0 +1,11 @@ +{ + "profiles": { + "Elsa.Samples.Rebus.AzureServiceBusWorker": { + "commandName": "Project", + "dotnetRunMessages": "true", + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + } + } + } +} diff --git a/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs b/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs new file mode 100644 index 000000000..dce443023 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/Workflows/ConsumerWorkflow.cs @@ -0,0 +1,21 @@ +using Elsa.Activities.Console; +using Elsa.Activities.Rebus; +using Elsa.Builders; +using Elsa.Samples.RebusWorker.Messages; + +namespace Elsa.Samples.RebusWorker.Workflows +{ + public class ConsumerWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder workflow) + { + workflow + .StartWith(messageReceived => messageReceived.Set(x => x.MessageType, typeof(Greeting))) + .WriteLine(context => + { + var greeting = context.GetInput(); + return $"Received a greeting from {greeting.From}, saying \"{greeting.Message}\" to {greeting.To}!"; + }); + } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs b/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs new file mode 100644 index 000000000..c8bda1957 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/Workflows/ProducerWorkflow.cs @@ -0,0 +1,60 @@ +using System; +using Elsa.Activities.Console; +using Elsa.Activities.Rebus; +using Elsa.Activities.Timers; +using Elsa.Builders; +using Elsa.Samples.RebusWorker.Messages; +using NodaTime; + +namespace Elsa.Samples.RebusWorker.Workflows +{ + public class ProducerWorkflow : IWorkflow + { + private readonly IClock _clock; + private readonly Random _random; + + public ProducerWorkflow(IClock clock) + { + _clock = clock; + _random = new Random(); + } + + public void Build(IWorkflowBuilder workflow) + { + workflow + .TimerEvent(Duration.FromSeconds(5)) + .WriteLine("Sending a random greeting to the \"greetings\" queue.") + .Then(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting)) + //.Then(sendMessage => sendMessage.Set(x => x.Message, GetRandomGreeting)) + .WriteLine(() => $"Message sent at {_clock.GetCurrentInstant()}"); + } + + private Greeting GetRandomGreeting() + { + var greetings = new[] + { + new Greeting + { + From = "John", + To = "Jill", + Message = "Hello!" + }, + new Greeting + { + From = "Julia", + To = "Miriam", + Message = "Happy Monday!" + }, + new Greeting + { + From = "Jack", + To = "Bob", + Message = "How do you do?" + } + }; + + var index = _random.Next(0, greetings.Length); + return greetings[index]; + } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.RebusWorker/appsettings.Development.json b/src/samples/Elsa.Samples.RebusWorker/appsettings.Development.json new file mode 100644 index 000000000..8983e0fc1 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/appsettings.Development.json @@ -0,0 +1,9 @@ +{ + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + } +} diff --git a/src/samples/Elsa.Samples.RebusWorker/appsettings.json b/src/samples/Elsa.Samples.RebusWorker/appsettings.json new file mode 100644 index 000000000..8983e0fc1 --- /dev/null +++ b/src/samples/Elsa.Samples.RebusWorker/appsettings.json @@ -0,0 +1,9 @@ +{ + "Logging": { + "LogLevel": { + "Default": "Information", + "Microsoft": "Warning", + "Microsoft.Hosting.Lifetime": "Information" + } + } +} diff --git a/test/component/Elsa.ComponentTests/Endpoints/WorkflowDefinitions/PostTests.cs b/test/component/Elsa.ComponentTests/Endpoints/WorkflowDefinitions/PostTests.cs index cebeaedc8..a1d15edd8 100644 --- a/test/component/Elsa.ComponentTests/Endpoints/WorkflowDefinitions/PostTests.cs +++ b/test/component/Elsa.ComponentTests/Endpoints/WorkflowDefinitions/PostTests.cs @@ -6,7 +6,6 @@ using AutoFixture; using Elsa.Activities.Console; using Elsa.ComponentTests.Helpers; using Elsa.Models; -using Elsa.Server.Api.Endpoints.WorkflowDefinitions; using Elsa.Testing.Shared.AutoFixture; using Elsa.Testing.Shared.Helpers; using Xunit;