Add commit state behavior support in workflows

Introduced `ActivityCommitStateBehavior` and `WorkflowCommitStateOptions` to enable flexible state commit handling in workflows. Integrated commit logic into activity and workflow execution contexts and middleware. This improves control over when state is committed during workflow execution.
This commit is contained in:
Sipke Schoorstra 2025-01-25 23:26:37 +01:00
parent adea7fa666
commit 9365328f0b
16 changed files with 258 additions and 36 deletions

View file

@ -1,4 +1,5 @@
using System.Text.Json.Nodes;
using Elsa.Api.Client.Resources.WorkflowDefinitions.Models;
using Elsa.Api.Client.Shared.Models;
namespace Elsa.Api.Client.Extensions;
@ -182,4 +183,14 @@ public static class ActivityExtensions
/// Sets a value indicating whether the specified activity can trigger the workflow.
/// </summary>
public static void SetRunAsynchronously(this JsonObject activity, bool value) => activity.SetProperty(JsonValue.Create(value), "customProperties", "runAsynchronously");
/// <summary>
/// Gets the commit state behavior for the specified activity.
/// </summary>
public static ActivityCommitStateBehavior GetCommitStateBehavior(this JsonObject activity) => activity.TryGetProperty<ActivityCommitStateBehavior?>("customProperties", "commitStateBehavior") ?? ActivityCommitStateBehavior.Default;
/// <summary>
/// Sets the commit state behavior for the specified activity.
/// </summary>
public static void SetCommitStateBehavior(this JsonObject activity, ActivityCommitStateBehavior value) => activity.SetProperty(JsonValue.Create(value.ToString()), "customProperties", "commitStateBehavior");
}

View file

@ -15,52 +15,43 @@ public static class JsonObjectExtensions
{
return obj.ContainsKey("type") && obj.ContainsKey("id") && obj.ContainsKey("version");
}
/// <summary>
/// Serializes the specified value to a <see cref="JsonObject"/>.
/// </summary>
/// <param name="value">The value to serialize.</param>
/// <param name="options">The <see cref="JsonSerializerOptions"/> to use.</param>
/// <returns>A <see cref="JsonObject"/> representing the specified value.</returns>
public static JsonNode SerializeToNode(this object value, JsonSerializerOptions? options = default)
public static JsonNode SerializeToNode(this object value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase
};
options ??= new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
return JsonSerializer.SerializeToNode(value, options)!;
}
/// <summary>
/// Serializes the specified value to a <see cref="JsonArray"/>.
/// </summary>
/// <param name="value">The value to serialize.</param>
/// <param name="options">The <see cref="JsonSerializerOptions"/> to use.</param>
/// <returns>A <see cref="JsonObject"/> representing the specified value.</returns>
public static JsonArray SerializeToArray(this object value, JsonSerializerOptions? options = default)
public static JsonArray SerializeToArray(this object value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase
};
options ??= new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
return JsonSerializer.SerializeToNode(value, options)!.AsArray();
}
/// <summary>
/// Serializes the specified value to a <see cref="JsonArray"/>.
/// </summary>
/// <param name="value">The value to serialize.</param>
/// <param name="options">The <see cref="JsonSerializerOptions"/> to use.</param>
/// <returns>A <see cref="JsonObject"/> representing the specified value.</returns>
public static JsonArray SerializeToArray<T>(this IEnumerable<T> value, JsonSerializerOptions? options = default)
public static JsonArray SerializeToArray<T>(this IEnumerable<T> value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase
};
options ??= new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
return JsonSerializer.SerializeToNode(value, options)!.AsArray();
}
@ -71,19 +62,23 @@ public static class JsonObjectExtensions
/// <param name="options">The <see cref="JsonSerializerOptions"/> to use.</param>
/// <typeparam name="T">The type to deserialize to.</typeparam>
/// <returns>The deserialized value.</returns>
public static T Deserialize<T>(this JsonNode value, JsonSerializerOptions? options = default)
public static T Deserialize<T>(this JsonNode value, JsonSerializerOptions? options = null)
{
options ??= new JsonSerializerOptions
{
PropertyNamingPolicy = JsonNamingPolicy.CamelCase
};
options ??= new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
if (value is JsonObject jsonObject)
return JsonSerializer.Deserialize<T>(jsonObject, options)!;
if (value is JsonArray jsonArray)
return JsonSerializer.Deserialize<T>(jsonArray, options)!;
if (typeof(T).IsEnum || (Nullable.GetUnderlyingType(typeof(T))?.IsEnum ?? false))
{
if (value.GetValueKind() == JsonValueKind.Null)
return default!;
return (T)Enum.Parse(Nullable.GetUnderlyingType(typeof(T)) ?? typeof(T), value.ToString());
}
if (value is JsonValue jsonValue)
return jsonValue.GetValue<T>();
@ -101,7 +96,7 @@ public static class JsonObjectExtensions
model = GetPropertyContainer(model, path);
model[path.Last()] = value?.SerializeToNode();
}
/// <summary>
/// Sets the property value of the specified model.
/// </summary>
@ -113,7 +108,7 @@ public static class JsonObjectExtensions
model = GetPropertyContainer(model, path);
model[path.Last()] = value?.SerializeToNode();
}
/// <summary>
/// Sets the property value of the specified model.
/// </summary>
@ -125,7 +120,7 @@ public static class JsonObjectExtensions
model = GetPropertyContainer(model, path);
model[path.Last()] = new JsonArray(value.Select(x => x.SerializeToNode()).ToArray());
}
/// <summary>
/// Gets the property value of the specified model.
/// </summary>
@ -139,7 +134,7 @@ public static class JsonObjectExtensions
foreach (var prop in path.SkipLast(1))
{
if (currentModel[prop] is not JsonObject value)
return default;
return null;
currentModel = value;
}
@ -147,6 +142,25 @@ public static class JsonObjectExtensions
return currentModel[path.Last()];
}
/// <summary>
/// Gets the property value of the specified model.
/// </summary>
/// <param name="model">The model to get the property value from.</param>
/// <param name="path">The path to the property.</param>
/// <typeparam name="T">The type to deserialize to.</typeparam>
/// <returns>The property value.</returns>
public static T? TryGetProperty<T>(this JsonObject model, params string[] path)
{
try
{
return model.GetProperty<T>(path);
}
catch (Exception e)
{
return default;
}
}
/// <summary>
/// Gets the property value of the specified model.
/// </summary>
@ -173,7 +187,7 @@ public static class JsonObjectExtensions
var property = GetProperty(model, path);
return property != null ? property.Deserialize<T>(options) : default;
}
/// <summary>
/// Returns the property container of the specified model.
/// </summary>
@ -190,5 +204,4 @@ public static class JsonObjectExtensions
return model;
}
}

