Exportable & importable workflow state + execution (#3317)

* Incremental work on externalizing workflow state execution

* Finish resumable WriteHttpResponse
This commit is contained in:
Sipke Schoorstra 2022-10-04 11:57:01 +02:00 committed by GitHub
parent a2cc77d952
commit 4b0af64cee
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
53 changed files with 878 additions and 328 deletions

View file

@ -2,6 +2,7 @@
<s:Boolean x:Key="/Default/UserDictionary/Words/=Comparers/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Computables/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Configurer/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Dtos/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=materializer/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Materializers/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Pomelo/@EntryIndexedValue">True</s:Boolean>

View file

@ -19,6 +19,7 @@ using Elsa.Persistence.EntityFrameworkCore.Sqlite.Modules.ActivityDefinitions;
using Elsa.Persistence.EntityFrameworkCore.Sqlite.Modules.Labels;
using Elsa.Persistence.EntityFrameworkCore.Sqlite.Modules.Management;
using Elsa.Persistence.EntityFrameworkCore.Sqlite.Modules.Runtime;
using Elsa.ProtoActor.Extensions;
using Elsa.Requirements;
using Elsa.Scheduling.Extensions;
using Elsa.WorkflowContexts.Extensions;
@ -34,6 +35,9 @@ using Elsa.WorkflowServer.Web.Jobs;
using FastEndpoints;
using Microsoft.AspNetCore.Authentication.JwtBearer;
using Microsoft.AspNetCore.Authorization;
using Microsoft.Data.Sqlite;
using Proto.Persistence.Sqlite;
using Event = Elsa.Workflows.Core.Activities.Event;
var builder = WebApplication.CreateBuilder(args);
var services = builder.Services;
@ -68,9 +72,9 @@ services
identity.CreateDefaultUser = true;
identity.IdentityOptions = options => identitySection.Bind(options);
})
//.UseRuntime(runtime => runtime.UseProtoActor(proto => proto.PersistenceProvider = _ => new SqliteProvider(new SqliteConnectionStringBuilder(sqliteConnectionString))))
.UseRuntime(runtime =>
{
runtime.UseProtoActor(proto => proto.PersistenceProvider = _ => new SqliteProvider(new SqliteConnectionStringBuilder(sqliteConnectionString)));
runtime.UseEntityFrameworkCore(ef => ef.UseSqlite(sqliteConnectionString));
runtime.WorkflowStateExporter = sp => sp.GetRequiredService<AsyncWorkflowStateExporter>();
})

View file

@ -498,6 +498,7 @@ export class FlowchartComponent implements ContainerActivityComponent {
};
private onNodeContextMenu = async (e: PositionEventArgs<JQuery.ContextMenuEvent>) => {
debugger;
const node = e.node as ActivityNodeShape;
const activity = e.node.data as Activity;

View file

@ -33,10 +33,10 @@ public static class ObjectConverter
options.ReferenceHandler = ReferenceHandler.Preserve;
options.PropertyNameCaseInsensitive = true;
options.Converters.Add(new JsonStringEnumConverter());
if (value is DahomeyJsonNode { ValueKind: JsonValueKind.Object } dahomyJsonObject)
return ToObject(dahomyJsonObject, targetType, options);
if (value is JsonElement { ValueKind: JsonValueKind.Object } jsonObject)
return jsonObject.Deserialize(targetType, options);
@ -56,7 +56,9 @@ public static class ObjectConverter
var targetTypeConverter = TypeDescriptor.GetConverter(underlyingTargetType);
if (targetTypeConverter.CanConvertFrom(underlyingSourceType))
return targetTypeConverter.IsValid(value) ? targetTypeConverter.ConvertFrom(value) : targetType.GetDefaultValue();
return targetTypeConverter.IsValid(value)
? targetTypeConverter.ConvertFrom(value)
: targetType.GetDefaultValue();
var sourceTypeConverter = TypeDescriptor.GetConverter(underlyingSourceType);
@ -65,8 +67,14 @@ public static class ObjectConverter
if (underlyingTargetType.IsEnum)
{
if (underlyingSourceType != typeof(string))
if (underlyingSourceType == typeof(string))
return Enum.Parse(underlyingTargetType, (string)value);
if (underlyingSourceType == typeof(int))
return Enum.ToObject(underlyingTargetType, value);
if (underlyingSourceType == typeof(double))
return Enum.ToObject(underlyingTargetType, Convert.ChangeType(value, typeof(int)));
}
try
@ -75,7 +83,9 @@ public static class ObjectConverter
}
catch (InvalidCastException e)
{
throw new TypeConversionException($"Failed to convert an object of type {sourceType} to {underlyingTargetType}", value, underlyingTargetType, e);
throw new TypeConversionException(
$"Failed to convert an object of type {sourceType} to {underlyingTargetType}", value,
underlyingTargetType, e);
}
}

View file

@ -62,6 +62,6 @@ public class HttpEndpoint : Trigger<HttpRequestModel>
// Generate bookmark data for path and selected methods.
var path = context.Get(Path);
var methods = context.Get(SupportedMethods);
return methods!.Select(x => new HttpEndpointBookmarkData(path!, x.ToLowerInvariant())).Cast<object>().ToArray();
return methods!.Select(x => new HttpEndpointBookmarkPayload(path!, x.ToLowerInvariant())).Cast<object>().ToArray();
}
}

View file

