parent
cbf609760e
commit
f6c4ffb934
|
|
@ -41,19 +41,26 @@ public class ParallelForEach<T> : Activity
|
|||
{
|
||||
var items = context.Get(Items)!.ToList();
|
||||
var tags = new List<Guid>();
|
||||
var currentIndex = 0;
|
||||
|
||||
foreach (var item in items)
|
||||
{
|
||||
// For each item, declare a new variable the work to be scheduled.
|
||||
var variable = new Variable<T>("CurrentValue", item)
|
||||
var currentValueVariable = new Variable<T>("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<Variable<T>>
|
||||
|
||||
var currentIndexVariable = new Variable<int>("CurrentIndex", currentIndex++)
|
||||
{
|
||||
variable
|
||||
StorageDriverType = typeof(WorkflowStorageDriver)
|
||||
};
|
||||
|
||||
var variables = new List<Variable>
|
||||
{
|
||||
currentValueVariable,
|
||||
currentIndexVariable
|
||||
};
|
||||
|
||||
// Schedule a body of work for each item.
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -20,7 +20,8 @@ namespace Elsa.Workflows.Core;
|
|||
/// <param name="Owner">The activity scheduling the <see cref="Child"/> activity.</param>
|
||||
/// <param name="Child">The child <see cref="IActivity"/> being scheduled.</param>
|
||||
/// <param name="CompletionCallback">The <see cref="ActivityCompletionCallback"/> delegate to invoke when the scheduled <see cref="Child"/> activity completes.</param>
|
||||
public record ActivityCompletionCallbackEntry(ActivityExecutionContext Owner, ActivityNode Child, ActivityCompletionCallback? CompletionCallback);
|
||||
/// <param name="Tag">An optional tag.</param>
|
||||
public record ActivityCompletionCallbackEntry(ActivityExecutionContext Owner, ActivityNode Child, ActivityCompletionCallback? CompletionCallback, object? Tag = default);
|
||||
|
||||
/// <summary>
|
||||
/// Provides context to the currently executing workflow.
|
||||
|
|
@ -282,9 +283,9 @@ public class WorkflowExecutionContext : IExecutionContext
|
|||
/// <summary>
|
||||
/// Registers a completion callback for the specified activity.
|
||||
/// </summary>
|
||||
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);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -8,9 +8,31 @@ namespace Elsa.Workflows.Core.Contracts;
|
|||
/// </summary>
|
||||
public interface IActivityScheduler
|
||||
{
|
||||
/// <summary>
|
||||
/// Returns true if there are any work items in the scheduler.
|
||||
/// </summary>
|
||||
bool HasAny { get; }
|
||||
|
||||
/// <summary>
|
||||
/// Schedules a work item.
|
||||
/// </summary>
|
||||
/// <param name="workItem"></param>
|
||||
void Schedule(ActivityWorkItem workItem);
|
||||
|
||||
/// <summary>
|
||||
/// Takes the next work item from the scheduler.
|
||||
/// </summary>
|
||||
/// <returns></returns>
|
||||
ActivityWorkItem Take();
|
||||
|
||||
/// <summary>
|
||||
/// Returns a list of all work items in the scheduler.
|
||||
/// </summary>
|
||||
/// <returns></returns>
|
||||
IEnumerable<ActivityWorkItem> List();
|
||||
|
||||
/// <summary>
|
||||
/// Clears all work items from the scheduler.
|
||||
/// </summary>
|
||||
void Clear();
|
||||
}
|
||||
|
|
@ -11,19 +11,6 @@ namespace Elsa.Extensions;
|
|||
/// </summary>
|
||||
public static class WorkflowExecutionContextExtensions
|
||||
{
|
||||
// /// <summary>
|
||||
// /// Remove the specified set of <see cref="ActivityExecutionContext"/> from the workflow execution context.
|
||||
// /// </summary>
|
||||
// public static async Task RemoveActivityExecutionContextsAsync(this WorkflowExecutionContext workflowExecutionContext, IEnumerable<ActivityExecutionContext> 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);
|
||||
// }
|
||||
|
||||
/// <summary>
|
||||
/// Schedules the workflow for execution.
|
||||
/// </summary>
|
||||
|
|
@ -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);
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
|
|
|
|||
|
|
@ -1,15 +1,29 @@
|
|||
using Elsa.Workflows.Core.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using JetBrains.Annotations;
|
||||
|
||||
namespace Elsa.Workflows.Core.Services;
|
||||
|
||||
/// <summary>
|
||||
/// A FIFO queue based activity scheduler.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public class QueueBasedActivityScheduler : IActivityScheduler
|
||||
{
|
||||
private readonly Queue<ActivityWorkItem> _queue = new();
|
||||
|
||||
/// <inheritdoc />
|
||||
public bool HasAny => _queue.Any();
|
||||
public void Schedule(ActivityWorkItem activity) => _queue.Enqueue(activity);
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Schedule(ActivityWorkItem workItem) => _queue.Enqueue(workItem);
|
||||
|
||||
/// <inheritdoc />
|
||||
public ActivityWorkItem Take() => _queue.Dequeue();
|
||||
|
||||
/// <inheritdoc />
|
||||
public IEnumerable<ActivityWorkItem> List() => _queue.ToList();
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Clear() => _queue.Clear();
|
||||
}
|
||||
|
|
@ -1,15 +1,29 @@
|
|||
using Elsa.Workflows.Core.Contracts;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using JetBrains.Annotations;
|
||||
|
||||
namespace Elsa.Workflows.Core.Services;
|
||||
|
||||
/// <summary>
|
||||
/// A LIFO stack based activity scheduler.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
public class StackBasedActivityScheduler : IActivityScheduler
|
||||
{
|
||||
private readonly Stack<ActivityWorkItem> _stack = new();
|
||||
|
||||
/// <inheritdoc />
|
||||
public bool HasAny => _stack.Any();
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Schedule(ActivityWorkItem activity) => _stack.Push(activity);
|
||||
|
||||
/// <inheritdoc />
|
||||
public ActivityWorkItem Take() => _stack.Pop();
|
||||
|
||||
/// <inheritdoc />
|
||||
public IEnumerable<ActivityWorkItem> List() => _stack.ToList();
|
||||
|
||||
/// <inheritdoc />
|
||||
public void Clear() => _stack.Clear();
|
||||
}
|
||||
|
|
@ -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();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
||||
/// <summary>
|
||||
/// Represents a serializable completion callback that is registered by an activity.
|
||||
/// </summary>
|
||||
public class CompletionCallbackState
|
||||
{
|
||||
// ReSharper disable once UnusedMember.Global
|
||||
// Required for JSON serialization configured with reference handling.
|
||||
/// <summary>
|
||||
/// Creates a new instance of the <see cref="CompletionCallbackState"/> class.
|
||||
/// </summary>
|
||||
public CompletionCallbackState()
|
||||
{
|
||||
}
|
||||
|
||||
public CompletionCallbackState(string ownerInstanceId, string childNodeId, string? methodName)
|
||||
/// <summary>
|
||||
/// Creates a new instance of the <see cref="CompletionCallbackState"/> class.
|
||||
/// </summary>
|
||||
/// <param name="ownerInstanceId">The ID of the activity instance that registered the callback.</param>
|
||||
/// <param name="childNodeId">The ID of the child node that the callback is registered for.</param>
|
||||
/// <param name="methodName">The name of the method to invoke when the child node completes.</param>
|
||||
/// <param name="tag">An optional tag.</param>
|
||||
public CompletionCallbackState(string ownerInstanceId, string childNodeId, string? methodName, object? tag = default)
|
||||
{
|
||||
OwnerInstanceId = ownerInstanceId;
|
||||
ChildNodeId = childNodeId;
|
||||
MethodName = methodName;
|
||||
Tag = tag;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Gets the ID of the activity instance that registered the callback.
|
||||
/// </summary>
|
||||
public string OwnerInstanceId { get; init; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// Gets the ID of the child node that the callback is registered for.
|
||||
/// </summary>
|
||||
public string ChildNodeId { get; init; } = default!;
|
||||
public string? MethodName { get; init; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// Gets the name of the method to invoke when the child node completes.
|
||||
/// </summary>
|
||||
public string? MethodName { get; init; }
|
||||
|
||||
/// <summary>
|
||||
/// Gets the tag.
|
||||
/// </summary>
|
||||
public object? Tag { get; init; }
|
||||
}
|
||||
|
|
@ -99,6 +99,10 @@
|
|||
<None Update="Scenarios\ExplicitJoins\Workflows\flow-join-any.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
<None Update="Scenarios\FlowchartCompletion\Workflows\workflow6.json">
|
||||
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
|
||||
</None>
|
||||
|
||||
|
||||
|
||||
</ItemGroup>
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue