diff --git a/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs b/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs index 1e3c080f4..882dda36c 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs @@ -1,4 +1,6 @@ using System.Runtime.CompilerServices; +using System.Text.Json; +using System.Text.Json.Nodes; using Elsa.Expressions.Helpers; using Elsa.Extensions; using Elsa.Workflows.Attributes; @@ -42,8 +44,8 @@ public class ParallelForEach : Activity var tags = new List(); var currentIndex = 0; - context.SetProperty(ScheduledTagsProperty, tags); - context.SetProperty(CompletedTagsProperty, new List()); + SetTagList(context, ScheduledTagsProperty, tags); + SetTagList(context, CompletedTagsProperty, new List()); await foreach (var item in items) { @@ -51,10 +53,10 @@ public class ParallelForEach : Activity var currentValueVariable = new Variable("CurrentValue", item) { // TODO: This should be configurable, because this won't work for e.g. file streams and other non-serializable types. - StorageDriverType = typeof(WorkflowStorageDriver) + StorageDriverType = typeof(WorkflowInstanceStorageDriver) }; - var currentIndexVariable = new Variable("CurrentIndex", currentIndex++) { StorageDriverType = typeof(WorkflowStorageDriver) }; + var currentIndexVariable = new Variable("CurrentIndex", currentIndex++) { StorageDriverType = typeof(WorkflowInstanceStorageDriver) }; var variables = new List { currentValueVariable, currentIndexVariable }; // Schedule a body of work for each item. @@ -71,15 +73,13 @@ public class ParallelForEach : Activity private async ValueTask OnChildCompleted(ActivityCompletedContext context) { var targetContext = context.TargetContext; - var scheduledTags = targetContext.GetProperty>(ScheduledTagsProperty)!; + var scheduledTags = GetTagList(targetContext, ScheduledTagsProperty); var completedTag = targetContext.Tag.ConvertTo(); - - var completedTags = new HashSet(targetContext.UpdateProperty>(CompletedTagsProperty, completedTags => - { - completedTags!.Add(completedTag); - return completedTags; - })); - + var completedTags = GetTagList(targetContext, CompletedTagsProperty); + + completedTags.Add(completedTag); + SetTagList(targetContext, CompletedTagsProperty, completedTags); + // If not all scheduled activities have completed yet, we're not done yet. if (!scheduledTags.IsEqualTo(completedTags)) return; @@ -87,4 +87,18 @@ public class ParallelForEach : Activity // We're done, so complete the activity. await targetContext.CompleteActivityAsync(); } + + private ICollection GetTagList(ActivityExecutionContext context, string propertyName) + { + // Read the list of tags from the context using the specified property name. The value is stored as JsonArray, so we need to deserialize it. + var jsonArray = context.GetProperty(propertyName); + return jsonArray.Select(x => x.ConvertTo()).ToList(); + } + + private void SetTagList(ActivityExecutionContext context, string propertyName, ICollection tags) + { + // Serialize the list of tags to a JsonArray and store it in the context using the specified property name. + var jsonArray = JsonSerializer.SerializeToNode(tags) as JsonArray; + context.SetProperty(propertyName, jsonArray); + } } \ No newline at end of file