Refactor back-paddling of branching activities

This commit is contained in:
Sipke Schoorstra 2021-01-12 13:22:40 +01:00
parent fcfe7f08b4
commit 8613133add
14 changed files with 48 additions and 130 deletions

View file

@ -13,6 +13,7 @@ namespace Elsa.Models
{
Variables = new Variables();
ScheduledActivities = new Stack<ScheduledActivity>();
Scopes = new Stack<string>();
}
public string DefinitionId { get; set; } = default!;
@ -38,8 +39,9 @@ namespace Elsa.Models
get => _blockingActivities;
set => _blockingActivities = new HashSet<BlockingActivity>(value, BlockingActivityEqualityComparer.Instance);
}
public WorkflowFault? Fault { get; set; }
public Stack<ScheduledActivity> ScheduledActivities { get; set; }
public Stack<string> Scopes { get; set; }
}
}

View file

@ -25,6 +25,7 @@ namespace Elsa.Services
{
if (!IsScheduled)
{
context.WorkflowInstance.Scopes.Push(Id);
IsScheduled = true;
await OnEnterAsync(context);
return Outcome(Enter);

View file

@ -1,7 +0,0 @@
namespace Elsa.Services.Models
{
public interface IBranchingActivity
{
void Unwind(ActivityExecutionContext context);
}
}

View file

@ -1,6 +1,7 @@
using System;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
@ -11,7 +12,7 @@ namespace Elsa.Activities.ControlFlow
Description = "Iterate between two numbers.",
Outcomes = new[] { OutcomeNames.Iterate, OutcomeNames.Done }
)]
public class For : IteratingActivity
public class For : Activity
{
[ActivityProperty(Hint = "An expression that evaluates to the starting number.")]
public long Start { get; set; }
@ -46,6 +47,7 @@ namespace Elsa.Activities.ControlFlow
if (loop)
{
context.WorkflowInstance.Scopes.Push(Id);
CurrentValue = currentValue + Step;
return Outcome(OutcomeNames.Iterate, currentValue);
}

View file

@ -14,7 +14,7 @@ namespace Elsa.Activities.ControlFlow
Description = "Iterate over a collection.",
Outcomes = new[] { OutcomeNames.Iterate, OutcomeNames.Done }
)]
public class ForEach : IteratingActivity
public class ForEach : Activity
{
[ActivityProperty(Hint = "Enter an expression that evaluates to a collection of items to iterate over.")]
public ICollection<object> Items { get; set; } = new Collection<object>();
@ -44,6 +44,7 @@ namespace Elsa.Activities.ControlFlow
{
var output = collection[currentIndex];
CurrentIndex = currentIndex + 1;
context.WorkflowInstance.Scopes.Push(Id);
return Combine(Outcome(OutcomeNames.Iterate, output));
}

View file

@ -1,5 +1,3 @@
using System.Collections.Generic;
using System.Linq;
using Elsa.ActivityResults;
using Elsa.Attributes;
using Elsa.Services;
@ -14,58 +12,35 @@ namespace Elsa.Activities.ControlFlow
Description = "Evaluate a Boolean expression and continue execution depending on the result.",
Outcomes = new[] { True, False, OutcomeNames.Done }
)]
public class IfElse : Activity, IBranchingActivity
public class IfElse : Activity
{
public const string True = "True";
public const string False = "False";
[ActivityProperty(Hint = "The condition to evaluate.")]
public bool Condition { get; set; }
private bool EnteredScope
{
get => GetState<bool>();
set => SetState(value);
}
protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context)
{
var outcome = Condition ? True : False;
return Outcome(outcome);
}
public void Unwind(ActivityExecutionContext context)
{
var workflowExecutionContext = context.WorkflowExecutionContext;
var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint;
var activityExecutionContext = context;
var currentActivityBlueprint = activityExecutionContext.ActivityBlueprint;
var currentActivityId = currentActivityBlueprint.Id;
// An IfElse activity completed within a burst of execution, which means its outcome (true or false) did not yield child activities to execute (an empty branch).
// Schedule its child activities connected to the "Done" outcome.
if (currentActivityBlueprint.Type == nameof(IfElse))
if (!EnteredScope)
{
var nextActivities = GetNextActivities(workflowExecutionContext, currentActivityId);
workflowExecutionContext.ScheduleActivities(nextActivities);
context.WorkflowInstance.Scopes.Push(Id);
EnteredScope = true;
}
else
{
// Get all incoming connections.
var inboundConnections = workflowBlueprint.GetInboundConnectionPath(currentActivityId).ToList();
// Filter out those connections who have a source of IfElse.
var query =
from inboundConnection in inboundConnections
let parentActivityBlueprint = inboundConnection.Source.Activity
where inboundConnection.Source.Activity.Type == nameof(IfElse)
select inboundConnection;
var firstMatch = query.FirstOrDefault();
if (firstMatch != null && firstMatch.Source.Outcome != OutcomeNames.Done)
{
var parentActivityBlueprint = firstMatch.Source.Activity;
var nextActivities = OutcomeResult.GetNextActivities(workflowExecutionContext, parentActivityBlueprint.Id, new[] { OutcomeNames.Done }).ToList();
workflowExecutionContext.ScheduleActivities(nextActivities);
}
EnteredScope = false;
return Done();
}
var outcome = Condition ? True : False;
return Outcome(outcome);
}
private IEnumerable<string> GetNextActivities(WorkflowExecutionContext workflowExecutionContext, string currentActivityId) => OutcomeResult.GetNextActivities(workflowExecutionContext, currentActivityId, new[] { OutcomeNames.Done });
}
}

View file

@ -1,29 +0,0 @@
using System.Linq;
using Elsa.Services;
using Elsa.Services.Models;
// ReSharper disable once CheckNamespace
namespace Elsa.Activities.ControlFlow
{
public abstract class IteratingActivity : Activity, IBranchingActivity
{
public void Unwind(ActivityExecutionContext context)
{
var workflowExecutionContext = context.WorkflowExecutionContext;
var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint;
var currentActivityId = context.ActivityBlueprint.Id;
var inboundConnections = workflowBlueprint.GetInboundConnectionPath(currentActivityId).ToList();
var query =
from inboundConnection in inboundConnections
let parentActivityBlueprint = inboundConnection.Source.Activity
where inboundConnection.Source.Activity.Type == Type
select inboundConnection;
var firstMatch = query.FirstOrDefault();
if(firstMatch != null && firstMatch.Source.Outcome == OutcomeNames.Iterate)
workflowExecutionContext.ScheduleActivity(firstMatch.Source.Activity.Id);
}
}
}

View file

@ -11,7 +11,7 @@ namespace Elsa.Activities.ControlFlow
Description = "Execute while a given condition is true.",
Outcomes = new[] { OutcomeNames.Iterate, OutcomeNames.Done }
)]
public class While : IteratingActivity
public class While : Activity
{
[ActivityProperty(Hint = "The condition to evaluate.")]
public bool Condition { get; set; }
@ -21,8 +21,12 @@ namespace Elsa.Activities.ControlFlow
var loop = Condition;
if (loop)
{
context.WorkflowInstance.Scopes.Push(Id);
return Outcome(OutcomeNames.Iterate);
}
return Done();
}
}

View file

@ -55,11 +55,6 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton(options.StorageFactory)
.AddStartupTask<CreateSubscriptions>();
services
.AddTransient<IBranchingActivity, IfElse>()
.AddTransient<IBranchingActivity, For>()
.AddTransient<IBranchingActivity, While>();
options
.AddWorkflowsCore()
.AddCoreActivities();
@ -163,7 +158,6 @@ namespace Microsoft.Extensions.DependencyInjection
.AddActivity<ParallelForEach>()
.AddActivity<Fork>()
.AddActivity<IfElse>()
.AddActivity<Switch>()
.AddActivity<Join>()
.AddActivity<Switch>()
.AddActivity<While>()

View file

@ -59,6 +59,7 @@ namespace Elsa.Handlers
private async ValueTask SaveWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken)
{
workflowExecutionContext.PruneActivityData();
var workflowInstance = workflowExecutionContext.WorkflowInstance;
await _workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken);
_logger.LogDebug("Committed workflow {WorkflowInstanceId} to storage", workflowInstance.Id);

View file

