Add execution context affinity and prevent duplicate scheduling for Join

This fixes the issue where the Join activity executes twice and potentially one instance never completes.
This commit is contained in:
Sipke Schoorstra 2023-09-12 20:27:36 +02:00
parent 287442ca7e
commit 6fddd5223d
15 changed files with 81 additions and 24 deletions

View file

@ -19,7 +19,7 @@ using Proto.Persistence.SqlServer;
const bool useMongoDb = false;
const bool useSqlServer = false;
const bool useDapper = false;
const bool useProtoActor = true;
const bool useProtoActor = false;
const bool useHangfire = false;
var builder = WebApplication.CreateBuilder(args);

View file

@ -43,20 +43,24 @@ public class FlowJoin : Activity, IJoinNode
switch (mode)
{
case FlowJoinMode.WaitAll:
{
// If all left-inbound activities have executed, complete & continue.
var haveAllInboundActivitiesExecuted = inboundActivities.All(x => flowScope.GetExecutionCount(x) > executionCount);
if (haveAllInboundActivitiesExecuted)
{
await context.CompleteActivityAsync();
break;
case FlowJoinMode.WaitAny:
// If this activity was already executed, complete it with the AlreadyCompleted result. This will prevent the Flowchart activity from scheduling its siblings again.
var alreadyExecuted = inboundActivities.Max(x => flowScope.GetExecutionCount(x)) == executionCount;
var result = alreadyExecuted ? new AlreadyCompleted() : default;
await ClearBookmarksAsync(flowchart, context);
}
await context.CompleteActivityAsync(result);
break;
}
case FlowJoinMode.WaitAny:
{
await context.CompleteActivityAsync();
await ClearBookmarksAsync(flowchart, context);
break;
}
}
}

View file

@ -6,6 +6,7 @@ using Elsa.Workflows.Core.Activities.Flowchart.Extensions;
using Elsa.Workflows.Core.Activities.Flowchart.Models;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Contracts;
using Elsa.Workflows.Core.Options;
using Elsa.Workflows.Core.Signals;
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
@ -202,7 +203,15 @@ public class Flowchart : Container
}
else
{
await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync);
// Select an existing activity execution context for this activity, if any.
var joinContext = flowchartContext.WorkflowExecutionContext.ActiveActivityExecutionContexts.FirstOrDefault(x => x.ParentActivityExecutionContext == flowchartContext && x.Activity == activity);
var scheduleWorkOptions = new ScheduleWorkOptions
{
CompletionCallback = OnChildCompletedAsync,
ReuseActivityExecutionContextId = joinContext?.Id,
PreventDuplicateScheduling = true
};
await flowchartContext.ScheduleActivityAsync(activity, scheduleWorkOptions);
}
}
}

View file

@ -16,21 +16,25 @@ public interface IActivityScheduler
/// <summary>
/// Schedules a work item.
/// </summary>
/// <param name="workItem"></param>
/// <param name="workItem">The work item to schedule.</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>
/// Returns true if there are any work items matching the specified predicate.
/// </summary>
/// <param name="predicate">The predicate to match.</param>
bool Any(Func<ActivityWorkItem, bool> predicate);
/// <summary>
/// Clears all work items from the scheduler.
/// </summary>

View file

@ -87,15 +87,22 @@ public static class WorkflowExecutionContextExtensions
ActivityExecutionContext owner,
ScheduleWorkOptions? options = default)
{
var scheduler = workflowExecutionContext.Scheduler;
if(options?.PreventDuplicateScheduling == true && scheduler.Any(x => x.ActivityId == activityNode.NodeId))
return;
var activityInvoker = workflowExecutionContext.GetRequiredService<IActivityInvoker>();
var toolVersion = workflowExecutionContext.Workflow.WorkflowMetadata.ToolVersion;
var activityId = toolVersion?.Major >= 3 ? activityNode.Activity.Id : activityNode.NodeId;
var tag = options?.Tag;
var activityInvocationOptions = new ActivityInvocationOptions(owner, tag, options?.Variables);
var activityInvocationOptions = new ActivityInvocationOptions(owner, tag, options?.Variables, options?.ReuseActivityExecutionContextId);
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, tag);
scheduler.Schedule(workItem);
}
/// <summary>

View file

