Merge remote-tracking branch 'origin/main'

This commit is contained in:
Sipke Schoorstra 2025-04-16 08:41:52 +02:00
commit 9c05a4247c
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
8 changed files with 255 additions and 226 deletions

View file

@ -24,7 +24,7 @@ public static class JsonObjectExtensions
/// <returns>A <see cref="JsonObject"/> representing the specified value.</returns>
public static JsonNode SerializeToNode(this object value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
options ??= new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
return JsonSerializer.SerializeToNode(value, options)!;
}
@ -37,7 +37,7 @@ public static class JsonObjectExtensions
/// <returns>A <see cref="JsonObject"/> representing the specified value.</returns>
public static JsonArray SerializeToArray(this object value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
options ??= new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
return JsonSerializer.SerializeToNode(value, options)!.AsArray();
}
@ -50,7 +50,7 @@ public static class JsonObjectExtensions
/// <returns>A <see cref="JsonObject"/> representing the specified value.</returns>
public static JsonArray SerializeToArray<T>(this IEnumerable<T> value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
options ??= new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
return JsonSerializer.SerializeToNode(value, options)!.AsArray();
}

View file

@ -23,7 +23,7 @@ public class ActivityDesignerMetadata
/// <param name="y">The y coordinate of the activity.</param>
/// <param name="width">The width of the activity.</param>
/// <param name="height">The height of the activity.</param>
public ActivityDesignerMetadata(double x, double y, double width, double height) : this(new Position(x, y), new Size(width, height))
public ActivityDesignerMetadata(double x, double y, double width, double height) : this(new(x, y), new(width, height))
{
}

View file

@ -8,10 +8,15 @@ public class Connection
/// <summary>
/// Gets or sets the source endpoint.
/// </summary>
public Endpoint Source { get; set; } = default!;
public Endpoint Source { get; set; } = null!;
/// <summary>
/// Gets or sets the target endpoint.
/// </summary>
public Endpoint Target { get; set; } = default!;
public Endpoint Target { get; set; } = null!;
/// <summary>
/// Gets or sets the vertices of the connection.
/// </summary>
public ICollection<Position> Vertices { get; set; } = [];
}

View file