@ -1,3 +1,4 @@
using System;
using System.Net;
using System.Threading.Tasks;
using Elsa.Workflows.Core.Attributes;
@ -16,6 +17,37 @@ public class WriteHttpResponse : Activity
{
var httpContextAccessor = context.GetRequiredService<IHttpContextAccessor>();
var httpContext = httpContextAccessor.HttpContext;
if (httpContext == null)
{
// We're executing in a non-HTTP context (like in a virtual actor).
// Create a bookmark to allow the invoker to get the export state and resume execution from there.
context.CreateBookmark(OnResumeAsync);
return;
}
var response = httpContext.Response;
response.StatusCode = (int)context.Get(StatusCode);
var content = context.Get(Content);
if (content != null)
await response.WriteAsync(content, context.CancellationToken);
}
private async ValueTask OnResumeAsync(ActivityExecutionContext context)
{
var httpContextAccessor = context.GetRequiredService<IHttpContextAccessor>();
var httpContext = httpContextAccessor.HttpContext;
if (httpContext == null)
{
// We're not in an HTTP context, so let's fail.
throw new Exception("Cannot execute in a non-HTTP context");
}
var response = httpContext.Response;
response.StatusCode = (int)context.Get(StatusCode);

View file

@ -12,7 +12,7 @@ namespace Elsa.Http.Extensions;
public static class RouteTableExtensions
{
public static void AddRoutes(this IRouteTable routeTable, IEnumerable<WorkflowTrigger> triggers)
public static void AddRoutes(this IRouteTable routeTable, IEnumerable<StoredTrigger> triggers)
{
var paths = Filter(triggers).Select(Deserialize).Select(x => x.Path).ToList();
routeTable.AddRange(paths);
@ -24,7 +24,7 @@ public static class RouteTableExtensions
routeTable.AddRange(paths);
}
public static void RemoveRoutes(this IRouteTable routeTable, IEnumerable<WorkflowTrigger> triggers)
public static void RemoveRoutes(this IRouteTable routeTable, IEnumerable<StoredTrigger> triggers)
{
var paths = Filter(triggers).Select(Deserialize).Select(x => x.Path).ToList();
routeTable.RemoveRange(paths);
@ -36,9 +36,9 @@ public static class RouteTableExtensions
routeTable.RemoveRange(paths);
}
private static IEnumerable<WorkflowTrigger> Filter(IEnumerable<WorkflowTrigger> triggers) => triggers.Where(x => x.Name == ActivityTypeNameHelper.GenerateTypeName<HttpEndpoint>());
private static IEnumerable<StoredTrigger> Filter(IEnumerable<StoredTrigger> triggers) => triggers.Where(x => x.Name == ActivityTypeNameHelper.GenerateTypeName<HttpEndpoint>());
private static IEnumerable<Bookmark> Filter(IEnumerable<Bookmark> triggers) => triggers.Where(x => x.Name == ActivityTypeNameHelper.GenerateTypeName<HttpEndpoint>());
private static HttpEndpointBookmarkData Deserialize(WorkflowTrigger trigger) => Deserialize(trigger.Data!);
private static HttpEndpointBookmarkData Deserialize(Bookmark bookmark) => Deserialize(bookmark.Data!);
private static HttpEndpointBookmarkData Deserialize(string model) => JsonSerializer.Deserialize<HttpEndpointBookmarkData>(model)!;
private static HttpEndpointBookmarkPayload Deserialize(StoredTrigger trigger) => Deserialize(trigger.Data!);
private static HttpEndpointBookmarkPayload Deserialize(Bookmark bookmark) => Deserialize(bookmark.Data!);
private static HttpEndpointBookmarkPayload Deserialize(string model) => JsonSerializer.Deserialize<HttpEndpointBookmarkPayload>(model)!;
}

View file

@ -5,10 +5,13 @@ using System.Net.Mime;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Common.Models;
using Elsa.Http.Models;
using Elsa.Http.Options;
using Elsa.Http.Services;
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Runtime.Services;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.Http.Extensions;
using Microsoft.AspNetCore.Routing;
@ -20,13 +23,26 @@ namespace Elsa.Http.Middleware;
public class HttpTriggerMiddleware
{
private readonly RequestDelegate _next;
private readonly IHasher _hasher;
private readonly IBookmarkHasher _hasher;
private readonly IWorkflowRuntime _workflowRuntime;
private readonly IWorkflowHostFactory _workflowHostFactory;
private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly HttpActivityOptions _options;
private readonly string _activityTypeName = ActivityTypeNameHelper.GenerateTypeName<HttpEndpoint>();
public HttpTriggerMiddleware(RequestDelegate next, IHasher hasher, IOptions<HttpActivityOptions> options)
public HttpTriggerMiddleware(
RequestDelegate next,
IBookmarkHasher hasher,
IWorkflowRuntime workflowRuntime,
IWorkflowHostFactory workflowHostFactory,
IWorkflowDefinitionService workflowDefinitionService,
IOptions<HttpActivityOptions> options)
{
_next = next;
_hasher = hasher;
_workflowRuntime = workflowRuntime;
_workflowHostFactory = workflowHostFactory;
_workflowDefinitionService = workflowDefinitionService;
_options = options.Value;
}
@ -50,8 +66,7 @@ public class HttpTriggerMiddleware
var request = httpContext.Request;
var method = request.Method!.ToLowerInvariant();
var abortToken = httpContext.RequestAborted;
var hash = _hasher.Hash(new HttpEndpointBookmarkData(path, method));
var cancellationToken = httpContext.RequestAborted;
var routeData = GetRouteData(httpContext, routeMatcher, path);
var requestModel = new HttpRequestModel(
@ -63,7 +78,70 @@ public class HttpTriggerMiddleware
request.Headers.ToDictionary(x => x.Key, x => x.Value.ToString())
);
var input = new Dictionary<string, object>() { [HttpEndpoint.InputKey] = requestModel };
var input = new Dictionary<string, object> { [HttpEndpoint.InputKey] = requestModel };
// TODO: Get correlation ID from query string or header etc.
var correlationId = default(string);
// Trigger the workflow.
var bookmarkPayload = new HttpEndpointBookmarkPayload(path, method);
var triggerOptions = new TriggerWorkflowsOptions(correlationId, input);
var triggerResult = await _workflowRuntime.TriggerWorkflowsAsync(
_activityTypeName,
bookmarkPayload,
triggerOptions,
cancellationToken);
// Check to see if we received any WriteHttpResponse activity bookmarks. If we do, acquire a lock on the workflow instance and resume it from here within an actual HTTP context so that the activity can complete its HTTP response.
var writeHttpResponseTypeName = ActivityTypeNameHelper.GenerateTypeName<WriteHttpResponse>();
var query =
from triggeredWorkflow in triggerResult.TriggeredWorkflows
from bookmark in triggeredWorkflow.Bookmarks
where bookmark.Name == writeHttpResponseTypeName
select (triggeredWorkflow.InstanceId, bookmark.Id);
var workflowExecutionResults = new Stack<(string InstanceId, string BookmarkId)>(query);
while (workflowExecutionResults.TryPop(out var result))
{
// Resume the workflow "in-process".
var workflowState = await _workflowRuntime.ExportWorkflowStateAsync(
result.InstanceId,
cancellationToken);
if (workflowState == null)
{
// TODO: log this, shouldn't normally happen.
continue;
}
var workflowDefinition = await _workflowDefinitionService.FindAsync(
workflowState.DefinitionId,
VersionOptions.SpecificVersion(workflowState.DefinitionVersion),
cancellationToken);
if (workflowDefinition == null)
{
// TODO: Log this, shouldn't normally happen.
continue;
}
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(
workflowDefinition,
cancellationToken);
var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
await workflowHost.ResumeWorkflowAsync(
result.BookmarkId,
null,
cancellationToken);
// Import the updated workflow state into the runtime.
await _workflowRuntime.ImportWorkflowStateAsync(workflowState, cancellationToken);
}
}
private static async Task WriteResponseAsync(HttpContext httpContext, CancellationToken cancellationToken)

View file

@ -1,11 +1,11 @@
namespace Elsa.Http.Models;
public record HttpEndpointBookmarkData
public record HttpEndpointBookmarkPayload
{
private readonly string _path = default!;
private readonly string _method = default!;
public HttpEndpointBookmarkData(string path, string method)
public HttpEndpointBookmarkPayload(string path, string method)
{
Path = path;
Method = method;

View file

@ -1,26 +0,0 @@
using Elsa.Persistence.EntityFrameworkCore.Common;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime;
public class EFCoreBookmarkStore : IBookmarkStore
{
private readonly Store<RuntimeDbContext, StoredBookmark> _store;
public EFCoreBookmarkStore(Store<RuntimeDbContext, StoredBookmark> store) => _store = store;
public async ValueTask SaveAsync(string activityTypeName, string hash, string workflowInstanceId, IEnumerable<string> bookmarkIds, CancellationToken cancellationToken = default)
{
var storedBookmarks = bookmarkIds.Select(x => new StoredBookmark(activityTypeName, hash, workflowInstanceId, x)).ToList();
await _store.SaveManyAsync(storedBookmarks, cancellationToken);
}
public async ValueTask<IEnumerable<StoredBookmark>> LoadAsync(string workflowInstanceId, CancellationToken cancellationToken = default) =>
await _store.FindManyAsync(x => x.WorkflowInstanceId == workflowInstanceId, cancellationToken);
public async ValueTask<IEnumerable<StoredBookmark>> LoadAsync(string activityTypeName, string hash, CancellationToken cancellationToken = default) =>
await _store.FindManyAsync(x => x.ActivityTypeName == activityTypeName && x.Hash == hash, cancellationToken);
public async ValueTask DeleteAsync(string activityTypeName, string hash, string workflowInstanceId, CancellationToken cancellationToken = default) =>
await _store.DeleteWhereAsync(x => x.WorkflowInstanceId == workflowInstanceId && x.ActivityTypeName == activityTypeName && x.Hash == hash, cancellationToken);
}

View file

@ -11,7 +11,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime
{
public class Configurations :
IEntityTypeConfiguration<WorkflowState>,
IEntityTypeConfiguration<WorkflowTrigger>,
IEntityTypeConfiguration<StoredTrigger>,
IEntityTypeConfiguration<WorkflowExecutionLogRecord>,
IEntityTypeConfiguration<StoredBookmark>
{
@ -38,11 +38,11 @@ namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime
builder.HasIndex("UpdatedAt").HasDatabaseName($"IX_{nameof(WorkflowState)}_UpdatedAt");
}
public void Configure(EntityTypeBuilder<WorkflowTrigger> builder)
public void Configure(EntityTypeBuilder<StoredTrigger> builder)
{
builder.HasIndex(x => x.WorkflowDefinitionId).HasDatabaseName($"IX_{nameof(WorkflowTrigger)}_{nameof(WorkflowTrigger.WorkflowDefinitionId)}");
builder.HasIndex(x => x.Name).HasDatabaseName($"IX_{nameof(WorkflowTrigger)}_{nameof(WorkflowTrigger.Name)}");
builder.HasIndex(x => x.Hash).HasDatabaseName($"IX_{nameof(WorkflowTrigger)}_{nameof(WorkflowTrigger.Hash)}");
builder.HasIndex(x => x.WorkflowDefinitionId).HasDatabaseName($"IX_{nameof(StoredTrigger)}_{nameof(StoredTrigger.WorkflowDefinitionId)}");
builder.HasIndex(x => x.Name).HasDatabaseName($"IX_{nameof(StoredTrigger)}_{nameof(StoredTrigger.Name)}");
builder.HasIndex(x => x.Hash).HasDatabaseName($"IX_{nameof(StoredTrigger)}_{nameof(StoredTrigger.Hash)}");
}
public void Configure(EntityTypeBuilder<WorkflowExecutionLogRecord> builder)

View file

@ -0,0 +1,39 @@
using Elsa.Persistence.EntityFrameworkCore.Common;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime;
public class EFCoreBookmarkStore : IBookmarkStore
{
private readonly Store<RuntimeDbContext, StoredBookmark> _store;
public EFCoreBookmarkStore(Store<RuntimeDbContext, StoredBookmark> store) => _store = store;
public async ValueTask SaveAsync(
string activityTypeName,
string hash,
string workflowInstanceId,
IEnumerable<string> bookmarkIds, CancellationToken cancellationToken = default)
{
var storedBookmarks = bookmarkIds.Select(x => new StoredBookmark(activityTypeName, hash, workflowInstanceId, x))
.ToList();
await _store.SaveManyAsync(storedBookmarks, cancellationToken);
}
public async ValueTask<IEnumerable<StoredBookmark>> FindByWorkflowInstanceAsync(
string workflowInstanceId,
CancellationToken cancellationToken = default) =>
await _store.FindManyAsync(x => x.WorkflowInstanceId == workflowInstanceId, cancellationToken);
public async ValueTask<IEnumerable<StoredBookmark>> FindByHashAsync(string hash,
CancellationToken cancellationToken = default) =>
await _store.FindManyAsync(x => x.Hash == hash, cancellationToken);
public async ValueTask DeleteAsync(
string hash,
string workflowInstanceId,
CancellationToken cancellationToken = default) =>
await _store.DeleteWhereAsync(
x => x.WorkflowInstanceId == workflowInstanceId && x.Hash == hash,
cancellationToken);
}

View file

@ -21,7 +21,7 @@ public class EFCoreRuntimePersistenceFeature : PersistenceFeatureBase<RuntimeDbC
Module.Configure<WorkflowRuntimeFeature>(feature =>
{
feature.WorkflowStateStore = sp => sp.GetRequiredService<EFCoreWorkflowStateStore>();
feature.WorkflowTriggerStore = sp => sp.GetRequiredService<EFCoreWorkflowTriggerStore>();
feature.WorkflowTriggerStore = sp => sp.GetRequiredService<EFCoreTriggerStore>();
feature.BookmarkStore = sp => sp.GetRequiredService<EFCoreBookmarkStore>();
feature.WorkflowExecutionLogStore = sp => sp.GetRequiredService<EFCoreWorkflowExecutionLogStore>();
});
@ -32,7 +32,7 @@ public class EFCoreRuntimePersistenceFeature : PersistenceFeatureBase<RuntimeDbC
base.Apply();
AddStore<WorkflowState, EFCoreWorkflowStateStore>();
AddStore<WorkflowTrigger, EFCoreWorkflowTriggerStore>();
AddStore<StoredTrigger, EFCoreTriggerStore>();
AddStore<StoredBookmark, EFCoreBookmarkStore>();
AddStore<WorkflowExecutionLogRecord, EFCoreWorkflowExecutionLogStore>();
}

View file

@ -0,0 +1,38 @@
using Elsa.Persistence.EntityFrameworkCore.Common;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime;
public class EFCoreTriggerStore : ITriggerStore
{
private readonly Store<RuntimeDbContext, StoredTrigger> _store;
public EFCoreTriggerStore(Store<RuntimeDbContext, StoredTrigger> store)
{
_store = store;
}
public async Task SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default) => await _store.SaveAsync(record, cancellationToken);
public async Task SaveManyAsync(IEnumerable<StoredTrigger> records, CancellationToken cancellationToken = default) => await _store.SaveManyAsync(records, cancellationToken);
public async Task<IEnumerable<StoredTrigger>> FindAsync(
string hash,
CancellationToken cancellationToken = default) =>
await _store.QueryAsync(query => query.Where(x => x.Hash == hash), cancellationToken);
public async Task<IEnumerable<StoredTrigger>> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) =>
await _store.QueryAsync(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId), cancellationToken);
public async Task ReplaceAsync(IEnumerable<StoredTrigger> removed, IEnumerable<StoredTrigger> added, CancellationToken cancellationToken = default)
{
await _store.DeleteManyAsync(removed, cancellationToken);
await _store.SaveManyAsync(added, cancellationToken);
}
public async Task DeleteManyAsync(IEnumerable<string> ids, CancellationToken cancellationToken = default)
{
var idList = ids.ToList();
await _store.DeleteWhereAsync(x => idList.Contains(x.Id), cancellationToken);
}
}