@ -8,4 +8,4 @@ namespace Elsa.Workflows.Core.Options;
/// <param name="Owner">The activity execution context that owns this invocation.</param>
/// <param name="Tag">An optional tag that can be used to identify the invocation.</param>
/// <param name="Variables">The variables to declare in the activity execution context that will be created for this invocation.</param>
public record ActivityInvocationOptions(ActivityExecutionContext? Owner, object? Tag, IEnumerable<Variable>? Variables);
public record ActivityInvocationOptions(ActivityExecutionContext? Owner, object? Tag, IEnumerable<Variable>? Variables, string? ReuseActivityExecutionContextId = default);

View file

@ -9,4 +9,9 @@ namespace Elsa.Workflows.Core.Options;
/// <param name="CompletionCallback">The callback to invoke when the work item has completed.</param>
/// <param name="Tag">A tag that can be used to identify the work item.</param>
/// <param name="Variables">A collection of variables to declare in the activity execution context that will be created for this work item.</param>
public record ScheduleWorkOptions(ActivityCompletionCallback? CompletionCallback = default, object? Tag = default, ICollection<Variable>? Variables = default);
public record ScheduleWorkOptions(
ActivityCompletionCallback? CompletionCallback = default,
object? Tag = default,
ICollection<Variable>? Variables = default,
string? ReuseActivityExecutionContextId = default,
bool PreventDuplicateScheduling = false);

View file

@ -19,11 +19,21 @@ public class ActivityInvoker : IActivityInvoker
/// <inheritdoc />
public async Task InvokeAsync(WorkflowExecutionContext workflowExecutionContext, IActivity activity, ActivityInvocationOptions? options = default)
{
// Setup an activity execution context.
var activityExecutionContext = workflowExecutionContext.CreateActivityExecutionContext(activity, options);
// Setup an activity execution context, potentially reusing an existing one if requested.
var reuseActivityExecutionContextId = options?.ReuseActivityExecutionContextId;
// Add the activity context to the workflow context.
workflowExecutionContext.AddActivityExecutionContext(activityExecutionContext);
var activityExecutionContext = reuseActivityExecutionContextId != null
? workflowExecutionContext.ActiveActivityExecutionContexts.FirstOrDefault(x => x.Id == reuseActivityExecutionContextId)
: default;
if (activityExecutionContext == null)
{
// Create a new activity execution context.
activityExecutionContext = workflowExecutionContext.CreateActivityExecutionContext(activity, options);
// Add the activity context to the workflow context.
workflowExecutionContext.AddActivityExecutionContext(activityExecutionContext);
}
// Execute the activity execution pipeline.
await InvokeAsync(activityExecutionContext);

View file

@ -24,6 +24,9 @@ public class QueueBasedActivityScheduler : IActivityScheduler
/// <inheritdoc />
public IEnumerable<ActivityWorkItem> List() => _queue.ToList();
/// <inheritdoc />
public bool Any(Func<ActivityWorkItem, bool> predicate) => _queue.Any(predicate);
/// <inheritdoc />
public void Clear() => _queue.Clear();
}

View file

@ -24,6 +24,9 @@ public class StackBasedActivityScheduler : IActivityScheduler
/// <inheritdoc />
public IEnumerable<ActivityWorkItem> List() => _stack.ToList();
/// <inheritdoc />
public bool Any(Func<ActivityWorkItem, bool> predicate) => _stack.Any(predicate);
/// <inheritdoc />
public void Clear() => _stack.Clear();
}

View file

@ -96,12 +96,18 @@
<None Update="Scenarios\ParentChildInputs\Workflows\parent2.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\ExplicitJoins\Workflows\flow-join-any.json">
<None Update="Scenarios\ExplicitJoins\Workflows\join-any-1.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\FlowchartCompletion\Workflows\workflow6.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\ExplicitJoins\Workflows\join-all-1.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>
<None Update="Scenarios\ExplicitJoins\Workflows\join-all-2.json">
<CopyToOutputDirectory>Always</CopyToOutputDirectory>
</None>

View file

