Implicit fork & join (#3331)
* Initial experiment with implicit fork and join construct * Implement flow execution counters * Incremental work on FlowJoin * Implement bookmark removal behavior * Set icon for FlowJoin
This commit is contained in:
parent
c19f24c1d2
commit
9a772af3c8
|
|
@ -1,4 +1,9 @@
|
|||
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
|
||||
<s:Int64 x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/WRAP_LIMIT/@EntryValue">300</s:Int64>
|
||||
<s:Boolean x:Key="/Default/Environment/SettingsMigration/IsMigratorApplied/=JetBrains_002EReSharper_002EPsi_002ECSharp_002ECodeStyle_002ECSharpKeepExistingMigration/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/Environment/SettingsMigration/IsMigratorApplied/=JetBrains_002EReSharper_002EPsi_002ECSharp_002ECodeStyle_002ECSharpPlaceEmbeddedOnSameLineMigration/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/Environment/SettingsMigration/IsMigratorApplied/=JetBrains_002EReSharper_002EPsi_002ECSharp_002ECodeStyle_002ECSharpUseContinuousIndentInsideBracesMigration/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/Environment/SettingsMigration/IsMigratorApplied/=JetBrains_002EReSharper_002EPsi_002ECSharp_002ECodeStyle_002ESettingsUpgrade_002EMigrateBlankLinesAroundFieldToBlankLinesAroundProperty/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Comparers/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Computables/@EntryIndexedValue">True</s:Boolean>
|
||||
<s:Boolean x:Key="/Default/UserDictionary/Words/=Configurer/@EntryIndexedValue">True</s:Boolean>
|
||||
|
|
|
|||
|
|
@ -60,6 +60,7 @@ services
|
|||
.AddActivity<Flowchart>()
|
||||
.AddActivity<FlowDecision>()
|
||||
.AddActivity<FlowSwitch>()
|
||||
.AddActivity<FlowJoin>()
|
||||
.AddActivity<Elsa.Scheduling.Activities.Delay>()
|
||||
.AddActivity<Elsa.Scheduling.Activities.Timer>()
|
||||
.AddActivity<ForEach>()
|
||||
|
|
|
|||
|
|
@ -0,0 +1,13 @@
|
|||
import {FunctionalComponent, h} from '@stencil/core';
|
||||
import {ActivityIconSettings, getActivityIconCssClass} from "./models";
|
||||
|
||||
export const FlowJoinIcon: FunctionalComponent<ActivityIconSettings> = (settings) => (
|
||||
<svg class={getActivityIconCssClass(settings)} width="24" height="24" viewBox="0 0 24 24" stroke-width="2" stroke="currentColor" fill="none" stroke-linecap="round" stroke-linejoin="round">
|
||||
<path stroke="none" d="M0 0h24v24H0z"/>
|
||||
<circle cx="12" cy="18" r="2"/>
|
||||
<circle cx="7" cy="6" r="2"/>
|
||||
<circle cx="17" cy="6" r="2"/>
|
||||
<path d="M7 8v2a2 2 0 0 0 2 2h6a2 2 0 0 0 2 -2v-2"/>
|
||||
<line x1="12" y1="12" x2="12" y2="16"/>
|
||||
</svg>
|
||||
);
|
||||
|
|
@ -15,6 +15,7 @@ import {
|
|||
WriteLineIcon
|
||||
} from "../components/icons/activities";
|
||||
import {WriteHttpResponseIcon} from "../components/icons/activities/write-http-response";
|
||||
import {FlowJoinIcon} from "../components/icons/activities/flow-join";
|
||||
|
||||
export type ActivityType = string;
|
||||
export type ActivityIcon = (ActivityIconSettings?) => any;
|
||||
|
|
@ -37,6 +38,7 @@ export class ActivityIconRegistry {
|
|||
this.add('Elsa.FlowDecision', settings => <FlowDecisionIcon size={settings?.size}/>);
|
||||
this.add('Elsa.Event', settings => <EventIcon size={settings?.size}/>);
|
||||
this.add('Elsa.RunJavaScript', settings => <RunJavaScriptIcon size={settings?.size}/>);
|
||||
this.add('Elsa.FlowJoin', settings => <FlowJoinIcon size={settings?.size}/>);
|
||||
}
|
||||
|
||||
public add(activityType: ActivityType, icon: ActivityIcon) {
|
||||
|
|
|
|||
|
|
@ -0,0 +1,53 @@
|
|||
using Elsa.Common.Extensions;
|
||||
using Elsa.Workflows.Core.Activities.Flowchart.Contracts;
|
||||
using Elsa.Workflows.Core.Activities.Flowchart.Extensions;
|
||||
using Elsa.Workflows.Core.Activities.Flowchart.Models;
|
||||
using Elsa.Workflows.Core.Attributes;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
|
||||
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
|
||||
|
||||
[Activity("Elsa", "Flow", "Merge multiple branches into a single branch of execution.")]
|
||||
public class FlowJoin : ActivityBase, IJoinNode
|
||||
{
|
||||
[Input] public Input<JoinMode> Mode { get; set; } = new(JoinMode.WaitAll);
|
||||
|
||||
protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
|
||||
{
|
||||
var flowchartExecutionContext = context.ParentActivityExecutionContext!;
|
||||
var flowchart = (Flowchart)flowchartExecutionContext.Activity;
|
||||
var inboundActivities = flowchart.Connections.LeftInboundActivities(this).ToList();
|
||||
var flowScope = flowchartExecutionContext.GetProperty<FlowScope>(Flowchart.ScopeProperty)!;
|
||||
var executionCount = flowScope.GetExecutionCount(this);
|
||||
var mode = context.Get(Mode);
|
||||
|
||||
switch (mode)
|
||||
{
|
||||
case JoinMode.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 JoinMode.WaitAny:
|
||||
// Only complete if we haven't already executed.
|
||||
var alreadyExecuted = inboundActivities.Max(x => flowScope.GetExecutionCount(x)) == executionCount;
|
||||
|
||||
if (!alreadyExecuted)
|
||||
{
|
||||
await context.CompleteActivityAsync();
|
||||
ClearBookmarks(flowchart, context);
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
private void ClearBookmarks(Flowchart flowchart, ActivityExecutionContext context)
|
||||
{
|
||||
// Clear any bookmarks created between this join and its most recent fork.
|
||||
var connections = flowchart.Connections;
|
||||
var inboundActivities = connections.LeftAncestorActivities(this).Select(x => x.Id).ToList();
|
||||
context.WorkflowExecutionContext.Bookmarks.RemoveWhere(x => inboundActivities.Contains(x.ActivityId));
|
||||
}
|
||||
}
|
||||
|
|
@ -1,3 +1,5 @@
|
|||
using Elsa.Workflows.Core.Activities.Flowchart.Contracts;
|
||||
using Elsa.Workflows.Core.Activities.Flowchart.Extensions;
|
||||
using Elsa.Workflows.Core.Activities.Flowchart.Models;
|
||||
using Elsa.Workflows.Core.Attributes;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
|
|
@ -9,6 +11,8 @@ namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
|
|||
[Activity("Elsa", "Flow", "A flowchart is a collection of activities and connections between them.")]
|
||||
public class Flowchart : Container
|
||||
{
|
||||
internal const string ScopeProperty = "Scope";
|
||||
|
||||
public Flowchart()
|
||||
{
|
||||
OnSignalReceived<ActivityCompleted>(OnDescendantCompletedAsync);
|
||||
|
|
@ -32,33 +36,72 @@ public class Flowchart : Container
|
|||
{
|
||||
await ScheduleChildrenAsync(signal, context);
|
||||
}
|
||||
|
||||
|
||||
private async Task ScheduleChildrenAsync(ActivityCompleted signal, SignalContext context)
|
||||
{
|
||||
var activityExecutionContext = context.ReceiverActivityExecutionContext;
|
||||
var parent = context.SenderActivityExecutionContext.Activity;
|
||||
var flowchartActivityExecutionContext = context.ReceiverActivityExecutionContext;
|
||||
var completedActivity = context.SenderActivityExecutionContext.Activity;
|
||||
|
||||
if (parent == null!)
|
||||
if (completedActivity == null!)
|
||||
return;
|
||||
|
||||
// Ignore completed activities that are not immediate children.
|
||||
var isDirectChild = Activities.Contains(parent);
|
||||
|
||||
var isDirectChild = Activities.Contains(completedActivity);
|
||||
|
||||
if (!isDirectChild)
|
||||
return;
|
||||
|
||||
// If a specific outcome was provided by the completed activity, use it to find the connection to the next activity.
|
||||
Func<Connection, bool> outboundConnectionsQuery = signal.Result is Outcome outcome
|
||||
? connection => connection.Source == parent && connection.SourcePort == outcome.Name
|
||||
: connection => connection.Source == parent;
|
||||
Func<Connection, bool> outboundConnectionsQuery = signal.Result is Outcome outcome
|
||||
? connection => connection.Source == completedActivity && connection.SourcePort == outcome.Name
|
||||
: connection => connection.Source == completedActivity;
|
||||
|
||||
var outboundConnections = Connections.Where(outboundConnectionsQuery).ToList();
|
||||
var children = outboundConnections.Select(x => x.Target).ToList();
|
||||
var scope = flowchartActivityExecutionContext.GetProperty(ScopeProperty, () => new FlowScope());
|
||||
|
||||
scope.RegisterActivityExecution(completedActivity);
|
||||
|
||||
if (children.Any())
|
||||
activityExecutionContext.ScheduleActivities(children);
|
||||
else
|
||||
await activityExecutionContext.CompleteActivityAsync();
|
||||
{
|
||||
scope.AddActivities(children);
|
||||
|
||||
// Schedule each child, but only if all of its left inbound activities have already executed.
|
||||
foreach (var activity in children)
|
||||
{
|
||||
var inboundActivities = Connections.LeftInboundActivities(activity).ToList();
|
||||
|
||||
// If the completed activity is not part of the left inbound path, always allow its children to be scheduled.
|
||||
if (!inboundActivities.Contains(completedActivity))
|
||||
{
|
||||
flowchartActivityExecutionContext.ScheduleActivity(activity);
|
||||
continue;
|
||||
}
|
||||
|
||||
// If the activity is anything but a join activity, only schedule it if all of its left-inbound activities have executed, effectively implementing a "wait all" join.
|
||||
if (activity is not IJoinNode)
|
||||
{
|
||||
var executionCount = scope.GetExecutionCount(activity);
|
||||
var haveInboundActivitiesExecuted = inboundActivities.All(x => scope.GetExecutionCount(x) > executionCount);
|
||||
|
||||
if (haveInboundActivitiesExecuted)
|
||||
flowchartActivityExecutionContext.ScheduleActivity(activity);
|
||||
}
|
||||
else
|
||||
{
|
||||
flowchartActivityExecutionContext.ScheduleActivity(activity);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (!children.Any())
|
||||
{
|
||||
// If there are no more pending activities in any of the scopes, mark this activity as completed.
|
||||
var hasPendingChildren = scope.HasPendingActivities();
|
||||
|
||||
if (!hasPendingChildren)
|
||||
await flowchartActivityExecutionContext.CompleteActivityAsync();
|
||||
}
|
||||
|
||||
context.StopPropagation();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,10 @@
|
|||
using Elsa.Workflows.Core.Services;
|
||||
|
||||
namespace Elsa.Workflows.Core.Activities.Flowchart.Contracts;
|
||||
|
||||
/// <summary>
|
||||
/// Gives implementing activities a chance to customize certain flowchart execution behaviors.
|
||||
/// </summary>
|
||||
public interface IJoinNode : IActivity
|
||||
{
|
||||
}
|
||||
|
|
@ -5,18 +5,75 @@ namespace Elsa.Workflows.Core.Activities.Flowchart.Extensions;
|
|||
|
||||
public static class ConnectionsExtensions
|
||||
{
|
||||
public static IEnumerable<Connection> Descendants(this ICollection<Connection> allConnections, IActivity parent)
|
||||
public static IEnumerable<Connection> Descendants(this ICollection<Connection> connections, IActivity parent)
|
||||
{
|
||||
var children = allConnections.Where(x => parent == x.Source).ToList();
|
||||
var visitedActivities = new HashSet<IActivity>();
|
||||
return connections.Descendants(parent, visitedActivities);
|
||||
}
|
||||
|
||||
public static IEnumerable<Connection> Ancestors(this ICollection<Connection> connections, IActivity activity)
|
||||
{
|
||||
var visitedActivities = new HashSet<IActivity>();
|
||||
return connections.Ancestors(activity, visitedActivities);
|
||||
}
|
||||
|
||||
public static IEnumerable<Connection> InboundConnections(this ICollection<Connection> connections, IActivity activity) => connections.Where(x => x.Target == activity).ToList();
|
||||
|
||||
public static IEnumerable<Connection> LeftInboundConnections(this ICollection<Connection> connections, IActivity activity)
|
||||
{
|
||||
// We only take "left" inbound connections, which means we exclude descendent connections looping back.
|
||||
var descendantConnections = connections.Descendants(activity).ToList();
|
||||
var filteredConnections = connections.InboundConnections(activity).Except(descendantConnections).ToList();
|
||||
|
||||
return filteredConnections;
|
||||
}
|
||||
|
||||
public static IEnumerable<Connection> LeftAncestorConnections(this ICollection<Connection> connections, IActivity activity)
|
||||
{
|
||||
// We only take "left" inbound connections, which means we exclude descendent connections looping back.
|
||||
var descendantConnections = connections.Descendants(activity).ToList();
|
||||
var filteredConnections = connections.Ancestors(activity).Except(descendantConnections).ToList();
|
||||
|
||||
return filteredConnections;
|
||||
}
|
||||
|
||||
public static IEnumerable<IActivity> InboundActivities(this ICollection<Connection> connections, IActivity activity) => connections.InboundConnections(activity).Select(x => x.Source);
|
||||
public static IEnumerable<IActivity> LeftInboundActivities(this ICollection<Connection> connections, IActivity activity) => connections.LeftInboundConnections(activity).Select(x => x.Source);
|
||||
public static IEnumerable<IActivity> LeftAncestorActivities(this ICollection<Connection> connections, IActivity activity) => connections.LeftAncestorConnections(activity).Select(x => x.Source);
|
||||
|
||||
private static IEnumerable<Connection> Descendants(this ICollection<Connection> connections, IActivity parent, ISet<IActivity> visitedActivities)
|
||||
{
|
||||
var children = connections.Where(x => parent == x.Source && !visitedActivities.Contains(x.Target)).ToList();
|
||||
|
||||
foreach (var child in children)
|
||||
{
|
||||
visitedActivities.Add(child.Target);
|
||||
yield return child;
|
||||
|
||||
var descendants = allConnections.Descendants(child.Target).ToList();
|
||||
var descendants = connections.Descendants(child.Target, visitedActivities).ToList();
|
||||
|
||||
foreach (var descendant in descendants)
|
||||
{
|
||||
yield return descendant;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static IEnumerable<Connection> Ancestors(this ICollection<Connection> connections, IActivity activity, ISet<IActivity> visitedActivities)
|
||||
{
|
||||
var parents = connections.Where(x => activity == x.Target && !visitedActivities.Contains(x.Source)).ToList();
|
||||
|
||||
foreach (var parent in parents)
|
||||
{
|
||||
visitedActivities.Add(parent.Source);
|
||||
yield return parent;
|
||||
|
||||
var ancestors = connections.Ancestors(parent.Source, visitedActivities).ToList();
|
||||
|
||||
foreach (var ancestor in ancestors)
|
||||
{
|
||||
yield return ancestor;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,22 @@
|
|||
using System.Diagnostics;
|
||||
using System.Text.Json.Serialization;
|
||||
|
||||
namespace Elsa.Workflows.Core.Activities.Flowchart.Models;
|
||||
|
||||
[DebuggerDisplay("ActivityId = {ActivityId}, ExecutionCount = {ExecutionCount}")]
|
||||
public class ActivityFlowState
|
||||
{
|
||||
[JsonConstructor]
|
||||
public ActivityFlowState()
|
||||
{
|
||||
}
|
||||
|
||||
public ActivityFlowState(string activityId, long executionCount = 0)
|
||||
{
|
||||
ActivityId = activityId;
|
||||
ExecutionCount = executionCount;
|
||||
}
|
||||
|
||||
public string ActivityId { get; set; } = default!;
|
||||
public long ExecutionCount { get; set; }
|
||||
}
|
||||
|
|
@ -0,0 +1,75 @@
|
|||
using System.Text.Json.Serialization;
|
||||
using Elsa.Workflows.Core.Services;
|
||||
|
||||
namespace Elsa.Workflows.Core.Activities.Flowchart.Models;
|
||||
|
||||
public class FlowScope
|
||||
{
|
||||
[JsonConstructor]
|
||||
public FlowScope()
|
||||
{
|
||||
}
|
||||
|
||||
public FlowScope(string ownerActivityId)
|
||||
{
|
||||
OwnerActivityId = ownerActivityId;
|
||||
Activities.Add(ownerActivityId, new ActivityFlowState(ownerActivityId, 1));
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// The activity from which the scope was created.
|
||||
/// </summary>
|
||||
public string OwnerActivityId { get; set; } = default!;
|
||||
|
||||
/// <summary>
|
||||
/// A list of scheduled activity IDs and a flag whether they executed or not.
|
||||
/// </summary>
|
||||
public IDictionary<string, ActivityFlowState> Activities { get; set; } =
|
||||
new Dictionary<string, ActivityFlowState>();
|
||||
|
||||
public void AddActivities(IEnumerable<IActivity> activities, long executionCount = 0)
|
||||
{
|
||||
foreach (var activity in activities)
|
||||
AddActivity(activity, executionCount);
|
||||
}
|
||||
|
||||
public void AddActivity(IActivity activity, long executionCount = 0) => EnsureActivity(activity, executionCount);
|
||||
|
||||
public ActivityFlowState EnsureActivity(IActivity activity, long executionCount = 0)
|
||||
{
|
||||
if (Activities.ContainsKey(activity.Id))
|
||||
return Activities[activity.Id];
|
||||
|
||||
var state = new ActivityFlowState(activity.Id, executionCount);
|
||||
Activities.Add(activity.Id, state);
|
||||
return state;
|
||||
|
||||
}
|
||||
|
||||
public void RegisterActivityExecution(IActivity activity)
|
||||
{
|
||||
var state = Activities.TryGetValue(activity.Id, out var s) ? s : default;
|
||||
|
||||
if (state == null)
|
||||
{
|
||||
state = new ActivityFlowState(activity.Id);
|
||||
Activities[activity.Id] = state;
|
||||
}
|
||||
|
||||
state.ExecutionCount++;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Return a list excluding any activities that already executed.
|
||||
/// </summary>
|
||||
public IEnumerable<IActivity> ExcludeExecutedActivities(IEnumerable<IActivity> activities) =>
|
||||
activities.Where(x => !Activities.ContainsKey(x.Id) || Activities[x.Id].ExecutionCount == 0);
|
||||
|
||||
public bool HasPendingActivities()
|
||||
{
|
||||
var sample = Activities.Values.First().ExecutionCount;
|
||||
return Activities.Values.Any(x => x.ExecutionCount != sample);
|
||||
}
|
||||
|
||||
public long GetExecutionCount(IActivity activity) => EnsureActivity(activity).ExecutionCount;
|
||||
}
|
||||
|
|
@ -0,0 +1,7 @@
|
|||
namespace Elsa.Workflows.Core.Activities.Flowchart.Activities;
|
||||
|
||||
public enum JoinMode
|
||||
{
|
||||
WaitAll,
|
||||
WaitAny
|
||||
}
|
||||
|
|
@ -21,22 +21,10 @@ public class ActivityInvoker : IActivityInvoker
|
|||
ActivityExecutionContext? owner,
|
||||
IEnumerable<MemoryBlockReference>? memoryReferences = default)
|
||||
{
|
||||
var cancellationToken = workflowExecutionContext.CancellationToken;
|
||||
|
||||
// Get a handle to the parent execution context.
|
||||
var parentActivityExecutionContext = owner;
|
||||
var parentExpressionExecutionContext = parentActivityExecutionContext?.ExpressionExecutionContext;
|
||||
|
||||
// Setup an activity execution context.
|
||||
var workflowMemory = workflowExecutionContext.MemoryRegister;
|
||||
var workflow = workflowExecutionContext.Workflow;
|
||||
var transientProperties = workflowExecutionContext.TransientProperties;
|
||||
var input = workflowExecutionContext.Input;
|
||||
var applicationProperties = ExpressionExecutionContextExtensions.CreateApplicationPropertiesFrom(workflow, transientProperties, input);
|
||||
var parentMemory = parentActivityExecutionContext?.ExpressionExecutionContext.Memory ?? workflowMemory;
|
||||
var activityMemory = new MemoryRegister(workflowMemory);
|
||||
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, parentMemory, parentExpressionExecutionContext, applicationProperties, cancellationToken);
|
||||
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, parentActivityExecutionContext, expressionExecutionContext, activity, cancellationToken);
|
||||
var activityExecutionContext = workflowExecutionContext.CreateActivityExecutionContext(activity, owner);
|
||||
|
||||
// Declare memory.
|
||||
if (memoryReferences != null)
|
||||
|
|
|
|||
|
|
@ -78,7 +78,7 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
|
|||
foreach (var entry in state.ActivityOutput)
|
||||
{
|
||||
var activityId = entry.Key;
|
||||
var node = workflowExecutionContext.FindActivityNodeById(activityId);
|
||||
var node = workflowExecutionContext.FindNodeById(activityId);
|
||||
var activityType = node.Activity.GetType();
|
||||
|
||||
foreach (var outputEntry in entry.Value)
|
||||
|
|
@ -96,7 +96,7 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
|
|||
foreach (var completionCallbackEntry in state.CompletionCallbacks)
|
||||
{
|
||||
var owner = workflowExecutionContext.ActivityExecutionContexts.First(x => x.Id == completionCallbackEntry.OwnerId);
|
||||
var child = workflowExecutionContext.FindActivityNodeById(completionCallbackEntry.ChildId).Activity;
|
||||
var child = workflowExecutionContext.FindNodeById(completionCallbackEntry.ChildId).Activity;
|
||||
var callbackName = completionCallbackEntry.MethodName;
|
||||
var callbackDelegate = owner.Activity.GetActivityCompletionCallback(callbackName);
|
||||
workflowExecutionContext.AddCompletionCallback(owner, child, callbackDelegate);
|
||||
|
|
@ -133,21 +133,11 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
|
|||
{
|
||||
ActivityExecutionContext CreateActivityExecutionContext(ActivityExecutionContextState activityExecutionContextState)
|
||||
{
|
||||
var cancellationToken = workflowExecutionContext.CancellationToken;
|
||||
var activity = workflowExecutionContext.FindActivityById(activityExecutionContextState.ScheduledActivityId);
|
||||
var workflowMemory = workflowExecutionContext.MemoryRegister;
|
||||
var workflow = workflowExecutionContext.Workflow;
|
||||
var expressionInput = workflowExecutionContext.Input;
|
||||
var transientProperties = workflowExecutionContext.TransientProperties;
|
||||
var applicationProperties = ExpressionExecutionContextExtensions.CreateApplicationPropertiesFrom(workflow, transientProperties, expressionInput);
|
||||
var activityMemory = new MemoryRegister(workflowMemory);
|
||||
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, activityMemory, null, applicationProperties, cancellationToken);
|
||||
var properties = activityExecutionContextState.Properties;
|
||||
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, default, expressionExecutionContext, activity, cancellationToken)
|
||||
{
|
||||
Id = activityExecutionContextState.Id,
|
||||
ApplicationProperties = properties
|
||||
};
|
||||
var activityExecutionContext = workflowExecutionContext.CreateActivityExecutionContext(activity);
|
||||
activityExecutionContext.Id = activityExecutionContextState.Id;
|
||||
activityExecutionContext.ApplicationProperties = properties;
|
||||
return activityExecutionContext;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
using System.Diagnostics;
|
||||
using System.Text.Json.Serialization;
|
||||
using Elsa.Workflows.Core.Behaviors;
|
||||
using Elsa.Workflows.Core.Helpers;
|
||||
|
|
@ -5,6 +6,7 @@ using Elsa.Workflows.Core.Services;
|
|||
|
||||
namespace Elsa.Workflows.Core.Models;
|
||||
|
||||
[DebuggerDisplay("{Type} - {Id}")]
|
||||
public abstract class ActivityBase : IActivity, ISignalHandler
|
||||
{
|
||||
private readonly ICollection<SignalHandlerRegistration> _signalHandlers = new List<SignalHandlerRegistration>();
|
||||
|
|
|
|||
|
|
@ -135,6 +135,18 @@ public class ActivityExecutionContext
|
|||
public void ClearBookmarks() => _bookmarks.Clear();
|
||||
|
||||
public T? GetProperty<T>(string key) => ApplicationProperties!.TryGetValue<T?>(key, out var value) ? value : default;
|
||||
|
||||
public T GetProperty<T>(string key, Func<T> defaultValue)
|
||||
{
|
||||
if (ApplicationProperties.TryGetValue<T?>(key, out var value))
|
||||
return value!;
|
||||
|
||||
value = defaultValue();
|
||||
ApplicationProperties[key] = value!;
|
||||
|
||||
return value!;
|
||||
}
|
||||
|
||||
public void SetProperty<T>(string key, T? value) => ApplicationProperties[key] = value!;
|
||||
|
||||
public T UpdateProperty<T>(string key, Func<T?, T> updater) where T : notnull
|
||||
|
|
|
|||
|
|
@ -22,27 +22,27 @@ public class ActivityNode
|
|||
foreach (var child in Children)
|
||||
{
|
||||
yield return child;
|
||||
|
||||
|
||||
var descendants = child.Descendants();
|
||||
|
||||
foreach (var descendant in descendants)
|
||||
yield return descendant;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public IEnumerable<ActivityNode> Ancestors()
|
||||
{
|
||||
foreach (var parent in Parents)
|
||||
{
|
||||
yield return parent;
|
||||
|
||||
|
||||
var ancestors = parent.Ancestors();
|
||||
|
||||
foreach (var ancestor in ancestors)
|
||||
yield return ancestor;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public IEnumerable<ActivityNode> Siblings() => Parents.SelectMany(parent => parent.Children);
|
||||
public IEnumerable<ActivityNode> SiblingsAndCousins() => Parents.SelectMany(parent => parent.Descendants());
|
||||
}
|
||||
|
|
@ -109,9 +109,9 @@ public class WorkflowExecutionContext
|
|||
_completionCallbackEntries.Remove(entry);
|
||||
}
|
||||
|
||||
public ActivityNode FindActivityNodeById(string nodeId) => NodeIdLookup[nodeId];
|
||||
public ActivityNode FindNodeById(string nodeId) => NodeIdLookup[nodeId];
|
||||
public ActivityNode FindNodeByActivity(IActivity activity) => NodeActivityLookup[activity];
|
||||
public IActivity FindActivityById(string activityId) => FindActivityNodeById(activityId).Activity;
|
||||
public IActivity FindActivityById(string activityId) => FindNodeById(activityId).Activity;
|
||||
public T? GetProperty<T>(string key) => Properties.TryGetValue(key, out var value) ? (T?)value : default(T);
|
||||
public void SetProperty<T>(string key, T value) => Properties[key] = value!;
|
||||
|
||||
|
|
@ -134,6 +134,16 @@ public class WorkflowExecutionContext
|
|||
SubStatus = subStatus;
|
||||
}
|
||||
|
||||
public ActivityExecutionContext CreateActivityExecutionContext(IActivity activity, ActivityExecutionContext? parentContext = default)
|
||||
{
|
||||
var parentExpressionExecutionContext = parentContext?.ExpressionExecutionContext;
|
||||
var applicationProperties = ExpressionExecutionContextExtensions.CreateApplicationPropertiesFrom(Workflow, TransientProperties, Input);
|
||||
var parentMemory = parentContext?.ExpressionExecutionContext.Memory ?? MemoryRegister;
|
||||
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, parentMemory, parentExpressionExecutionContext, applicationProperties, CancellationToken);
|
||||
var activityExecutionContext = new ActivityExecutionContext(this, parentContext, expressionExecutionContext, activity, CancellationToken);
|
||||
return activityExecutionContext;
|
||||
}
|
||||
|
||||
private WorkflowStatus GetMainStatus(WorkflowSubStatus subStatus) =>
|
||||
subStatus switch
|
||||
{
|
||||
|
|
|
|||
|
|
@ -98,7 +98,8 @@ public class ActivityDescriber : IActivityDescriber
|
|||
{
|
||||
var inputAttribute = propertyInfo.GetCustomAttribute<InputAttribute>();
|
||||
var descriptionAttribute = propertyInfo.GetCustomAttribute<DescriptionAttribute>();
|
||||
var wrappedPropertyType = propertyInfo.PropertyType.GenericTypeArguments[0];
|
||||
var propertyType = propertyInfo.PropertyType;
|
||||
var wrappedPropertyType = !typeof(Input).IsAssignableFrom(propertyType) ? propertyType : propertyInfo.PropertyType.GenericTypeArguments[0];
|
||||
|
||||
yield return new InputDescriptor
|
||||
(
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
using System.Reflection;
|
||||
using Elsa.Workflows.Core.Attributes;
|
||||
using Elsa.Workflows.Core.Models;
|
||||
using Elsa.Workflows.Management.Extensions;
|
||||
using Elsa.Workflows.Management.Models;
|
||||
using Elsa.Workflows.Management.Services;
|
||||
|
|
@ -38,13 +39,14 @@ public class PropertyOptionsResolver : IPropertyOptionsResolver
|
|||
{
|
||||
var isNullable = activityPropertyInfo.PropertyType.IsNullableType();
|
||||
var propertyType = isNullable ? activityPropertyInfo.PropertyType.GetTypeOfNullable() : activityPropertyInfo.PropertyType;
|
||||
var wrappedPropertyType = !typeof(Input).IsAssignableFrom(propertyType) ? propertyType : activityPropertyInfo.PropertyType.GenericTypeArguments[0];
|
||||
|
||||
items = null;
|
||||
|
||||
if (!propertyType.IsEnum)
|
||||
if (!wrappedPropertyType.IsEnum)
|
||||
return false;
|
||||
|
||||
items = propertyType.GetEnumNames().Select(x => new SelectListItem(x.Humanize(LetterCasing.Title), x)).ToList();
|
||||
items = wrappedPropertyType.GetEnumNames().Select(x => new SelectListItem(x.Humanize(LetterCasing.Title), x)).ToList();
|
||||
|
||||
if (isNullable)
|
||||
items.Insert(0, new SelectListItem("-", ""));
|
||||
|
|
|
|||
Loading…
Reference in a new issue