View file

@ -14,7 +14,7 @@ public class RuntimeDbContext : DbContextBase
}
public DbSet<WorkflowState> WorkflowStates { get; set; } = default!;
public DbSet<WorkflowTrigger> WorkflowTriggers { get; set; } = default!;
public DbSet<StoredTrigger> WorkflowTriggers { get; set; } = default!;
public DbSet<WorkflowExecutionLogRecord> WorkflowExecutionLogRecords { get; set; } = default!;
public DbSet<StoredBookmark> Bookmarks { get; set; } = default!;
@ -22,7 +22,7 @@ public class RuntimeDbContext : DbContextBase
{
var config = new Configurations();
modelBuilder.ApplyConfiguration<WorkflowState>(config);
modelBuilder.ApplyConfiguration<WorkflowTrigger>(config);
modelBuilder.ApplyConfiguration<StoredTrigger>(config);
modelBuilder.ApplyConfiguration<WorkflowExecutionLogRecord>(config);
modelBuilder.ApplyConfiguration<StoredBookmark>(config);
}

View file

@ -1,44 +0,0 @@
using Elsa.Persistence.EntityFrameworkCore.Common;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime;
public class EFCoreWorkflowTriggerStore : IWorkflowTriggerStore
{
private readonly Store<RuntimeDbContext, WorkflowTrigger> _store;
public EFCoreWorkflowTriggerStore(Store<RuntimeDbContext, WorkflowTrigger> store)
{
_store = store;
}
public async Task SaveAsync(WorkflowTrigger record, CancellationToken cancellationToken = default) => await _store.SaveAsync(record, cancellationToken);
public async Task SaveManyAsync(IEnumerable<WorkflowTrigger> records, CancellationToken cancellationToken = default) => await _store.SaveManyAsync(records, cancellationToken);
public async Task<IEnumerable<WorkflowTrigger>> FindManyByNameAsync(string name, string? hash, CancellationToken cancellationToken = default) =>
await _store.QueryAsync(query =>
{
query = query.Where(x => x.Name == name);
if (hash != null)
query = query.Where(x => x.Hash == hash);
return query;
}, cancellationToken);
public async Task<IEnumerable<WorkflowTrigger>> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) =>
await _store.QueryAsync(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId), cancellationToken);
public async Task ReplaceAsync(IEnumerable<WorkflowTrigger> removed, IEnumerable<WorkflowTrigger> added, CancellationToken cancellationToken = default)
{
await _store.DeleteManyAsync(removed, cancellationToken);
await _store.SaveManyAsync(added, cancellationToken);
}
public async Task DeleteManyAsync(IEnumerable<string> ids, CancellationToken cancellationToken = default)
{
var idList = ids.ToList();
await _store.DeleteWhereAsync(x => idList.Contains(x.Id), cancellationToken);
}
}

View file

@ -0,0 +1,7 @@
namespace Elsa.ProtoActor.Extensions;
public static class ProtoStringExtensions
{
public static string EmptyIfNull(this string? value) => value ?? "";
public static string? NullIfEmpty(this string value) => value == "" ? default : value;
}

View file

@ -1,6 +1,7 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Common.Models;
@ -9,6 +10,7 @@ using Elsa.ProtoActor.Extensions;
using Elsa.Runtime.Protos;
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Serialization;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Core.State;
using Elsa.Workflows.Runtime.Models;
@ -17,39 +19,38 @@ using Elsa.Workflows.Runtime.Services;
using Proto;
using Proto.Cluster;
using Proto.Persistence;
using Bookmark = Elsa.Workflows.Core.Models.Bookmark;
using RunWorkflowResult = Elsa.Runtime.Protos.RunWorkflowResult;
namespace Elsa.ProtoActor.Grains;
using Persistence = Proto.Persistence.Persistence;
/// <summary>
/// Executes a workflow.
/// </summary>
public class WorkflowGrain : WorkflowGrainBase
{
private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly IWorkflowRunner _workflowRunner;
private readonly IEventPublisher _eventPublisher;
private readonly IWorkflowHostFactory _workflowHostFactory;
private readonly SerializerOptionsProvider _serializerOptionsProvider;
private readonly Persistence _persistence;
private string _definitionId = default!;
private int _version;
private IDictionary<string, object>? _input;
private Workflow _workflow = default!;
private IWorkflowHost _workflowHost = default!;
private WorkflowState _workflowState = default!;
private ICollection<Bookmark> _bookmarks = new List<Bookmark>();
//private ICollection<Bookmark> _bookmarks = default!;
public WorkflowGrain(
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowRunner workflowRunner,
IEventPublisher eventPublisher,
IWorkflowHostFactory workflowHostFactory,
SerializerOptionsProvider serializerOptionsProvider,
IProvider provider,
IContext context) : base(context)
{
_workflowDefinitionService = workflowDefinitionService;
_workflowRunner = workflowRunner;
_eventPublisher = eventPublisher;
_workflowHostFactory = workflowHostFactory;
_serializerOptionsProvider = serializerOptionsProvider;
_persistence = Persistence.WithSnapshotting(provider, WorkflowInstanceId, ApplySnapshot);
}
@ -71,7 +72,21 @@ public class WorkflowGrain : WorkflowGrainBase
throw new Exception("Workflow definition is no longer available");
// Materialize the workflow.
_workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
// Create an initial workflow state.
if (_workflowState == null!)
{
_workflowState = new WorkflowState
{
DefinitionId = workflow.Identity.DefinitionId,
DefinitionVersion = workflow.Identity.Version,
//Bookmarks = _bookmarks
};
}
// Create a workflow host.
_workflowHost = await _workflowHostFactory.CreateAsync(workflow, _workflowState, cancellationToken);
}
public override async Task<StartWorkflowResponse> Start(StartWorkflowRequest request)
@ -86,61 +101,80 @@ public class WorkflowGrain : WorkflowGrainBase
if (workflowDefinition == null)
throw new Exception("Specified workflow definition and version does not exist");
_workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var version = workflowDefinition.Version;
await _eventPublisher.PublishAsync(new WorkflowExecuting(_workflow), cancellationToken);
var workflowResult = await _workflowRunner.RunAsync(_workflow, WorkflowInstanceId, _input, cancellationToken);
var finished = workflowResult.WorkflowState.Status == WorkflowStatus.Finished;
_workflowState = workflowResult.WorkflowState;
_version = version;
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
_version = workflow.Version;
_definitionId = definitionId;
_input = input;
// Create a workflow host.
_workflowHost = await _workflowHostFactory.CreateAsync(workflow, cancellationToken);
await UpdateBookmarksAsync(_workflowState.Bookmarks, cancellationToken);
var startWorkflowResult = await _workflowHost.StartWorkflowAsync(WorkflowInstanceId, input, cancellationToken);
_workflowState = _workflowHost.WorkflowState;
await UpdateBookmarksAsync(startWorkflowResult.BookmarksDiff, cancellationToken);
await SaveSnapshotAsync();
await _eventPublisher.PublishAsync(new WorkflowExecuted(_workflow, _workflowState), cancellationToken);
return new StartWorkflowResponse
{
Result = finished ? RunWorkflowResult.Finished : RunWorkflowResult.Suspended
Result = _workflowHost.WorkflowState.Status == WorkflowStatus.Finished ? RunWorkflowResult.Finished : RunWorkflowResult.Suspended,
Bookmarks = { Map(_workflowHost.WorkflowState.Bookmarks) }
};
}
public override async Task<ResumeWorkflowResponse> Resume(ResumeWorkflowRequest request)
{
var input = request.Input?.Deserialize();
_input = request.Input?.Deserialize();
var bookmarkId = request.BookmarkId;
var cancellationToken = Context.CancellationToken;
await _eventPublisher.PublishAsync(new WorkflowExecuting(_workflow), cancellationToken);
var workflowResult = await _workflowRunner.RunAsync(_workflow, _workflowState, bookmarkId, input, cancellationToken);
var finished = workflowResult.WorkflowState.Status == WorkflowStatus.Finished;
_workflowState = workflowResult.WorkflowState;
_input = input;
var resumeWorkflowResult = await _workflowHost.ResumeWorkflowAsync(bookmarkId, _input, cancellationToken);
var finished = _workflowHost.WorkflowState.Status == WorkflowStatus.Finished;
await UpdateBookmarksAsync(_workflowState.Bookmarks, cancellationToken);
_workflowState = _workflowHost.WorkflowState;
await UpdateBookmarksAsync(resumeWorkflowResult.BookmarksDiff, cancellationToken);
await SaveSnapshotAsync();
await _eventPublisher.PublishAsync(new WorkflowExecuted(_workflow, _workflowState), cancellationToken);
return new ResumeWorkflowResponse
{
Result = finished ? RunWorkflowResult.Finished : RunWorkflowResult.Suspended
Result = finished ? RunWorkflowResult.Finished : RunWorkflowResult.Suspended,
Bookmarks = { Map(_workflowHost.WorkflowState.Bookmarks) }
};
}
private async Task UpdateBookmarksAsync(ICollection<Bookmark> bookmarks, CancellationToken cancellationToken)
public override Task<ExportWorkflowStateResponse> ExportState(ExportWorkflowStateRequest request)
{
var originalBookmarks = _bookmarks;
_bookmarks = bookmarks;
var options = _serializerOptionsProvider.CreatePersistenceOptions();
var json = JsonSerializer.Serialize(_workflowHost.WorkflowState, options);
var response = new ExportWorkflowStateResponse
{
SerializedWorkflowState = new Json
{
Text = json
}
};
await RemoveBookmarksAsync(originalBookmarks, cancellationToken);
await StoreBookmarksAsync(bookmarks, cancellationToken);
await PublishChangedBookmarksAsync(originalBookmarks, bookmarks, cancellationToken);
return Task.FromResult(response);
}
public override Task<ImportWorkflowStateResponse> ImportState(ImportWorkflowStateRequest request)
{
var options = _serializerOptionsProvider.CreatePersistenceOptions();
var workflowState = JsonSerializer.Deserialize<WorkflowState>(request.SerializedWorkflowState.Text, options)!;
_workflowState = workflowState;
_workflowHost.WorkflowState = workflowState;
return Task.FromResult(new ImportWorkflowStateResponse());
}
private async Task UpdateBookmarksAsync(Diff<Bookmark> bookmarksDiff, CancellationToken cancellationToken)
{
await RemoveBookmarksAsync(bookmarksDiff.Removed, cancellationToken);
await StoreBookmarksAsync(bookmarksDiff.Added, cancellationToken);
}
private async Task StoreBookmarksAsync(ICollection<Bookmark> bookmarks, CancellationToken cancellationToken)
@ -175,15 +209,19 @@ public class WorkflowGrain : WorkflowGrainBase
}
}
private async Task PublishChangedBookmarksAsync(ICollection<Bookmark> originalBookmarks, ICollection<Bookmark> updatedBookmarks, CancellationToken cancellationToken)
{
var diff = Diff.For(originalBookmarks, updatedBookmarks);
var removedBookmarks = diff.Removed;
var createdBookmarks = diff.Added;
await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(new IndexedWorkflowBookmarks(_workflowState, createdBookmarks, removedBookmarks)), cancellationToken);
}
private void ApplySnapshot(Snapshot snapshot) => (_definitionId, _version, _workflowState, _bookmarks, _input) = (WorkflowSnapshot)snapshot.State;
private void ApplySnapshot(Snapshot snapshot) => (_definitionId, _version, _workflowState, _input) = (WorkflowSnapshot)snapshot.State;
private async Task SaveSnapshotAsync() => await _persistence.PersistSnapshotAsync(GetState());
private object GetState() => new WorkflowSnapshot(_definitionId, _version, _workflowState, _bookmarks, _input);
private object GetState() => new WorkflowSnapshot(_definitionId, _version, _workflowState, _input);
private static IEnumerable<BookmarkDto> Map(IEnumerable<Bookmark> bookmarks) =>
bookmarks.Select(x => new BookmarkDto
{
Id = x.Id,
Name = x.Name,
ActivityId = x.ActivityId,
ActivityInstanceId = x.ActivityInstanceId,
Hash = x.Hash,
Data = x.Data.EmptyIfNull(),
CallbackMethodName = x.CallbackMethodName.EmptyIfNull()
});
}

