Fix persistent variable serialization

This commit is contained in:
Sipke Schoorstra 2022-12-16 12:11:45 +01:00
parent 764cd6187b
commit 594970bfc6
7 changed files with 67 additions and 47 deletions

View file

@ -20,7 +20,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime
builder.Ignore(x => x.Properties);
builder.Ignore(x => x.ActivityOutput);
builder.Ignore(x => x.CompletionCallbacks);
builder.Ignore(x => x.PersistentVariables);
builder.Ignore(x => x.ActivityExecutionContexts);
builder.Property<string>("Data");
builder.Property<DateTimeOffset>("CreatedAt");

View file

@ -1,4 +1,3 @@
using System.Diagnostics;
using System.Reflection;
using Elsa.Workflows.Core.Attributes;
using Elsa.Workflows.Core.Models;
@ -9,13 +8,6 @@ namespace Elsa.Workflows.Core.Implementations;
public class WorkflowStateSerializer : IWorkflowStateSerializer
{
private readonly IServiceProvider _serviceProvider;
public WorkflowStateSerializer(IServiceProvider serviceProvider)
{
_serviceProvider = serviceProvider;
}
public WorkflowState SerializeState(WorkflowExecutionContext workflowExecutionContext)
{
var state = new WorkflowState
@ -32,7 +24,6 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
SerializeProperties(state, workflowExecutionContext);
SerializeCompletionCallbacks(state, workflowExecutionContext);
SerializeActivityExecutionContexts(state, workflowExecutionContext);
SerializePersistentVariables(state, workflowExecutionContext);
return state;
}
@ -46,17 +37,16 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
DeserializeProperties(state, workflowExecutionContext);
DeserializeActivityExecutionContexts(state, workflowExecutionContext);
DeserializeCompletionCallbacks(state, workflowExecutionContext);
//DeserializePersistentVariables(state, workflowExecutionContext);
}
private void SerializeProperties(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
{
state.Properties = workflowExecutionContext.Properties;
state.Properties = new PropertyBag(workflowExecutionContext.Properties);
}
private void DeserializeProperties(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
{
workflowExecutionContext.Properties = state.Properties;
workflowExecutionContext.Properties = state.Properties.Properties;
}
private void GetOutput(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
@ -181,16 +171,6 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
workflowExecutionContext.ActivityExecutionContexts = activityExecutionContexts;
}
private void SerializePersistentVariables(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
{
var workflow = workflowExecutionContext.Workflow;
state.PersistentVariables = workflow.Variables
.Where(x => x.StorageDriverId != null)
.Select(x => new PersistentVariableState(x.Name, x.StorageDriverId!))
.ToList();
}
private Dictionary<string, object> GetOutputFrom(ActivityNode activityNode) =>
activityNode.GetType().GetProperties(BindingFlags.Public).Where(x => x.GetCustomAttribute<OutputAttribute>() != null).ToDictionary(x => x.Name, x => x.GetValue(activityNode)!);
}

View file

@ -0,0 +1,20 @@
using System.Text.Json.Serialization;
using Elsa.Workflows.Core.Serialization.Converters;
namespace Elsa.Workflows.Core.Models;
[JsonConverter(typeof(PropertyBagConverter))]
public class PropertyBag
{
[JsonConstructor]
public PropertyBag() : this(new Dictionary<string, object>())
{
}
public PropertyBag(IDictionary<string, object> properties)
{
Properties = properties;
}
public IDictionary<string, object> Properties { get; init; }
}

View file

@ -0,0 +1,19 @@
using System.Text.Json;
using System.Text.Json.Serialization;
using Elsa.Workflows.Core.Models;
namespace Elsa.Workflows.Core.Serialization.Converters;
public class PropertyBagConverter : JsonConverter<PropertyBag>
{
public override PropertyBag? Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
{
var dictionary = JsonSerializer.Deserialize<Dictionary<string, object>>(ref reader);
return new PropertyBag(dictionary);
}
public override void Write(Utf8JsonWriter writer, PropertyBag value, JsonSerializerOptions options)
{
JsonSerializer.Serialize(writer, value.Properties);
}
}

View file

@ -1,4 +1,6 @@
using System.Text.Json.Serialization;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Serialization.Converters;
namespace Elsa.Workflows.Core.State;
@ -56,14 +58,10 @@ public class WorkflowState
/// A flattened list of <see cref="ActivityExecutionContextState"/> objects, representing the various active "call stacks" of the workflow.
/// </summary>
public ICollection<ActivityExecutionContextState> ActivityExecutionContexts { get; set; } = new List<ActivityExecutionContextState>();
/// <summary>
/// A global property bag that contains properties set by application code and/or activities.
/// </summary>
public IDictionary<string, object> Properties { get; set; } = new Dictionary<string, object>();
/// <summary>
/// A list of variables that can be persisted.
/// </summary>
public ICollection<PersistentVariableState> PersistentVariables { get; set; } = new List<PersistentVariableState>();
[JsonConverter(typeof(PropertyBagConverter))]
public PropertyBag Properties { get; set; } = new();
}

View file

@ -24,8 +24,4 @@
<PackageReference Include="Open.Linq.AsyncExtensions" Version="1.2.0" />
</ItemGroup>
<ItemGroup>
<Folder Include="Consumers" />
</ItemGroup>
</Project>

View file

@ -1,3 +1,5 @@
using Elsa.Expressions.Helpers;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Pipelines.WorkflowExecution;
using Elsa.Workflows.Core.Services;
@ -24,35 +26,41 @@ public class PersistentVariablesMiddleware : WorkflowExecutionMiddleware
// Load persistent variables.
var dataDriveContext = new DataDriveContext(context, cancellationToken);
var persistentVariables = context.Workflow.Variables
.Where(x => x.StorageDriverId != null)
.Select(x => new PersistentVariableState(x.Name, x.StorageDriverId!))
.ToList();
var persistentVariables = context.Workflow.Variables.Where(x => x.StorageDriverId != null).ToList();
foreach (var variableState in persistentVariables)
foreach (var variable in persistentVariables)
{
var drive = _storageDriverManager.GetDriveById(variableState.StorageDriverId);
var drive = _storageDriverManager.GetDriveById(variable.StorageDriverId!);
if (drive == null) continue;
var id = $"{context.Id}:{variableState.Name}";
var id = $"{context.Id}:{variable.Name}";
var value = await drive.ReadAsync(id, dataDriveContext);
if (value == null) continue;
var variable = new Variable(variableState.Name, value);
var parsedValue = ParseVariableValue(variable, value);
context.MemoryRegister.Declare(variable);
variable.Set(context.MemoryRegister, parsedValue);
}
// Invoke next middleware.
await Next(context);
// Persist variables.
foreach (var variableState in persistentVariables)
foreach (var variable in persistentVariables)
{
var drive = _storageDriverManager.GetDriveById(variableState.StorageDriverId);
var drive = _storageDriverManager.GetDriveById(variable.StorageDriverId!);
if (drive == null) continue;
if (!context.MemoryRegister.TryGetBlock(variableState.Name, out var block)) continue;
if (!context.MemoryRegister.TryGetBlock(variable.Name, out var block)) continue;
if (block.Value == null) continue;
var id = $"{context.Id}:{variableState.Name}";
var id = $"{context.Id}:{variable.Name}";
await drive.WriteAsync(id, block.Value, dataDriveContext);
}
}
private object ParseVariableValue(Variable variable, object value)
{
if (!variable.GetType().GenericTypeArguments.Any())
return value;
var type = variable.GetType().GenericTypeArguments.First();
return value.ConvertTo(type)!;
}
}