Add WorkflowContexts module

This commit is contained in:
Sipke Schoorstra 2022-03-08 14:40:47 +01:00
parent fc5f14a18e
commit c33549c291
28 changed files with 342 additions and 133 deletions

View file

@ -90,6 +90,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.JavaScript", "
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.Hangfire", "src\modules\Elsa.Modules.Hangfire\Elsa.Modules.Hangfire.csproj", "{0601A2A6-2C62-418B-9104-8CDE497E5283}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Modules.WorkflowContexts", "src\modules\Elsa.Modules.WorkflowContexts\Elsa.Modules.WorkflowContexts.csproj", "{302BFC43-ED2F-43AE-8AD4-FCD481B0AC67}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -200,6 +202,10 @@ Global
{0601A2A6-2C62-418B-9104-8CDE497E5283}.Debug|Any CPU.Build.0 = Debug|Any CPU
{0601A2A6-2C62-418B-9104-8CDE497E5283}.Release|Any CPU.ActiveCfg = Release|Any CPU
{0601A2A6-2C62-418B-9104-8CDE497E5283}.Release|Any CPU.Build.0 = Release|Any CPU
{302BFC43-ED2F-43AE-8AD4-FCD481B0AC67}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{302BFC43-ED2F-43AE-8AD4-FCD481B0AC67}.Debug|Any CPU.Build.0 = Debug|Any CPU
{302BFC43-ED2F-43AE-8AD4-FCD481B0AC67}.Release|Any CPU.ActiveCfg = Release|Any CPU
{302BFC43-ED2F-43AE-8AD4-FCD481B0AC67}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(NestedProjects) = preSolution
{155227F0-A33B-40AA-A4B4-06F813EB921B} = {61017E64-6D00-49CB-9E81-5002DC8F7D5F}
@ -240,5 +246,6 @@ Global
{49716A83-239C-4913-BC11-E379ED2F676E} = {C6658DE0-2B2F-47F0-BB61-2CA66D435C09}
{D31581AB-A6C1-4B73-AB63-45667F6C82AE} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79}
{0601A2A6-2C62-418B-9104-8CDE497E5283} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79}
{302BFC43-ED2F-43AE-8AD4-FCD481B0AC67} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79}
EndGlobalSection
EndGlobal

View file

@ -11,6 +11,7 @@ public class WorkflowDefinitionBuilder : IWorkflowDefinitionBuilder
public int Version { get; private set; } = 1;
public IActivity? Root { get; private set; }
public ICollection<Variable> Variables { get; set; } = new List<Variable>();
public IDictionary<string, object> ApplicationProperties { get; set; } = new Dictionary<string, object>();
public IWorkflowDefinitionBuilder WithId(string id)
{
@ -36,6 +37,18 @@ public class WorkflowDefinitionBuilder : IWorkflowDefinitionBuilder
return this;
}
public IWorkflowDefinitionBuilder WithVariable(Variable variable)
{
Variables.Add(variable);
return this;
}
public IWorkflowDefinitionBuilder WithApplicationProperty(string name, object value)
{
ApplicationProperties[name] = value;
return this;
}
public Workflow BuildWorkflow()
{
var definitionId = DefinitionId ?? Guid.NewGuid().ToString("N");
@ -44,6 +57,6 @@ public class WorkflowDefinitionBuilder : IWorkflowDefinitionBuilder
var identity = new WorkflowIdentity(definitionId, Version, id);
var publication = WorkflowPublication.LatestAndPublished;
var metadata = new WorkflowMetadata();
return new Workflow(identity, publication, metadata, root, Variables);
return new Workflow(identity, publication, metadata, root, Variables, ApplicationProperties);
}
}

View file

@ -8,8 +8,11 @@ public interface IWorkflowDefinitionBuilder
int Version { get; }
IActivity? Root { get; }
ICollection<Variable> Variables { get; set; }
IDictionary<string, object> ApplicationProperties { get; }
IWorkflowDefinitionBuilder WithDefinitionId(string definitionId);
IWorkflowDefinitionBuilder WithVersion(int version);
IWorkflowDefinitionBuilder WithRoot(IActivity root);
IWorkflowDefinitionBuilder WithVariable(Variable variable);
IWorkflowDefinitionBuilder WithApplicationProperty(string name, object value);
Workflow BuildWorkflow();
}

View file

