From 3a9eaeb7997968f054df9fbacfe9d64611f88a78 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 3 Nov 2020 16:58:00 +0100 Subject: [PATCH] fix Signal activities, triggers and add sample --- Samples.sln | 7 +++ .../Builders/IActivityBuilder.cs | 1 - .../Builders/ICompositeActivityBuilder.cs | 2 +- .../Builders/IOutcomeBuilder.cs | 2 + .../Services/IActivityActivator.cs | 2 + .../Models/ActivityExecutionContext.cs | 11 +++- .../Models/WorkflowExecutionContext.cs | 13 +++-- .../TriggeredSignal.cs => Models/Signal.cs} | 7 ++- .../ReceiveSignal.cs} | 17 ++++--- .../ReceiveSignal/ReceiveSignalExtensions.cs | 16 ++++++ .../ReceiveSignal/ReceiveSignalTrigger.cs | 23 +++++++++ .../Signaling/Services/ISignaler.cs | 10 ++++ .../Activities/Signaling/Services/Signaler.cs | 25 +++++++++ .../Signaled/SignaledBuilderExtensions.cs | 16 ------ .../Signaling/Signaled/SignaledTrigger.cs | 21 -------- .../Signaling/TriggerEvent/TriggerEvent.cs | 48 ----------------- .../{TriggerSignal.cs => SendSignal.cs} | 25 +++------ .../Elsa.Core/Builders/ActivityBuilder.cs | 3 +- src/core/Elsa.Core/Builders/OutcomeBuilder.cs | 20 +++++--- .../Elsa.Core/Elsa.Core.csproj.DotSettings | 1 + .../ElsaServiceCollectionExtensions.cs | 9 ++-- .../Elsa.Core/Services/ActivityActivator.cs | 13 +++++ .../Services/WorkflowBlueprintMaterializer.cs | 2 +- .../Workflows/WalkAroundWorkflow.cs | 1 - .../Elsa.Samples.SignalingConsole.csproj | 16 ++++++ .../Elsa.Samples.SignalingConsole/Program.cs | 51 +++++++++++++++++++ .../TrafficLightWorkflow.cs | 20 ++++++++ .../Workflows/ForEachWorkflow.cs | 2 +- .../Workflows/ForkJoinWorkflow.cs | 6 +-- .../Workflows/ForkJoinWorkflowTests.cs | 15 ++++-- .../Workflows/ParallelForEachWorkflow.cs | 2 +- 31 files changed, 263 insertions(+), 144 deletions(-) rename src/core/Elsa.Core/Activities/Signaling/{TriggerSignal/TriggeredSignal.cs => Models/Signal.cs} (53%) rename src/core/Elsa.Core/Activities/Signaling/{Signaled/Signaled.cs => ReceiveSignal/ReceiveSignal.cs} (62%) create mode 100644 src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalExtensions.cs create mode 100644 src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalTrigger.cs create mode 100644 src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs create mode 100644 src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs delete mode 100644 src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledBuilderExtensions.cs delete mode 100644 src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledTrigger.cs delete mode 100644 src/core/Elsa.Core/Activities/Signaling/TriggerEvent/TriggerEvent.cs rename src/core/Elsa.Core/Activities/Signaling/TriggerSignal/{TriggerSignal.cs => SendSignal.cs} (59%) create mode 100644 src/samples/Elsa.Samples.SignalingConsole/Elsa.Samples.SignalingConsole.csproj create mode 100644 src/samples/Elsa.Samples.SignalingConsole/Program.cs create mode 100644 src/samples/Elsa.Samples.SignalingConsole/TrafficLightWorkflow.cs diff --git a/Samples.sln b/Samples.sln index 26cff5fa2..c4f4d74e5 100644 --- a/Samples.sln +++ b/Samples.sln @@ -116,6 +116,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.DeclarativeCom EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.GoBackConsole", "src\samples\Elsa.Samples.GoBackConsole\Elsa.Samples.GoBackConsole.csproj", "{F4454030-8295-4575-A54C-D816235173C5}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.SignalingConsole", "src\samples\Elsa.Samples.SignalingConsole\Elsa.Samples.SignalingConsole.csproj", "{3529F81C-9256-4C24-A1E6-ADC27093AB49}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -279,6 +281,10 @@ Global {F4454030-8295-4575-A54C-D816235173C5}.Debug|Any CPU.Build.0 = Debug|Any CPU {F4454030-8295-4575-A54C-D816235173C5}.Release|Any CPU.ActiveCfg = Release|Any CPU {F4454030-8295-4575-A54C-D816235173C5}.Release|Any CPU.Build.0 = Release|Any CPU + {3529F81C-9256-4C24-A1E6-ADC27093AB49}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {3529F81C-9256-4C24-A1E6-ADC27093AB49}.Debug|Any CPU.Build.0 = Debug|Any CPU + {3529F81C-9256-4C24-A1E6-ADC27093AB49}.Release|Any CPU.ActiveCfg = Release|Any CPU + {3529F81C-9256-4C24-A1E6-ADC27093AB49}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -333,6 +339,7 @@ Global {CC39E9F8-72D2-4A5B-846B-E2764DFE1C19} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} {E422689B-9F82-45B6-BE14-62B4EEF43B57} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} {F4454030-8295-4575-A54C-D816235173C5} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} + {3529F81C-9256-4C24-A1E6-ADC27093AB49} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158} diff --git a/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs b/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs index c13740a71..3aa059a18 100644 --- a/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/IActivityBuilder.cs @@ -23,6 +23,5 @@ namespace Elsa.Builders IActivityBuilder WithId(string? id); IActivityBuilder WithName(string? name); Func> BuildActivityAsync(); - //IWorkflowBlueprint Build(); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Builders/ICompositeActivityBuilder.cs b/src/core/Elsa.Abstractions/Builders/ICompositeActivityBuilder.cs index 2f689af30..0bf3a826a 100644 --- a/src/core/Elsa.Abstractions/Builders/ICompositeActivityBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/ICompositeActivityBuilder.cs @@ -27,7 +27,7 @@ namespace Elsa.Builders where T : class, IActivity; IActivityBuilder Add( - Action>? setup = default, + Action>? setup, Action? branch = default) where T : class, IActivity; IActivityBuilder Add( diff --git a/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs b/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs index 82440ee99..1bcd2054c 100644 --- a/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs +++ b/src/core/Elsa.Abstractions/Builders/IOutcomeBuilder.cs @@ -1,3 +1,4 @@ +using System; using Elsa.Services.Models; namespace Elsa.Builders @@ -8,6 +9,7 @@ namespace Elsa.Builders IActivityBuilder Source { get; } string? Outcome { get; } IConnectionBuilder Then(string activityName); + IConnectionBuilder Then(IActivityBuilder targetActivity, Action? branch = default); IWorkflowBlueprint Build(); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IActivityActivator.cs b/src/core/Elsa.Abstractions/Services/IActivityActivator.cs index b3eca2421..0f922d8a6 100644 --- a/src/core/Elsa.Abstractions/Services/IActivityActivator.cs +++ b/src/core/Elsa.Abstractions/Services/IActivityActivator.cs @@ -1,12 +1,14 @@ using System; using System.Collections.Generic; using Elsa.Models; +using Elsa.Services.Models; namespace Elsa.Services { public interface IActivityActivator { IActivity ActivateActivity(string activityTypeName, Action? setup = default); + IActivity ActivateActivity(IActivityBlueprint activityBlueprint); T ActivateActivity(Action? configure = default) where T : class, IActivity; IActivity ActivateActivity(ActivityDefinition activityDefinition); IEnumerable GetActivityTypes(); diff --git a/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs index a8885776b..e2ceeab39 100644 --- a/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/ActivityExecutionContext.cs @@ -43,6 +43,13 @@ namespace Elsa.Services.Models public T GetVariable() => GetVariable(typeof(T).Name); public T GetService() => WorkflowExecutionContext.ServiceProvider.GetService(); + public async ValueTask ActivateActivityAsync(CancellationToken cancellationToken = default) + { + var activity = ActivateActivity(); + await SetActivityPropertiesAsync(activity, cancellationToken); + return activity; + } + public async ValueTask SetActivityPropertiesAsync( IActivity activity, CancellationToken cancellationToken = default) => @@ -51,10 +58,10 @@ namespace Elsa.Services.Models this, cancellationToken); - public IActivity ActivateActivity(string activityType, Action? setupActivity = default) + public IActivity ActivateActivity() { var activityActivator = ServiceProvider.GetRequiredService(); - var activity = activityActivator.ActivateActivity(activityType, setupActivity); + var activity = activityActivator.ActivateActivity(ActivityBlueprint); activity.Data = ActivityInstance.Data; return activity; } diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index 0f396cd32..b50444631 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -1,11 +1,12 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Threading; +using System.Threading.Tasks; using Elsa.Models; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Localization; using Newtonsoft.Json; -using Newtonsoft.Json.Linq; namespace Elsa.Services.Models { @@ -15,8 +16,8 @@ namespace Elsa.Services.Models IServiceProvider serviceProvider, IWorkflowBlueprint workflowBlueprint, WorkflowInstance workflowInstance, - object? input, - object? workflowContext + object? input = default, + object? workflowContext = default ) { ServiceProvider = serviceProvider; @@ -114,5 +115,11 @@ namespace Elsa.Services.Models public T GetOutputFrom(string activityName) => (T)GetOutputFrom(activityName)!; public void SetWorkflowContext(object? value) => WorkflowContext = value; public T GetWorkflowContext() => (T)WorkflowContext!; + + public async ValueTask> ActivateActivitiesAsync(CancellationToken cancellationToken = default) + { + var activityExecutionContexts = WorkflowBlueprint.Activities.Select(x => new ActivityExecutionContext(this, ServiceProvider, x)); + return await Task.WhenAll(activityExecutionContexts.Select(async x => await x.ActivateActivityAsync(cancellationToken))); + } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/TriggerSignal/TriggeredSignal.cs b/src/core/Elsa.Core/Activities/Signaling/Models/Signal.cs similarity index 53% rename from src/core/Elsa.Core/Activities/Signaling/TriggerSignal/TriggeredSignal.cs rename to src/core/Elsa.Core/Activities/Signaling/Models/Signal.cs index d601423ed..fd9f79fa4 100644 --- a/src/core/Elsa.Core/Activities/Signaling/TriggerSignal/TriggeredSignal.cs +++ b/src/core/Elsa.Core/Activities/Signaling/Models/Signal.cs @@ -1,9 +1,8 @@ -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.Signaling +namespace Elsa.Activities.Signaling.Models { - public class TriggeredSignal + public class Signal { - public TriggeredSignal(string signalName, object? input) + public Signal(string signalName, object? input = default) { SignalName = signalName; Input = input; diff --git a/src/core/Elsa.Core/Activities/Signaling/Signaled/Signaled.cs b/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignal.cs similarity index 62% rename from src/core/Elsa.Core/Activities/Signaling/Signaled/Signaled.cs rename to src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignal.cs index 049881eec..b79751013 100644 --- a/src/core/Elsa.Core/Activities/Signaling/Signaled/Signaled.cs +++ b/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignal.cs @@ -1,4 +1,5 @@ using System; +using Elsa.Activities.Signaling.Models; using Elsa.ActivityResults; using Elsa.Attributes; using Elsa.Services; @@ -8,29 +9,31 @@ using Elsa.Services.Models; namespace Elsa.Activities.Signaling { /// - /// Halts workflow execution until the specified signal is received. + /// Suspends workflow execution until the specified signal is received. /// [ActivityDefinition( Category = "Workflows", Description = "Halt workflow execution until the specified signal is received.", Icon = "fas fa-traffic-light" )] - public class Signaled : Activity + public class ReceiveSignal : Activity { - [ActivityProperty(Hint = "An expression that evaluates to the name of the signal to wait for.")] + [ActivityProperty(Hint = "The name of the signal to wait for.")] public string Signal { get; set; } = default!; protected override bool OnCanExecute(ActivityExecutionContext context) { - var signal = Signal; - var triggeredSignal = (TriggeredSignal)context.Input!; - return string.Equals(triggeredSignal.SignalName, signal, StringComparison.OrdinalIgnoreCase); + if (context.Input is Signal triggeredSignal) + return string.Equals(triggeredSignal.SignalName, Signal, StringComparison.OrdinalIgnoreCase); + + return false; } protected override IActivityExecutionResult OnExecute() => Suspend(); + protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) { - var triggeredSignal = (TriggeredSignal)context.Input!; + var triggeredSignal = context.GetInput(); return Done(triggeredSignal.Input); } } diff --git a/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalExtensions.cs b/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalExtensions.cs new file mode 100644 index 000000000..cf42447c4 --- /dev/null +++ b/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalExtensions.cs @@ -0,0 +1,16 @@ +using System; +using Elsa.Activities.Signaling; +using Elsa.Builders; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.ControlFlow +{ + public static class ReceiveSignalExtensions + { + public static IActivityBuilder ReceiveSignal(this IBuilder builder, Action> setup) => builder.Then(setup); + public static IActivityBuilder ReceiveSignal(this IBuilder builder, Func signal) => builder.ReceiveSignal(activity => activity.Set(x => x.Signal, signal)); + public static IActivityBuilder ReceiveSignal(this IBuilder builder, Func signal) => builder.ReceiveSignal(activity => activity.Set(x => x.Signal, signal)); + public static IActivityBuilder ReceiveSignal(this IBuilder builder, string signal) => builder.ReceiveSignal(activity => activity.Set(x => x.Signal, signal)); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalTrigger.cs b/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalTrigger.cs new file mode 100644 index 000000000..33f536513 --- /dev/null +++ b/src/core/Elsa.Core/Activities/Signaling/ReceiveSignal/ReceiveSignalTrigger.cs @@ -0,0 +1,23 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Triggers; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.Signaling +{ + public class ReceiveSignalTrigger : Trigger + { + public string Signal { get; set; } = default!; + public string? CorrelationId { get; set; } + } + + public class ReceiveSignalTriggerProvider : TriggerProvider + { + public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => + new ReceiveSignalTrigger + { + Signal = await context.Activity.GetPropertyValueAsync(x => x.Signal, cancellationToken), + CorrelationId = context.ActivityExecutionContext.WorkflowExecutionContext.CorrelationId + }; + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs b/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs new file mode 100644 index 000000000..95542942f --- /dev/null +++ b/src/core/Elsa.Core/Activities/Signaling/Services/ISignaler.cs @@ -0,0 +1,10 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Elsa.Activities.Signaling.Services +{ + public interface ISignaler + { + Task SendSignal(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs new file mode 100644 index 000000000..8d423f01b --- /dev/null +++ b/src/core/Elsa.Core/Activities/Signaling/Services/Signaler.cs @@ -0,0 +1,25 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.Signaling.Models; +using Elsa.Services; + +namespace Elsa.Activities.Signaling.Services +{ + public class Signaler : ISignaler + { + private readonly IWorkflowScheduler _workflowScheduler; + + public Signaler(IWorkflowScheduler workflowScheduler) + { + _workflowScheduler = workflowScheduler; + } + + public async Task SendSignal(string signal, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default) => + await _workflowScheduler.TriggerWorkflowsAsync( + x => x.Signal == signal && (x.CorrelationId == null || x.CorrelationId == correlationId), + new Signal(signal, input), + correlationId, + cancellationToken: cancellationToken + ); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledBuilderExtensions.cs b/src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledBuilderExtensions.cs deleted file mode 100644 index 069bd95c3..000000000 --- a/src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledBuilderExtensions.cs +++ /dev/null @@ -1,16 +0,0 @@ -using System; -using Elsa.Activities.Signaling; -using Elsa.Builders; -using Elsa.Services.Models; - -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.ControlFlow -{ - public static class SignaledBuilderExtensions - { - public static IActivityBuilder Signaled(this IBuilder builder, Action> setup) => builder.Then(setup); - public static IActivityBuilder Signaled(this IBuilder builder, Func signal) => builder.Signaled(activity => activity.Set(x => x.Signal, signal)); - public static IActivityBuilder Signaled(this IBuilder builder, Func signal) => builder.Signaled(activity => activity.Set(x => x.Signal, signal)); - public static IActivityBuilder Signaled(this IBuilder builder, string signal) => builder.Signaled(activity => activity.Set(x => x.Signal, signal)); - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledTrigger.cs b/src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledTrigger.cs deleted file mode 100644 index 65670e84a..000000000 --- a/src/core/Elsa.Core/Activities/Signaling/Signaled/SignaledTrigger.cs +++ /dev/null @@ -1,21 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Triggers; - -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.Signaling -{ - public class SignaledTrigger : Trigger - { - public string Signal { get; set; } = default!; - } - - public class SignaledTriggerProvider : TriggerProvider - { - public override async ValueTask GetTriggerAsync(TriggerProviderContext context, CancellationToken cancellationToken) => - new SignaledTrigger - { - Signal = await context.Activity.GetPropertyValueAsync(x => x.Signal, cancellationToken) - }; - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/TriggerEvent/TriggerEvent.cs b/src/core/Elsa.Core/Activities/Signaling/TriggerEvent/TriggerEvent.cs deleted file mode 100644 index 3bec148a0..000000000 --- a/src/core/Elsa.Core/Activities/Signaling/TriggerEvent/TriggerEvent.cs +++ /dev/null @@ -1,48 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.ActivityResults; -using Elsa.Attributes; -using Elsa.Services; -using Elsa.Services.Models; - -// ReSharper disable once CheckNamespace -namespace Elsa.Activities.Signaling -{ - [ActivityDefinition( - Category = "Workflows", - Description = "Trigger all workflows that start with or are blocked on the specified activity type.", - Icon = "fas fa-sitemap" - )] - public class TriggerEvent : Activity - { - private readonly IWorkflowScheduler _workflowScheduler; - - public TriggerEvent(IWorkflowScheduler workflowScheduler) - { - _workflowScheduler = workflowScheduler; - } - - [ActivityProperty(Hint = "An expression that evaluates to the activity type to use when triggering workflows.")] - public string ActivityType { get; set; } = default!; - - [ActivityProperty( - Hint = "An expression that evaluates to a dictionary to be provided as input when triggering workflows." - )] - public object? Input { get; set; } - - [ActivityProperty(Hint = "An expression that evaluates to the correlation ID to use when triggering workflows.")] - public string? CorrelationId { get; set; } - - protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) - { - await _workflowScheduler.TriggerWorkflowsAsync( - ActivityType, - Input, - CorrelationId, - cancellationToken: cancellationToken - ); - - return Done(); - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Signaling/TriggerSignal/TriggerSignal.cs b/src/core/Elsa.Core/Activities/Signaling/TriggerSignal/SendSignal.cs similarity index 59% rename from src/core/Elsa.Core/Activities/Signaling/TriggerSignal/TriggerSignal.cs rename to src/core/Elsa.Core/Activities/Signaling/TriggerSignal/SendSignal.cs index 04a9257b5..977cd05e4 100644 --- a/src/core/Elsa.Core/Activities/Signaling/TriggerSignal/TriggerSignal.cs +++ b/src/core/Elsa.Core/Activities/Signaling/TriggerSignal/SendSignal.cs @@ -1,5 +1,6 @@ using System.Threading; using System.Threading.Tasks; +using Elsa.Activities.Signaling.Services; using Elsa.ActivityResults; using Elsa.Attributes; using Elsa.Services; @@ -13,16 +14,16 @@ namespace Elsa.Activities.Signaling /// [ActivityDefinition( Category = "Workflows", - Description = "Trigger the specified signal.", + Description = "Sends the specified signal.", Icon = "fas fa-broadcast-tower" )] - public class TriggerSignal : Activity + public class SendSignal : Activity { - private readonly IWorkflowScheduler _workflowScheduler; + private readonly ISignaler _signaler; - public TriggerSignal(IWorkflowScheduler workflowScheduler) + public SendSignal(ISignaler signaler) { - _workflowScheduler = workflowScheduler; + _signaler = signaler; } [ActivityProperty(Hint = "An expression that evaluates to the name of the signal to trigger.")] @@ -34,19 +35,9 @@ namespace Elsa.Activities.Signaling [ActivityProperty(Hint = "An expression that evaluates to an input value when triggering the signal.")] public object? Input { get; set; } - protected override async ValueTask OnExecuteAsync( - ActivityExecutionContext context, - CancellationToken cancellationToken) + protected override async ValueTask OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) { - var triggeredSignal = new TriggeredSignal(Signal, Input); - - await _workflowScheduler.TriggerWorkflowsAsync( - nameof(Signaled), - triggeredSignal, - CorrelationId, - cancellationToken: cancellationToken - ); - + await _signaler.SendSignal(Signal, Input, CorrelationId, cancellationToken); return Done(); } } diff --git a/src/core/Elsa.Core/Builders/ActivityBuilder.cs b/src/core/Elsa.Core/Builders/ActivityBuilder.cs index 323ed14a5..e83c6fe91 100644 --- a/src/core/Elsa.Core/Builders/ActivityBuilder.cs +++ b/src/core/Elsa.Core/Builders/ActivityBuilder.cs @@ -75,11 +75,10 @@ namespace Elsa.Builders public Func> BuildActivityAsync() => async (context, cancellationToken) => { - var activity = context.ActivateActivity(context.ActivityBlueprint.Type); + var activity = await context.ActivateActivityAsync(cancellationToken); activity.Id = ActivityId; activity.Name = Name; activity.Description = Description; - await context.SetActivityPropertiesAsync(activity, cancellationToken); return activity; }; } diff --git a/src/core/Elsa.Core/Builders/OutcomeBuilder.cs b/src/core/Elsa.Core/Builders/OutcomeBuilder.cs index 10f664b78..2749fbd4f 100644 --- a/src/core/Elsa.Core/Builders/OutcomeBuilder.cs +++ b/src/core/Elsa.Core/Builders/OutcomeBuilder.cs @@ -20,11 +20,20 @@ namespace Elsa.Builders public IActivityBuilder Then( Action>? setup = default, - Action? branch = default) where T : class, IActivity => - Then(WorkflowBuilder.Add(setup), branch); + Action? branch = default) where T : class, IActivity + { + var activityBuilder = WorkflowBuilder.Add(setup); + Then(activityBuilder, branch); + return activityBuilder; + } public IActivityBuilder Then(Action? branch = default) - where T : class, IActivity => Then(WorkflowBuilder.Add(branch)); + where T : class, IActivity + { + var activityBuilder = WorkflowBuilder.Add(branch); + Then(activityBuilder); + return activityBuilder; + } public IConnectionBuilder Then(string activityName) { @@ -34,11 +43,10 @@ namespace Elsa.Builders Outcome); } - private IActivityBuilder Then(IActivityBuilder activityBuilder, Action? branch = default) + public IConnectionBuilder Then(IActivityBuilder activityBuilder, Action? branch = default) { branch?.Invoke(activityBuilder); - WorkflowBuilder.Connect(Source, activityBuilder, Outcome); - return activityBuilder; + return WorkflowBuilder.Connect(Source, activityBuilder, Outcome); } public IWorkflowBlueprint Build() => ((IWorkflowBuilder)WorkflowBuilder).BuildBlueprint(); diff --git a/src/core/Elsa.Core/Elsa.Core.csproj.DotSettings b/src/core/Elsa.Core/Elsa.Core.csproj.DotSettings index 25daf1ffb..af314d94e 100644 --- a/src/core/Elsa.Core/Elsa.Core.csproj.DotSettings +++ b/src/core/Elsa.Core/Elsa.Core.csproj.DotSettings @@ -1,3 +1,4 @@  True + True True \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index eae9b8bfa..07c6bc633 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -4,6 +4,7 @@ using Elsa; using Elsa.Activities.ControlFlow; using Elsa.Activities.Primitives; using Elsa.Activities.Signaling; +using Elsa.Activities.Signaling.Services; using Elsa.Activities.Workflows; using Elsa.Builders; using Elsa.Consumers; @@ -149,10 +150,10 @@ namespace Microsoft.Extensions.DependencyInjection .AddActivity() .AddActivity() .AddActivity() - .AddActivity() - .AddTriggerProvider() - .AddActivity() - .AddActivity() + .AddActivity() + .AddTriggerProvider() + .AddActivity() + .AddScoped() .AddActivity() .AddTriggerProvider(); diff --git a/src/core/Elsa.Core/Services/ActivityActivator.cs b/src/core/Elsa.Core/Services/ActivityActivator.cs index 4bbe67c2c..f50e7e352 100644 --- a/src/core/Elsa.Core/Services/ActivityActivator.cs +++ b/src/core/Elsa.Core/Services/ActivityActivator.cs @@ -2,6 +2,7 @@ using System; using System.Collections.Generic; using System.Linq; using Elsa.Models; +using Elsa.Services.Models; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Services @@ -48,6 +49,18 @@ namespace Elsa.Services return activity; } + public IActivity ActivateActivity(IActivityBlueprint activityBlueprint) + { + return ActivateActivity( + activityBlueprint.Type, + activity => + { + activity.Id = activityBlueprint.Id; + activity.Name = activityBlueprint.Name; + activity.PersistWorkflow = activityBlueprint.PersistWorkflow; + }); + } + public T ActivateActivity(Action? setup = null) where T : class, IActivity { var activity = ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider); diff --git a/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs b/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs index 5f821f84a..d0e941592 100644 --- a/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs +++ b/src/core/Elsa.Core/Services/WorkflowBlueprintMaterializer.cs @@ -98,7 +98,7 @@ namespace Elsa.Services private static async ValueTask CreateActivityAsync(ActivityDefinition activityDefinition, ActivityExecutionContext context, CancellationToken cancellationToken) { - var activity = context.ActivateActivity(activityDefinition.Type); + var activity = context.ActivateActivity(); activity.Description = activityDefinition.Description; activity.Id = activityDefinition.ActivityId; activity.Name = activityDefinition.Name; diff --git a/src/samples/Elsa.Samples.GoBackConsole/Workflows/WalkAroundWorkflow.cs b/src/samples/Elsa.Samples.GoBackConsole/Workflows/WalkAroundWorkflow.cs index 7057e97d6..1d4e28cdb 100644 --- a/src/samples/Elsa.Samples.GoBackConsole/Workflows/WalkAroundWorkflow.cs +++ b/src/samples/Elsa.Samples.GoBackConsole/Workflows/WalkAroundWorkflow.cs @@ -2,7 +2,6 @@ using Elsa.Activities.ControlFlow; using Elsa.Builders; using Elsa.Samples.GoBackConsole.Activities; -using Elsa.Services.Models; namespace Elsa.Samples.GoBackConsole.Workflows { diff --git a/src/samples/Elsa.Samples.SignalingConsole/Elsa.Samples.SignalingConsole.csproj b/src/samples/Elsa.Samples.SignalingConsole/Elsa.Samples.SignalingConsole.csproj new file mode 100644 index 000000000..7fbd68615 --- /dev/null +++ b/src/samples/Elsa.Samples.SignalingConsole/Elsa.Samples.SignalingConsole.csproj @@ -0,0 +1,16 @@ + + + + Exe + netcoreapp3.1 + + + + + + + + + + + diff --git a/src/samples/Elsa.Samples.SignalingConsole/Program.cs b/src/samples/Elsa.Samples.SignalingConsole/Program.cs new file mode 100644 index 000000000..eea1d0d21 --- /dev/null +++ b/src/samples/Elsa.Samples.SignalingConsole/Program.cs @@ -0,0 +1,51 @@ +using System; +using System.Threading.Tasks; +using Elsa.Activities.Signaling.Services; +using Elsa.Services; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Samples.SignalingConsole +{ + /// + /// Demonstrates a workflow with a While looping construct. + /// + static class Program + { + private static async Task Main() + { + // Create a service container with Elsa services. + var services = new ServiceCollection() + .AddElsa() + .AddConsoleActivities() + .AddWorkflow() + .BuildServiceProvider(); + + // Run startup actions (not needed when registering Elsa with a Host). + var startupRunner = services.GetRequiredService(); + await startupRunner.StartupAsync(); + + // Get a workflow runner. + var workflowRunner = services.GetRequiredService(); + + // Define a couple of cars so we can correlate workflows with them. + var cars = new[] { "Car 1", "Car 2" }; + + // Execute a workflow for each car. + foreach (var car in cars) + await workflowRunner.RunWorkflowAsync(correlationId: car); + + Console.WriteLine("Hit enter to signal green light for Car 2."); + Console.ReadLine(); + + // The workflows are now suspended at the red light. + // Trigger a green light signal for the first car. + var signaler = services.GetRequiredService(); + await signaler.SendSignal("Green", correlationId: "Car 2"); + + // Notice that only the workflow correlated to the second car executed. + + // Keep the application alive for the workflow scheduler to have enough time to resume the workflow. + Console.ReadLine(); + } + } +} \ No newline at end of file diff --git a/src/samples/Elsa.Samples.SignalingConsole/TrafficLightWorkflow.cs b/src/samples/Elsa.Samples.SignalingConsole/TrafficLightWorkflow.cs new file mode 100644 index 000000000..8bec79e56 --- /dev/null +++ b/src/samples/Elsa.Samples.SignalingConsole/TrafficLightWorkflow.cs @@ -0,0 +1,20 @@ +using Elsa.Activities.Console; +using Elsa.Activities.ControlFlow; +using Elsa.Builders; +using Elsa.Services.Models; + +namespace Elsa.Samples.SignalingConsole +{ + public class TrafficLightWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder workflow) + { + workflow + .WriteLine(context => $"{GetCarName(context)} is approaching red traffic light...") + .ReceiveSignal("Green") + .WriteLine(context => $"Light turned green for {GetCarName(context)}. Hit that power pedal!"); + } + + private string GetCarName(ActivityExecutionContext context) => context.WorkflowExecutionContext.CorrelationId; + } +} \ No newline at end of file diff --git a/test/integration/Elsa.Core.IntegrationTests/Workflows/ForEachWorkflow.cs b/test/integration/Elsa.Core.IntegrationTests/Workflows/ForEachWorkflow.cs index 8013f94df..eeafcccca 100644 --- a/test/integration/Elsa.Core.IntegrationTests/Workflows/ForEachWorkflow.cs +++ b/test/integration/Elsa.Core.IntegrationTests/Workflows/ForEachWorkflow.cs @@ -23,7 +23,7 @@ namespace Elsa.Core.IntegrationTests.Workflows _items, iterate => iterate .Then(activity => activity.Set(x => x.Text, context => $"{context.Input}")).WithId("WriteLine") - .Then() /* Block workflow.*/ + .Then() /* Block workflow.*/ .WriteLine("Resumed")) .WriteLine("One iterations executing, rest is blocked"); } diff --git a/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflow.cs b/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflow.cs index be1197d6d..af25404a5 100644 --- a/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflow.cs +++ b/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflow.cs @@ -20,9 +20,9 @@ namespace Elsa.Core.IntegrationTests.Workflows activity => activity.Set(x => x.Branches, new HashSet(new[] { "Branch 1", "Branch 2", "Branch 3" })), fork => { - fork.When("Branch 1").Signaled("Signal1").WriteLine("Branch 1 executed", "WriteLine1").Then("Join"); - fork.When("Branch 2").Signaled("Signal2").WriteLine("Branch 2 executed", "WriteLine2").Then("Join"); - fork.When("Branch 3").Signaled("Signal3").WriteLine("Branch 3 executed", "WriteLine3").Then("Join"); + fork.When("Branch 1").ReceiveSignal("Signal1").WriteLine("Branch 1 executed", "WriteLine1").Then("Join"); + fork.When("Branch 2").ReceiveSignal("Signal2").WriteLine("Branch 2 executed", "WriteLine2").Then("Join"); + fork.When("Branch 3").ReceiveSignal("Signal3").WriteLine("Branch 3 executed", "WriteLine3").Then("Join"); }) .Add(join => join.Set(x => x.Mode, _joinMode)).WithName("Join") .WriteLine("Finished", "Finished"); diff --git a/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflowTests.cs b/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflowTests.cs index 5acc13c8d..59e15555c 100644 --- a/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflowTests.cs +++ b/test/integration/Elsa.Core.IntegrationTests/Workflows/ForkJoinWorkflowTests.cs @@ -2,10 +2,10 @@ using System.Linq; using System.Threading.Tasks; using Elsa.Activities.ControlFlow; using Elsa.Activities.Signaling; +using Elsa.Activities.Signaling.Models; using Elsa.Models; using Elsa.Services.Models; using Elsa.Testing.Shared.Helpers; -using Open.Linq.AsyncExtensions; using Xunit; using Xunit.Abstractions; @@ -58,7 +58,7 @@ namespace Elsa.Core.IntegrationTests.Workflows var workflow = new ForkJoinWorkflow(Join.JoinMode.WaitAny); var workflowBlueprint = WorkflowBuilder.Build(workflow); var workflowInstance = await WorkflowRunner.RunWorkflowAsync(workflowBlueprint); - + bool GetActivityHasExecuted(string name) => (from entry in workflowInstance.ExecutionLog let activity = workflowBlueprint.Activities.First(x => x.Id == entry.ActivityId) where activity.Name == name select activity).Any(); bool GetIsFinished() => GetActivityHasExecuted("Finished"); @@ -74,9 +74,14 @@ namespace Elsa.Core.IntegrationTests.Workflows private async Task TriggerSignalAsync(IWorkflowBlueprint workflowBlueprint, WorkflowInstance workflowInstance, string signal) { - var signaled = await WorkflowSelector.GetTriggersAsync(x => x.Signal == signal).First(); - var triggeredSignal = new TriggeredSignal(signal, null); - return await WorkflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, signaled.Id, triggeredSignal); + var workflowExecutionContext = new WorkflowExecutionContext(ServiceProvider, workflowBlueprint, workflowInstance); + var activities = await workflowExecutionContext.ActivateActivitiesAsync(); + var blockingActivityIds = workflowInstance.BlockingActivities.Where(x => x.ActivityType == nameof(ReceiveSignal)).Select(x => x.ActivityId); + var receiveSignalActivities = activities.Where(x => blockingActivityIds.Contains(x.Id)).Cast(); + var receiveSignal = receiveSignalActivities.Single(x => x.Signal == signal); + + var triggeredSignal = new Signal(signal); + return await WorkflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance, receiveSignal.Id, triggeredSignal); } } } \ No newline at end of file diff --git a/test/integration/Elsa.Core.IntegrationTests/Workflows/ParallelForEachWorkflow.cs b/test/integration/Elsa.Core.IntegrationTests/Workflows/ParallelForEachWorkflow.cs index 62fec8bb6..1d1f81d2d 100644 --- a/test/integration/Elsa.Core.IntegrationTests/Workflows/ParallelForEachWorkflow.cs +++ b/test/integration/Elsa.Core.IntegrationTests/Workflows/ParallelForEachWorkflow.cs @@ -23,7 +23,7 @@ namespace Elsa.Core.IntegrationTests.Workflows _items, iterate => iterate .Then(activity => activity.Set(x => x.Text, context => $"{context.Input}")).WithId("WriteLine") - .Then() /* Block workflow.*/ + .Then() /* Block workflow.*/ .WriteLine("Resumed")) .WriteLine("All iterations executing in parallel"); }