View file

@ -1,9 +1,16 @@
using System.Collections.Generic;
using System.Linq;
using System.Text.Json;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Common.Models;
using Elsa.ProtoActor.Extensions;
using Elsa.Runtime.Protos;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Serialization;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Core.State;
using Elsa.Workflows.Runtime.Services;
using Proto.Cluster;
@ -12,17 +19,29 @@ namespace Elsa.ProtoActor.Implementations;
public class ProtoActorWorkflowRuntime : IWorkflowRuntime
{
private readonly Cluster _cluster;
private readonly SerializerOptionsProvider _serializerOptionsProvider;
private readonly ITriggerStore _triggerStore;
private readonly IIdentityGenerator _identityGenerator;
private readonly IHasher _hasher;
private readonly IBookmarkHasher _hasher;
public ProtoActorWorkflowRuntime(Cluster cluster, IIdentityGenerator identityGenerator, IHasher hasher)
public ProtoActorWorkflowRuntime(
Cluster cluster,
SerializerOptionsProvider serializerOptionsProvider,
ITriggerStore triggerStore,
IIdentityGenerator identityGenerator,
IBookmarkHasher hasher)
{
_cluster = cluster;
_serializerOptionsProvider = serializerOptionsProvider;
_triggerStore = triggerStore;
_identityGenerator = identityGenerator;
_hasher = hasher;
}
public async Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowOptions options, CancellationToken cancellationToken = default)
public async Task<StartWorkflowResult> StartWorkflowAsync(
string definitionId,
StartWorkflowOptions options,
CancellationToken cancellationToken = default)
{
var versionOptions = options.VersionOptions;
var correlationId = options.CorrelationId;
@ -39,11 +58,16 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
var workflowInstanceId = _identityGenerator.GenerateId();
var client = _cluster.GetWorkflowGrain(workflowInstanceId);
var response = await client.Start(request, cancellationToken);
var bookmarks = Map(response!.Bookmarks).ToList();
return new StartWorkflowResult(workflowInstanceId);
return new StartWorkflowResult(workflowInstanceId, bookmarks);
}
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string instanceId, string bookmarkId, ResumeWorkflowOptions options, CancellationToken cancellationToken = default)
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(
string instanceId,
string bookmarkId,
ResumeWorkflowOptions options,
CancellationToken cancellationToken = default)
{
var request = new ResumeWorkflowRequest
{
@ -54,13 +78,34 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
var client = _cluster.GetWorkflowGrain(instanceId);
var response = await client.Resume(request, cancellationToken);
var bookmarks = Map(response!.Bookmarks).ToList();
return new ResumeWorkflowResult();
return new ResumeWorkflowResult(bookmarks);
}
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions options, CancellationToken cancellationToken = default)
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(
string activityTypeName,
object bookmarkPayload,
TriggerWorkflowsOptions options,
CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(bookmarkPayload);
var triggeredWorkflows = new List<TriggeredWorkflow>();
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
// Start new workflows.
var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
foreach (var trigger in triggers)
{
var startResult = await StartWorkflowAsync(
trigger.WorkflowDefinitionId,
new StartWorkflowOptions(options.CorrelationId, options.Input, VersionOptions.Published),
cancellationToken);
triggeredWorkflows.Add(new TriggeredWorkflow(startResult.InstanceId, startResult.Bookmarks));
}
// Resume existing workflow instances.
var client = _cluster.GetBookmarkGrain(hash);
var request = new ResolveBookmarksRequest() { BookmarkName = activityTypeName };
var bookmarksResponse = await client.Resolve(request, cancellationToken);
@ -69,9 +114,57 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
foreach (var bookmark in bookmarks)
{
var workflowInstanceId = bookmark.WorkflowInstanceId;
var resumeResult = await ResumeWorkflowAsync(workflowInstanceId, bookmark.BookmarkId, new ResumeWorkflowOptions(options.Input), cancellationToken);
var resumeResult = await ResumeWorkflowAsync(
workflowInstanceId,
bookmark.BookmarkId,
new ResumeWorkflowOptions(options.Input),
cancellationToken);
triggeredWorkflows.Add(new TriggeredWorkflow(workflowInstanceId, resumeResult.Bookmarks));
}
return new TriggerWorkflowsResult();
return new TriggerWorkflowsResult(triggeredWorkflows);
}
public async Task<WorkflowState?> ExportWorkflowStateAsync(
string instanceId,
CancellationToken cancellationToken = default)
{
var client = _cluster.GetWorkflowGrain(instanceId);
var response = await client.ExportState(new ExportWorkflowStateRequest(), cancellationToken);
var json = response!.SerializedWorkflowState.Text;
var options = _serializerOptionsProvider.CreatePersistenceOptions();
var workflowState = JsonSerializer.Deserialize<WorkflowState>(json, options);
return workflowState;
}
public async Task ImportWorkflowStateAsync(WorkflowState workflowState,
CancellationToken cancellationToken = default)
{
var options = _serializerOptionsProvider.CreatePersistenceOptions();
var client = _cluster.GetWorkflowGrain(workflowState.Id);
var json = JsonSerializer.Serialize(workflowState, options);
var request = new ImportWorkflowStateRequest
{
SerializedWorkflowState = new Json
{
Text = json
}
};
await client.ImportState(request, cancellationToken);
}
private static IEnumerable<Bookmark> Map(IEnumerable<BookmarkDto> bookmarkDtos) =>
bookmarkDtos.Select(x =>
new Bookmark(
x.Id,
x.Name,
x.Hash,
x.Data.NullIfEmpty(),
x.ActivityId,
x.ActivityInstanceId,
x.CallbackMethodName.NullIfEmpty()));
}

View file

@ -28,10 +28,10 @@ public class ProtoActorSystem
ClusterConfigurationSettings = clusterConfigurationSettings;
}
public IClusterProvider ClusterProvider { get; set; }
public GrpcNetRemoteConfig RemoteConfig { get; set; }
public IClusterProvider ClusterProvider { get; set; } = default!;
public GrpcNetRemoteConfig RemoteConfig { get; set; } = default!;
public ActorSystemConfig ActorSystemConfig { get; set; } = ActorSystemConfig.Setup();
public IIdentityLookup IdentityLookup { get; set; }
public IIdentityLookup IdentityLookup { get; set; }= default!;
public ClusterConfigurationSettings ClusterConfigurationSettings { get; set; } = new();
public string Name { get; set; }
public string Name { get; set; }= default!;
}

View file