@ -4,9 +4,9 @@ namespace Elsa.Extensions;
public static class DictionaryExtensions
{
public static bool TryGetValue<T>(this IReadOnlyDictionary<string, object> dictionary, string key, out T value) => TryGetValue((IDictionary<string, object>)dictionary, key, out value);
public static bool TryGetValue<T>(this IReadOnlyDictionary<string, object> dictionary, string key, out T? value) => TryGetValue((IDictionary<string, object?>)dictionary, key, out value);
public static bool TryGetValue<T>(this IDictionary<string, object> dictionary, string key, out T value)
public static bool TryGetValue<T>(this IDictionary<string, object?> dictionary, string key, out T? value)
{
if (!dictionary.TryGetValue(key, out var item))
{
@ -18,15 +18,26 @@ public static class DictionaryExtensions
return true;
}
public static T GetValue<T>(this IDictionary<string, object> dictionary, string key)
public static T? GetValue<T>(this IDictionary<string, object?> dictionary, string key)
{
return ConvertValue<T>(dictionary[key]);
}
private static T ConvertValue<T>(object value)
public static T GetOrAdd<T>(this IDictionary<string, object?> dictionary, string key, Func<T> valueFactory) where T : notnull
{
if (dictionary.TryGetValue<T>(key, out var value))
return value;
value = valueFactory();
dictionary.Add(key, value);
return value;
}
private static T? ConvertValue<T>(object? value)
{
return value switch
{
null => default,
T v => v,
JsonElement jsonElement => jsonElement.Deserialize<T>()!,
_ => throw new InvalidOperationException()

View file

@ -6,15 +6,19 @@ public class ExpressionExecutionContext
{
private readonly IServiceProvider _serviceProvider;
public ExpressionExecutionContext(IServiceProvider serviceProvider, Register register, ExpressionExecutionContext? parentContext, CancellationToken cancellationToken)
public ExpressionExecutionContext(IServiceProvider serviceProvider, Register register, Workflow workflow, IDictionary<string, object?> transientProperties, ExpressionExecutionContext? parentContext, CancellationToken cancellationToken)
{
_serviceProvider = serviceProvider;
Register = register;
Workflow = workflow;
TransientProperties = transientProperties;
ParentContext = parentContext;
CancellationToken = cancellationToken;
}
public Register Register { get; }
public Workflow Workflow { get; }
public IDictionary<string, object?> TransientProperties { get; }
public ExpressionExecutionContext? ParentContext { get; set; }
public CancellationToken CancellationToken { get; }

View file

@ -7,9 +7,10 @@ public record Workflow(
WorkflowPublication Publication,
WorkflowMetadata Metadata,
IActivity Root,
ICollection<Variable> Variables)
ICollection<Variable> Variables,
IDictionary<string, object> ApplicationProperties)
{
public static Workflow FromActivity(IActivity root) => new(WorkflowIdentity.VersionOne, WorkflowPublication.LatestDraft, new WorkflowMetadata(), root, new List<Variable>());
public static Workflow FromActivity(IActivity root) => new(WorkflowIdentity.VersionOne, WorkflowPublication.LatestDraft, new WorkflowMetadata(), root, new List<Variable>(), new Dictionary<string, object>());
public Workflow WithVersion(int version) => this with { Identity = Identity with { Version = version } };
public Workflow IncrementVersion() => WithVersion(Identity.Version + 1);

View file

@ -48,7 +48,18 @@ public class WorkflowExecutionContext
public IActivityScheduler Scheduler { get; }
public Bookmark? Bookmark { get; }
public IReadOnlyDictionary<string, object> Input { get; }
/// <summary>
/// A dictionary that can be used by application code and activities to store information. Values need to be serializable, since this dictionary will be persisted alongside the workflow instance.
/// </summary>
public IDictionary<string, object?> Properties { get; set; } = new Dictionary<string, object?>();
/// <summary>
/// A dictionary that can be used by application code and middleware to store information and even services. Values do not need to be serializable, since this dictionary will not be persisted.
/// All data will be gone once workflow execution completes.
/// </summary>
public IDictionary<string, object?> TransientProperties { get; set; } = new Dictionary<string, object?>();
public ExecuteActivityDelegate? ExecuteDelegate { get; set; }
public CancellationToken CancellationToken { get; }
public IReadOnlyCollection<Bookmark> Bookmarks => new ReadOnlyCollection<Bookmark>(_bookmarks);

View file

@ -27,8 +27,10 @@ public class ActivityInvoker : IActivityInvoker
// Setup an activity execution context.
var register = workflowExecutionContext.Workflow.CreateRegister();
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, parentActivityExecutionContext?.ExpressionExecutionContext, cancellationToken);
var workflow = workflowExecutionContext.Workflow;
var transientProperties = workflowExecutionContext.TransientProperties;
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, workflow, transientProperties, parentActivityExecutionContext?.ExpressionExecutionContext, cancellationToken);
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, parentActivityExecutionContext, expressionExecutionContext, activity, cancellationToken);
// Declare locations.

View file

@ -15,7 +15,7 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
{
_serviceProvider = serviceProvider;
}
public WorkflowState ReadState(WorkflowExecutionContext workflowExecutionContext)
{
var state = new WorkflowState
@ -44,7 +44,7 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
{
state.Properties = workflowExecutionContext.Properties;
}
private void SetProperties(WorkflowState state, WorkflowExecutionContext workflowExecutionContext)
{
workflowExecutionContext.Properties = state.Properties;
@ -127,7 +127,9 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
var cancellationToken = workflowExecutionContext.CancellationToken;
var activity = workflowExecutionContext.FindActivityById(activityExecutionContextState.ScheduledActivityId);
var register = new Register(activityExecutionContextState.Register.Locations);
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, default, cancellationToken);
var workflow = workflowExecutionContext.Workflow;
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, workflow, new Dictionary<string, object?>(), default, cancellationToken);
var properties = activityExecutionContextState.Properties;
var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, default, expressionExecutionContext, activity, cancellationToken)
{
@ -136,7 +138,7 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
};
return activityExecutionContext;
}
var activityExecutionContexts = state.ActivityExecutionContexts.Select(CreateActivityExecutionContext).ToList();
var lookup = activityExecutionContexts.ToDictionary(x => x.Id);
@ -146,10 +148,10 @@ public class WorkflowStateSerializer : IWorkflowStateSerializer
var parentContext = lookup[contextState.ParentContextId!];
var contextId = contextState.Id;
var context = lookup[contextId];
context.ExpressionExecutionContext.ParentContext = parentContext.ExpressionExecutionContext;
context.ExpressionExecutionContext.ParentContext = parentContext.ExpressionExecutionContext;
context.ParentActivityExecutionContext = parentContext;
}
workflowExecutionContext.ActivityExecutionContexts = activityExecutionContexts;
}