@ -18,14 +18,17 @@ public class ExplicitJoinWaitAnyTests
_services = new TestApplicationBuilder(testOutputHelper).WithCapturingTextWriter(_capturingTextWriter).Build();
}
[Fact(DisplayName = "Workflows with explicit joins using WaitAny wait for any of the specified activities to complete and complete the workflow.")]
public async Task Test1()
[Theory(DisplayName = "Workflows with explicit joins complete the workflow.")]
[InlineData("join-any-1.json", "Start; End")]
[InlineData("join-all-1.json", "Start; Line 1; Line 2; End")]
[InlineData("join-all-2.json", "Start; Line 1; Line 2; End")]
public async Task Test1(string workflowFileName, string expectedLines)
{
// Populate registries.
await _services.PopulateRegistriesAsync();
// Import workflow.
var fileName = "Scenarios/ExplicitJoins/Workflows/flow-join-any.json";
var fileName = $"Scenarios/ExplicitJoins/Workflows/{workflowFileName}";
var workflowDefinition = await _services.ImportWorkflowDefinitionAsync(fileName);
// Execute.
@ -33,7 +36,8 @@ public class ExplicitJoinWaitAnyTests
// Assert expected output.
var lines = _capturingTextWriter.Lines.ToList();
Assert.Equal(new[] { "Start", "End" }, lines);
var expectedLinesArray = expectedLines.Split(";", StringSplitOptions.TrimEntries | StringSplitOptions.RemoveEmptyEntries).ToList();
Assert.Equal(expectedLinesArray, lines);
// Assert expected workflow status.
Assert.Equal(WorkflowStatus.Finished, workflowState.Status);

View file

@ -0,0 +1 @@
{"id":"9948e541b0e540c384c6eb9fdcf1d544","definitionId":"296428cc22bd4f309e858f4b461dbb3d","name":"Join 1","createdAt":"2023-09-12T18:12:54.370149+00:00","version":2,"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":"sL581qWiIEOpRncv5Ez33A","metadata":{},"customProperties":{"source":"FlowchartJsonConverter.cs:46","notFoundConnections":[],"canStartWorkflow":false,"runAsynchronously":false},"activities":[{"mode":{"typeName":"Elsa.Workflows.Core.Activities.Flowchart.Models.FlowJoinMode, Elsa.Workflows.Core","expression":{"type":"Literal","value":"WaitAll"},"memoryReference":{"id":"aA4pfJOzyEqmLOY3U3vCoQ:input-0"}},"id":"aA4pfJOzyEqmLOY3U3vCoQ","name":"FlowJoin1","type":"Elsa.FlowJoin","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":198,"y":-266},"size":{"width":136.859375,"height":50}}}},{"text":{"typeName":"String","expression":{"type":"Literal","value":"Start"},"memoryReference":{"id":"FRXh0PlHekejZiWANAplUg:input-0"}},"id":"FRXh0PlHekejZiWANAplUg","name":"WriteLine1","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-287,"y":-266},"size":{"width":102.21875,"height":50}},"displayText":"Start"}},{"text":{"typeName":"String","expression":{"type":"Literal","value":"Line 1"},"memoryReference":{"id":"GaVrCWqBDkS6YrvydxKvFA:input-0"}},"id":"GaVrCWqBDkS6YrvydxKvFA","name":"WriteLine2","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-40,"y":-316},"size":{"width":110.65625,"height":50}},"displayText":"Line 1"}},{"text":{"typeName":"String","expression":{"type":"Literal","value":"Line 2"},"memoryReference":{"id":"9Z7bt98j00uRMt9TS2FD0g:input-0"}},"id":"9Z7bt98j00uRMt9TS2FD0g","name":"WriteLine3","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-40,"y":-207},"size":{"width":110.65625,"height":50}},"displayText":"Line 2"}},{"text":{"typeName":"String","expression":{"type":"Literal","value":"End"},"memoryReference":{"id":"MplDrS_bUE-8td-XBUvWmA:input-0"}},"id":"MplDrS_bUE-8td-XBUvWmA","name":"WriteLine4","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":430,"y":-266},"size":{"width":139.296875,"height":50}}}}],"connections":[{"source":{"activity":"FRXh0PlHekejZiWANAplUg","port":"Done"},"target":{"activity":"GaVrCWqBDkS6YrvydxKvFA","port":"In"}},{"source":{"activity":"FRXh0PlHekejZiWANAplUg","port":"Done"},"target":{"activity":"9Z7bt98j00uRMt9TS2FD0g","port":"In"}},{"source":{"activity":"9Z7bt98j00uRMt9TS2FD0g","port":"Done"},"target":{"activity":"aA4pfJOzyEqmLOY3U3vCoQ","port":"In"}},{"source":{"activity":"GaVrCWqBDkS6YrvydxKvFA","port":"Done"},"target":{"activity":"aA4pfJOzyEqmLOY3U3vCoQ","port":"In"}},{"source":{"activity":"aA4pfJOzyEqmLOY3U3vCoQ","port":"Done"},"target":{"activity":"MplDrS_bUE-8td-XBUvWmA","port":"In"}}]}}

View file

@ -0,0 +1 @@
{"id":"14984c048db34cceaf5e15ece80f9703","definitionId":"64f8f2c852f04dc28f60707034d7a906","name":"Join 2","createdAt":"2023-09-12T18:12:43.392272+00:00","version":1,"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":"B2jpPDriPEe-zSOhVuAGeQ","metadata":{},"customProperties":{"source":"FlowchartJsonConverter.cs:46","notFoundConnections":[],"canStartWorkflow":false,"runAsynchronously":false},"activities":[{"text":{"typeName":"String","expression":{"type":"Literal","value":"Start"},"memoryReference":{"id":"2yPLgUudQEGPaschdLdkGg:input-0"}},"id":"2yPLgUudQEGPaschdLdkGg","name":"WriteLine1","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-358.5,"y":-229},"size":{"width":139.296875,"height":50}},"displayText":"Start"}},{"body":{"type":"Elsa.Flowchart","version":1,"id":"zaAIlUtQDEaPqeITeRJ4_w","metadata":{},"customProperties":{"source":"FlowchartJsonConverter.cs:46","notFoundConnections":[],"canStartWorkflow":false,"runAsynchronously":false},"activities":[{"text":{"typeName":"String","expression":{"type":"Literal","value":"Line 2"},"memoryReference":{"id":"g3jVTdvU3EK7V5JqRqfqjQ:input-0"}},"id":"g3jVTdvU3EK7V5JqRqfqjQ","name":"WriteLine1","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-279.5,"y":-294},"size":{"width":139.296875,"height":50}},"displayText":"Line 2"}}],"connections":[]},"id":"eHAEEQ-Iu0GwM3mRxZin0Q","name":"FlowNode1","type":"Elsa.FlowNode","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-120,"y":-179},"size":{"width":211.15625,"height":120}},"displayText":"Line 2"}},{"text":{"typeName":"String","expression":{"type":"Literal","value":"Line 1"},"memoryReference":{"id":"7RQ2sbVHo0SxsyK0dvTKEA:input-0"}},"id":"7RQ2sbVHo0SxsyK0dvTKEA","name":"WriteLine2","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":-120,"y":-279},"size":{"width":139.296875,"height":50}},"displayText":"Line 1"}},{"mode":{"typeName":"Elsa.Workflows.Core.Activities.Flowchart.Models.FlowJoinMode, Elsa.Workflows.Core","expression":{"type":"Literal","value":"WaitAll"},"memoryReference":{"id":"dVUuwMm5XUiYpiYoGo5JKA:input-0"}},"id":"dVUuwMm5XUiYpiYoGo5JKA","name":"FlowJoin1","type":"Elsa.FlowJoin","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":180.5,"y":-229},"size":{"width":136.859375,"height":50}}}},{"text":{"typeName":"String","expression":{"type":"Literal","value":"End"},"memoryReference":{"id":"NpZ_yT_irUuE9ASEfYLxbw:input-0"}},"id":"NpZ_yT_irUuE9ASEfYLxbw","name":"WriteLine3","type":"Elsa.WriteLine","version":1,"customProperties":{"canStartWorkflow":false,"runAsynchronously":false},"metadata":{"designer":{"position":{"x":400,"y":-229},"size":{"width":139.296875,"height":50}},"displayText":"End"}}],"connections":[{"source":{"activity":"2yPLgUudQEGPaschdLdkGg","port":"Done"},"target":{"activity":"7RQ2sbVHo0SxsyK0dvTKEA","port":"In"}},{"source":{"activity":"2yPLgUudQEGPaschdLdkGg","port":"Done"},"target":{"activity":"eHAEEQ-Iu0GwM3mRxZin0Q","port":"In"}},{"source":{"activity":"7RQ2sbVHo0SxsyK0dvTKEA","port":"Done"},"target":{"activity":"dVUuwMm5XUiYpiYoGo5JKA","port":"In"}},{"source":{"activity":"eHAEEQ-Iu0GwM3mRxZin0Q","port":"Done"},"target":{"activity":"dVUuwMm5XUiYpiYoGo5JKA","port":"In"}},{"source":{"activity":"dVUuwMm5XUiYpiYoGo5JKA","port":"Done"},"target":{"activity":"NpZ_yT_irUuE9ASEfYLxbw","port":"In"}}]}}