@ -8,6 +8,8 @@ import "Messages.proto";
service WorkflowGrain {
rpc Start (StartWorkflowRequest) returns (StartWorkflowResponse);
rpc Resume (ResumeWorkflowRequest) returns (ResumeWorkflowResponse);
rpc ExportState(ExportWorkflowStateRequest) returns (ExportWorkflowStateResponse);
rpc ImportState(ImportWorkflowStateRequest) returns (ImportWorkflowStateResponse);
}
service BookmarkGrain {

View file

@ -15,6 +15,7 @@ message StartWorkflowRequest {
message StartWorkflowResponse {
RunWorkflowResult Result = 1;
repeated BookmarkDto Bookmarks = 2;
}
message ResumeWorkflowRequest {
@ -25,8 +26,21 @@ message ResumeWorkflowRequest {
message ResumeWorkflowResponse {
RunWorkflowResult Result = 1;
repeated BookmarkDto Bookmarks = 2;
}
message ExportWorkflowStateRequest {}
message ExportWorkflowStateResponse {
Json SerializedWorkflowState = 1;
}
message ImportWorkflowStateRequest {
Json SerializedWorkflowState = 1;
}
message ImportWorkflowStateResponse {}
enum RunWorkflowResult {
Finished = 0;
Suspended = 1;
@ -54,6 +68,16 @@ message StoredBookmark {
string BookmarkId = 2;
}
message BookmarkDto {
string Id = 1;
string Name = 2;
string Hash = 3;
optional string Data = 4;
string ActivityId = 5;
string ActivityInstanceId = 6;
optional string CallbackMethodName = 7;
}
message Json {
string text = 1;
}

View file

@ -5,5 +5,5 @@ using Elsa.Workflows.Core.State;
namespace Elsa.ProtoActor;
public record WorkflowSnapshot(string DefinitionId, int Version, WorkflowState WorkflowState, ICollection<Bookmark> Bookmarks, IDictionary<string, object>? Input);
public record WorkflowSnapshot(string DefinitionId, int Version, WorkflowState WorkflowState, IDictionary<string, object>? Input);
public record BookmarkSnapshot(ICollection<StoredBookmark> Bookmarks);

View file

@ -25,7 +25,7 @@ public class WorkflowTriggerScheduler : IWorkflowTriggerScheduler
_jobScheduler = jobScheduler;
}
public async Task ScheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default)
public async Task ScheduleTriggersAsync(IEnumerable<StoredTrigger> triggers, CancellationToken cancellationToken = default)
{
var triggerList = triggers.ToList();
@ -51,7 +51,7 @@ public class WorkflowTriggerScheduler : IWorkflowTriggerScheduler
}
}
public async Task UnscheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default)
public async Task UnscheduleTriggersAsync(IEnumerable<StoredTrigger> triggers, CancellationToken cancellationToken = default)
{
var triggerList = triggers.ToList();

View file

@ -10,6 +10,6 @@ namespace Elsa.Scheduling.Services;
/// </summary>
public interface IWorkflowTriggerScheduler
{
Task ScheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default);
Task UnscheduleTriggersAsync(IEnumerable<WorkflowTrigger> triggers, CancellationToken cancellationToken = default);
Task ScheduleTriggersAsync(IEnumerable<StoredTrigger> triggers, CancellationToken cancellationToken = default);
Task UnscheduleTriggersAsync(IEnumerable<StoredTrigger> triggers, CancellationToken cancellationToken = default);
}

View file

@ -28,7 +28,7 @@ public class Trigger : ElsaEndpoint<Request, Response>
public override async Task HandleAsync(Request request, CancellationToken cancellationToken)
{
var eventBookmark = new EventBookmarkData(request.EventName);
var eventBookmark = new EventBookmarkPayload(request.EventName);
var bookmarkName = ActivityTypeNameHelper.GenerateTypeName<Event>();
var result = await _workflowRuntime.TriggerWorkflowsAsync(bookmarkName, eventBookmark, new TriggerWorkflowsOptions(), cancellationToken);

View file

@ -35,6 +35,6 @@ public class Event : Trigger<object?>
protected override void Execute(ActivityExecutionContext context)
{
var eventName = context.Get(EventName)!;
context.CreateBookmark(new EventBookmarkData(eventName));
context.CreateBookmark(new EventBookmarkPayload(eventName));
}
}

View file

@ -21,6 +21,11 @@
<ProjectReference Include="..\Elsa.Common\Elsa.Common.csproj" />
<ProjectReference Include="..\Elsa.Expressions\Elsa.Expressions.csproj" />
<ProjectReference Include="..\..\common\Elsa.Features\Elsa.Features.csproj" />
<ProjectReference Include="..\Elsa.Mediator\Elsa.Mediator.csproj" />
</ItemGroup>
<ItemGroup>
<Folder Include="Notifications" />
</ItemGroup>
</Project>

View file

@ -65,6 +65,7 @@ public class WorkflowsFeature : FeatureBase
.AddSingleton<IWorkflowStateSerializer, WorkflowStateSerializer>()
.AddSingleton<IActivitySchedulerFactory, ActivitySchedulerFactory>()
.AddSingleton<IHasher, Hasher>()
.AddSingleton<IBookmarkHasher, BookmarkHasher>()
.AddSingleton<IIdentityGenerator, RandomIdentityGenerator>()
.AddSingleton<IBookmarkPayloadSerializer, BookmarkPayloadSerializer>()
.AddTransient<WorkflowBuilder>()

View file

@ -0,0 +1,29 @@
using Elsa.Workflows.Core.Services;
namespace Elsa.Workflows.Core.Implementations;
public class BookmarkHasher : IBookmarkHasher
{
private readonly IBookmarkPayloadSerializer _bookmarkPayloadSerializer;
private readonly IHasher _hasher;
public BookmarkHasher(IBookmarkPayloadSerializer bookmarkPayloadSerializer, IHasher hasher)
{
_bookmarkPayloadSerializer = bookmarkPayloadSerializer;
_hasher = hasher;
}
public string Hash(string activityTypeName, object? payload)
{
var json = payload != null ? _bookmarkPayloadSerializer.Serialize(payload) : null;
return Hash(activityTypeName, json);
}
public string Hash(string activityTypeName, string? serializedPayload)
{
var input = $"{activityTypeName}{(!string.IsNullOrWhiteSpace(serializedPayload) ? ":" + serializedPayload : "")}";
var hash = _hasher.Hash(input);
return hash;
}
}

View file

@ -113,11 +113,11 @@ public class ActivityExecutionContext
public Bookmark CreateBookmark(object? payload = default, ExecuteActivityDelegate? callback = default)
{
var hasher = GetRequiredService<IHasher>();
var bookmarkHasher = GetRequiredService<IBookmarkHasher>();
var identityGenerator = GetRequiredService<IIdentityGenerator>();
var payloadSerializer = GetRequiredService<IBookmarkPayloadSerializer>();
var payloadJson = payload != null ? payloadSerializer.Serialize(payload) : default;
var hash = payloadJson != null ? hasher.Hash(payloadJson) : default;
var hash = bookmarkHasher.Hash(Activity.Type, payloadJson);
var bookmark = new Bookmark(
identityGenerator.GenerateId(),

View file

@ -1,10 +1,10 @@
namespace Elsa.Workflows.Core.Models;
public class EventBookmarkData
public class EventBookmarkPayload
{
private readonly string _eventName = default!;
public EventBookmarkData(string eventName)
public EventBookmarkPayload(string eventName)
{
EventName = eventName;
}

View file

@ -0,0 +1,7 @@
namespace Elsa.Workflows.Core.Services;
public interface IBookmarkHasher
{
string Hash(string activityTypeName, object? payload);
string Hash(string activityTypeName, string? serializedPayload);
}

View file

@ -2,8 +2,8 @@ using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Comparers;
public class WorkflowTriggerHashEqualityComparer : IEqualityComparer<WorkflowTrigger>
public class WorkflowTriggerHashEqualityComparer : IEqualityComparer<StoredTrigger>
{
public bool Equals(WorkflowTrigger? x, WorkflowTrigger? y) => x?.Hash?.Equals(y?.Hash) ?? false;
public int GetHashCode(WorkflowTrigger obj) => obj.Hash?.GetHashCode() ?? "".GetHashCode();
public bool Equals(StoredTrigger? x, StoredTrigger? y) => x?.Hash?.Equals(y?.Hash) ?? false;
public int GetHashCode(StoredTrigger obj) => obj.Hash?.GetHashCode() ?? "".GetHashCode();
}

View file

@ -5,7 +5,7 @@ namespace Elsa.Workflows.Runtime.Entities;
/// <summary>
/// Represents a trigger associated with a workflow definition.
/// </summary>
public class WorkflowTrigger : Entity
public class StoredTrigger : Entity
{
public string WorkflowDefinitionId { get; set; } = default!;
public string Name { get; set; } = default!;

View file

@ -6,7 +6,7 @@ namespace Elsa.Workflows.Runtime.Extensions;
public static class WorkflowTriggerExtensions
{
public static IEnumerable<WorkflowTrigger> Filter<T>(this IEnumerable<WorkflowTrigger> triggers) where T : ITrigger
public static IEnumerable<StoredTrigger> Filter<T>(this IEnumerable<StoredTrigger> triggers) where T : ITrigger
{
var triggerName = ActivityTypeNameHelper.GenerateTypeName<T>();
return triggers.Where(x => x.Name == triggerName);

View file

@ -50,7 +50,7 @@ public class WorkflowRuntimeFeature : FeatureBase
/// </summary>
public Func<IServiceProvider, IBookmarkStore> BookmarkStore { get; set; } = sp => sp.GetRequiredService<MemoryBookmarkStore>();
public Func<IServiceProvider, IWorkflowTriggerStore> WorkflowTriggerStore { get; set; } = sp => sp.GetRequiredService<MemoryWorkflowTriggerStore>();
public Func<IServiceProvider, ITriggerStore> WorkflowTriggerStore { get; set; } = sp => sp.GetRequiredService<MemoryTriggerStore>();
public Func<IServiceProvider, IWorkflowExecutionLogStore> WorkflowExecutionLogStore { get; set; } = sp => sp.GetRequiredService<MemoryWorkflowExecutionLogStore>();
public Func<IServiceProvider, IWorkflowStateExporter> WorkflowStateExporter { get; set; } =
@ -75,6 +75,7 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddSingleton<ITriggerIndexer, TriggerIndexer>()
.AddSingleton<IWorkflowInstanceFactory, WorkflowInstanceFactory>()
.AddSingleton<IWorkflowDefinitionService, WorkflowDefinitionService>()
.AddSingleton<IWorkflowHostFactory, WorkflowHostFactory>()
.AddSingleton(WorkflowRuntime)
.AddSingleton(WorkflowDispatcher)
.AddSingleton(WorkflowStateStore)
@ -85,7 +86,7 @@ public class WorkflowRuntimeFeature : FeatureBase
// Memory Stores
.AddMemoryStore<WorkflowState, MemoryWorkflowStateStore>()
.AddMemoryStore<StoredBookmark, MemoryBookmarkStore>()
.AddMemoryStore<WorkflowTrigger, MemoryWorkflowTriggerStore>()
.AddMemoryStore<StoredTrigger, MemoryTriggerStore>()
.AddMemoryStore<WorkflowExecutionLogRecord, MemoryWorkflowExecutionLogStore>()
// Workflow definition providers.

View file

@ -13,54 +13,57 @@ namespace Elsa.Workflows.Runtime.Implementations;
public class DefaultWorkflowRuntime : IWorkflowRuntime
{
private readonly IWorkflowRunner _workflowRunner;
private readonly IWorkflowHostFactory _workflowHostFactory;
private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly IWorkflowStateStore _workflowStateStore;
private readonly ITriggerStore _triggerStore;
private readonly IBookmarkStore _bookmarkStore;
private readonly IHasher _hasher;
private readonly IBookmarkHasher _hasher;
private readonly IEventPublisher _eventPublisher;
public DefaultWorkflowRuntime(
IWorkflowRunner workflowRunner,
IWorkflowHostFactory workflowHostFactory,
IWorkflowDefinitionService workflowDefinitionService,
IWorkflowStateStore workflowStateStore,
ITriggerStore triggerStore,
IBookmarkStore bookmarkStore,
IHasher hasher,
IBookmarkHasher hasher,
IEventPublisher eventPublisher)
{
_workflowRunner = workflowRunner;
_workflowHostFactory = workflowHostFactory;
_workflowDefinitionService = workflowDefinitionService;
_workflowStateStore = workflowStateStore;
_triggerStore = triggerStore;
_bookmarkStore = bookmarkStore;
_hasher = hasher;
_eventPublisher = eventPublisher;
}
public async Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowOptions options, CancellationToken cancellationToken = default)
public async Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowOptions options,
CancellationToken cancellationToken = default)
{
var input = options.Input;
var versionOptions = options.VersionOptions;
var workflowDefinition = await _workflowDefinitionService.FindAsync(definitionId, versionOptions, cancellationToken);
var workflowDefinition =
await _workflowDefinitionService.FindAsync(definitionId, versionOptions, cancellationToken);
if (workflowDefinition == null)
throw new Exception("Specified workflow definition and version does not exist");
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
var workflowHost = await _workflowHostFactory.CreateAsync(workflow, cancellationToken);
await _eventPublisher.PublishAsync(new WorkflowExecuting(workflow), cancellationToken);
var workflowResult = await _workflowRunner.RunAsync(workflow, input, cancellationToken);
var workflowState = workflowResult.WorkflowState;
var finished = workflowResult.WorkflowState.Status == WorkflowStatus.Finished;
await workflowHost.StartWorkflowAsync(input, cancellationToken);
var workflowState = workflowHost.WorkflowState;
await SaveWorkflowStateAsync(workflowState, cancellationToken);
await UpdateBookmarksAsync(workflowState, new List<Bookmark>(), workflowResult.WorkflowState.Bookmarks, cancellationToken);
await _eventPublisher.PublishAsync(new WorkflowExecuted(workflow, workflowState), cancellationToken);
await UpdateBookmarksAsync(workflowState, new List<Bookmark>(), workflowState.Bookmarks, cancellationToken);
return new StartWorkflowResult(workflowState.Id);
return new StartWorkflowResult(workflowState.Id, workflowHost.WorkflowState.Bookmarks);
}
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string instanceId, string bookmarkId, ResumeWorkflowOptions options, CancellationToken cancellationToken = default)
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string instanceId, string bookmarkId,
ResumeWorkflowOptions options, CancellationToken cancellationToken = default)
{
var workflowState = await _workflowStateStore.LoadAsync(instanceId, cancellationToken);
@ -69,58 +72,109 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
var definitionId = workflowState.DefinitionId;
var version = workflowState.DefinitionVersion;
var workflowDefinition = await _workflowDefinitionService.FindAsync(definitionId, VersionOptions.SpecificVersion(version), cancellationToken);
var workflowDefinition = await _workflowDefinitionService.FindAsync(
definitionId,
VersionOptions.SpecificVersion(version),
cancellationToken);
if (workflowDefinition == null)
throw new Exception("Specified workflow definition and version does not exist");
var input = options.Input;
var existingStoredBookmarks = await _bookmarkStore.LoadAsync(workflowState.Id, cancellationToken).AsTask().ToList();
var existingStoredBookmarks = await _bookmarkStore.FindByWorkflowInstanceAsync(
workflowState.Id,
cancellationToken)
.AsTask()
.ToList();
var existingBookmarks = existingStoredBookmarks
.Select(storedBookmark => workflowState.Bookmarks.FirstOrDefault(bookmark => bookmark.Id == storedBookmark.BookmarkId))
.Select(storedBookmark =>
workflowState.Bookmarks.FirstOrDefault(bookmark => bookmark.Id == storedBookmark.BookmarkId))
.Where(x => x != null)
.Select(x => x!)
.ToList();
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken);
await _eventPublisher.PublishAsync(new WorkflowExecuting(workflow), cancellationToken);
var workflowResult = await _workflowRunner.RunAsync(workflow, workflowState, bookmarkId, input, cancellationToken);
var finished = workflowResult.WorkflowState.Status == WorkflowStatus.Finished;
var workflowHost = await _workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken);
await workflowHost.ResumeWorkflowAsync(bookmarkId, input, cancellationToken);
workflowState = workflowHost.WorkflowState;
await SaveWorkflowStateAsync(workflowState, cancellationToken);
await UpdateBookmarksAsync(workflowState, existingBookmarks, workflowResult.WorkflowState.Bookmarks, cancellationToken);
await _eventPublisher.PublishAsync(new WorkflowExecuted(workflow, workflowState), cancellationToken);
await UpdateBookmarksAsync(workflowState, existingBookmarks, workflowState.Bookmarks, cancellationToken);
return new ResumeWorkflowResult();
return new ResumeWorkflowResult(workflowState.Bookmarks);
}
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions options, CancellationToken cancellationToken = default)
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(
string activityTypeName,
object bookmarkPayload,
TriggerWorkflowsOptions options,
CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(bookmarkPayload);
var bookmarks = await _bookmarkStore.LoadAsync(activityTypeName, hash, cancellationToken);
var triggeredWorkflows = new List<TriggeredWorkflow>();
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
// Start new workflows.
var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
foreach (var trigger in triggers)
{
var startResult = await StartWorkflowAsync(
trigger.WorkflowDefinitionId,
new StartWorkflowOptions(options.CorrelationId, options.Input, VersionOptions.Published),
cancellationToken);
triggeredWorkflows.Add(new TriggeredWorkflow(startResult.InstanceId, startResult.Bookmarks));
}
// Resume bookmarks.
var bookmarks = await _bookmarkStore.FindByHashAsync(hash, cancellationToken);
foreach (var bookmark in bookmarks)
{
var workflowInstanceId = bookmark.WorkflowInstanceId;
var resumeResult = await ResumeWorkflowAsync(workflowInstanceId, bookmark.BookmarkId, new ResumeWorkflowOptions(options.Input), cancellationToken);
var resumeResult = await ResumeWorkflowAsync(
workflowInstanceId,
bookmark.BookmarkId,
new ResumeWorkflowOptions(options.Input),
cancellationToken);
triggeredWorkflows.Add(new TriggeredWorkflow(workflowInstanceId, resumeResult.Bookmarks));
}
return new TriggerWorkflowsResult();
return new TriggerWorkflowsResult(triggeredWorkflows);
}
private async Task SaveWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken) => await _workflowStateStore.SaveAsync(workflowState.Id, workflowState, cancellationToken);
public async Task<WorkflowState?> ExportWorkflowStateAsync(string instanceId,
CancellationToken cancellationToken = default)
{
return await _workflowStateStore.LoadAsync(instanceId, cancellationToken);
}
private async Task UpdateBookmarksAsync(WorkflowState workflowState, ICollection<Bookmark> previousBookmarks, ICollection<Bookmark> newBookmarks, CancellationToken cancellationToken)
public async Task ImportWorkflowStateAsync(WorkflowState workflowState,
CancellationToken cancellationToken = default)
{
await _workflowStateStore.SaveAsync(workflowState.Id, workflowState, cancellationToken);
}
private async Task SaveWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken) =>
await _workflowStateStore.SaveAsync(workflowState.Id, workflowState, cancellationToken);
private async Task UpdateBookmarksAsync(WorkflowState workflowState, ICollection<Bookmark> previousBookmarks,
ICollection<Bookmark> newBookmarks, CancellationToken cancellationToken)
{
await RemoveBookmarksAsync(workflowState.Id, previousBookmarks, cancellationToken);
await StoreBookmarksAsync(workflowState.Id, newBookmarks, cancellationToken);
await PublishChangedBookmarksAsync(workflowState, previousBookmarks, newBookmarks, cancellationToken);
}
private async Task StoreBookmarksAsync(string workflowInstanceId, ICollection<Bookmark> bookmarks, CancellationToken cancellationToken)
private async Task StoreBookmarksAsync(string workflowInstanceId, ICollection<Bookmark> bookmarks,
CancellationToken cancellationToken)
{
var groupedBookmarks = bookmarks.GroupBy(x => (x.Name, x.Hash));
@ -132,22 +186,29 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
}
}
private async Task RemoveBookmarksAsync(string workflowInstanceId, IEnumerable<Bookmark> bookmarks, CancellationToken cancellationToken)
private async Task RemoveBookmarksAsync(
string workflowInstanceId,
IEnumerable<Bookmark> bookmarks,
CancellationToken cancellationToken)
{
var groupedBookmarks = bookmarks.GroupBy(x => (x.Name, x.Hash));
var groupedBookmarks = bookmarks.GroupBy(x => x.Hash);
foreach (var groupedBookmark in groupedBookmarks)
{
var key = groupedBookmark.Key;
await _bookmarkStore.DeleteAsync(key.Name, key.Hash, workflowInstanceId, cancellationToken);
await _bookmarkStore.DeleteAsync(key, workflowInstanceId, cancellationToken);
}
}
private async Task PublishChangedBookmarksAsync(WorkflowState workflowState, ICollection<Bookmark> originalBookmarks, ICollection<Bookmark> updatedBookmarks, CancellationToken cancellationToken)
private async Task PublishChangedBookmarksAsync(WorkflowState workflowState,
ICollection<Bookmark> originalBookmarks, ICollection<Bookmark> updatedBookmarks,
CancellationToken cancellationToken)
{
var diff = Diff.For(originalBookmarks, updatedBookmarks);
var removedBookmarks = diff.Removed;
var createdBookmarks = diff.Added;
await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(new IndexedWorkflowBookmarks(workflowState, createdBookmarks, removedBookmarks)), cancellationToken);
await _eventPublisher.PublishAsync(
new WorkflowBookmarksIndexed(
new IndexedWorkflowBookmarks(workflowState, createdBookmarks, removedBookmarks)), cancellationToken);
}
}