View file

@ -34,7 +34,8 @@ namespace Elsa.Management.Services
WorkflowPublication.LatestDraft,
new WorkflowMetadata(CreatedAt: _systemClock.UtcNow),
new Sequence(),
new List<Variable>());
new List<Variable>(),
new Dictionary<string, object>());
}
public async Task<Workflow?> PublishAsync(string definitionId, CancellationToken cancellationToken = default)

View file

@ -0,0 +1,24 @@
using Elsa.Models;
using Elsa.Modules.WorkflowContexts.Contracts;
namespace Elsa.Modules.WorkflowContexts.Abstractions;
public abstract class WorkflowContextProvider<T> : IWorkflowContextProvider
{
async ValueTask<object?> IWorkflowContextProvider.LoadAsync(WorkflowExecutionContext workflowExecutionContext) => await LoadAsync(workflowExecutionContext);
async ValueTask IWorkflowContextProvider.SaveAsync(WorkflowExecutionContext workflowExecutionContext, object? context) => await SaveAsync(workflowExecutionContext, (T?)context);
protected virtual ValueTask<T?> LoadAsync(WorkflowExecutionContext workflowExecutionContext) => new(Load(workflowExecutionContext));
protected virtual T? Load(WorkflowExecutionContext workflowExecutionContext) => default;
protected virtual ValueTask SaveAsync(WorkflowExecutionContext workflowExecutionContext, T? context)
{
Save(workflowExecutionContext, context);
return ValueTask.CompletedTask;
}
protected virtual void Save(WorkflowExecutionContext workflowExecutionContext, T? context)
{
}
}

View file

@ -0,0 +1,21 @@
using System.Threading.Tasks.Sources;
using Elsa.Models;
namespace Elsa.Modules.WorkflowContexts.Contracts;
/// <summary>
/// Implement this interface to implement a workflow context provider that loads application-specific objects into the workflow.
/// These providers can then be configured on a given workflow.
/// </summary>
public interface IWorkflowContextProvider
{
/// <summary>
/// Implement this method to load an object into memory that is accessible throughout the lifetime of the workflow's current execution.
/// </summary>
ValueTask<object?> LoadAsync(WorkflowExecutionContext workflowExecutionContext);
/// <summary>
/// Implement this method to save an object that was loaded previously
/// </summary>
ValueTask SaveAsync(WorkflowExecutionContext workflowExecutionContext, object? context);
}

View file

