Add initial support for persisting workflow input

This commit is contained in:
Sipke Schoorstra 2023-08-31 21:04:19 +02:00
parent fe02011c8b
commit 4fe7518976
4 changed files with 46 additions and 11 deletions

View file

@ -166,7 +166,7 @@ public class WorkflowExecutionContext : IExecutionContext
/// <summary>
/// A dictionary of inputs provided at the start of the current workflow execution.
/// </summary>
public IDictionary<string, object> Input { get; }
public IDictionary<string, object> Input { get; set; }
/// <summary>
/// A dictionary of outputs provided by the current workflow execution.

View file

@ -22,6 +22,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
SubStatus = workflowExecutionContext.SubStatus,
Bookmarks = workflowExecutionContext.Bookmarks,
ExecutionLogSequence = workflowExecutionContext.ExecutionLogSequence,
Input = GetPersistableInput(workflowExecutionContext),
Output = workflowExecutionContext.Output,
Fault = MapFault(workflowExecutionContext.Fault),
CreatedAt = workflowExecutionContext.CreatedAt
@ -34,6 +35,22 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
return state;
}
private IDictionary<string, object> GetPersistableInput(WorkflowExecutionContext workflowExecutionContext)
{
// TODO: This is a temporary solution. We need to find a better way to handle this.
var persistableInput = workflowExecutionContext.Workflow.Inputs.Where(x => x.StorageDriverType == typeof(WorkflowStorageDriver)).ToList();
var input = workflowExecutionContext.Input;
var filteredInput = new Dictionary<string, object>();
foreach (var inputDefinition in persistableInput)
{
if (input.TryGetValue(inputDefinition.Name, out var value))
filteredInput.Add(inputDefinition.Name, value);
}
return filteredInput;
}
/// <inheritdoc />
public WorkflowExecutionContext Apply(WorkflowExecutionContext workflowExecutionContext, WorkflowState state)
{
@ -41,6 +58,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
workflowExecutionContext.CorrelationId = state.CorrelationId;
workflowExecutionContext.SubStatus = state.SubStatus;
workflowExecutionContext.Bookmarks = state.Bookmarks;
workflowExecutionContext.Input = state.Input;
workflowExecutionContext.Output = state.Output;
workflowExecutionContext.ExecutionLogSequence = state.ExecutionLogSequence;
workflowExecutionContext.CreatedAt = state.CreatedAt;
@ -66,11 +84,11 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
{
var ownerActivityExecutionContext = workflowExecutionContext.ActiveActivityExecutionContexts.First(x => x.Id == completionCallbackEntry.OwnerInstanceId);
var childNode = workflowExecutionContext.ActiveActivityExecutionContexts.FirstOrDefault(x => x.NodeId == completionCallbackEntry.ChildNodeId)?.ActivityNode;
// If the child node is null, it means the completion callback was registered for an activity instance that has already completed or was canceled.
if(childNode == null)
if (childNode == null)
continue;
var callbackName = completionCallbackEntry.MethodName;
var callbackDelegate = !string.IsNullOrEmpty(callbackName) ? ownerActivityExecutionContext.Activity.GetActivityCompletionCallback(callbackName) : default;
workflowExecutionContext.AddCompletionCallback(ownerActivityExecutionContext, childNode, callbackDelegate);
@ -87,7 +105,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
if (ownerContext == null)
throw new Exception("Lost an owner context");
}
var completionCallbacks = workflowExecutionContext.CompletionCallbacks.Select(x => new CompletionCallbackState(x.Owner.Id, x.Child.NodeId, x.CompletionCallback?.Method.Name));
state.CompletionCallbacks = completionCallbacks.ToList();
}
@ -105,7 +123,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
if (parentContext == null)
throw new Exception("We lost a context");
}
var activityExecutionContextState = new ActivityExecutionContextState
{
Id = activityExecutionContext.Id,
@ -119,7 +137,6 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
CompletedAt = activityExecutionContext.CompletedAt,
Tag = activityExecutionContext.Tag,
DynamicVariables = activityExecutionContext.DynamicVariables
};
return activityExecutionContextState;
}
@ -159,16 +176,16 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor
context.ExpressionExecutionContext.ParentContext = parentContext.ExpressionExecutionContext;
context.ParentActivityExecutionContext = parentContext;
}
// Assign root expression execution context.
var rootActivityExecutionContexts = activityExecutionContexts.Where(x => x.ExpressionExecutionContext.ParentContext == null);
foreach (var rootActivityExecutionContext in rootActivityExecutionContexts)
foreach (var rootActivityExecutionContext in rootActivityExecutionContexts)
rootActivityExecutionContext.ExpressionExecutionContext.ParentContext = workflowExecutionContext.ExpressionExecutionContext;
workflowExecutionContext.ActivityExecutionContexts = activityExecutionContexts;
}
private static WorkflowFaultState? MapFault(WorkflowFault? fault)
{
if (fault == null)

View file

@ -69,6 +69,11 @@ public class WorkflowState
/// </summary>
public long ExecutionLogSequence { get; set; }
/// <summary>
/// A dictionary of inputs sent to the workflow.
/// </summary>
public IDictionary<string, object> Input { get; set; } = new Dictionary<string, object>();
/// <summary>
/// A dictionary of outputs produced by the workflow.
/// </summary>

View file

@ -2,4 +2,17 @@ using Elsa.Common.Models;
namespace Elsa.Workflows.Runtime.Contracts;
public record StartWorkflowRuntimeOptions(string? CorrelationId = default, IDictionary<string, object>? Input = default, VersionOptions VersionOptions = default, string? TriggerActivityId = default, string? InstanceId = default);
/// <summary>
/// Represents options for starting a workflow.
/// </summary>
/// <param name="CorrelationId"></param>
/// <param name="Input"></param>
/// <param name="VersionOptions"></param>
/// <param name="TriggerActivityId"></param>
/// <param name="InstanceId"></param>
public record StartWorkflowRuntimeOptions(
string? CorrelationId = default,
IDictionary<string, object>? Input = default,
VersionOptions VersionOptions = default,
string? TriggerActivityId = default,
string? InstanceId = default);