View file

@ -6,30 +6,30 @@ namespace Elsa.Workflows.Runtime.Implementations;
public class MemoryBookmarkStore : IBookmarkStore
{
private readonly ConcurrentDictionary<(string ActivityTypeName, string Hash), ICollection<StoredBookmark>> _bookmarks = new();
private readonly ConcurrentDictionary<string, ICollection<StoredBookmark>> _bookmarks = new();
public ValueTask SaveAsync(string activityTypeName, string hash, string workflowInstanceId, IEnumerable<string> bookmarkIds, CancellationToken cancellationToken = default)
{
var storedBookmarks = bookmarkIds.Select(x => new StoredBookmark(activityTypeName, hash, workflowInstanceId, x)).ToList();
_bookmarks.AddOrUpdate((activityTypeName, hash), new List<StoredBookmark>(storedBookmarks), (s, bookmarks) => bookmarks.Concat(storedBookmarks).ToList());
_bookmarks.AddOrUpdate(hash, new List<StoredBookmark>(storedBookmarks), (s, bookmarks) => bookmarks.Concat(storedBookmarks).ToList());
return ValueTask.CompletedTask;
}
public ValueTask<IEnumerable<StoredBookmark>> LoadAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
public ValueTask<IEnumerable<StoredBookmark>> FindByWorkflowInstanceAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
{
var bookmarks = _bookmarks.Values.SelectMany(x => x).Where(x => x.WorkflowInstanceId == workflowInstanceId).ToList();
return new(bookmarks);
}
public ValueTask<IEnumerable<StoredBookmark>> LoadAsync(string activityTypeName, string hash, CancellationToken cancellationToken = default)
public ValueTask<IEnumerable<StoredBookmark>> FindByHashAsync(string hash, CancellationToken cancellationToken = default)
{
var bookmarks = _bookmarks.TryGetValue((activityTypeName, hash), out var value) ? value : Enumerable.Empty<StoredBookmark>();
var bookmarks = _bookmarks.TryGetValue(hash, out var value) ? value : Enumerable.Empty<StoredBookmark>();
return new(bookmarks);
}
public ValueTask DeleteAsync(string activityTypeName, string hash, string workflowInstanceId, CancellationToken cancellationToken = default)
public ValueTask DeleteAsync(string hash, string workflowInstanceId, CancellationToken cancellationToken = default)
{
var key = (activityTypeName, hash);
var key = hash;
if(!_bookmarks.TryGetValue(key, out var bookmarks))
return ValueTask.CompletedTask;

View file

@ -0,0 +1,56 @@
using Elsa.Common.Implementations;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Workflows.Runtime.Implementations;
public class MemoryTriggerStore : ITriggerStore
{
private readonly MemoryStore<StoredTrigger> _store;
public MemoryTriggerStore(MemoryStore<StoredTrigger> store)
{
_store = store;
}
public Task SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default)
{
_store.Save(record, x => x.Id);
return Task.CompletedTask;
}
public Task SaveManyAsync(IEnumerable<StoredTrigger> records, CancellationToken cancellationToken = default)
{
_store.SaveMany(records, x => x.Id);
return Task.CompletedTask;
}
public Task<IEnumerable<StoredTrigger>> FindAsync(string hash, CancellationToken cancellationToken = default)
{
var triggers = _store.Query(query => query.Where(x => x.Hash == hash));
return Task.FromResult<IEnumerable<StoredTrigger>>(triggers.ToList());
}
public Task<IEnumerable<StoredTrigger>> FindManyByWorkflowDefinitionIdAsync(
string workflowDefinitionId,
CancellationToken cancellationToken = default)
{
var triggers = _store.Query(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId));
return Task.FromResult<IEnumerable<StoredTrigger>>(triggers.ToList());
}
public Task ReplaceAsync(IEnumerable<StoredTrigger> removed, IEnumerable<StoredTrigger> added,
CancellationToken cancellationToken = default)
{
_store.DeleteMany(removed, x => x.Id);
_store.SaveMany(added, x => x.Id);
return Task.CompletedTask;
}
public Task DeleteManyAsync(IEnumerable<string> ids, CancellationToken cancellationToken = default)
{
_store.DeleteMany(ids);
return Task.CompletedTask;
}
}