@ -0,0 +1,13 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net6.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,13 @@
using Elsa.Contracts;
using Elsa.Modules.WorkflowContexts.Middleware;
using Elsa.Pipelines.WorkflowExecution;
namespace Elsa.Modules.WorkflowContexts.Extensions;
public static class WorkflowExecutionBuilderExtensions
{
/// <summary>
/// Installs the <see cref="WorkflowContextMiddleware"/>.
/// </summary>
public static IWorkflowExecutionBuilder UseWorkflowContexts(this IWorkflowExecutionBuilder builder) => builder.UseMiddleware<WorkflowContextMiddleware>();
}

View file

@ -0,0 +1,22 @@
using Elsa.Contracts;
using Elsa.Extensions;
using Elsa.Modules.WorkflowContexts.Contracts;
using Elsa.Modules.WorkflowContexts.Models;
namespace Elsa.Modules.WorkflowContexts.Extensions;
public static class WorkflowExtensions
{
/// <summary>
/// Installs the specified workflow context provider type into the specified workflow.
/// </summary>
public static IWorkflowDefinitionBuilder AddWorkflowContext<T, TProvider>(this IWorkflowDefinitionBuilder workflow, WorkflowContext<T, TProvider> workflowContext) where TProvider : IWorkflowContextProvider
{
var providerTypes = workflow.ApplicationProperties!.GetOrAdd("Elsa:WorkflowContexts", () => new List<WorkflowContext>());
providerTypes.Add(workflowContext);
return workflow;
}
}

View file

@ -0,0 +1,57 @@
using Elsa.Extensions;
using Elsa.Models;
using Elsa.Modules.WorkflowContexts.Contracts;
using Elsa.Modules.WorkflowContexts.Models;
using Elsa.Pipelines.WorkflowExecution;
using Elsa.Pipelines.WorkflowExecution.Components;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Modules.WorkflowContexts.Middleware;
/// <summary>
/// Middleware that loads & save workflow context into the currently executing workflow using installed workflow context providers.
/// </summary>
public class WorkflowContextMiddleware : WorkflowExecutionMiddleware
{
private readonly IServiceProvider _serviceProvider;
public WorkflowContextMiddleware(WorkflowMiddlewareDelegate next, IServiceProvider serviceProvider) : base(next)
{
_serviceProvider = serviceProvider;
}
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
// Check if the workflow contains any workflow contexts.
if (!context.Workflow.ApplicationProperties!.TryGetValue<ICollection<WorkflowContext>>("Elsa:WorkflowContexts", out var workflowContexts))
{
await Next(context);
return;
}
// For each workflow context, invoke its provider.
foreach (var workflowContext in workflowContexts!)
{
var provider = (IWorkflowContextProvider)ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider, workflowContext.ProviderType);
var value = await provider.LoadAsync(context);
// Store the loaded value into the workflow execution context.
var contextDictionary = context.TransientProperties.GetOrAdd("WorkflowContexts", () => new Dictionary<WorkflowContext, object?>());
contextDictionary.Add(workflowContext, value);
}
// Invoke the next middleware.
await Next(context);
// For each workflow context, invoke its provider to update the context.
foreach (var workflowContext in workflowContexts!)
{
// Get the loaded value from the workflow execution context.
var contextDictionary = context.TransientProperties.GetOrAdd("WorkflowContexts", () => new Dictionary<WorkflowContext, object?>());
var value = contextDictionary[workflowContext];
var provider = (IWorkflowContextProvider)ActivatorUtilities.GetServiceOrCreateInstance(_serviceProvider, workflowContext.ProviderType);
await provider.SaveAsync(context, value);
}
}
}

View file

@ -0,0 +1,28 @@
using Elsa.Models;
using Elsa.Modules.WorkflowContexts.Contracts;
namespace Elsa.Modules.WorkflowContexts.Models;
public class WorkflowContext
{
public WorkflowContext(Type providerType)
{
ProviderType = providerType;
}
public Type ProviderType { get; }
}
public class WorkflowContext<T, TProvider> : WorkflowContext where TProvider:IWorkflowContextProvider
{
public WorkflowContext() : base(typeof(TProvider))
{
}
public T? Get(ExpressionExecutionContext context)
{
var workflowContexts = (IDictionary<WorkflowContext, object?>)context.TransientProperties["WorkflowContexts"]!;
return (T?)workflowContexts[this];
}
}

View file

@ -15,6 +15,7 @@ public class WorkflowDefinition : Entity
public int Version { get; set; } = 1;
public IActivity Root { get; set; } = default!;
public ICollection<Variable> Variables { get; set; } = new List<Variable>();
public IDictionary<string, object> ApplicationProperties { get; set; } = new Dictionary<string, object>();
public bool IsLatest { get; set; }
public bool IsPublished { get; set; }

View file

