diff --git a/Elsa.sln.DotSettings b/Elsa.sln.DotSettings index c882dee0a..21965872e 100644 --- a/Elsa.sln.DotSettings +++ b/Elsa.sln.DotSettings @@ -2,6 +2,7 @@ True True True + True True True True diff --git a/src/bundles/Elsa.WorkflowServer.Web/Program.cs b/src/bundles/Elsa.WorkflowServer.Web/Program.cs index c187a5f36..f36d9d744 100644 --- a/src/bundles/Elsa.WorkflowServer.Web/Program.cs +++ b/src/bundles/Elsa.WorkflowServer.Web/Program.cs @@ -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(); }) diff --git a/src/designer/elsa-workflows-designer/src/modules/flowchart/flowchart.tsx b/src/designer/elsa-workflows-designer/src/modules/flowchart/flowchart.tsx index 949a73fbe..4f356c161 100644 --- a/src/designer/elsa-workflows-designer/src/modules/flowchart/flowchart.tsx +++ b/src/designer/elsa-workflows-designer/src/modules/flowchart/flowchart.tsx @@ -498,6 +498,7 @@ export class FlowchartComponent implements ContainerActivityComponent { }; private onNodeContextMenu = async (e: PositionEventArgs) => { + debugger; const node = e.node as ActivityNodeShape; const activity = e.node.data as Activity; diff --git a/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs b/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs index af901425b..ce2fe5c8b 100644 --- a/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs +++ b/src/modules/Elsa.Expressions/Helpers/ObjectConverter.cs @@ -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); } } diff --git a/src/modules/Elsa.Http/Activities/HttpEndpoint.cs b/src/modules/Elsa.Http/Activities/HttpEndpoint.cs index f2eef3b1b..27ada6c23 100644 --- a/src/modules/Elsa.Http/Activities/HttpEndpoint.cs +++ b/src/modules/Elsa.Http/Activities/HttpEndpoint.cs @@ -62,6 +62,6 @@ public class HttpEndpoint : Trigger // 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().ToArray(); + return methods!.Select(x => new HttpEndpointBookmarkPayload(path!, x.ToLowerInvariant())).Cast().ToArray(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Http/Activities/WriteHttpResponse.cs b/src/modules/Elsa.Http/Activities/WriteHttpResponse.cs index 7a76135b7..a43cf622f 100644 --- a/src/modules/Elsa.Http/Activities/WriteHttpResponse.cs +++ b/src/modules/Elsa.Http/Activities/WriteHttpResponse.cs @@ -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(); 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(); + 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); diff --git a/src/modules/Elsa.Http/Extensions/RouteTableExtensions.cs b/src/modules/Elsa.Http/Extensions/RouteTableExtensions.cs index c5bf28236..b950489bc 100644 --- a/src/modules/Elsa.Http/Extensions/RouteTableExtensions.cs +++ b/src/modules/Elsa.Http/Extensions/RouteTableExtensions.cs @@ -12,7 +12,7 @@ namespace Elsa.Http.Extensions; public static class RouteTableExtensions { - public static void AddRoutes(this IRouteTable routeTable, IEnumerable triggers) + public static void AddRoutes(this IRouteTable routeTable, IEnumerable 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 triggers) + public static void RemoveRoutes(this IRouteTable routeTable, IEnumerable 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 Filter(IEnumerable triggers) => triggers.Where(x => x.Name == ActivityTypeNameHelper.GenerateTypeName()); + private static IEnumerable Filter(IEnumerable triggers) => triggers.Where(x => x.Name == ActivityTypeNameHelper.GenerateTypeName()); private static IEnumerable Filter(IEnumerable triggers) => triggers.Where(x => x.Name == ActivityTypeNameHelper.GenerateTypeName()); - 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(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(model)!; } \ No newline at end of file diff --git a/src/modules/Elsa.Http/Middleware/HttpTriggerMiddleware.cs b/src/modules/Elsa.Http/Middleware/HttpTriggerMiddleware.cs index 2469449b7..cc2b7b5f1 100644 --- a/src/modules/Elsa.Http/Middleware/HttpTriggerMiddleware.cs +++ b/src/modules/Elsa.Http/Middleware/HttpTriggerMiddleware.cs @@ -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(); - public HttpTriggerMiddleware(RequestDelegate next, IHasher hasher, IOptions options) + public HttpTriggerMiddleware( + RequestDelegate next, + IBookmarkHasher hasher, + IWorkflowRuntime workflowRuntime, + IWorkflowHostFactory workflowHostFactory, + IWorkflowDefinitionService workflowDefinitionService, + IOptions 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() { [HttpEndpoint.InputKey] = requestModel }; + var input = new Dictionary { [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(); + + 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) diff --git a/src/modules/Elsa.Http/Models/HttpEndpointBookmarkData.cs b/src/modules/Elsa.Http/Models/HttpEndpointBookmarkPayload.cs similarity index 78% rename from src/modules/Elsa.Http/Models/HttpEndpointBookmarkData.cs rename to src/modules/Elsa.Http/Models/HttpEndpointBookmarkPayload.cs index 968d42eab..f546aa859 100644 --- a/src/modules/Elsa.Http/Models/HttpEndpointBookmarkData.cs +++ b/src/modules/Elsa.Http/Models/HttpEndpointBookmarkPayload.cs @@ -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; diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/BookmarkStore.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/BookmarkStore.cs deleted file mode 100644 index e34c7ead4..000000000 --- a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/BookmarkStore.cs +++ /dev/null @@ -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 _store; - public EFCoreBookmarkStore(Store store) => _store = store; - - public async ValueTask SaveAsync(string activityTypeName, string hash, string workflowInstanceId, IEnumerable bookmarkIds, CancellationToken cancellationToken = default) - { - var storedBookmarks = bookmarkIds.Select(x => new StoredBookmark(activityTypeName, hash, workflowInstanceId, x)).ToList(); - await _store.SaveManyAsync(storedBookmarks, cancellationToken); - } - - public async ValueTask> LoadAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => - await _store.FindManyAsync(x => x.WorkflowInstanceId == workflowInstanceId, cancellationToken); - - public async ValueTask> 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); -} \ No newline at end of file diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Configurations.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Configurations.cs index 22fe80bc1..255cab388 100644 --- a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Configurations.cs +++ b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Configurations.cs @@ -11,7 +11,7 @@ namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime { public class Configurations : IEntityTypeConfiguration, - IEntityTypeConfiguration, + IEntityTypeConfiguration, IEntityTypeConfiguration, IEntityTypeConfiguration { @@ -38,11 +38,11 @@ namespace Elsa.Persistence.EntityFrameworkCore.Modules.Runtime builder.HasIndex("UpdatedAt").HasDatabaseName($"IX_{nameof(WorkflowState)}_UpdatedAt"); } - public void Configure(EntityTypeBuilder builder) + public void Configure(EntityTypeBuilder 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 builder) diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreBookmarkStore.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreBookmarkStore.cs new file mode 100644 index 000000000..7d8c05e86 --- /dev/null +++ b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreBookmarkStore.cs @@ -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 _store; + public EFCoreBookmarkStore(Store store) => _store = store; + + public async ValueTask SaveAsync( + string activityTypeName, + string hash, + string workflowInstanceId, + IEnumerable bookmarkIds, CancellationToken cancellationToken = default) + { + var storedBookmarks = bookmarkIds.Select(x => new StoredBookmark(activityTypeName, hash, workflowInstanceId, x)) + .ToList(); + await _store.SaveManyAsync(storedBookmarks, cancellationToken); + } + + public async ValueTask> FindByWorkflowInstanceAsync( + string workflowInstanceId, + CancellationToken cancellationToken = default) => + await _store.FindManyAsync(x => x.WorkflowInstanceId == workflowInstanceId, cancellationToken); + + public async ValueTask> 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); +} \ No newline at end of file diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Feature.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreRuntimePersistenceFeature.cs similarity index 92% rename from src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Feature.cs rename to src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreRuntimePersistenceFeature.cs index c97aa25bd..0d9160c5c 100644 --- a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/Feature.cs +++ b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreRuntimePersistenceFeature.cs @@ -21,7 +21,7 @@ public class EFCoreRuntimePersistenceFeature : PersistenceFeatureBase(feature => { feature.WorkflowStateStore = sp => sp.GetRequiredService(); - feature.WorkflowTriggerStore = sp => sp.GetRequiredService(); + feature.WorkflowTriggerStore = sp => sp.GetRequiredService(); feature.BookmarkStore = sp => sp.GetRequiredService(); feature.WorkflowExecutionLogStore = sp => sp.GetRequiredService(); }); @@ -32,7 +32,7 @@ public class EFCoreRuntimePersistenceFeature : PersistenceFeatureBase(); - AddStore(); + AddStore(); AddStore(); AddStore(); } diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreTriggerStore.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreTriggerStore.cs new file mode 100644 index 000000000..c5d451188 --- /dev/null +++ b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreTriggerStore.cs @@ -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 _store; + + public EFCoreTriggerStore(Store store) + { + _store = store; + } + + public async Task SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default) => await _store.SaveAsync(record, cancellationToken); + public async Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default) => await _store.SaveManyAsync(records, cancellationToken); + + public async Task> FindAsync( + string hash, + CancellationToken cancellationToken = default) => + await _store.QueryAsync(query => query.Where(x => x.Hash == hash), cancellationToken); + + public async Task> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) => + await _store.QueryAsync(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId), cancellationToken); + + public async Task ReplaceAsync(IEnumerable removed, IEnumerable added, CancellationToken cancellationToken = default) + { + await _store.DeleteManyAsync(removed, cancellationToken); + await _store.SaveManyAsync(added, cancellationToken); + } + + public async Task DeleteManyAsync(IEnumerable ids, CancellationToken cancellationToken = default) + { + var idList = ids.ToList(); + await _store.DeleteWhereAsync(x => idList.Contains(x.Id), cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreWorkflowExecutionLogStore.cs similarity index 100% rename from src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs rename to src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreWorkflowExecutionLogStore.cs diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowStateStore.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreWorkflowStateStore.cs similarity index 100% rename from src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowStateStore.cs rename to src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/EFCoreWorkflowStateStore.cs diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/DbContext.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/RuntimeDbContext.cs similarity index 91% rename from src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/DbContext.cs rename to src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/RuntimeDbContext.cs index 2ddb47e7d..f51297d00 100644 --- a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/DbContext.cs +++ b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/RuntimeDbContext.cs @@ -14,7 +14,7 @@ public class RuntimeDbContext : DbContextBase } public DbSet WorkflowStates { get; set; } = default!; - public DbSet WorkflowTriggers { get; set; } = default!; + public DbSet WorkflowTriggers { get; set; } = default!; public DbSet WorkflowExecutionLogRecords { get; set; } = default!; public DbSet Bookmarks { get; set; } = default!; @@ -22,7 +22,7 @@ public class RuntimeDbContext : DbContextBase { var config = new Configurations(); modelBuilder.ApplyConfiguration(config); - modelBuilder.ApplyConfiguration(config); + modelBuilder.ApplyConfiguration(config); modelBuilder.ApplyConfiguration(config); modelBuilder.ApplyConfiguration(config); } diff --git a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowTriggerStore.cs b/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowTriggerStore.cs deleted file mode 100644 index 3e917e779..000000000 --- a/src/modules/Elsa.Persistence.EntityFrameworkCore/Modules/Runtime/WorkflowTriggerStore.cs +++ /dev/null @@ -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 _store; - - public EFCoreWorkflowTriggerStore(Store store) - { - _store = store; - } - - public async Task SaveAsync(WorkflowTrigger record, CancellationToken cancellationToken = default) => await _store.SaveAsync(record, cancellationToken); - public async Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default) => await _store.SaveManyAsync(records, cancellationToken); - - public async Task> 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> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) => - await _store.QueryAsync(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId), cancellationToken); - - public async Task ReplaceAsync(IEnumerable removed, IEnumerable added, CancellationToken cancellationToken = default) - { - await _store.DeleteManyAsync(removed, cancellationToken); - await _store.SaveManyAsync(added, cancellationToken); - } - - public async Task DeleteManyAsync(IEnumerable ids, CancellationToken cancellationToken = default) - { - var idList = ids.ToList(); - await _store.DeleteWhereAsync(x => idList.Contains(x.Id), cancellationToken); - } -} \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Extensions/ProtoStringExtensions.cs b/src/modules/Elsa.ProtoActor/Extensions/ProtoStringExtensions.cs new file mode 100644 index 000000000..ccd3747a5 --- /dev/null +++ b/src/modules/Elsa.ProtoActor/Extensions/ProtoStringExtensions.cs @@ -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; +} \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Grains/WorkflowGrain.cs b/src/modules/Elsa.ProtoActor/Grains/WorkflowGrain.cs index 67dfd0884..e485cad8e 100644 --- a/src/modules/Elsa.ProtoActor/Grains/WorkflowGrain.cs +++ b/src/modules/Elsa.ProtoActor/Grains/WorkflowGrain.cs @@ -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; - /// /// Executes a workflow. /// 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? _input; - private Workflow _workflow = default!; + private IWorkflowHost _workflowHost = default!; private WorkflowState _workflowState = default!; - private ICollection _bookmarks = new List(); + //private ICollection _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 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 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 bookmarks, CancellationToken cancellationToken) + public override Task 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 ImportState(ImportWorkflowStateRequest request) + { + var options = _serializerOptionsProvider.CreatePersistenceOptions(); + var workflowState = JsonSerializer.Deserialize(request.SerializedWorkflowState.Text, options)!; + + _workflowState = workflowState; + _workflowHost.WorkflowState = workflowState; + + return Task.FromResult(new ImportWorkflowStateResponse()); + } + + private async Task UpdateBookmarksAsync(Diff bookmarksDiff, CancellationToken cancellationToken) + { + await RemoveBookmarksAsync(bookmarksDiff.Removed, cancellationToken); + await StoreBookmarksAsync(bookmarksDiff.Added, cancellationToken); } private async Task StoreBookmarksAsync(ICollection bookmarks, CancellationToken cancellationToken) @@ -175,15 +209,19 @@ public class WorkflowGrain : WorkflowGrainBase } } - private async Task PublishChangedBookmarksAsync(ICollection originalBookmarks, ICollection 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 Map(IEnumerable 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() + }); } \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs index d25576493..6ae8935fc 100644 --- a/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Implementations/ProtoActorWorkflowRuntime.cs @@ -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 StartWorkflowAsync(string definitionId, StartWorkflowOptions options, CancellationToken cancellationToken = default) + public async Task 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 ResumeWorkflowAsync(string instanceId, string bookmarkId, ResumeWorkflowOptions options, CancellationToken cancellationToken = default) + public async Task 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 TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions options, CancellationToken cancellationToken = default) + public async Task TriggerWorkflowsAsync( + string activityTypeName, + object bookmarkPayload, + TriggerWorkflowsOptions options, + CancellationToken cancellationToken = default) { - var hash = _hasher.Hash(bookmarkPayload); + var triggeredWorkflows = new List(); + 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 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(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 Map(IEnumerable bookmarkDtos) => + bookmarkDtos.Select(x => + new Bookmark( + x.Id, + x.Name, + x.Hash, + x.Data.NullIfEmpty(), + x.ActivityId, + x.ActivityInstanceId, + x.CallbackMethodName.NullIfEmpty())); } \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/ProtoActorSystem.cs b/src/modules/Elsa.ProtoActor/ProtoActorSystem.cs index 2998779a5..2bc5e0c38 100644 --- a/src/modules/Elsa.ProtoActor/ProtoActorSystem.cs +++ b/src/modules/Elsa.ProtoActor/ProtoActorSystem.cs @@ -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!; } \ No newline at end of file diff --git a/src/modules/Elsa.ProtoActor/Protos/Grains.proto b/src/modules/Elsa.ProtoActor/Protos/Grains.proto index f5252a7b8..3c6e84482 100644 --- a/src/modules/Elsa.ProtoActor/Protos/Grains.proto +++ b/src/modules/Elsa.ProtoActor/Protos/Grains.proto @@ -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 { diff --git a/src/modules/Elsa.ProtoActor/Protos/Messages.proto b/src/modules/Elsa.ProtoActor/Protos/Messages.proto index ab42f9eb7..ff8355ad8 100644 --- a/src/modules/Elsa.ProtoActor/Protos/Messages.proto +++ b/src/modules/Elsa.ProtoActor/Protos/Messages.proto @@ -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; } diff --git a/src/modules/Elsa.ProtoActor/Snapshots.cs b/src/modules/Elsa.ProtoActor/Snapshots.cs index 38e768acc..73c06b7a4 100644 --- a/src/modules/Elsa.ProtoActor/Snapshots.cs +++ b/src/modules/Elsa.ProtoActor/Snapshots.cs @@ -5,5 +5,5 @@ using Elsa.Workflows.Core.State; namespace Elsa.ProtoActor; -public record WorkflowSnapshot(string DefinitionId, int Version, WorkflowState WorkflowState, ICollection Bookmarks, IDictionary? Input); +public record WorkflowSnapshot(string DefinitionId, int Version, WorkflowState WorkflowState, IDictionary? Input); public record BookmarkSnapshot(ICollection Bookmarks); \ No newline at end of file diff --git a/src/modules/Elsa.Scheduling/Implementations/WorkflowTriggerScheduler.cs b/src/modules/Elsa.Scheduling/Implementations/WorkflowTriggerScheduler.cs index 978cff721..d854d247a 100644 --- a/src/modules/Elsa.Scheduling/Implementations/WorkflowTriggerScheduler.cs +++ b/src/modules/Elsa.Scheduling/Implementations/WorkflowTriggerScheduler.cs @@ -25,7 +25,7 @@ public class WorkflowTriggerScheduler : IWorkflowTriggerScheduler _jobScheduler = jobScheduler; } - public async Task ScheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default) + public async Task ScheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default) { var triggerList = triggers.ToList(); @@ -51,7 +51,7 @@ public class WorkflowTriggerScheduler : IWorkflowTriggerScheduler } } - public async Task UnscheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default) + public async Task UnscheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default) { var triggerList = triggers.ToList(); diff --git a/src/modules/Elsa.Scheduling/Services/IWorkflowTriggerScheduler.cs b/src/modules/Elsa.Scheduling/Services/IWorkflowTriggerScheduler.cs index 76ec2305f..ac4440a1a 100644 --- a/src/modules/Elsa.Scheduling/Services/IWorkflowTriggerScheduler.cs +++ b/src/modules/Elsa.Scheduling/Services/IWorkflowTriggerScheduler.cs @@ -10,6 +10,6 @@ namespace Elsa.Scheduling.Services; /// public interface IWorkflowTriggerScheduler { - Task ScheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default); - Task UnscheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default); + Task ScheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default); + Task UnscheduleTriggersAsync(IEnumerable triggers, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/Events/Trigger/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/Events/Trigger/Endpoint.cs index dd0c62ff1..f8b3531be 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/Events/Trigger/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/Events/Trigger/Endpoint.cs @@ -28,7 +28,7 @@ public class Trigger : ElsaEndpoint 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(); var result = await _workflowRuntime.TriggerWorkflowsAsync(bookmarkName, eventBookmark, new TriggerWorkflowsOptions(), cancellationToken); diff --git a/src/modules/Elsa.Workflows.Core/Activities/Event.cs b/src/modules/Elsa.Workflows.Core/Activities/Event.cs index 1b8f66bfe..7ae2ec64a 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Event.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Event.cs @@ -35,6 +35,6 @@ public class Event : Trigger protected override void Execute(ActivityExecutionContext context) { var eventName = context.Get(EventName)!; - context.CreateBookmark(new EventBookmarkData(eventName)); + context.CreateBookmark(new EventBookmarkPayload(eventName)); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Elsa.Workflows.Core.csproj b/src/modules/Elsa.Workflows.Core/Elsa.Workflows.Core.csproj index eebe78a67..29803b09e 100644 --- a/src/modules/Elsa.Workflows.Core/Elsa.Workflows.Core.csproj +++ b/src/modules/Elsa.Workflows.Core/Elsa.Workflows.Core.csproj @@ -21,6 +21,11 @@ + + + + + diff --git a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs index 413c90767..8b5c6b17e 100644 --- a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs +++ b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs @@ -65,6 +65,7 @@ public class WorkflowsFeature : FeatureBase .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton() .AddTransient() diff --git a/src/modules/Elsa.Workflows.Core/Implementations/BookmarkHasher.cs b/src/modules/Elsa.Workflows.Core/Implementations/BookmarkHasher.cs new file mode 100644 index 000000000..c279bb4c4 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Implementations/BookmarkHasher.cs @@ -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; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs index cdfbb7400..8f77b6a75 100644 --- a/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Models/ActivityExecutionContext.cs @@ -113,11 +113,11 @@ public class ActivityExecutionContext public Bookmark CreateBookmark(object? payload = default, ExecuteActivityDelegate? callback = default) { - var hasher = GetRequiredService(); + var bookmarkHasher = GetRequiredService(); var identityGenerator = GetRequiredService(); var payloadSerializer = GetRequiredService(); 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(), diff --git a/src/modules/Elsa.Workflows.Core/Models/EventBookmarkData.cs b/src/modules/Elsa.Workflows.Core/Models/EventBookmarkPayload.cs similarity index 75% rename from src/modules/Elsa.Workflows.Core/Models/EventBookmarkData.cs rename to src/modules/Elsa.Workflows.Core/Models/EventBookmarkPayload.cs index c41dc6060..5e30a8fe1 100644 --- a/src/modules/Elsa.Workflows.Core/Models/EventBookmarkData.cs +++ b/src/modules/Elsa.Workflows.Core/Models/EventBookmarkPayload.cs @@ -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; } diff --git a/src/modules/Elsa.Workflows.Core/Services/IBookmarkHasher.cs b/src/modules/Elsa.Workflows.Core/Services/IBookmarkHasher.cs new file mode 100644 index 000000000..7eeda0360 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Services/IBookmarkHasher.cs @@ -0,0 +1,7 @@ +namespace Elsa.Workflows.Core.Services; + +public interface IBookmarkHasher +{ + string Hash(string activityTypeName, object? payload); + string Hash(string activityTypeName, string? serializedPayload); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerHashEqualityComparer.cs b/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerHashEqualityComparer.cs index df191e136..96d082529 100644 --- a/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerHashEqualityComparer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Comparers/WorkflowTriggerHashEqualityComparer.cs @@ -2,8 +2,8 @@ using Elsa.Workflows.Runtime.Entities; namespace Elsa.Workflows.Runtime.Comparers; -public class WorkflowTriggerHashEqualityComparer : IEqualityComparer +public class WorkflowTriggerHashEqualityComparer : IEqualityComparer { - 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(); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowTrigger.cs b/src/modules/Elsa.Workflows.Runtime/Entities/StoredTrigger.cs similarity index 90% rename from src/modules/Elsa.Workflows.Runtime/Entities/WorkflowTrigger.cs rename to src/modules/Elsa.Workflows.Runtime/Entities/StoredTrigger.cs index 1fa7f704c..e708be610 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowTrigger.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/StoredTrigger.cs @@ -5,7 +5,7 @@ namespace Elsa.Workflows.Runtime.Entities; /// /// Represents a trigger associated with a workflow definition. /// -public class WorkflowTrigger : Entity +public class StoredTrigger : Entity { public string WorkflowDefinitionId { get; set; } = default!; public string Name { get; set; } = default!; diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowTriggerExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowTriggerExtensions.cs index d88b8df0d..b890e2e11 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowTriggerExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowTriggerExtensions.cs @@ -6,7 +6,7 @@ namespace Elsa.Workflows.Runtime.Extensions; public static class WorkflowTriggerExtensions { - public static IEnumerable Filter(this IEnumerable triggers) where T : ITrigger + public static IEnumerable Filter(this IEnumerable triggers) where T : ITrigger { var triggerName = ActivityTypeNameHelper.GenerateTypeName(); return triggers.Where(x => x.Name == triggerName); diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index da0183aa9..ecbcbb5cf 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -50,7 +50,7 @@ public class WorkflowRuntimeFeature : FeatureBase /// public Func BookmarkStore { get; set; } = sp => sp.GetRequiredService(); - public Func WorkflowTriggerStore { get; set; } = sp => sp.GetRequiredService(); + public Func WorkflowTriggerStore { get; set; } = sp => sp.GetRequiredService(); public Func WorkflowExecutionLogStore { get; set; } = sp => sp.GetRequiredService(); public Func WorkflowStateExporter { get; set; } = @@ -75,6 +75,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddSingleton(WorkflowRuntime) .AddSingleton(WorkflowDispatcher) .AddSingleton(WorkflowStateStore) @@ -85,7 +86,7 @@ public class WorkflowRuntimeFeature : FeatureBase // Memory Stores .AddMemoryStore() .AddMemoryStore() - .AddMemoryStore() + .AddMemoryStore() .AddMemoryStore() // Workflow definition providers. diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs index d7278c4c7..646c3e14d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/DefaultWorkflowRuntime.cs @@ -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 StartWorkflowAsync(string definitionId, StartWorkflowOptions options, CancellationToken cancellationToken = default) + public async Task 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(), workflowResult.WorkflowState.Bookmarks, cancellationToken); - await _eventPublisher.PublishAsync(new WorkflowExecuted(workflow, workflowState), cancellationToken); + await UpdateBookmarksAsync(workflowState, new List(), workflowState.Bookmarks, cancellationToken); - return new StartWorkflowResult(workflowState.Id); + return new StartWorkflowResult(workflowState.Id, workflowHost.WorkflowState.Bookmarks); } - public async Task ResumeWorkflowAsync(string instanceId, string bookmarkId, ResumeWorkflowOptions options, CancellationToken cancellationToken = default) + public async Task 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 TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions options, CancellationToken cancellationToken = default) + public async Task 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(); + 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 ExportWorkflowStateAsync(string instanceId, + CancellationToken cancellationToken = default) + { + return await _workflowStateStore.LoadAsync(instanceId, cancellationToken); + } - private async Task UpdateBookmarksAsync(WorkflowState workflowState, ICollection previousBookmarks, ICollection 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 previousBookmarks, + ICollection 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 bookmarks, CancellationToken cancellationToken) + private async Task StoreBookmarksAsync(string workflowInstanceId, ICollection 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 bookmarks, CancellationToken cancellationToken) + private async Task RemoveBookmarksAsync( + string workflowInstanceId, + IEnumerable 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 originalBookmarks, ICollection updatedBookmarks, CancellationToken cancellationToken) + private async Task PublishChangedBookmarksAsync(WorkflowState workflowState, + ICollection originalBookmarks, ICollection 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); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryBookmarkStore.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryBookmarkStore.cs index 84e89957d..db961f8ae 100644 --- a/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryBookmarkStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryBookmarkStore.cs @@ -6,30 +6,30 @@ namespace Elsa.Workflows.Runtime.Implementations; public class MemoryBookmarkStore : IBookmarkStore { - private readonly ConcurrentDictionary<(string ActivityTypeName, string Hash), ICollection> _bookmarks = new(); + private readonly ConcurrentDictionary> _bookmarks = new(); public ValueTask SaveAsync(string activityTypeName, string hash, string workflowInstanceId, IEnumerable bookmarkIds, CancellationToken cancellationToken = default) { var storedBookmarks = bookmarkIds.Select(x => new StoredBookmark(activityTypeName, hash, workflowInstanceId, x)).ToList(); - _bookmarks.AddOrUpdate((activityTypeName, hash), new List(storedBookmarks), (s, bookmarks) => bookmarks.Concat(storedBookmarks).ToList()); + _bookmarks.AddOrUpdate(hash, new List(storedBookmarks), (s, bookmarks) => bookmarks.Concat(storedBookmarks).ToList()); return ValueTask.CompletedTask; } - public ValueTask> LoadAsync(string workflowInstanceId, CancellationToken cancellationToken = default) + public ValueTask> FindByWorkflowInstanceAsync(string workflowInstanceId, CancellationToken cancellationToken = default) { var bookmarks = _bookmarks.Values.SelectMany(x => x).Where(x => x.WorkflowInstanceId == workflowInstanceId).ToList(); return new(bookmarks); } - public ValueTask> LoadAsync(string activityTypeName, string hash, CancellationToken cancellationToken = default) + public ValueTask> FindByHashAsync(string hash, CancellationToken cancellationToken = default) { - var bookmarks = _bookmarks.TryGetValue((activityTypeName, hash), out var value) ? value : Enumerable.Empty(); + var bookmarks = _bookmarks.TryGetValue(hash, out var value) ? value : Enumerable.Empty(); 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; diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryTriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryTriggerStore.cs new file mode 100644 index 000000000..15bdf8682 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryTriggerStore.cs @@ -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 _store; + + public MemoryTriggerStore(MemoryStore store) + { + _store = store; + } + + public Task SaveAsync(StoredTrigger record, CancellationToken cancellationToken = default) + { + _store.Save(record, x => x.Id); + return Task.CompletedTask; + } + + public Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default) + { + _store.SaveMany(records, x => x.Id); + return Task.CompletedTask; + } + + public Task> FindAsync(string hash, CancellationToken cancellationToken = default) + { + var triggers = _store.Query(query => query.Where(x => x.Hash == hash)); + return Task.FromResult>(triggers.ToList()); + } + + public Task> FindManyByWorkflowDefinitionIdAsync( + string workflowDefinitionId, + CancellationToken cancellationToken = default) + { + var triggers = _store.Query(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId)); + return Task.FromResult>(triggers.ToList()); + } + + public Task ReplaceAsync(IEnumerable removed, IEnumerable added, + CancellationToken cancellationToken = default) + { + _store.DeleteMany(removed, x => x.Id); + _store.SaveMany(added, x => x.Id); + + return Task.CompletedTask; + } + + public Task DeleteManyAsync(IEnumerable ids, CancellationToken cancellationToken = default) + { + _store.DeleteMany(ids); + return Task.CompletedTask; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryWorkflowTriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryWorkflowTriggerStore.cs deleted file mode 100644 index 28954422b..000000000 --- a/src/modules/Elsa.Workflows.Runtime/Implementations/MemoryWorkflowTriggerStore.cs +++ /dev/null @@ -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 _store; - - public MemoryWorkflowTriggerStore(MemoryStore store) - { - _store = store; - } - - public Task SaveAsync(WorkflowTrigger record, CancellationToken cancellationToken = default) - { - _store.Save(record, x => x.Id); - return Task.CompletedTask; - } - - public Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default) - { - _store.SaveMany(records, x => x.Id); - return Task.CompletedTask; - } - - public Task> 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>(triggers.ToList()); - } - - public Task> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) - { - var triggers = _store.Query(query => query.Where(x => x.WorkflowDefinitionId == workflowDefinitionId)); - return Task.FromResult>(triggers.ToList()); - } - - public Task ReplaceAsync(IEnumerable removed, IEnumerable added, CancellationToken cancellationToken = default) - { - _store.DeleteMany(removed, x => x.Id); - _store.SaveMany(added, x => x.Id); - - return Task.CompletedTask; - } - - public Task DeleteManyAsync(IEnumerable ids, CancellationToken cancellationToken = default) - { - _store.DeleteMany(ids); - return Task.CompletedTask; - } -} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/TriggerIndexer.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/TriggerIndexer.cs index 8e3a56f63..dc52a27ca 100644 --- a/src/modules/Elsa.Workflows.Runtime/Implementations/TriggerIndexer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/TriggerIndexer.cs @@ -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 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(0); + : new List(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> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) => - await _workflowTriggerStore.FindManyByWorkflowDefinitionIdAsync(workflowDefinitionId, cancellationToken); + private async Task> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) => + await _triggerStore.FindManyByWorkflowDefinitionIdAsync(workflowDefinitionId, cancellationToken); - private async IAsyncEnumerable GetTriggersAsync(Workflow workflow, [EnumeratorCancellation] CancellationToken cancellationToken = default) + private async IAsyncEnumerable 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> GetTriggersAsync(WorkflowIndexingContext context, IActivity activity) + private async Task> 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> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger) + private async Task> 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) }); diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowHost.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowHost.cs new file mode 100644 index 000000000..8a7e13689 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowHost.cs @@ -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 StartWorkflowAsync(IDictionary? input = default, CancellationToken cancellationToken = default) + { + var instanceId = _identityGenerator.GenerateId(); + return await StartWorkflowAsync(instanceId, input, cancellationToken); + } + + public async Task StartWorkflowAsync(string instanceId, IDictionary? 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 ResumeWorkflowAsync(string bookmarkId, IDictionary? 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 diff, CancellationToken cancellationToken) + { + var removedBookmarks = diff.Removed; + var createdBookmarks = diff.Added; + await _eventPublisher.PublishAsync(new WorkflowBookmarksIndexed(new IndexedWorkflowBookmarks(WorkflowState, createdBookmarks, removedBookmarks)), cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowHostFactory.cs b/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowHostFactory.cs new file mode 100644 index 000000000..483eddea8 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Implementations/WorkflowHostFactory.cs @@ -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 CreateAsync(Workflow workflow, WorkflowState workflowState, CancellationToken cancellationToken = default) + { + var workflowHost = (IWorkflowHost)ActivatorUtilities.CreateInstance(_serviceProvider, workflow, workflowState); + return Task.FromResult(workflowHost); + } + + public Task CreateAsync(Workflow workflow, CancellationToken cancellationToken = default) + { + var workflowState = new WorkflowState + { + DefinitionId = workflow.Identity.DefinitionId, + DefinitionVersion = workflow.Identity.Version + }; + + return CreateAsync(workflow, workflowState, cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/IndexedWorkflowTriggers.cs b/src/modules/Elsa.Workflows.Runtime/Models/IndexedWorkflowTriggers.cs index b2f7cad0b..77caf9cf2 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/IndexedWorkflowTriggers.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/IndexedWorkflowTriggers.cs @@ -3,4 +3,4 @@ using Elsa.Workflows.Runtime.Entities; namespace Elsa.Workflows.Runtime.Models; -public record IndexedWorkflowTriggers(Workflow Workflow, ICollection AddedTriggers, ICollection RemovedTriggers, ICollection UnchangedTriggers); \ No newline at end of file +public record IndexedWorkflowTriggers(Workflow Workflow, ICollection AddedTriggers, ICollection RemovedTriggers, ICollection UnchangedTriggers); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IBookmarkStore.cs b/src/modules/Elsa.Workflows.Runtime/Services/IBookmarkStore.cs index fcb80b662..a6a266f58 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/IBookmarkStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/IBookmarkStore.cs @@ -5,7 +5,7 @@ namespace Elsa.Workflows.Runtime.Services; public interface IBookmarkStore { ValueTask SaveAsync(string activityTypeName, string hash, string workflowInstanceId, IEnumerable bookmarkIds, CancellationToken cancellationToken = default); - ValueTask> LoadAsync(string workflowInstanceId, CancellationToken cancellationToken = default); - ValueTask> LoadAsync(string activityTypeName, string hash, CancellationToken cancellationToken = default); - ValueTask DeleteAsync(string activityTypeName, string hash, string workflowInstanceId, CancellationToken cancellationToken = default); + ValueTask> FindByWorkflowInstanceAsync(string workflowInstanceId, CancellationToken cancellationToken = default); + ValueTask> FindByHashAsync(string hash, CancellationToken cancellationToken = default); + ValueTask DeleteAsync(string hash, string workflowInstanceId, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/ITriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Services/ITriggerStore.cs new file mode 100644 index 000000000..cbddd122e --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/ITriggerStore.cs @@ -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 records, CancellationToken cancellationToken = default); + Task> FindAsync(string hash, CancellationToken cancellationToken = default); + Task> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default); + Task ReplaceAsync(IEnumerable removed, IEnumerable added, CancellationToken cancellationToken = default); + Task DeleteManyAsync(IEnumerable ids, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowHost.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowHost.cs new file mode 100644 index 000000000..b1c9b6c2a --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowHost.cs @@ -0,0 +1,21 @@ +using Elsa.Workflows.Core.Helpers; +using Elsa.Workflows.Core.Models; +using Elsa.Workflows.Core.State; + +namespace Elsa.Workflows.Runtime.Services; + +/// +/// Represents a single workflow instance that can be executed and takes care of publishing various lifecycle events. +/// +public interface IWorkflowHost +{ + Workflow Workflow { get; set; } + WorkflowState WorkflowState { get; set; } + Task StartWorkflowAsync(IDictionary? input = default, CancellationToken cancellationToken = default); + Task StartWorkflowAsync(string instanceId, IDictionary? input = default, CancellationToken cancellationToken = default); + Task ResumeWorkflowAsync(string bookmarkId, IDictionary? input = default, CancellationToken cancellationToken = default); +} + +public record StartWorkflowHostResult(Diff BookmarksDiff); + +public record ResumeWorkflowHostResult(Diff BookmarksDiff); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowHostFactory.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowHostFactory.cs new file mode 100644 index 000000000..9617ee01e --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowHostFactory.cs @@ -0,0 +1,19 @@ +using Elsa.Workflows.Core.Models; +using Elsa.Workflows.Core.State; + +namespace Elsa.Workflows.Runtime.Services; + +/// +/// Creates objects. +/// +public interface IWorkflowHostFactory +{ + Task CreateAsync( + Workflow workflow, + WorkflowState workflowState, + CancellationToken cancellationToken = default); + + Task CreateAsync( + Workflow workflow, + CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs index 76e9ee359..2c52eee3a 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowRuntime.cs @@ -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 StartWorkflowAsync(string definitionId, StartWorkflowOptions options, CancellationToken cancellationToken = default); Task ResumeWorkflowAsync(string instanceId, string bookmarkId, ResumeWorkflowOptions options, CancellationToken cancellationToken = default); Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions options, CancellationToken cancellationToken = default); + Task ExportWorkflowStateAsync(string instanceId, CancellationToken cancellationToken = default); + Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default); } public record StartWorkflowOptions(string? CorrelationId = default, IDictionary? Input = default, VersionOptions VersionOptions = default); public record ResumeWorkflowOptions(IDictionary? Input = default); -public record StartWorkflowResult(string InstanceId); -public record ResumeWorkflowResult; -public record TriggerWorkflowsOptions(IDictionary? Input = default); -public record TriggerWorkflowsResult; \ No newline at end of file +public record StartWorkflowResult(string InstanceId, ICollection Bookmarks); +public record ResumeWorkflowResult(ICollection Bookmarks); +public record TriggerWorkflowsOptions(string? CorrelationId = default, IDictionary? Input = default); +public record TriggerWorkflowsResult(ICollection TriggeredWorkflows); +public record TriggeredWorkflow(string InstanceId, ICollection Bookmarks); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowTriggerStore.cs b/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowTriggerStore.cs deleted file mode 100644 index 9f9fd8521..000000000 --- a/src/modules/Elsa.Workflows.Runtime/Services/IWorkflowTriggerStore.cs +++ /dev/null @@ -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 records, CancellationToken cancellationToken = default); - Task> FindManyByNameAsync(string name, string? hash, CancellationToken cancellationToken = default); - Task> FindManyByWorkflowDefinitionIdAsync(string workflowDefinitionId, CancellationToken cancellationToken = default); - Task ReplaceAsync(IEnumerable removed, IEnumerable added, CancellationToken cancellationToken = default); - Task DeleteManyAsync(IEnumerable ids, CancellationToken cancellationToken = default); -} \ No newline at end of file