View file

@ -1,62 +0,0 @@
using Elsa.Common.Implementations;
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Workflows.Runtime.Implementations;
public class MemoryWorkflowTriggerStore : IWorkflowTriggerStore
{
private readonly MemoryStore<WorkflowTrigger> _store;
public MemoryWorkflowTriggerStore(MemoryStore<WorkflowTrigger> store)
{
_store = store;
}
public Task SaveAsync(WorkflowTrigger record, CancellationToken cancellationToken = default)
{
_store.Save(record, x => x.Id);
return Task.CompletedTask;
}
public Task SaveManyAsync(IEnumerable<WorkflowTrigger> records, CancellationToken cancellationToken = default)
{
_store.SaveMany(records, x => x.Id);
return Task.CompletedTask;
}
public Task<IEnumerable<WorkflowTrigger>> FindManyByNameAsync(string name, string? hash, CancellationToken cancellationToken = default)
{
var triggers = _store.Query(query =>
{
query = query.Where(x => x.Name == name);
if (hash != null)
query = query.Where(x => x.Hash == hash);
return query;
});
return Task.FromResult<IEnumerable<WorkflowTrigger>>(triggers.ToList());
}
public Task<IEnumerable<WorkflowTrigger>> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default)
{
var triggers = _store.Query(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId));
return Task.FromResult<IEnumerable<WorkflowTrigger>>(triggers.ToList());
}
public Task ReplaceAsync(IEnumerable<WorkflowTrigger> removed, IEnumerable<WorkflowTrigger> added, CancellationToken cancellationToken = default)
{
_store.DeleteMany(removed, x => x.Id);
_store.SaveMany(added, x => x.Id);
return Task.CompletedTask;
}
public Task DeleteManyAsync(IEnumerable<string> ids, CancellationToken cancellationToken = default)
{
_store.DeleteMany(ids);
return Task.CompletedTask;
}
}

View file

@ -24,10 +24,10 @@ public class TriggerIndexer : ITriggerIndexer
private readonly IWorkflowDefinitionService _workflowDefinitionService;
private readonly IExpressionEvaluator _expressionEvaluator;
private readonly IIdentityGenerator _identityGenerator;
private readonly IWorkflowTriggerStore _workflowTriggerStore;
private readonly ITriggerStore _triggerStore;
private readonly IEventPublisher _eventPublisher;
private readonly IServiceProvider _serviceProvider;
private readonly IHasher _hasher;
private readonly IBookmarkHasher _hasher;
private readonly ILogger _logger;
public TriggerIndexer(
@ -35,16 +35,16 @@ public class TriggerIndexer : ITriggerIndexer
IWorkflowDefinitionService workflowDefinitionService,
IExpressionEvaluator expressionEvaluator,
IIdentityGenerator identityGenerator,
IWorkflowTriggerStore workflowTriggerStore,
ITriggerStore triggerStore,
IEventPublisher eventPublisher,
IServiceProvider serviceProvider,
IHasher hasher,
IBookmarkHasher hasher,
ILogger<TriggerIndexer> logger)
{
_activityWalker = activityWalker;
_expressionEvaluator = expressionEvaluator;
_identityGenerator = identityGenerator;
_workflowTriggerStore = workflowTriggerStore;
_triggerStore = triggerStore;
_eventPublisher = eventPublisher;
_serviceProvider = serviceProvider;
_hasher = hasher;
@ -66,13 +66,13 @@ public class TriggerIndexer : ITriggerIndexer
// Collect new triggers **if workflow is published**.
var newTriggers = workflow.Publication.IsPublished
? await GetTriggersAsync(workflow, cancellationToken).ToListAsync(cancellationToken)
: new List<WorkflowTrigger>(0);
: new List<StoredTrigger>(0);
// Diff triggers.
var diff = Diff.For(currentTriggers, newTriggers, new WorkflowTriggerHashEqualityComparer());
// Replace triggers for the specified workflow.
await _workflowTriggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed, diff.Unchanged);
@ -81,10 +81,10 @@ public class TriggerIndexer : ITriggerIndexer
return indexedWorkflow;
}
private async Task<IEnumerable<WorkflowTrigger>> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) =>
await _workflowTriggerStore.FindManyByWorkflowDefinitionIdAsync(workflowDefinitionId, cancellationToken);
private async Task<IEnumerable<StoredTrigger>> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) =>
await _triggerStore.FindManyByWorkflowDefinitionIdAsync(workflowDefinitionId, cancellationToken);
private async IAsyncEnumerable<WorkflowTrigger> GetTriggersAsync(Workflow workflow, [EnumeratorCancellation] CancellationToken cancellationToken = default)
private async IAsyncEnumerable<StoredTrigger> GetTriggersAsync(Workflow workflow, [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
var context = new WorkflowIndexingContext(workflow, cancellationToken);
var nodes = await _activityWalker.WalkAsync(workflow.Root, cancellationToken);
@ -105,7 +105,7 @@ public class TriggerIndexer : ITriggerIndexer
}
}
private async Task<IEnumerable<WorkflowTrigger>> GetTriggersAsync(WorkflowIndexingContext context, IActivity activity)
private async Task<IEnumerable<StoredTrigger>> GetTriggersAsync(WorkflowIndexingContext context, IActivity activity)
{
// If the activity implements ITrigger, request its trigger data. Otherwise, create one trigger datum.
if (activity is ITrigger trigger)
@ -117,10 +117,10 @@ public class TriggerIndexer : ITriggerIndexer
return new[] { simpleTrigger };
}
private WorkflowTrigger CreateWorkflowTrigger(WorkflowIndexingContext context, IActivity activity)
private StoredTrigger CreateWorkflowTrigger(WorkflowIndexingContext context, IActivity activity)
{
var workflow = context.Workflow;
return new WorkflowTrigger
return new StoredTrigger
{
Id = _identityGenerator.GenerateId(),
WorkflowDefinitionId = workflow.Identity.DefinitionId,
@ -128,7 +128,7 @@ public class TriggerIndexer : ITriggerIndexer
};
}
private async Task<ICollection<WorkflowTrigger>> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger)
private async Task<ICollection<StoredTrigger>> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger)
{
var workflow = context.Workflow;
var cancellationToken = context.CancellationToken;
@ -138,12 +138,12 @@ public class TriggerIndexer : ITriggerIndexer
var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext);
var triggerTypeName = trigger.Type;
var triggers = triggerData.Select(x => new WorkflowTrigger
var triggers = triggerData.Select(x => new StoredTrigger
{
Id = _identityGenerator.GenerateId(),
WorkflowDefinitionId = workflow.Identity.DefinitionId,
Name = triggerTypeName,
Hash = _hasher.Hash(x),
Hash = _hasher.Hash(triggerTypeName, x),
Data = JsonSerializer.Serialize(x)
});

View file

@ -0,0 +1,77 @@
using Elsa.Mediator.Services;
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Services;
using Elsa.Workflows.Core.State;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Notifications;
using Elsa.Workflows.Runtime.Services;
namespace Elsa.Workflows.Runtime.Implementations;
public class WorkflowHost : IWorkflowHost
{
private readonly IWorkflowRunner _workflowRunner;
private readonly IEventPublisher _eventPublisher;
private readonly IIdentityGenerator _identityGenerator;
public WorkflowHost(Workflow workflow, WorkflowState workflowState, IWorkflowRunner workflowRunner, IEventPublisher eventPublisher, IIdentityGenerator identityGenerator)
{
Workflow = workflow;
WorkflowState = workflowState;
_workflowRunner = workflowRunner;
_eventPublisher = eventPublisher;
_identityGenerator = identityGenerator;
}
public Workflow Workflow { get; set; }
public WorkflowState WorkflowState { get; set; }
public async Task<StartWorkflowHostResult> StartWorkflowAsync(IDictionary<string, object>? input = default, CancellationToken cancellationToken = default)
{
var instanceId = _identityGenerator.GenerateId();
return await StartWorkflowAsync(instanceId, input, cancellationToken);
}
public async Task<StartWorkflowHostResult> StartWorkflowAsync(string instanceId, IDictionary<string, object>? input = default, CancellationToken cancellationToken = default)
{
await _eventPublisher.PublishAsync(new WorkflowExecuting(Workflow), cancellationToken);
var originalBookmarks = WorkflowState.Bookmarks.ToList();
var workflowResult = await _workflowRunner.RunAsync(instanceId, Workflow, input, cancellationToken);
WorkflowState = workflowResult.WorkflowState;
var updatedBookmarks = WorkflowState.Bookmarks;
var diff = Diff.For(originalBookmarks, updatedBookmarks);
await PublishChangedBookmarksAsync(diff, cancellationToken);
await _eventPublisher.PublishAsync(new WorkflowExecuted(Workflow, WorkflowState), cancellationToken);
return new StartWorkflowHostResult(diff);
}
public async Task<ResumeWorkflowHostResult> ResumeWorkflowAsync(string bookmarkId, IDictionary<string, object>? input = default, CancellationToken cancellationToken = default)
{
await _eventPublisher.PublishAsync(new WorkflowExecuting(Workflow), cancellationToken);
var originalBookmarks = WorkflowState.Bookmarks.ToList();
var workflowResult = await _workflowRunner.RunAsync(Workflow, WorkflowState, bookmarkId, input, cancellationToken);
WorkflowState = workflowResult.WorkflowState;
var updatedBookmarks = WorkflowState.Bookmarks;
var diff = Diff.For(originalBookmarks, updatedBookmarks);
await PublishChangedBookmarksAsync(diff, cancellationToken);
await _eventPublisher.PublishAsync(new WorkflowExecuted(Workflow, WorkflowState), cancellationToken);
return new ResumeWorkflowHostResult(diff);
}
private async Task PublishChangedBookmarksAsync(Diff<Bookmark> diff, CancellationToken cancellationToken)
{
var removedBookmarks = diff.Removed;
var createdBookmarks = diff.Added;
await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(new IndexedWorkflowBookmarks(WorkflowState, createdBookmarks, removedBookmarks)), cancellationToken);
}
}