@ -15,7 +15,8 @@ public class WorkflowDefinitionMapper
new WorkflowPublication(definition.IsLatest, definition.IsPublished),
new WorkflowMetadata(definition.Name, definition.Description, definition.CreatedAt),
definition.Root,
definition.Variables);
definition.Variables,
definition.ApplicationProperties);
}
public WorkflowDefinition? Map(Workflow? workflow) => workflow == null ? null : Map(workflow, new WorkflowDefinition());
@ -35,6 +36,8 @@ public class WorkflowDefinitionMapper
definition.CreatedAt = metadata.CreatedAt;
definition.IsLatest = publication.IsLatest;
definition.IsPublished = publication.IsPublished;
definition.Variables = workflow.Variables;
definition.ApplicationProperties = workflow.ApplicationProperties;
return definition;
}

View file

@ -11,7 +11,7 @@ using Microsoft.EntityFrameworkCore.Storage.ValueConversion;
namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
{
[DbContext(typeof(ElsaDbContext))]
[Migration("20220307110612_Initial")]
[Migration("20220308133708_Initial")]
partial class Initial
{
protected override void BuildTargetModel(ModelBuilder modelBuilder)
@ -19,24 +19,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
#pragma warning disable 612, 618
modelBuilder.HasAnnotation("ProductVersion", "6.0.1");
modelBuilder.Entity("Elsa.Models.Variable", b =>
{
b.Property<string>("Id")
.HasColumnType("TEXT");
b.Property<string>("Name")
.HasColumnType("TEXT");
b.Property<string>("WorkflowDefinitionId")
.HasColumnType("TEXT");
b.HasKey("Id");
b.HasIndex("WorkflowDefinitionId");
b.ToTable("Variable");
});
modelBuilder.Entity("Elsa.Persistence.Entities.WorkflowBookmark", b =>
{
b.Property<string>("Id")
@ -103,7 +85,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.Property<string>("Id")
.HasColumnType("TEXT");
b.Property<DateTime>("CreatedAt")
b.Property<DateTimeOffset>("CreatedAt")
.HasColumnType("TEXT");
b.Property<string>("Data")
@ -174,7 +156,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.Property<string>("Source")
.HasColumnType("TEXT");
b.Property<DateTime>("Timestamp")
b.Property<DateTimeOffset>("Timestamp")
.HasColumnType("TEXT");
b.Property<string>("WorkflowInstanceId")
@ -206,14 +188,14 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.Property<string>("Id")
.HasColumnType("TEXT");
b.Property<DateTime?>("CancelledAt")
b.Property<DateTimeOffset?>("CancelledAt")
.HasColumnType("TEXT");
b.Property<string>("CorrelationId")
.IsRequired()
.HasColumnType("TEXT");
b.Property<DateTime>("CreatedAt")
b.Property<DateTimeOffset>("CreatedAt")
.HasColumnType("TEXT");
b.Property<string>("Data")
@ -227,13 +209,13 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
.IsRequired()
.HasColumnType("TEXT");
b.Property<DateTime?>("FaultedAt")
b.Property<DateTimeOffset?>("FaultedAt")
.HasColumnType("TEXT");
b.Property<DateTime?>("FinishedAt")
b.Property<DateTimeOffset?>("FinishedAt")
.HasColumnType("TEXT");
b.Property<DateTime?>("LastExecutedAt")
b.Property<DateTimeOffset?>("LastExecutedAt")
.HasColumnType("TEXT");
b.Property<string>("Name")
@ -312,18 +294,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.ToTable("WorkflowTriggers");
});
modelBuilder.Entity("Elsa.Models.Variable", b =>
{
b.HasOne("Elsa.Persistence.Entities.WorkflowDefinition", null)
.WithMany("Variables")
.HasForeignKey("WorkflowDefinitionId");
});
modelBuilder.Entity("Elsa.Persistence.Entities.WorkflowDefinition", b =>
{
b.Navigation("Variables");
});
#pragma warning restore 612, 618
}
}

View file