View file

@ -0,0 +1,29 @@
namespace Elsa.Api.Client.Resources.WorkflowDefinitions.Models;
public enum ActivityCommitStateBehavior
{
/// <summary>
/// Never commit state, regardless of the workflow commit state options.
/// </summary>
Never,
/// <summary>
/// Look at the workflow commit state options to determine if state should be committed.
/// </summary>
Default,
/// <summary>
/// Commit state before the activity starts.
/// </summary>
Executing,
/// <summary>
/// Commit state after the activity executes.
/// </summary>
Executed,
/// <summary>
/// Commit state before the activity starts and after the activity executes.
/// </summary>
BeforeAndAfterExecution
}

View file

@ -0,0 +1,19 @@
namespace Elsa.Api.Client.Resources.WorkflowDefinitions.Models;
public class WorkflowCommitStateOptions
{
/// <summary>
/// Commit workflow state before the workflow starts.
/// </summary>
public bool Starting { get; set; }
/// <summary>
/// Commit workflow state before an activity executes, unless the activity is configured to not commit state.
/// </summary>
public bool ActivityExecuting { get; set; }
/// <summary>
/// Commit workflow state after an activity executes, unless the activity is configured to not commit state.
/// </summary>
public bool ActivityExecuted { get; set; }
}

View file

@ -29,4 +29,9 @@ public class WorkflowOptions
/// The type of <c>IIncidentStrategy</c> to use when a fault occurs in the workflow.
/// </summary>
public string? IncidentStrategyType { get; set; }
/// <summary>
/// The options for committing workflow state.
/// </summary>
public WorkflowCommitStateOptions CommitStateOptions { get; set; } = new();
}

View file