@ -1,62 +1,31 @@
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Events;
using Elsa.Models;
using Elsa.Services.Models;
using MediatR;
namespace Elsa.Handlers
{
/// <summary>
/// Walks up the tee of inbound connections along the "Iterate" outcome of the looping construct (While/For/ForEach) and re-schedules the looping activity.
/// Walks up the tree of inbound connections along the "Iterate" outcome of the looping construct (While/For/ForEach) and re-schedules the looping activity.
/// Also handles composite activity re-scheduling.
/// </summary>
public class RescheduleBranchingActivitiesAndContainers : INotificationHandler<WorkflowExecutionBurstCompleted>
{
private readonly IEnumerable<IBranchingActivity> _branchingActivities;
public RescheduleBranchingActivitiesAndContainers(IEnumerable<IBranchingActivity> branchingActivities)
{
_branchingActivities = branchingActivities;
}
public Task Handle(WorkflowExecutionBurstCompleted notification, CancellationToken cancellationToken)
{
var workflowExecutionContext = notification.WorkflowExecutionContext;
var activityExecutionContext = notification.ActivityExecutionContext;
foreach (var branchingActivity in _branchingActivities)
{
if (workflowExecutionContext.HasScheduledActivities || workflowExecutionContext.Status != WorkflowStatus.Running)
break;
branchingActivity.Unwind(activityExecutionContext);
}
ScheduleContainers(notification);
return Task.CompletedTask;
}
private static void ScheduleContainers(WorkflowExecutionBurstCompleted notification)
{
var workflowExecutionContext = notification.WorkflowExecutionContext;
var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint;
var workflowExecutionContext = activityExecutionContext.WorkflowExecutionContext;
// If no suspension has been instructed, re-schedule any container activities.
if (workflowExecutionContext.HasScheduledActivities || workflowExecutionContext.Status != WorkflowStatus.Running)
return;
var activityBlueprint = notification.ActivityExecutionContext.ActivityBlueprint;
// Re-schedule the parent activity, if any.
if (activityBlueprint.Parent != null && workflowBlueprint.GetActivity(activityBlueprint.Parent.Id) != null)
{
var output = GetFinishOutput(workflowExecutionContext);
workflowExecutionContext.ScheduleActivity(activityBlueprint.Parent.Id, output);
}
}
if (workflowExecutionContext.HasScheduledActivities || workflowExecutionContext.Status != WorkflowStatus.Running || !workflowExecutionContext.WorkflowInstance.Scopes.Any())
return Task.CompletedTask;
private static FinishOutput? GetFinishOutput(WorkflowExecutionContext workflowExecutionContext) => workflowExecutionContext.WorkflowInstance.Output as FinishOutput;
var parentActivityId = workflowExecutionContext.WorkflowInstance.Scopes.Pop();
workflowExecutionContext.ScheduleActivity(parentActivityId, activityExecutionContext.Output);
return Task.CompletedTask;
}
}
}

View file

@ -290,13 +290,13 @@ namespace Elsa.Services
await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken);
await result.ExecuteAsync(activityExecutionContext, cancellationToken);
workflowExecutionContext.WorkflowInstance.Output = activityExecutionContext.Output;
workflowExecutionContext.PruneActivityData();
activityOperation = Execute;
workflowExecutionContext.CompletePass();
await _mediator.Publish(new WorkflowExecutionPassCompleted(workflowExecutionContext, activityExecutionContext), cancellationToken);
if (!workflowExecutionContext.HasScheduledActivities)
await _mediator.Publish(new WorkflowExecutionBurstCompleted(workflowExecutionContext, activityExecutionContext), cancellationToken);
activityOperation = Execute;
}
if (workflowExecutionContext.HasBlockingActivities)

View file

@ -19,6 +19,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Configuration
builder.Ignore(x => x.BlockingActivities);
builder.Ignore(x => x.Fault);
builder.Ignore(x => x.ScheduledActivities);
builder.Ignore(x => x.Scopes);
builder.Ignore(x => x.Variables);
builder.Property<string>("Data");
}

View file

@ -1,4 +1,5 @@
using System;
using System.Collections.Generic;
using System.Linq.Expressions;
using Elsa.Models;
using Elsa.Persistence.Specifications;
@ -34,6 +35,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Stores
entity.ActivityOutput,
entity.BlockingActivities,
entity.ScheduledActivities,
entity.Scopes,
entity.Fault
};
@ -52,6 +54,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Stores
entity.ActivityOutput,
entity.BlockingActivities,
entity.ScheduledActivities,
entity.Scopes,
entity.Fault
};
@ -66,6 +69,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Stores
entity.ActivityOutput = data.ActivityOutput;
entity.BlockingActivities = data.BlockingActivities;
entity.ScheduledActivities = data.ScheduledActivities;
entity.Scopes = data.Scopes ?? new Stack<string>();
entity.Fault = data.Fault;
}
}