@ -37,7 +37,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
DefinitionId = table.Column<string>(type: "TEXT", nullable: false),
Name = table.Column<string>(type: "TEXT", nullable: true),
Description = table.Column<string>(type: "TEXT", nullable: true),
CreatedAt = table.Column<DateTime>(type: "TEXT", nullable: false),
CreatedAt = table.Column<DateTimeOffset>(type: "TEXT", nullable: false),
Version = table.Column<int>(type: "INTEGER", nullable: false),
IsLatest = table.Column<bool>(type: "INTEGER", nullable: false),
IsPublished = table.Column<bool>(type: "INTEGER", nullable: false),
@ -56,7 +56,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
WorkflowInstanceId = table.Column<string>(type: "TEXT", nullable: false),
ActivityId = table.Column<string>(type: "TEXT", nullable: false),
ActivityType = table.Column<string>(type: "TEXT", nullable: false),
Timestamp = table.Column<DateTime>(type: "TEXT", nullable: false),
Timestamp = table.Column<DateTimeOffset>(type: "TEXT", nullable: false),
EventName = table.Column<string>(type: "TEXT", nullable: true),
Message = table.Column<string>(type: "TEXT", nullable: true),
Source = table.Column<string>(type: "TEXT", nullable: true),
@ -78,11 +78,11 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
WorkflowStatus = table.Column<int>(type: "INTEGER", nullable: false),
CorrelationId = table.Column<string>(type: "TEXT", nullable: false),
Name = table.Column<string>(type: "TEXT", nullable: true),
CreatedAt = table.Column<DateTime>(type: "TEXT", nullable: false),
LastExecutedAt = table.Column<DateTime>(type: "TEXT", nullable: true),
FinishedAt = table.Column<DateTime>(type: "TEXT", nullable: true),
CancelledAt = table.Column<DateTime>(type: "TEXT", nullable: true),
FaultedAt = table.Column<DateTime>(type: "TEXT", nullable: true),
CreatedAt = table.Column<DateTimeOffset>(type: "TEXT", nullable: false),
LastExecutedAt = table.Column<DateTimeOffset>(type: "TEXT", nullable: true),
FinishedAt = table.Column<DateTimeOffset>(type: "TEXT", nullable: true),
CancelledAt = table.Column<DateTimeOffset>(type: "TEXT", nullable: true),
FaultedAt = table.Column<DateTimeOffset>(type: "TEXT", nullable: true),
Data = table.Column<string>(type: "TEXT", nullable: true)
},
constraints: table =>
@ -105,29 +105,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
table.PrimaryKey("PK_WorkflowTriggers", x => x.Id);
});
migrationBuilder.CreateTable(
name: "Variable",
columns: table => new
{
Id = table.Column<string>(type: "TEXT", nullable: false),
Name = table.Column<string>(type: "TEXT", nullable: true),
WorkflowDefinitionId = table.Column<string>(type: "TEXT", nullable: true)
},
constraints: table =>
{
table.PrimaryKey("PK_Variable", x => x.Id);
table.ForeignKey(
name: "FK_Variable_WorkflowDefinitions_WorkflowDefinitionId",
column: x => x.WorkflowDefinitionId,
principalTable: "WorkflowDefinitions",
principalColumn: "Id");
});
migrationBuilder.CreateIndex(
name: "IX_Variable_WorkflowDefinitionId",
table: "Variable",
column: "WorkflowDefinitionId");
migrationBuilder.CreateIndex(
name: "IX_WorkflowBookmark_ActivityId",
table: "WorkflowBookmarks",
@ -278,10 +255,10 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropTable(
name: "Variable");
name: "WorkflowBookmarks");
migrationBuilder.DropTable(
name: "WorkflowBookmarks");
name: "WorkflowDefinitions");
migrationBuilder.DropTable(
name: "WorkflowExecutionLogRecords");
@ -291,9 +268,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
migrationBuilder.DropTable(
name: "WorkflowTriggers");
migrationBuilder.DropTable(
name: "WorkflowDefinitions");
}
}
}

View file