@ -37,6 +37,7 @@ public partial class WorkflowExecutionContext : IExecutionContext
private readonly IList<ActivityCompletionCallbackEntry> _completionCallbackEntries = new List<ActivityCompletionCallbackEntry>();
private IList<ActivityExecutionContext> _activityExecutionContexts;
private readonly IHasher _hasher;
private readonly ICommitStateHandler _commitStateHandler;
/// <summary>
/// Initializes a new instance of <see cref="WorkflowExecutionContext"/>.
@ -61,6 +62,7 @@ public partial class WorkflowExecutionContext : IExecutionContext
ActivityRegistry = serviceProvider.GetRequiredService<IActivityRegistry>();
ActivityRegistryLookup = serviceProvider.GetRequiredService<IActivityRegistryLookupService>();
_hasher = serviceProvider.GetRequiredService<IHasher>();
_commitStateHandler = serviceProvider.GetRequiredService<ICommitStateHandler>();
SubStatus = WorkflowSubStatus.Pending;
Id = id;
CorrelationId = correlationId;
@ -238,6 +240,11 @@ public partial class WorkflowExecutionContext : IExecutionContext
/// The current sub status of the workflow.
public WorkflowSubStatus SubStatus { get; internal set; }
/// <summary>
/// The previous sub status of the workflow.
/// </summary>
public WorkflowSubStatus PreviousSubStatus { get; internal set; }
/// The root <see cref="MemoryRegister"/> associated with the execution context.
public MemoryRegister MemoryRegister { get; private set; } = null!;
@ -510,8 +517,9 @@ public partial class WorkflowExecutionContext : IExecutionContext
internal void TransitionTo(WorkflowSubStatus subStatus)
{
if (!ValidateStatusTransition())
throw new Exception($"Cannot transition from {SubStatus} to {subStatus}");
throw new($"Cannot transition from {SubStatus} to {subStatus}");
PreviousSubStatus = SubStatus;
SubStatus = subStatus;
UpdatedAt = SystemClock.UtcNow;
@ -614,4 +622,9 @@ public partial class WorkflowExecutionContext : IExecutionContext
var currentMainStatus = GetMainStatus(SubStatus);
return currentMainStatus != WorkflowStatus.Finished;
}
public Task CommitAsync()
{
return _commitStateHandler.CommitAsync(this, CancellationToken);
}
}

View file

@ -4,5 +4,6 @@ namespace Elsa.Workflows;
public interface ICommitStateHandler
{
Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default);
Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default);
}

View file

@ -7,6 +7,8 @@
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Cprimitives_005Csetname/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Csignaling_005Cactivities_005Csignalreceived/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=activities_005Cworkflows_005Cfreeflowchart/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=commitstates_005Ccontracts/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=commitstates_005Cmodels/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=constants/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=contexts/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=contracts/@EntryIndexedValue">True</s:Boolean>

View file

@ -0,0 +1,29 @@
namespace Elsa.Workflows;
public enum ActivityCommitStateBehavior
{
/// <summary>
/// Never commit state, regardless of the workflow commit state options.
/// </summary>
Never,
/// <summary>
/// Look at the workflow commit state options to determine if state should be committed.
/// </summary>
Default,
/// <summary>
/// Commit state before the activity starts.
/// </summary>
Executing,
/// <summary>
/// Commit state after the activity executes.
/// </summary>
Executed,
/// <summary>
/// Commit state before the activity starts and after the activity executes.
/// </summary>
BeforeAndAfterExecution
}

View file

@ -11,6 +11,7 @@ public static class ActivityPropertyExtensions
private static readonly string[] CanStartWorkflowPropertyName = ["canStartWorkflow", "CanStartWorkflow"];
private static readonly string[] RunAsynchronouslyPropertyName = ["runAsynchronously", "RunAsynchronously"];
private static readonly string[] SourcePropertyName = ["source", "Source"];
private static readonly string[] CommitStateBehaviorName = ["commitStateBehavior", "CommitStateBehavior"];
/// <summary>
/// Gets a flag indicating whether this activity can be used for starting a workflow.
@ -46,6 +47,16 @@ public static class ActivityPropertyExtensions
/// Sets the source file and line number where this activity was instantiated, if any.
/// </summary>
public static void SetSource(this IActivity activity, string value) => activity.CustomProperties[SourcePropertyName[0]] = value;
/// <summary>
/// Gets the commit state behavior for the specified activity.
/// </summary>
public static ActivityCommitStateBehavior GetCommitStateBehavior(this IActivity activity) => activity.CustomProperties.GetValueOrDefault(CommitStateBehaviorName, () => ActivityCommitStateBehavior.Default);
/// <summary>
/// Sets the commit state behavior for the specified activity.
/// </summary>
public static void SetCommitStateBehavior(this IActivity activity, ActivityCommitStateBehavior value) => activity.CustomProperties[CommitStateBehaviorName[0]] = value;
/// <summary>
/// Sets the source file and line number where this activity was instantiated, if any.
@ -62,7 +73,7 @@ public static class ActivityPropertyExtensions
/// <summary>
/// Gets the display text for the specified activity.
/// </summary>
public static string? GetDisplayText(this IActivity activity) => activity.Metadata.TryGetValue("displayText", out var value) ? value.ToString() : default;
public static string? GetDisplayText(this IActivity activity) => activity.Metadata.TryGetValue("displayText", out var value) ? value.ToString() : null;
/// <summary>
/// Sets the display text for the specified activity.
@ -72,7 +83,7 @@ public static class ActivityPropertyExtensions
/// <summary>
/// Gets the description for the specified activity.
/// </summary>
public static string? GetDescription(this IActivity activity) => activity.Metadata.TryGetValue("description", out var value) ? value.ToString() : default;
public static string? GetDescription(this IActivity activity) => activity.Metadata.TryGetValue("description", out var value) ? value.ToString() : null;
/// <summary>
/// Sets the description for the specified activity.

