From 8cb43db070efc0ed3e2b354a5cf975ead5d14f2f Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 27 Jan 2025 20:07:23 +0100 Subject: [PATCH] Refactor workflow state persistence for improved consistency. Removed obsolete middleware for persisting bookmarks, execution logs, and variables, integrating their functionality into the commit state handler. Added early exit checks for empty collections in persistence methods. Updated activity invoker logic for better state commit handling during execution. --- .../Elsa.EntityFrameworkCore.Common/Store.cs | 12 ++++++++++-- .../DefaultActivityInvokerMiddleware.cs | 4 ++-- ...rkflowExecutionPipelineBuilderExtensions.cs | 6 +++--- .../PersistActivityExecutionLogMiddleware.cs | 8 +++++--- .../Workflows/PersistBookmarkMiddleware.cs | 18 ++++++------------ .../PersistWorkflowExecutionLogMiddleware.cs | 7 ++++--- .../Workflows/PersistentVariablesMiddleware.cs | 17 +++-------------- .../Services/BookmarkUpdater.cs | 18 ++++++++++++------ .../Services/StoreActivityExecutionLogSink.cs | 4 ++++ .../Services/StoreCommitStateHandler.cs | 15 ++++++++++++++- .../Services/StoreWorkflowExecutionLogSink.cs | 4 ++++ 11 files changed, 67 insertions(+), 46 deletions(-) diff --git a/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs b/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs index d5b8fdd6c..3387e912a 100644 --- a/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs +++ b/src/modules/Elsa.EntityFrameworkCore.Common/Store.cs @@ -79,9 +79,13 @@ public class Store(IDbContextFactory dbContextF Func? onSaving = default, CancellationToken cancellationToken = default) { - await using var dbContext = await CreateDbContextAsync(cancellationToken); var entityList = entities.ToList(); + if (entityList.Count == 0) + return; + + await using var dbContext = await CreateDbContextAsync(cancellationToken); + if (onSaving != null) { var savingTasks = entityList.Select(entity => onSaving(dbContext, entity, cancellationToken).AsTask()).ToList(); @@ -162,9 +166,13 @@ public class Store(IDbContextFactory dbContextF Func? onSaving = default, CancellationToken cancellationToken = default) { - await using var dbContext = await CreateDbContextAsync(cancellationToken); var entityList = entities.ToList(); + if (entityList.Count == 0) + return; + + await using var dbContext = await CreateDbContextAsync(cancellationToken); + if (onSaving != null) { var savingTasks = entityList.Select(entity => onSaving(dbContext, entity, cancellationToken).AsTask()).ToList(); diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs index a7bd74c6b..a1d884d10 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs @@ -49,7 +49,7 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I } // Conditionally commit the workflow state. - if(ShouldCommitWhenStarting(context)) + if(ShouldCommitWhenExecuting(context)) await context.WorkflowExecutionContext.CommitAsync(); context.TransitionTo(ActivityStatus.Running); @@ -120,7 +120,7 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I await context.EvaluateInputPropertiesAsync(); } - private bool ShouldCommitWhenStarting(ActivityExecutionContext context) + private bool ShouldCommitWhenExecuting(ActivityExecutionContext context) { var behavior = context.Activity.GetCommitStateBehavior(); diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs index cd48e66f6..9affb2058 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs @@ -18,9 +18,6 @@ public static class WorkflowExecutionPipelineBuilderExtensions pipelineBuilder .Reset() .UseEngineExceptionHandling() - .UseBookmarkPersistence() - .UseActivityExecutionLogPersistence() - .UseWorkflowExecutionLogPersistence() .UsePersistentVariables() .UseExceptionHandling() .UseDefaultActivityScheduler(); @@ -33,15 +30,18 @@ public static class WorkflowExecutionPipelineBuilderExtensions /// /// Installs middleware that persists bookmarks after workflow execution. /// + [Obsolete("This middleware is no longer used and will be removed in a future version. Bookmarks are now persisted through the commit state handler.")] public static IWorkflowExecutionPipelineBuilder UseBookmarkPersistence(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); /// /// Installs middleware that persists the workflow execution journal. /// + [Obsolete("This middleware is no longer used and will be removed in a future version. Execution logs are now persisted through the commit state handler.")] public static IWorkflowExecutionPipelineBuilder UseWorkflowExecutionLogPersistence(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); /// /// Installs middleware that persists activity execution records. /// + [Obsolete("This middleware is no longer used and will be removed in a future version. Activity state is now persisted through the commit state handler.")] public static IWorkflowExecutionPipelineBuilder UseActivityExecutionLogPersistence(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs index e81f53fd7..b659c8635 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistActivityExecutionLogMiddleware.cs @@ -1,17 +1,19 @@ using Elsa.Workflows.Pipelines.WorkflowExecution; -using Elsa.Workflows.Runtime.Entities; namespace Elsa.Workflows.Runtime.Middleware.Workflows; /// /// Creates and updates activity execution records from activity execution contexts. /// -public class PersistActivityExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink sink) : WorkflowExecutionMiddleware(next) +[Obsolete("This middleware is no longer used and will be removed in a future version. Activity state is now persisted through the commit state handler")] +public class PersistActivityExecutionLogMiddleware(WorkflowMiddlewareDelegate next) : WorkflowExecutionMiddleware(next) { /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) { await Next(context); - await sink.PersistExecutionLogsAsync(context); + + // Not used anymore. + //await sink.PersistExecutionLogsAsync(context); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs index 0319baf07..c44b255ab 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistBookmarkMiddleware.cs @@ -1,26 +1,20 @@ using Elsa.Workflows.Pipelines.WorkflowExecution; -using Elsa.Workflows.Runtime.Requests; namespace Elsa.Workflows.Runtime.Middleware.Workflows; /// /// Takes care of loading and persisting bookmarks. /// -public class PersistBookmarkMiddleware : WorkflowExecutionMiddleware +[Obsolete("This middleware is no longer used and will be removed in a future version. Bookmarks are now persisted through the commit state handler")] +public class PersistBookmarkMiddleware(WorkflowMiddlewareDelegate next) : WorkflowExecutionMiddleware(next) { - private readonly IBookmarksPersister _bookmarksPersister; - - /// - public PersistBookmarkMiddleware(WorkflowMiddlewareDelegate next, IBookmarksPersister bookmarksPersister) : base(next) - { - _bookmarksPersister = bookmarksPersister; - } - /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) { await Next(context); - var bookmarkRequest = new UpdateBookmarksRequest(context, context.BookmarksDiff, context.CorrelationId); - await _bookmarksPersister.PersistBookmarksAsync(bookmarkRequest); + + // Not used anymore. + // var bookmarkRequest = new UpdateBookmarksRequest(context, context.BookmarksDiff, context.CorrelationId); + // await _bookmarksPersister.PersistBookmarksAsync(bookmarkRequest); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs index 25d863054..4c2d39534 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistWorkflowExecutionLogMiddleware.cs @@ -1,12 +1,12 @@ using Elsa.Workflows.Pipelines.WorkflowExecution; -using Elsa.Workflows.Runtime.Entities; namespace Elsa.Workflows.Runtime.Middleware.Workflows; /// /// Takes care of persisting workflow execution log entries. /// -public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink sink) : WorkflowExecutionMiddleware(next) +[Obsolete("This middleware is no longer used and will be removed in a future version. Execution logs are now persisted through the commit state handler.")] +public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next) : WorkflowExecutionMiddleware(next) { /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) @@ -14,6 +14,7 @@ public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate ne // Invoke next middleware. await Next(context); - await sink.PersistExecutionLogsAsync(context); + // Not used anymore. + //await sink.PersistExecutionLogsAsync(context); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistentVariablesMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistentVariablesMiddleware.cs index b743983da..ab464717d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistentVariablesMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/PersistentVariablesMiddleware.cs @@ -5,28 +5,17 @@ namespace Elsa.Workflows.Runtime.Middleware.Workflows; /// /// Takes care of loading and persisting workflow variables. /// -public class PersistentVariablesMiddleware : WorkflowExecutionMiddleware +public class PersistentVariablesMiddleware(WorkflowMiddlewareDelegate next, IVariablePersistenceManager variablePersistenceManager) : WorkflowExecutionMiddleware(next) { - private readonly IVariablePersistenceManager _variablePersistenceManager; - - /// - /// Constructor. - /// - public PersistentVariablesMiddleware(WorkflowMiddlewareDelegate next, IVariablePersistenceManager variablePersistenceManager) : base(next) - { - _variablePersistenceManager = variablePersistenceManager; - } - /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) { // Load variables into the workflow execution context. - await _variablePersistenceManager.LoadVariablesAsync(context); + await variablePersistenceManager.LoadVariablesAsync(context); // Invoke next middleware. await Next(context); - // Persist variables. - await _variablePersistenceManager.SaveVariablesAsync(context); + // Variables are persisted through the commit state handler. } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkUpdater.cs b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkUpdater.cs index d42bfab61..d7a9be23e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BookmarkUpdater.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BookmarkUpdater.cs @@ -11,12 +11,15 @@ public class BookmarkUpdater(IBookmarkManager bookmarkManager, IBookmarkStore bo public async Task UpdateBookmarksAsync(UpdateBookmarksRequest request, CancellationToken cancellationToken = default) { var instanceId = request.WorkflowExecutionContext.Id; - await RemoveBookmarksAsync(instanceId, request.Diff.Removed, cancellationToken); - await StoreBookmarksAsync(request.WorkflowExecutionContext, request.Diff.Added, cancellationToken); + await RemoveBookmarksAsync(instanceId, request.Diff.Removed.ToList(), cancellationToken); + await StoreBookmarksAsync(request.WorkflowExecutionContext, request.Diff.Added.ToList(), cancellationToken); } - - private async Task RemoveBookmarksAsync(string workflowInstanceId, IEnumerable bookmarks, CancellationToken cancellationToken) + + private async Task RemoveBookmarksAsync(string workflowInstanceId, ICollection bookmarks, CancellationToken cancellationToken) { + if (bookmarks.Count == 0) + return; + var matchingIds = bookmarks.Select(x => x.Id).ToList(); var filter = new BookmarkFilter { @@ -25,9 +28,12 @@ public class BookmarkUpdater(IBookmarkManager bookmarkManager, IBookmarkStore bo }; await bookmarkManager.DeleteManyAsync(filter, cancellationToken); } - - private async Task StoreBookmarksAsync(WorkflowExecutionContext context, IEnumerable bookmarks, CancellationToken cancellationToken) + + private async Task StoreBookmarksAsync(WorkflowExecutionContext context, ICollection bookmarks, CancellationToken cancellationToken) { + if (bookmarks.Count == 0) + return; + foreach (var bookmark in bookmarks) { var storedBookmark = context.MapBookmark(bookmark); diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs index 2c280a22c..23cab65a3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs @@ -15,6 +15,10 @@ public class StoreActivityExecutionLogSink(IActivityExecutionStore activityExecu public async Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken = default) { var records = await extractor.ExtractLogRecordsAsync(context).ToList(); + + if(records.Count == 0) + return; + await activityExecutionStore.SaveManyAsync(records, cancellationToken); await notificationSender.SendAsync(new ActivityExecutionLogUpdated(context, records), cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreCommitStateHandler.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreCommitStateHandler.cs index 5db9391ba..6479c0eb6 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreCommitStateHandler.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreCommitStateHandler.cs @@ -1,9 +1,16 @@ using Elsa.Workflows.Management; +using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.State; namespace Elsa.Workflows.Runtime; -public class StoreCommitStateHandler(IWorkflowInstanceManager workflowInstanceManager) : ICommitStateHandler +public class StoreCommitStateHandler( + IWorkflowInstanceManager workflowInstanceManager, + IBookmarksPersister bookmarkPersister, + IVariablePersistenceManager variablePersistenceManager, + ILogRecordSink activityExecutionLogRecordSink, + ILogRecordSink workflowExecutionLogRecordSink) : ICommitStateHandler { public async Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default) { @@ -13,7 +20,13 @@ public class StoreCommitStateHandler(IWorkflowInstanceManager workflowInstanceMa public async Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default) { + var updateBookmarksRequest = new UpdateBookmarksRequest(workflowExecutionContext, workflowExecutionContext.BookmarksDiff, workflowExecutionContext.CorrelationId); + await bookmarkPersister.PersistBookmarksAsync(updateBookmarksRequest); + await activityExecutionLogRecordSink.PersistExecutionLogsAsync(workflowExecutionContext, cancellationToken); + await workflowExecutionLogRecordSink.PersistExecutionLogsAsync(workflowExecutionContext, cancellationToken); + await variablePersistenceManager.SaveVariablesAsync(workflowExecutionContext); await workflowInstanceManager.SaveAsync(workflowState, cancellationToken); + workflowExecutionContext.ExecutionLog.Clear(); await workflowExecutionContext.ExecuteDeferredTasksAsync(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs index a91667858..72bce2ebb 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreWorkflowExecutionLogSink.cs @@ -14,6 +14,10 @@ public class StoreWorkflowExecutionLogSink(IWorkflowExecutionLogStore store, ILo public async Task PersistExecutionLogsAsync(WorkflowExecutionContext context, CancellationToken cancellationToken) { var records = await extractor.ExtractLogRecordsAsync(context).ToList(); + + if(records.Count == 0) + return; + await store.AddManyAsync(records, context.CancellationToken); await notificationSender.SendAsync(new WorkflowExecutionLogUpdated(context), context.CancellationToken); }