@ -17,24 +17,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
#pragma warning disable 612, 618
modelBuilder.HasAnnotation("ProductVersion", "6.0.1");
modelBuilder.Entity("Elsa.Models.Variable", b =>
{
b.Property<string>("Id")
.HasColumnType("TEXT");
b.Property<string>("Name")
.HasColumnType("TEXT");
b.Property<string>("WorkflowDefinitionId")
.HasColumnType("TEXT");
b.HasKey("Id");
b.HasIndex("WorkflowDefinitionId");
b.ToTable("Variable");
});
modelBuilder.Entity("Elsa.Persistence.Entities.WorkflowBookmark", b =>
{
b.Property<string>("Id")
@ -101,7 +83,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.Property<string>("Id")
.HasColumnType("TEXT");
b.Property<DateTime>("CreatedAt")
b.Property<DateTimeOffset>("CreatedAt")
.HasColumnType("TEXT");
b.Property<string>("Data")
@ -172,7 +154,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.Property<string>("Source")
.HasColumnType("TEXT");
b.Property<DateTime>("Timestamp")
b.Property<DateTimeOffset>("Timestamp")
.HasColumnType("TEXT");
b.Property<string>("WorkflowInstanceId")
@ -204,14 +186,14 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.Property<string>("Id")
.HasColumnType("TEXT");
b.Property<DateTime?>("CancelledAt")
b.Property<DateTimeOffset?>("CancelledAt")
.HasColumnType("TEXT");
b.Property<string>("CorrelationId")
.IsRequired()
.HasColumnType("TEXT");
b.Property<DateTime>("CreatedAt")
b.Property<DateTimeOffset>("CreatedAt")
.HasColumnType("TEXT");
b.Property<string>("Data")
@ -225,13 +207,13 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
.IsRequired()
.HasColumnType("TEXT");
b.Property<DateTime?>("FaultedAt")
b.Property<DateTimeOffset?>("FaultedAt")
.HasColumnType("TEXT");
b.Property<DateTime?>("FinishedAt")
b.Property<DateTimeOffset?>("FinishedAt")
.HasColumnType("TEXT");
b.Property<DateTime?>("LastExecutedAt")
b.Property<DateTimeOffset?>("LastExecutedAt")
.HasColumnType("TEXT");
b.Property<string>("Name")
@ -310,18 +292,6 @@ namespace Elsa.Persistence.EntityFrameworkCore.Sqlite.Migrations
b.ToTable("WorkflowTriggers");
});
modelBuilder.Entity("Elsa.Models.Variable", b =>
{
b.HasOne("Elsa.Persistence.Entities.WorkflowDefinition", null)
.WithMany("Variables")
.HasForeignKey("WorkflowDefinitionId");
});
modelBuilder.Entity("Elsa.Persistence.Entities.WorkflowDefinition", b =>
{
b.Navigation("Variables");
});
#pragma warning restore 612, 618
}
}

View file

@ -9,6 +9,8 @@ namespace Elsa.Persistence.EntityFrameworkCore.Configuration
public void Configure(EntityTypeBuilder<WorkflowDefinition> builder)
{
builder.Ignore(x => x.Root);
builder.Ignore(x => x.Variables);
builder.Ignore(x => x.ApplicationProperties);
builder.Property<string>("Data");
builder.HasIndex(x => new {x.DefinitionId, x.Version}).HasDatabaseName($"IX_{nameof(WorkflowDefinition)}_{nameof(WorkflowDefinition.DefinitionId)}_{nameof(WorkflowDefinition.Version)}").IsUnique();

View file

@ -1,6 +1,7 @@
using System.Text.Json;
using Elsa.Contracts;
using Elsa.Management.Serialization;
using Elsa.Models;
using Elsa.Persistence.Entities;
using Elsa.Persistence.EntityFrameworkCore.Contracts;
@ -20,6 +21,8 @@ public class WorkflowDefinitionSerializer : IEntitySerializer<WorkflowDefinition
var data = new
{
entity.Root,
entity.Variables,
entity.ApplicationProperties
};
var options = _workflowSerializerOptionsProvider.CreatePersistenceOptions();
@ -30,7 +33,7 @@ public class WorkflowDefinitionSerializer : IEntitySerializer<WorkflowDefinition
public void Deserialize(ElsaDbContext dbContext, WorkflowDefinition entity)
{
var data = new WorkflowDefinitionState(entity.Root);
var data = new WorkflowDefinitionState(entity.Root, entity.Variables, entity.ApplicationProperties);
var json = (string?) dbContext.Entry(entity).Property("Data").CurrentValue;
if (!string.IsNullOrWhiteSpace(json))
@ -40,6 +43,8 @@ public class WorkflowDefinitionSerializer : IEntitySerializer<WorkflowDefinition
}
entity.Root = data.Root;
entity.Variables = data.Variables;
entity.ApplicationProperties = data.ApplicationProperties;
}
// Can't use records when using System.Text.Json serialization and reference handling. Hence, using a class with default constructor.
@ -49,11 +54,15 @@ public class WorkflowDefinitionSerializer : IEntitySerializer<WorkflowDefinition
{
}
public WorkflowDefinitionState(IActivity root)
public WorkflowDefinitionState(IActivity root, ICollection<Variable> variables, IDictionary<string, object> applicationProperties)
{
Root = root;
Variables = variables;
ApplicationProperties = applicationProperties;
}
public IActivity Root { get; init; } = default!;
public ICollection<Variable> Variables { get; set; } = new List<Variable>();
public IDictionary<string, object> ApplicationProperties { get; set; } = new Dictionary<string, object>();
}
}

View file

