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); }