< Summary

Information
Class: Elsa.Workflows.Activities.ParallelForEach<T>
Assembly: Elsa.Workflows.Core
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs
Line coverage
92%
Covered lines: 46
Uncovered lines: 4
Coverable lines: 50
Total lines: 129
Line coverage: 92%
Branch coverage
100%
Covered branches: 8
Total branches: 8
Branch coverage: 100%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%210%
.ctor(...)100%210%
.ctor(...)100%11100%
.ctor(...)100%11100%
.ctor(...)100%11100%
get_Items()100%11100%
get_Body()100%11100%
ExecuteAsync()100%44100%
OnChildCompleted()100%22100%
GetTagList(...)100%22100%
SetTagList(...)100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs

#LineLine coverage
 1using System.Runtime.CompilerServices;
 2using System.Text.Json;
 3using System.Text.Json.Nodes;
 4using Elsa.Expressions.Helpers;
 5using Elsa.Expressions.Models;
 6using Elsa.Extensions;
 7using Elsa.Workflows.Attributes;
 8using Elsa.Workflows.Memory;
 9using Elsa.Workflows.Models;
 10
 11namespace Elsa.Workflows.Activities;
 12
 13/// <summary>
 14/// Schedule an activity for each item in parallel.
 15/// </summary>
 16/// <typeparam name="T"></typeparam>
 17[Activity("Elsa", "Looping", "Schedule an activity for each item in parallel.")]
 18public class ParallelForEach<T> : Activity
 19{
 20    private const string ScheduledTagsProperty = nameof(ScheduledTagsProperty);
 21    private const string CompletedTagsProperty = nameof(CompletedTagsProperty);
 22
 23    /// <inheritdoc />
 024    public ParallelForEach(Func<ExpressionExecutionContext, ICollection<T>> @delegate, [CallerFilePath] string? source =
 25    {
 026    }
 27
 28    /// <inheritdoc />
 029    public ParallelForEach(Func<ICollection<T>> @delegate, [CallerFilePath] string? source = null, [CallerLineNumber] in
 30    {
 031    }
 32
 33    /// <inheritdoc />
 1534    public ParallelForEach(ICollection<T> items, [CallerFilePath] string? source = null, [CallerLineNumber] int? line = 
 35    {
 1536    }
 37
 38    /// <inheritdoc />
 1539    public ParallelForEach(Input<object> items, [CallerFilePath] string? source = null, [CallerLineNumber] int? line = n
 40    {
 1541        Items = items;
 1542    }
 43
 44    /// <inheritdoc />
 1745    public ParallelForEach([CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : base(source, l
 46    {
 1747    }
 48
 49    /// <summary>
 50    /// The items to iterate.
 51    /// </summary>
 52    [Input(Description = "The items to iterate through.")]
 7153    public Input<object> Items { get; set; } = new(Array.Empty<T>());
 54
 55    /// <summary>
 56    /// The <see cref="IActivity"/> to execute each iteration.
 57    /// </summary>
 58    [Port]
 6459    public IActivity Body { get; set; } = null!;
 60
 61    /// <inheritdoc />
 62    protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
 63    {
 1664        var items = context.GetItemSource<T>(Items);
 1665        var tags = new List<Guid>();
 1666        var currentIndex = 0;
 67
 10068        await foreach (var item in items)
 69        {
 70            // Schedule a body of work for each item.
 3471            var tag = Guid.NewGuid();
 3472            tags.Add(tag);
 73
 74            // Give each iteration its own variable ids so nested loops cannot share a parent memory block.
 3475            var currentValueVariable = new Variable<T>("CurrentValue", item, $"{tag}:CurrentValue")
 3476            {
 3477                // TODO: This should be configurable, because this won't work for e.g. file streams and other non-serial
 3478                StorageDriverType = typeof(WorkflowInstanceStorageDriver)
 3479            };
 80
 3481            var currentIndexVariable = new Variable<int>("CurrentIndex", currentIndex++, $"{tag}:CurrentIndex")
 3482            {
 3483                StorageDriverType = typeof(WorkflowInstanceStorageDriver)
 3484            };
 3485            var variables = new List<Variable> { currentValueVariable, currentIndexVariable };
 86
 3487            await context.ScheduleActivityAsync(Body, OnChildCompleted, tag, variables);
 88        }
 89
 1690        SetTagList(context, ScheduledTagsProperty, tags);
 1691        SetTagList(context, CompletedTagsProperty, new List<Guid>());
 92
 93        // If there were no items, we're done.
 1694        if (tags.Count == 0)
 395            await context.CompleteActivityAsync();
 1696    }
 97
 98    private async ValueTask OnChildCompleted(ActivityCompletedContext context)
 99    {
 22100        var targetContext = context.TargetContext;
 22101        var scheduledTags = GetTagList(targetContext, ScheduledTagsProperty);
 22102        var completedTag = targetContext.Tag.ConvertTo<Guid>();
 22103        var completedTags = GetTagList(targetContext, CompletedTagsProperty);
 104
 22105        completedTags.Add(completedTag);
 22106        SetTagList(targetContext, CompletedTagsProperty, completedTags);
 107
 108        // If not all scheduled activities have completed yet, we're not done yet.
 22109        if (!scheduledTags.IsEqualTo(completedTags))
 15110            return;
 111
 112        // We're done, so complete the activity.
 7113        await targetContext.CompleteActivityAsync();
 22114    }
 115
 116    private ICollection<Guid> GetTagList(ActivityExecutionContext context, string propertyName)
 117    {
 118        // Read the list of tags from the context using the specified property name. The value is stored as JsonArray, s
 44119        var jsonArray = context.GetProperty<JsonArray>(propertyName)!;
 134120        return jsonArray.Select(x => x.ConvertTo<Guid>()).ToList();
 121    }
 122
 123    private void SetTagList(ActivityExecutionContext context, string propertyName, ICollection<Guid> tags)
 124    {
 125        // Serialize the list of tags to a JsonArray and store it in the context using the specified property name.
 54126        var jsonArray = JsonSerializer.SerializeToNode(tags) as JsonArray;
 54127        context.SetProperty(propertyName, jsonArray);
 54128    }
 129}