@ -1,134 +1,134 @@
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Extensions;
using Elsa.Workflows.Activities.Flowchart.Contracts;
using Elsa.Workflows.Activities.Flowchart.Models;
using Elsa.Workflows.Attributes;
using Elsa.Workflows.Options;
using Elsa.Workflows.Signals;
namespace Elsa.Workflows.Activities.Flowchart.Activities;
/// <summary>
/// A flowchart consists of a collection of activities and connections between them.
/// </summary>
[Activity("Elsa", "Flow", "A flowchart is a collection of activities and connections between them.")]
[Browsable(false)]
public class Flowchart : Container
{
using System.ComponentModel;
using System.Runtime.CompilerServices;
using Elsa.Extensions;
using Elsa.Workflows.Activities.Flowchart.Contracts;
using Elsa.Workflows.Activities.Flowchart.Models;
using Elsa.Workflows.Attributes;
using Elsa.Workflows.Options;
using Elsa.Workflows.Signals;
namespace Elsa.Workflows.Activities.Flowchart.Activities;
/// <summary>
/// A flowchart consists of a collection of activities and connections between them.
/// </summary>
[Activity("Elsa", "Flow", "A flowchart is a collection of activities and connections between them.")]
[Browsable(false)]
public class Flowchart : Container
{
private const string ScopeProperty = "FlowScope";
private const string GraphTransientProperty = "FlowGraph";
private const string BackwardConnectionActivityInput = "BackwardConnection";
/// <inheritdoc />
public Flowchart([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
{
OnSignalReceived<ScheduleActivityOutcomes>(OnScheduleOutcomesAsync);
OnSignalReceived<ScheduleChildActivity>(OnScheduleChildActivityAsync);
OnSignalReceived<CancelSignal>(OnActivityCanceledAsync);
}
/// <summary>
/// The activity to execute when the flowchart starts.
/// </summary>
[Port] [Browsable(false)] public IActivity? Start { get; set; }
/// <summary>
/// A list of connections between activities.
/// </summary>
public ICollection<Connection> Connections { get; set; } = new List<Connection>();
/// <inheritdoc />
protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
var startActivity = GetStartActivity(context);
if (startActivity == null)
{
// Nothing else to execute.
await context.CompleteActivityAsync();
return;
}
// Schedule the start activity.
await context.ScheduleActivityAsync(startActivity, OnChildCompletedAsync);
}
private IActivity? GetStartActivity(ActivityExecutionContext context)
{
// If there's a trigger that triggered this workflow, use that.
var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId;
var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : default;
if (triggerActivity != null)
return triggerActivity;
// If an explicit Start activity was provided, use that.
if (Start != null)
return Start;
// If there is a Start activity on the flowchart, use that.
var startActivity = Activities.FirstOrDefault(x => x is Start);
if (startActivity != null)
return startActivity;
// If there's an activity marked as "Can Start Workflow", use that.
var canStartWorkflowActivity = Activities.FirstOrDefault(x => x.GetCanStartWorkflow());
if (canStartWorkflowActivity != null)
return canStartWorkflowActivity;
// If there is a single activity that has no inbound connections, use that.
var root = GetRootActivity();
if (root != null)
return root;
// If no start activity found, return the first activity.
return Activities.FirstOrDefault();
}
/// <summary>
/// Checks if there is any pending work for the flowchart.
/// </summary>
private bool HasPendingWork(ActivityExecutionContext context)
{
var workflowExecutionContext = context.WorkflowExecutionContext;
var activityIds = Activities.Select(x => x.Id).ToList();
var children = context.Children;
var hasRunningActivityInstances = children.Where(x => activityIds.Contains(x.Activity.Id)).Any(x => x.Status == ActivityStatus.Running);
var hasPendingWork = workflowExecutionContext.Scheduler.List().Any(workItem =>
{
var ownerInstanceId = workItem.Owner?.Id;
if (ownerInstanceId == null)
return false;
if (ownerInstanceId == context.Id)
return true;
var ownerContext = context.WorkflowExecutionContext.ActivityExecutionContexts.First(x => x.Id == ownerInstanceId);
var ancestors = ownerContext.GetAncestors().ToList();
return ancestors.Any(x => x == context);
});
return hasRunningActivityInstances || hasPendingWork;
}
private IActivity? GetRootActivity()
{
// Get the first activity that has no inbound connections.
var query =
from activity in Activities
let inboundConnections = Connections.Any(x => x.Target.Activity == activity)
where !inboundConnections
select activity;
var rootActivity = query.FirstOrDefault();
return rootActivity;
/// <inheritdoc />
public Flowchart([CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : base(source, line)
{
OnSignalReceived<ScheduleActivityOutcomes>(OnScheduleOutcomesAsync);
OnSignalReceived<ScheduleChildActivity>(OnScheduleChildActivityAsync);
OnSignalReceived<CancelSignal>(OnActivityCanceledAsync);
}
/// <summary>
/// The activity to execute when the flowchart starts.
/// </summary>
[Port] [Browsable(false)] public IActivity? Start { get; set; }
/// <summary>
/// A list of connections between activities.
/// </summary>
public ICollection<Connection> Connections { get; set; } = new List<Connection>();
/// <inheritdoc />
protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
var startActivity = GetStartActivity(context);
if (startActivity == null)
{
// Nothing else to execute.
await context.CompleteActivityAsync();
return;
}
// Schedule the start activity.
await context.ScheduleActivityAsync(startActivity, OnChildCompletedAsync);
}
private IActivity? GetStartActivity(ActivityExecutionContext context)
{
// If there's a trigger that triggered this workflow, use that.
var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId;
var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : null;
if (triggerActivity != null)
return triggerActivity;
// If an explicit Start activity was provided, use that.
if (Start != null)
return Start;
// If there is a Start activity on the flowchart, use that.
var startActivity = Activities.FirstOrDefault(x => x is Start);
if (startActivity != null)
return startActivity;
// If there's an activity marked as "Can Start Workflow", use that.
var canStartWorkflowActivity = Activities.FirstOrDefault(x => x.GetCanStartWorkflow());
if (canStartWorkflowActivity != null)
return canStartWorkflowActivity;
// If there is a single activity that has no inbound connections, use that.
var root = GetRootActivity();
if (root != null)
return root;
// If no start activity found, return the first activity.
return Activities.FirstOrDefault();
}
/// <summary>
/// Checks if there is any pending work for the flowchart.
/// </summary>
private bool HasPendingWork(ActivityExecutionContext context)
{
var workflowExecutionContext = context.WorkflowExecutionContext;
var activityIds = Activities.Select(x => x.Id).ToList();
var children = context.Children;
var hasRunningActivityInstances = children.Where(x => activityIds.Contains(x.Activity.Id)).Any(x => x.Status == ActivityStatus.Running);
var hasPendingWork = workflowExecutionContext.Scheduler.List().Any(workItem =>
{
var ownerInstanceId = workItem.Owner?.Id;
if (ownerInstanceId == null)
return false;
if (ownerInstanceId == context.Id)
return true;
var ownerContext = context.WorkflowExecutionContext.ActivityExecutionContexts.First(x => x.Id == ownerInstanceId);
var ancestors = ownerContext.GetAncestors().ToList();
return ancestors.Any(x => x == context);
});
return hasRunningActivityInstances || hasPendingWork;
}
private IActivity? GetRootActivity()
{
// Get the first activity that has no inbound connections.
var query =
from activity in Activities
let inboundConnections = Connections.Any(x => x.Target.Activity == activity)
where !inboundConnections
select activity;
var rootActivity = query.FirstOrDefault();
return rootActivity;
}
private FlowGraph GetFlowGraph(ActivityExecutionContext context)
@ -142,16 +142,16 @@ public class Flowchart : Container
return context.GetProperty(ScopeProperty, () => new FlowScope());
}
private async ValueTask OnChildCompletedAsync(ActivityCompletedContext context)
private async ValueTask OnChildCompletedAsync(ActivityCompletedContext context)
{
var flowchartContext = context.TargetContext;
var completedActivityContext = context.ChildContext;
var completedActivity = completedActivityContext.Activity;
var result = context.Result;
if (flowchartContext.Activity != this)
{
throw new Exception("Target context activity must be this flowchart");
if (flowchartContext.Activity != this)
{
throw new Exception("Target context activity must be this flowchart");
}
// If the completed activity's status is anything but "Completed", do not schedule its outbound activities.
@ -174,13 +174,13 @@ public class Flowchart : Container
var flowGraph = GetFlowGraph(flowchartContext);
var flowScope = GetFlowScope(flowchartContext);
var completedActivityExcecutedByBackwardConnection = completedActivityContext.ActivityInput.GetValueOrDefault<bool>(BackwardConnectionActivityInput);
bool hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExcecutedByBackwardConnection);
bool hasScheduledActivity = await ScheduleOutboundActivitiesAsync(flowGraph, flowScope, flowchartContext, completedActivity, outcomes, completedActivityExcecutedByBackwardConnection);
// If there are not any outbound connections, complete the flowchart activity if there is no other pending work
if (!hasScheduledActivity)
{
await CompleteIfNoPendingWorkAsync(flowchartContext);
}
}
}
/// <summary>
@ -355,65 +355,65 @@ public class Flowchart : Container
await activityExecutionContext.CancelActivityAsync();
}
}
private async Task CompleteIfNoPendingWorkAsync(ActivityExecutionContext context)
{
var hasPendingWork = HasPendingWork(context);
if (!hasPendingWork)
{
var hasFaultedActivities = context.Children.Any(x => x.Status == ActivityStatus.Faulted);
if (!hasFaultedActivities)
{
await context.CompleteActivityAsync();
}
}
}
private async ValueTask OnScheduleOutcomesAsync(ScheduleActivityOutcomes signal, SignalContext context)
{
var flowchartContext = context.ReceiverActivityExecutionContext;
var schedulingActivityContext = context.SenderActivityExecutionContext;
var schedulingActivity = schedulingActivityContext.Activity;
var outcomes = signal.Outcomes;
var outboundConnections = Connections.Where(connection => connection.Source.Activity == schedulingActivity && outcomes.Contains(connection.Source.Port!)).ToList();
var outboundActivities = outboundConnections.Select(x => x.Target.Activity).ToList();
if (outboundActivities.Any())
{
// Schedule each child.
foreach (var activity in outboundActivities) await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync);
}
}
private async ValueTask OnScheduleChildActivityAsync(ScheduleChildActivity signal, SignalContext context)
{
var flowchartContext = context.ReceiverActivityExecutionContext;
var activity = signal.Activity;
var activityExecutionContext = signal.ActivityExecutionContext;
if (activityExecutionContext != null)
{
await flowchartContext.ScheduleActivityAsync(activityExecutionContext.Activity, new ScheduleWorkOptions
{
ExistingActivityExecutionContext = activityExecutionContext,
CompletionCallback = OnChildCompletedAsync,
Input = signal.Input
});
}
else
{
await flowchartContext.ScheduleActivityAsync(activity, new ScheduleWorkOptions
{
CompletionCallback = OnChildCompletedAsync,
Input = signal.Input
});
}
}
private async ValueTask OnActivityCanceledAsync(CancelSignal signal, SignalContext context)
{
await CompleteIfNoPendingWorkAsync(context.ReceiverActivityExecutionContext);
}
private async Task CompleteIfNoPendingWorkAsync(ActivityExecutionContext context)
{
var hasPendingWork = HasPendingWork(context);
if (!hasPendingWork)
{
var hasFaultedActivities = context.Children.Any(x => x.Status == ActivityStatus.Faulted);
if (!hasFaultedActivities)
{
await context.CompleteActivityAsync();
}
}
}
private async ValueTask OnScheduleOutcomesAsync(ScheduleActivityOutcomes signal, SignalContext context)
{
var flowchartContext = context.ReceiverActivityExecutionContext;
var schedulingActivityContext = context.SenderActivityExecutionContext;
var schedulingActivity = schedulingActivityContext.Activity;
var outcomes = signal.Outcomes;
var outboundConnections = Connections.Where(connection => connection.Source.Activity == schedulingActivity && outcomes.Contains(connection.Source.Port!)).ToList();
var outboundActivities = outboundConnections.Select(x => x.Target.Activity).ToList();
if (outboundActivities.Any())
{
// Schedule each child.
foreach (var activity in outboundActivities) await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync);
}
}
private async ValueTask OnScheduleChildActivityAsync(ScheduleChildActivity signal, SignalContext context)
{
var flowchartContext = context.ReceiverActivityExecutionContext;
var activity = signal.Activity;
var activityExecutionContext = signal.ActivityExecutionContext;
if (activityExecutionContext != null)
{
await flowchartContext.ScheduleActivityAsync(activityExecutionContext.Activity, new ScheduleWorkOptions
{
ExistingActivityExecutionContext = activityExecutionContext,
CompletionCallback = OnChildCompletedAsync,
Input = signal.Input
});
}
else
{
await flowchartContext.ScheduleActivityAsync(activity, new ScheduleWorkOptions
{
CompletionCallback = OnChildCompletedAsync,
Input = signal.Input
});
}
}
private async ValueTask OnActivityCanceledAsync(CancelSignal signal, SignalContext context)
{
await CompleteIfNoPendingWorkAsync(context.ReceiverActivityExecutionContext);
}
}

