diff --git a/src/core/Elsa.Core/Activities/ControlFlow/ParallelForEach/ParallelForEach.cs b/src/core/Elsa.Core/Activities/ControlFlow/ParallelForEach/ParallelForEach.cs new file mode 100644 index 000000000..0f09b13bb --- /dev/null +++ b/src/core/Elsa.Core/Activities/ControlFlow/ParallelForEach/ParallelForEach.cs @@ -0,0 +1,31 @@ +using System.Collections.Generic; +using System.Collections.ObjectModel; +using System.Linq; +using Elsa.ActivityResults; +using Elsa.Attributes; +using Elsa.Services; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.ControlFlow +{ + [ActivityDefinition( + Category = "Control Flow", + Description = "Iterate over a collection in parallel.", + Icon = "far fa-circle", + Outcomes = new[] { OutcomeNames.Iterate, OutcomeNames.Done } + )] + public class ParallelForEach : Activity + { + [ActivityProperty(Hint = "A collection of items to iterate over.")] + public ICollection Items { get; set; } = new Collection(); + + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) + { + var items = Items.Reverse().ToList(); + var results = new List { Done() }; + results.AddRange(items.Select(x => Outcome(OutcomeNames.Iterate, x))); + return Combine(results); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/ControlFlow/ParallelForEach/ParallelForEachBuilderExtensions.cs b/src/core/Elsa.Core/Activities/ControlFlow/ParallelForEach/ParallelForEachBuilderExtensions.cs new file mode 100644 index 000000000..bf8a09304 --- /dev/null +++ b/src/core/Elsa.Core/Activities/ControlFlow/ParallelForEach/ParallelForEachBuilderExtensions.cs @@ -0,0 +1,37 @@ +using System; +using System.Collections.Generic; +using System.Threading.Tasks; +using Elsa.Builders; +using Elsa.Services; +using Elsa.Services.Models; + +// ReSharper disable once CheckNamespace +namespace Elsa.Activities.ControlFlow +{ + public static class ParallelForEachBuilderExtensions + { + public static IActivityBuilder ParallelForEach( + this IBuilder builder, + Action> setup, + Action iterate) => + builder.Then(setup, branch => iterate(branch.When(OutcomeNames.Iterate))); + + public static IActivityBuilder ParallelForEach( + this IBuilder builder, + Func> items, + Action iterate) => + builder.ParallelForEach(activity => activity.Set(x => x.Items, items), iterate); + + public static IActivityBuilder ParallelForEach( + this IBuilder builder, + Func> items, + Action iterate) => + builder.ParallelForEach(activity => activity.Set(x => x.Items, items), iterate); + + public static IActivityBuilder ParallelForEach( + this IBuilder builder, + ICollection items, + Action iterate) => + builder.ParallelForEach(activity => activity.Set(x => x.Items, items), iterate); + } +} \ 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 dcd3e1dc6..bf98c7ee1 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -130,6 +130,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddActivity() .AddActivity() .AddActivity() + .AddActivity() .AddActivity() .AddActivity() .AddActivity() diff --git a/test/integration/Elsa.Core.IntegrationTests/ParallelForEachWorkflowTests.cs b/test/integration/Elsa.Core.IntegrationTests/ParallelForEachWorkflowTests.cs new file mode 100644 index 000000000..7d730d8d7 --- /dev/null +++ b/test/integration/Elsa.Core.IntegrationTests/ParallelForEachWorkflowTests.cs @@ -0,0 +1,29 @@ +using System.Linq; +using System.Threading.Tasks; +using Elsa.Core.IntegrationTests.Workflows; +using Elsa.Models; +using Elsa.Testing.Shared.Helpers; +using Xunit; +using Xunit.Abstractions; + +namespace Elsa.Core.IntegrationTests +{ + public class ParallelForEachWorkflowTests : WorkflowsUnitTestBase + { + public ParallelForEachWorkflowTests(ITestOutputHelper testOutputHelper) : base(testOutputHelper) + { + } + + [Fact(DisplayName = "Runs Parallel For Each workflow.")] + public async Task Test01() + { + var items = Enumerable.Range(1, 10).Select(x => $"Item {x}").ToList(); + var workflow = new ParallelForEachWorkflow(items); + var workflowInstance = await WorkflowHost.RunWorkflowAsync(workflow); + var iterationLogs = workflowInstance.ExecutionLog.Where(x => x.ActivityId == "WriteLine").ToList(); + + Assert.Equal(WorkflowStatus.Suspended, workflowInstance.Status); + Assert.Equal(items.Count, iterationLogs.Count); + } + } +} \ 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 new file mode 100644 index 000000000..62fec8bb6 --- /dev/null +++ b/test/integration/Elsa.Core.IntegrationTests/Workflows/ParallelForEachWorkflow.cs @@ -0,0 +1,31 @@ +using System.Collections.Generic; +using System.Linq; +using Elsa.Activities.Console; +using Elsa.Activities.ControlFlow; +using Elsa.Activities.Signaling; +using Elsa.Builders; + +namespace Elsa.Core.IntegrationTests.Workflows +{ + public class ParallelForEachWorkflow : IWorkflow + { + private readonly List _items; + + public ParallelForEachWorkflow(IEnumerable items) + { + _items = items.Cast().ToList(); + } + + public void Build(IWorkflowBuilder workflow) + { + workflow + .ParallelForEach( + _items, + iterate => iterate + .Then(activity => activity.Set(x => x.Text, context => $"{context.Input}")).WithId("WriteLine") + .Then() /* Block workflow.*/ + .WriteLine("Resumed")) + .WriteLine("All iterations executing in parallel"); + } + } +} \ No newline at end of file