View file

@ -47,6 +47,10 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I
context.AddExecutionLogEntry("Precondition Failed", "Cannot execute at this time");
return;
}
// Conditionally commit the workflow state.
if(ShouldCommitWhenStarting(context))
await context.WorkflowExecutionContext.CommitAsync();
context.TransitionTo(ActivityStatus.Running);
@ -78,6 +82,10 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I
workflowExecutionContext.Bookmarks.AddRange(context.Bookmarks);
logger.LogDebug("Added {BookmarkCount} bookmarks to the workflow execution context", context.Bookmarks.Count);
}
// Conditionally commit the workflow state.
if(ShouldCommitWhenExecuted(context))
await context.WorkflowExecutionContext.CommitAsync();
}
/// <summary>
@ -111,4 +119,40 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I
// Evaluate input properties.
await context.EvaluateInputPropertiesAsync();
}
private bool ShouldCommitWhenStarting(ActivityExecutionContext context)
{
var behavior = context.Activity.GetCommitStateBehavior();
if (behavior == ActivityCommitStateBehavior.Executing)
return true;
if (behavior == ActivityCommitStateBehavior.Default)
{
var workflowOptions = context.WorkflowExecutionContext.Workflow.Options.CommitStateOptions;
if(workflowOptions.ActivityExecuting)
return true;
}
return false;
}
private bool ShouldCommitWhenExecuted(ActivityExecutionContext context)
{
var behavior = context.Activity.GetCommitStateBehavior();
if (behavior == ActivityCommitStateBehavior.Executed)
return true;
if (behavior == ActivityCommitStateBehavior.Default)
{
var workflowOptions = context.WorkflowExecutionContext.Workflow.Options.CommitStateOptions;
if(workflowOptions.ActivityExecuted)
return true;
}
return false;
}
}

View file

@ -36,6 +36,8 @@ public class DefaultActivitySchedulerMiddleware : WorkflowExecutionMiddleware
context.TransitionTo(WorkflowSubStatus.Executing);
await ConditionallyCommitStateAsync(context);
while (scheduler.HasAny)
{
// Do not start a workflow if cancellation has been requested.
@ -65,4 +67,12 @@ public class DefaultActivitySchedulerMiddleware : WorkflowExecutionMiddleware
await _activityInvoker.InvokeAsync(context, workItem.Activity, options);
}
private async Task ConditionallyCommitStateAsync(WorkflowExecutionContext context)
{
var shouldCommit = context.Workflow.Options.CommitStateOptions.Starting;
if (shouldCommit)
await context.CommitAsync();
}
}

View file

@ -0,0 +1,19 @@
namespace Elsa.Workflows.Models;
public class WorkflowCommitStateOptions
{
/// <summary>
/// Commit workflow state before the workflow starts.
/// </summary>
public bool Starting { get; set; }
/// <summary>
/// Commit workflow state before an activity executes, unless the activity is configured to not commit state.
/// </summary>
public bool ActivityExecuting { get; set; }
/// <summary>
/// Commit workflow state after an activity executes, unless the activity is configured to not commit state.
/// </summary>
public bool ActivityExecuted { get; set; }
}

View file

@ -29,4 +29,9 @@ public class WorkflowOptions
/// The type of <see cref="IIncidentStrategy"/> to use when a fault occurs in the workflow.
/// </summary>
public Type? IncidentStrategyType { get; set; }
/// <summary>
/// The options for committing workflow state.
/// </summary>
public WorkflowCommitStateOptions CommitStateOptions { get; set; } = new();
}

View file

@ -4,6 +4,11 @@ namespace Elsa.Workflows;
public class NoopCommitStateHandler : ICommitStateHandler
{
public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default)
{
return Task.CompletedTask;
}
public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default)
{
return Task.CompletedTask;

View file

@ -5,6 +5,12 @@ namespace Elsa.Workflows.Runtime;
public class StoreCommitStateHandler(IWorkflowInstanceManager workflowInstanceManager) : ICommitStateHandler
{
public async Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default)
{
var workflowState = workflowInstanceManager.ExtractWorkflowState(workflowExecutionContext);
await CommitAsync(workflowExecutionContext, workflowState, cancellationToken);
}
public async Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default)
{
await workflowInstanceManager.SaveAsync(workflowState, cancellationToken);