View file

@ -33,25 +33,30 @@ public class Connection : IEquatable<Connection>
/// <param name="target">The target endpoint.</param>
public Connection(IActivity source, IActivity target)
{
Source = new Endpoint(source);
Target = new Endpoint(target);
Source = new(source);
Target = new(target);
}
/// <summary>
/// The source endpoint.
/// </summary>
public Endpoint Source { get; set; } = default!;
public Endpoint Source { get; set; } = null!;
/// <summary>
/// The target endpoint.
/// </summary>
public Endpoint Target { get; set; } = default!;
public override string ToString() =>
$"{Source.Activity.Id}{(string.IsNullOrEmpty(Source.Port) ? "" : $":{Source.Port}")}->" +
$"{Target.Activity.Id}{(string.IsNullOrEmpty(Target.Port) ? "" : $":{Target.Port}")}";
// Implement equality logic
public Endpoint Target { get; set; } = null!;
/// <summary>
/// A collection of points representing the vertices of the connection.
/// </summary>
public ICollection<Position> Vertices { get; set; } = [];
public override string ToString() =>
$"{Source.Activity.Id}{(string.IsNullOrEmpty(Source.Port) ? "" : $":{Source.Port}")}->" +
$"{Target.Activity.Id}{(string.IsNullOrEmpty(Target.Port) ? "" : $":{Target.Port}")}";
// Implement equality logic
public bool Equals(Connection? other)
{
if (other == null) return false;
@ -72,7 +77,7 @@ public class Connection : IEquatable<Connection>
{
return e1.Activity.Equals(e2.Activity) && e1.Port == e2.Port;
}
private static int GetEndpointHashCode(Endpoint endpoint)
{
return HashCode.Combine(endpoint.Activity.GetHashCode(), endpoint.Port?.GetHashCode() ?? 0);

View file

@ -0,0 +1,10 @@
using JetBrains.Annotations;
namespace Elsa.Workflows.Activities.Flowchart.Models;
[UsedImplicitly]
public class Position
{
public double X { get; set; }
public double Y { get; set; }
}

View file

@ -29,15 +29,23 @@ public class ConnectionJsonConverter : JsonConverter<Connection>
var sourceElement = doc.RootElement.GetProperty("source");
var targetElement = doc.RootElement.GetProperty("target");
var sourceId = sourceElement.GetProperty("activity").GetString()!;
var targetId = targetElement.TryGetProperty("activity", out var targetIdValue) ? targetIdValue.GetString() : default;
var sourcePort = sourceElement.TryGetProperty("port", out var sourcePortValue) ? sourcePortValue.GetString() : default;
var targetPort = targetElement.TryGetProperty("port", out var targetPortValue) ? targetPortValue.GetString() : default;
var sourceActivity = _activities.TryGetValue(sourceId, out var s) ? s : default!;
var targetActivity = targetId != null ? _activities.TryGetValue(targetId, out var t) ? t : default! : default!;
var targetId = targetElement.TryGetProperty("activity", out var targetIdValue) ? targetIdValue.GetString() : null;
var sourcePort = sourceElement.TryGetProperty("port", out var sourcePortValue) ? sourcePortValue.GetString() : null;
var targetPort = targetElement.TryGetProperty("port", out var targetPortValue) ? targetPortValue.GetString() : null;
var sourceActivity = _activities.TryGetValue(sourceId, out var s) ? s : null!;
var targetActivity = targetId != null ? _activities.TryGetValue(targetId, out var t) ? t : null! : null!;
var source = new Endpoint(sourceActivity, sourcePort);
var target = new Endpoint(targetActivity, targetPort);
return new Connection(source, target);
var verticesElement = doc.RootElement.TryGetProperty("vertices", out var verticesValue) ? verticesValue : default;
var vertices = Array.Empty<Position>();
if (verticesElement.ValueKind == JsonValueKind.Array)
vertices = verticesElement.Deserialize<Position[]>(options)!;
return new(source, target)
{
Vertices = vertices
};
}
/// <inheritdoc />
@ -57,7 +65,8 @@ public class ConnectionJsonConverter : JsonConverter<Connection>
{
Activity = value.Target.Activity.Id,
Port = value.Target.Port
}
},
Vertices = value.Vertices
};
JsonSerializer.Serialize(writer, model, options);