@ -114,7 +114,7 @@ public class TriggerIndexer : ITriggerIndexer
private async IAsyncEnumerable<WorkflowTrigger> GetTriggersAsync(Workflow workflow, [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
var context = new WorkflowIndexingContext(workflow, cancellationToken);
// Get a list of activities that are configured as "startable".
var startableNodes = _activityWalker
.Walk(workflow.Root)
@ -173,7 +173,7 @@ public class TriggerIndexer : ITriggerIndexer
Hash = _hasher.Hash(x),
Data = JsonSerializer.Serialize(x)
});
return triggers.ToList();
}
@ -183,7 +183,7 @@ public class TriggerIndexer : ITriggerIndexer
var assignedInputs = inputs.Where(x => x.LocationReference != null!).ToList();
var register = context.GetOrCreateRegister(trigger);
var cancellationToken = context.CancellationToken;
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, default, cancellationToken);
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, context.Workflow, new Dictionary<string, object?>(), default, cancellationToken);
// Evaluate activity inputs before requesting trigger data.
foreach (var input in assignedInputs)

View file

@ -14,6 +14,7 @@
<ProjectReference Include="..\..\..\modules\Elsa.Modules.Quartz\Elsa.Modules.Quartz.csproj" />
<ProjectReference Include="..\..\..\api\Elsa.Api\Elsa.Api.csproj" />
<ProjectReference Include="..\..\..\core\Elsa\Elsa.csproj" />
<ProjectReference Include="..\..\..\modules\Elsa.Modules.WorkflowContexts\Elsa.Modules.WorkflowContexts.csproj" />
<ProjectReference Include="..\..\..\persistence\Elsa.Persistence.EntityFrameworkCore.Sqlite\Elsa.Persistence.EntityFrameworkCore.Sqlite.csproj" />
<ProjectReference Include="..\..\..\runtime\Elsa.Runtime.ProtoActor\Elsa.Runtime.ProtoActor.csproj" />
<ProjectReference Include="..\..\..\scripting\Elsa.Scripting.Liquid\Elsa.Scripting.Liquid.csproj" />

View file

@ -17,6 +17,7 @@ using Elsa.Modules.Quartz.Services;
using Elsa.Modules.Scheduling.Activities;
using Elsa.Modules.Scheduling.Extensions;
using Elsa.Modules.Scheduling.Triggers;
using Elsa.Modules.WorkflowContexts.Extensions;
using Elsa.Persistence.EntityFrameworkCore.Extensions;
using Elsa.Persistence.EntityFrameworkCore.Sqlite;
using Elsa.Pipelines.WorkflowExecution.Components;
@ -58,6 +59,7 @@ services
options.Workflows.Add(nameof(SendMessageWorkflow), new SendMessageWorkflow());
options.Workflows.Add(nameof(ReceiveMessageWorkflow), new ReceiveMessageWorkflow());
options.Workflows.Add(nameof(RunJavaScriptWorkflow), new RunJavaScriptWorkflow());
options.Workflows.Add(nameof(WorkflowContextsWorkflow), new WorkflowContextsWorkflow());
});
// Testing only: allow client app to connect from anywhere.
@ -101,6 +103,7 @@ serviceProvider.ConfigureDefaultWorkflowExecutionPipeline(pipeline =>
.UseWorkflowExecutionEvents()
.UseWorkflowExecutionLogPersistence()
.UsePersistence()
.UseWorkflowContexts()
.UseActivityScheduler()
);

View file

@ -0,0 +1,43 @@
using Elsa.Activities.Console;
using Elsa.Contracts;
using Elsa.Models;
using Elsa.Modules.WorkflowContexts.Abstractions;
using Elsa.Modules.WorkflowContexts.Extensions;
using Elsa.Modules.WorkflowContexts.Models;
using Elsa.Runtime.Contracts;
namespace Elsa.Samples.Web1.Workflows;
public class WorkflowContextsWorkflow : IWorkflow
{
public void Build(IWorkflowDefinitionBuilder workflow)
{
var documentContext = new WorkflowContext<Document, DocumentProvider>();
workflow
.AddWorkflowContext(documentContext)
.WithRoot(new WriteLine(context => $"Document title: {documentContext.Get(context)!.Title}"));
}
}
public class DocumentProvider : WorkflowContextProvider<Document>
{
protected override Document? Load(WorkflowExecutionContext workflowExecutionContext)
{
var idGenerator = workflowExecutionContext.GetRequiredService<IIdentityGenerator>();
return new Document
{
Id = idGenerator.GenerateId(),
Title = "Requirements",
Body = "Workflows should be able to load contextual data easily"
};
}
}
public class Document
{
public string Id { get; set; } = default!;
public string Title { get; set; } = default!;
public string Body { get; set; } = default!;
}