diff --git a/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs b/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs index 636103ba8..d5dde03fc 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/ParallelForEachT.cs @@ -41,19 +41,26 @@ public class ParallelForEach : Activity { var items = context.Get(Items)!.ToList(); var tags = new List(); + var currentIndex = 0; foreach (var item in items) { // For each item, declare a new variable the work to be scheduled. - var variable = new Variable("CurrentValue", item) + 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) }; - - var variables = new List> + + var currentIndexVariable = new Variable("CurrentIndex", currentIndex++) { - variable + StorageDriverType = typeof(WorkflowStorageDriver) + }; + + var variables = new List + { + currentValueVariable, + currentIndexVariable }; // Schedule a body of work for each item. diff --git a/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs b/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs index 2a4930d06..b2e3b1112 100644 --- a/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs +++ b/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs @@ -32,6 +32,8 @@ public class ScheduledChildCallbackBehavior : Behavior if (callbackEntry.CompletionCallback != null) { var completedContext = new ActivityCompletedContext(activityExecutionContext, childActivityExecutionContext, signal.Result); + var tag = callbackEntry.Tag; + completedContext.TargetContext.Tag = tag; await callbackEntry.CompletionCallback(completedContext); } } diff --git a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs index b88dd5f42..f137131c1 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs @@ -20,7 +20,8 @@ namespace Elsa.Workflows.Core; /// The activity scheduling the activity. /// The child being scheduled. /// The delegate to invoke when the scheduled activity completes. -public record ActivityCompletionCallbackEntry(ActivityExecutionContext Owner, ActivityNode Child, ActivityCompletionCallback? CompletionCallback); +/// An optional tag. +public record ActivityCompletionCallbackEntry(ActivityExecutionContext Owner, ActivityNode Child, ActivityCompletionCallback? CompletionCallback, object? Tag = default); /// /// Provides context to the currently executing workflow. @@ -282,9 +283,9 @@ public class WorkflowExecutionContext : IExecutionContext /// /// Registers a completion callback for the specified activity. /// - internal void AddCompletionCallback(ActivityExecutionContext owner, ActivityNode child, ActivityCompletionCallback? completionCallback = default) + internal void AddCompletionCallback(ActivityExecutionContext owner, ActivityNode child, ActivityCompletionCallback? completionCallback = default, object? tag = default) { - var entry = new ActivityCompletionCallbackEntry(owner, child, completionCallback); + var entry = new ActivityCompletionCallbackEntry(owner, child, completionCallback, tag); _completionCallbackEntries.Add(entry); } diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IActivityScheduler.cs b/src/modules/Elsa.Workflows.Core/Contracts/IActivityScheduler.cs index 93980dcd7..7ad3d9f6c 100644 --- a/src/modules/Elsa.Workflows.Core/Contracts/IActivityScheduler.cs +++ b/src/modules/Elsa.Workflows.Core/Contracts/IActivityScheduler.cs @@ -8,9 +8,31 @@ namespace Elsa.Workflows.Core.Contracts; /// public interface IActivityScheduler { + /// + /// Returns true if there are any work items in the scheduler. + /// bool HasAny { get; } + + /// + /// Schedules a work item. + /// + /// void Schedule(ActivityWorkItem workItem); + + /// + /// Takes the next work item from the scheduler. + /// + /// ActivityWorkItem Take(); + + /// + /// Returns a list of all work items in the scheduler. + /// + /// IEnumerable List(); + + /// + /// Clears all work items from the scheduler. + /// void Clear(); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs index d75cacedd..600af62a0 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs @@ -11,19 +11,6 @@ namespace Elsa.Extensions; /// public static class WorkflowExecutionContextExtensions { - // /// - // /// Remove the specified set of from the workflow execution context. - // /// - // public static async Task RemoveActivityExecutionContextsAsync(this WorkflowExecutionContext workflowExecutionContext, IEnumerable contexts) - // { - // // Copy each item into a new list to avoid changing the source enumerable while removing elements from it. - // var list = contexts.ToList(); - // - // // Remove each context. - // foreach (var context in list) - // await context.CompleteActivityExecutionContextAsync(context); - // } - /// /// Schedules the workflow for execution. /// @@ -108,7 +95,7 @@ public static class WorkflowExecutionContextExtensions var workItem = new ActivityWorkItem(activityId, owner.Id, async () => await activityInvoker.InvokeAsync(workflowExecutionContext, activityNode.Activity, activityInvocationOptions), tag); var completionCallback = options?.CompletionCallback; workflowExecutionContext.Scheduler.Schedule(workItem); - workflowExecutionContext.AddCompletionCallback(owner, activityNode, completionCallback); + workflowExecutionContext.AddCompletionCallback(owner, activityNode, completionCallback, tag); } /// diff --git a/src/modules/Elsa.Workflows.Core/Services/QueueBasedActivityScheduler.cs b/src/modules/Elsa.Workflows.Core/Services/QueueBasedActivityScheduler.cs index 6598adc79..d68611e17 100644 --- a/src/modules/Elsa.Workflows.Core/Services/QueueBasedActivityScheduler.cs +++ b/src/modules/Elsa.Workflows.Core/Services/QueueBasedActivityScheduler.cs @@ -1,15 +1,29 @@ using Elsa.Workflows.Core.Contracts; using Elsa.Workflows.Core.Models; +using JetBrains.Annotations; namespace Elsa.Workflows.Core.Services; +/// +/// A FIFO queue based activity scheduler. +/// +[PublicAPI] public class QueueBasedActivityScheduler : IActivityScheduler { private readonly Queue _queue = new(); + /// public bool HasAny => _queue.Any(); - public void Schedule(ActivityWorkItem activity) => _queue.Enqueue(activity); + + /// + public void Schedule(ActivityWorkItem workItem) => _queue.Enqueue(workItem); + + /// public ActivityWorkItem Take() => _queue.Dequeue(); + + /// public IEnumerable List() => _queue.ToList(); + + /// public void Clear() => _queue.Clear(); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/StackBasedActivityScheduler.cs b/src/modules/Elsa.Workflows.Core/Services/StackBasedActivityScheduler.cs index 44b769f35..ec8d92adf 100644 --- a/src/modules/Elsa.Workflows.Core/Services/StackBasedActivityScheduler.cs +++ b/src/modules/Elsa.Workflows.Core/Services/StackBasedActivityScheduler.cs @@ -1,15 +1,29 @@ using Elsa.Workflows.Core.Contracts; using Elsa.Workflows.Core.Models; +using JetBrains.Annotations; namespace Elsa.Workflows.Core.Services; +/// +/// A LIFO stack based activity scheduler. +/// +[PublicAPI] public class StackBasedActivityScheduler : IActivityScheduler { private readonly Stack _stack = new(); + /// public bool HasAny => _stack.Any(); + + /// public void Schedule(ActivityWorkItem activity) => _stack.Push(activity); + + /// public ActivityWorkItem Take() => _stack.Pop(); + + /// public IEnumerable List() => _stack.ToList(); + + /// public void Clear() => _stack.Clear(); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs index 4510c5d73..fd79ff9a6 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs @@ -91,7 +91,8 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor var callbackName = completionCallbackEntry.MethodName; var callbackDelegate = !string.IsNullOrEmpty(callbackName) ? ownerActivityExecutionContext.Activity.GetActivityCompletionCallback(callbackName) : default; - workflowExecutionContext.AddCompletionCallback(ownerActivityExecutionContext, childNode, callbackDelegate); + var tag = completionCallbackEntry.Tag; + workflowExecutionContext.AddCompletionCallback(ownerActivityExecutionContext, childNode, callbackDelegate, tag); } } @@ -106,7 +107,10 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor throw new Exception("Lost an owner context"); } - var completionCallbacks = workflowExecutionContext.CompletionCallbacks.Select(x => new CompletionCallbackState(x.Owner.Id, x.Child.NodeId, x.CompletionCallback?.Method.Name)); + var completionCallbacks = workflowExecutionContext + .CompletionCallbacks + .Select(x => new CompletionCallbackState(x.Owner.Id, x.Child.NodeId, x.CompletionCallback?.Method.Name, x.Tag)); + state.CompletionCallbacks = completionCallbacks.ToList(); } diff --git a/src/modules/Elsa.Workflows.Core/State/CompletionCallbackState.cs b/src/modules/Elsa.Workflows.Core/State/CompletionCallbackState.cs index 31477ba6a..c344df492 100644 --- a/src/modules/Elsa.Workflows.Core/State/CompletionCallbackState.cs +++ b/src/modules/Elsa.Workflows.Core/State/CompletionCallbackState.cs @@ -5,22 +5,52 @@ namespace Elsa.Workflows.Core.State; // Can't use records when using System.Text.Json serialization and reference handling. Hence, using a class with default constructor. //public record CompletionCallbackState(string OwnerId, string ChildId, string MethodName); +/// +/// Represents a serializable completion callback that is registered by an activity. +/// public class CompletionCallbackState { // ReSharper disable once UnusedMember.Global // Required for JSON serialization configured with reference handling. + /// + /// Creates a new instance of the class. + /// public CompletionCallbackState() { } - public CompletionCallbackState(string ownerInstanceId, string childNodeId, string? methodName) + /// + /// Creates a new instance of the class. + /// + /// The ID of the activity instance that registered the callback. + /// The ID of the child node that the callback is registered for. + /// The name of the method to invoke when the child node completes. + /// An optional tag. + public CompletionCallbackState(string ownerInstanceId, string childNodeId, string? methodName, object? tag = default) { OwnerInstanceId = ownerInstanceId; ChildNodeId = childNodeId; MethodName = methodName; + Tag = tag; } + /// + /// Gets the ID of the activity instance that registered the callback. + /// public string OwnerInstanceId { get; init; } = default!; + + /// + /// Gets the ID of the child node that the callback is registered for. + /// public string ChildNodeId { get; init; } = default!; - public string? MethodName { get; init; } = default!; + + /// + /// Gets the name of the method to invoke when the child node completes. + /// + public string? MethodName { get; init; } + + /// + /// Gets the tag. + /// + public object? Tag { get; init; } } \ No newline at end of file diff --git a/test/integration/Elsa.IntegrationTests/Elsa.IntegrationTests.csproj b/test/integration/Elsa.IntegrationTests/Elsa.IntegrationTests.csproj index 321126b3a..600d87b1d 100644 --- a/test/integration/Elsa.IntegrationTests/Elsa.IntegrationTests.csproj +++ b/test/integration/Elsa.IntegrationTests/Elsa.IntegrationTests.csproj @@ -99,6 +99,10 @@ Always + + Always + + diff --git a/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Tests.cs b/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Tests.cs index 61e25ffb9..bc2ad8fe8 100644 --- a/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Tests.cs +++ b/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Tests.cs @@ -26,6 +26,7 @@ public class Tests [InlineData("workflow3.json")] [InlineData("workflow4.json")] [InlineData("workflow5.json")] + [InlineData("workflow6.json")] public async Task Test1(string workflowFileName) { // Populate registries. diff --git a/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Workflows/workflow6.json b/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Workflows/workflow6.json new file mode 100644 index 000000000..151d3cd6f --- /dev/null +++ b/test/integration/Elsa.IntegrationTests/Scenarios/FlowchartCompletion/Workflows/workflow6.json @@ -0,0 +1,198 @@ +{ + "id": "e5190066597846cc926f8ae644d706f4", + "definitionId": "b07725956291433abee7ff229b415194", + "name": "Parallel For Each", + "createdAt": "2023-09-10T16:32:58.422674+00:00", + "version": 3, + "toolVersion": "3.0.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": {}, + "isReadonly": false, + "isLatest": true, + "isPublished": true, + "options": { + "autoUpdateConsumingWorkflows": false + }, + "root": { + "type": "Elsa.Flowchart", + "version": 1, + "id": "oplDmRa_zEqH8PCdTcFh8A", + "metadata": {}, + "customProperties": { + "source": "FlowchartJsonConverter.cs:45", + "NotFoundConnectionsKey": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "activities": [ + { + "text": { + "typeName": "String", + "expression": { + "type": "Literal", + "value": "Start" + }, + "memoryReference": { + "id": "NTW8X70bTEaBtlYyU44PcA:input-0" + } + }, + "id": "NTW8X70bTEaBtlYyU44PcA", + "name": "WriteLine1", + "type": "Elsa.WriteLine", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -352, + "y": -286.9921875 + }, + "size": { + "width": 139.296875, + "height": 50 + } + } + } + }, + { + "text": { + "typeName": "String", + "expression": { + "type": "Literal", + "value": "End" + }, + "memoryReference": { + "id": "6Bf011ZOvUWOfQ1rKJEDFA:input-0" + } + }, + "id": "6Bf011ZOvUWOfQ1rKJEDFA", + "name": "WriteLine2", + "type": "Elsa.WriteLine", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": 169, + "y": -286.9921875 + }, + "size": { + "width": 139.296875, + "height": 50 + } + } + } + }, + { + "items": { + "typeName": "Object[]", + "expression": { + "type": "JavaScript", + "value": "[\u0022Apple\u0022, \u0022Banana\u0022, \u0022Cherry\u0022]" + }, + "memoryReference": { + "id": "By4rVs_-2kSDENFO-bniRg:input-0" + } + }, + "body": { + "type": "Elsa.Flowchart", + "version": 1, + "id": "ylZLDh7FQ0WmnbdnRJHutg", + "metadata": {}, + "customProperties": { + "source": "FlowchartJsonConverter.cs:45", + "NotFoundConnectionsKey": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "activities": [ + { + "text": { + "typeName": "String", + "expression": { + "type": "JavaScript", + "value": "\u0060Current fruit: ${getCurrentValue()}\u0060" + }, + "memoryReference": { + "id": "fht2JVY0G0KURaRYB9jqlw:input-0" + } + }, + "id": "fht2JVY0G0KURaRYB9jqlw", + "name": "WriteLine1", + "type": "Elsa.WriteLine", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -459.5, + "y": -374 + }, + "size": { + "width": 139.296875, + "height": 50 + } + } + } + } + ], + "connections": [] + }, + "id": "By4rVs_-2kSDENFO-bniRg", + "name": "ParallelForEach1", + "type": "Elsa.ParallelForEach", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -100, + "y": -320 + }, + "size": { + "width": 149.203125, + "height": 116.015625 + } + } + } + } + ], + "connections": [ + { + "source": { + "activity": "NTW8X70bTEaBtlYyU44PcA", + "port": "Done" + }, + "target": { + "activity": "By4rVs_-2kSDENFO-bniRg", + "port": "In" + } + }, + { + "source": { + "activity": "By4rVs_-2kSDENFO-bniRg", + "port": "Done" + }, + "target": { + "activity": "6Bf011ZOvUWOfQ1rKJEDFA", + "port": "In" + } + } + ] + } +} \ No newline at end of file