View file

@ -0,0 +1,33 @@
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.State;
using Elsa.Workflows.Runtime.Services;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.Runtime.Implementations;
public class WorkflowHostFactory : IWorkflowHostFactory
{
private readonly IServiceProvider _serviceProvider;
public WorkflowHostFactory(IServiceProvider serviceProvider)
{
_serviceProvider = serviceProvider;
}
public Task<IWorkflowHost> CreateAsync(Workflow workflow, WorkflowState workflowState, CancellationToken cancellationToken = default)
{
var workflowHost = (IWorkflowHost)ActivatorUtilities.CreateInstance<WorkflowHost>(_serviceProvider, workflow, workflowState);
return Task.FromResult(workflowHost);
}
public Task<IWorkflowHost> CreateAsync(Workflow workflow, CancellationToken cancellationToken = default)
{
var workflowState = new WorkflowState
{
DefinitionId = workflow.Identity.DefinitionId,
DefinitionVersion = workflow.Identity.Version
};
return CreateAsync(workflow, workflowState, cancellationToken);
}
}

View file

@ -3,4 +3,4 @@ using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Models;
public record IndexedWorkflowTriggers(Workflow Workflow, ICollection<WorkflowTrigger> AddedTriggers, ICollection<WorkflowTrigger> RemovedTriggers, ICollection<WorkflowTrigger> UnchangedTriggers);
public record IndexedWorkflowTriggers(Workflow Workflow, ICollection<StoredTrigger> AddedTriggers, ICollection<StoredTrigger> RemovedTriggers, ICollection<StoredTrigger> UnchangedTriggers);

View file

@ -5,7 +5,7 @@ namespace Elsa.Workflows.Runtime.Services;
public interface IBookmarkStore
{
ValueTask SaveAsync(string activityTypeName, string hash, string workflowInstanceId, IEnumerable<string> bookmarkIds, CancellationToken cancellationToken = default);
ValueTask<IEnumerable<StoredBookmark>> LoadAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
ValueTask<IEnumerable<StoredBookmark>> LoadAsync(string activityTypeName, string hash, CancellationToken cancellationToken = default);
ValueTask DeleteAsync(string activityTypeName, string hash, string workflowInstanceId, CancellationToken cancellationToken = default);
ValueTask<IEnumerable<StoredBookmark>> FindByWorkflowInstanceAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
ValueTask<IEnumerable<StoredBookmark>> FindByHashAsync(string hash, CancellationToken cancellationToken = default);
ValueTask DeleteAsync(string hash, string workflowInstanceId, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,13 @@
using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Services;
public interface ITriggerStore
{
Task SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default);
Task SaveManyAsync(IEnumerable<StoredTrigger> records, CancellationToken cancellationToken = default);
Task<IEnumerable<StoredTrigger>> FindAsync(string hash, CancellationToken cancellationToken = default);
Task<IEnumerable<StoredTrigger>> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default);
Task ReplaceAsync(IEnumerable<StoredTrigger> removed, IEnumerable<StoredTrigger> added, CancellationToken cancellationToken = default);
Task DeleteManyAsync(IEnumerable<string> ids, CancellationToken cancellationToken = default);
}

View file

@ -0,0 +1,21 @@
using Elsa.Workflows.Core.Helpers;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.State;
namespace Elsa.Workflows.Runtime.Services;
/// <summary>
/// Represents a single workflow instance that can be executed and takes care of publishing various lifecycle events.
/// </summary>
public interface IWorkflowHost
{
Workflow Workflow { get; set; }
WorkflowState WorkflowState { get; set; }
Task<StartWorkflowHostResult> StartWorkflowAsync(IDictionary<string, object>? input = default, CancellationToken cancellationToken = default);
Task<StartWorkflowHostResult> StartWorkflowAsync(string instanceId, IDictionary<string, object>? input = default, CancellationToken cancellationToken = default);
Task<ResumeWorkflowHostResult> ResumeWorkflowAsync(string bookmarkId, IDictionary<string, object>? input = default, CancellationToken cancellationToken = default);
}
public record StartWorkflowHostResult(Diff<Bookmark> BookmarksDiff);
public record ResumeWorkflowHostResult(Diff<Bookmark> BookmarksDiff);

View file

@ -0,0 +1,19 @@
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.State;
namespace Elsa.Workflows.Runtime.Services;
/// <summary>
/// Creates <see cref="IWorkflowHost"/> objects.
/// </summary>
public interface IWorkflowHostFactory
{
Task<IWorkflowHost> CreateAsync(
Workflow workflow,
WorkflowState workflowState,
CancellationToken cancellationToken = default);
Task<IWorkflowHost> CreateAsync(
Workflow workflow,
CancellationToken cancellationToken = default);
}

View file

@ -1,4 +1,6 @@
using Elsa.Common.Models;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.State;
namespace Elsa.Workflows.Runtime.Services;
@ -7,11 +9,14 @@ public interface IWorkflowRuntime
Task<StartWorkflowResult> StartWorkflowAsync(string definitionId, StartWorkflowOptions options, CancellationToken cancellationToken = default);
Task<ResumeWorkflowResult> ResumeWorkflowAsync(string instanceId, string bookmarkId, ResumeWorkflowOptions options, CancellationToken cancellationToken = default);
Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions options, CancellationToken cancellationToken = default);
Task<WorkflowState?> ExportWorkflowStateAsync(string instanceId, CancellationToken cancellationToken = default);
Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default);
}
public record StartWorkflowOptions(string? CorrelationId = default, IDictionary<string, object>? Input = default, VersionOptions VersionOptions = default);
public record ResumeWorkflowOptions(IDictionary<string, object>? Input = default);
public record StartWorkflowResult(string InstanceId);
public record ResumeWorkflowResult;
public record TriggerWorkflowsOptions(IDictionary<string, object>? Input = default);
public record TriggerWorkflowsResult;
public record StartWorkflowResult(string InstanceId, ICollection<Bookmark> Bookmarks);
public record ResumeWorkflowResult(ICollection<Bookmark> Bookmarks);
public record TriggerWorkflowsOptions(string? CorrelationId = default, IDictionary<string, object>? Input = default);
public record TriggerWorkflowsResult(ICollection<TriggeredWorkflow> TriggeredWorkflows);
public record TriggeredWorkflow(string InstanceId, ICollection<Bookmark> Bookmarks);

View file

@ -1,13 +0,0 @@
using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Services;
public interface IWorkflowTriggerStore
{
Task SaveAsync(WorkflowTrigger record, CancellationToken cancellationToken = default);
Task SaveManyAsync(IEnumerable<WorkflowTrigger> records, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowTrigger>> FindManyByNameAsync(string name, string? hash, CancellationToken cancellationToken = default);
Task<IEnumerable<WorkflowTrigger>> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default);
Task ReplaceAsync(IEnumerable<WorkflowTrigger> removed, IEnumerable<WorkflowTrigger> added, CancellationToken cancellationToken = default);
Task DeleteManyAsync(IEnumerable<string> ids, CancellationToken cancellationToken = default);
}