View file

@ -24,9 +24,9 @@ public class FlowchartJsonConverter(IIdentityGenerator identityGenerator, IWellK
throw new JsonException("Failed to parse JsonDocument");
var id = doc.RootElement.TryGetProperty("id", out var idAttribute) ? idAttribute.GetString()! : identityGenerator.GenerateId();
var nodeId = doc.RootElement.TryGetProperty("nodeId", out var nodeIdAttribute) ? nodeIdAttribute.GetString() : default;
var name = doc.RootElement.TryGetProperty("name", out var nameElement) ? nameElement.GetString() : default;
var type = doc.RootElement.TryGetProperty("type", out var typeElement) ? typeElement.GetString() : default;
var nodeId = doc.RootElement.TryGetProperty("nodeId", out var nodeIdAttribute) ? nodeIdAttribute.GetString() : null;
var name = doc.RootElement.TryGetProperty("name", out var nameElement) ? nameElement.GetString() : null;
var type = doc.RootElement.TryGetProperty("type", out var typeElement) ? typeElement.GetString() : null;
var version = doc.RootElement.TryGetProperty("version", out var versionElement) ? versionElement.GetInt32() : 1;
var connectionsElement = doc.RootElement.TryGetProperty("connections", out var connectionsEl) ? connectionsEl : default;
@ -165,8 +165,8 @@ public class FlowchartJsonConverter(IIdentityGenerator identityGenerator, IWellK
connectionSerializerOptions.Converters.Add(new ObsoleteConnectionJsonConverter(activityDictionary));
var obsoleteConnections = connectionsElement.ValueKind != JsonValueKind.Undefined
? connectionsElement.Deserialize<ICollection<ObsoleteConnection>>(connectionSerializerOptions)?.Where(x => x.Source != null! && x.Target != null!).ToList() ?? new List<ObsoleteConnection>()
: new List<ObsoleteConnection>();
? connectionsElement.Deserialize<ICollection<ObsoleteConnection>>(connectionSerializerOptions)?.Where(x => x.Source != null! && x.Target != null!).ToList() ?? []
: [];
return obsoleteConnections.Select(x => new Connection(new Endpoint(x.Source, x.SourcePort), new Endpoint(x.Target, x.TargetPort))).ToList();
}
@ -174,7 +174,7 @@ public class FlowchartJsonConverter(IIdentityGenerator identityGenerator, IWellK
connectionSerializerOptions.Converters.Add(new ConnectionJsonConverter(activityDictionary));
return connectionsElement.ValueKind != JsonValueKind.Undefined
? connectionsElement.Deserialize<ICollection<Connection>>(connectionSerializerOptions)?.Where(x => x.Source != null! && x.Target != null!).ToList() ?? new List<Connection>()
: new List<Connection>();
? connectionsElement.Deserialize<ICollection<Connection>>(connectionSerializerOptions)?.Where(x => x.Source != null! && x.Target != null!).ToList() ?? []
: [];
}
}