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.
This commit is contained in:
Sipke Schoorstra 2025-01-27 20:07:23 +01:00
parent 879e546942
commit 8cb43db070
11 changed files with 67 additions and 46 deletions

View file

@ -79,9 +79,13 @@ public class Store<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
Func<TDbContext, TEntity, CancellationToken, ValueTask>? 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<TDbContext, TEntity>(IDbContextFactory<TDbContext> dbContextF
Func<TDbContext, TEntity, CancellationToken, ValueTask>? 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();

View file

@ -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();

View file

@ -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
/// <summary>
/// Installs middleware that persists bookmarks after workflow execution.
/// </summary>
[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<PersistBookmarkMiddleware>();
/// <summary>
/// Installs middleware that persists the workflow execution journal.
/// </summary>
[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<PersistWorkflowExecutionLogMiddleware>();
/// <summary>
/// Installs middleware that persists activity execution records.
/// </summary>
[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<PersistActivityExecutionLogMiddleware>();
}

View file

@ -1,17 +1,19 @@
using Elsa.Workflows.Pipelines.WorkflowExecution;
using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Middleware.Workflows;
/// <summary>
/// Creates and updates activity execution records from activity execution contexts.
/// </summary>
public class PersistActivityExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink<ActivityExecutionRecord> 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)
{
/// <inheritdoc />
public override async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
await Next(context);
await sink.PersistExecutionLogsAsync(context);
// Not used anymore.
//await sink.PersistExecutionLogsAsync(context);
}
}

View file

@ -1,26 +1,20 @@
using Elsa.Workflows.Pipelines.WorkflowExecution;
using Elsa.Workflows.Runtime.Requests;
namespace Elsa.Workflows.Runtime.Middleware.Workflows;
/// <summary>
/// Takes care of loading and persisting bookmarks.
/// </summary>
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;
/// <inheritdoc />
public PersistBookmarkMiddleware(WorkflowMiddlewareDelegate next, IBookmarksPersister bookmarksPersister) : base(next)
{
_bookmarksPersister = bookmarksPersister;
}
/// <inheritdoc />
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);
}
}

View file

@ -1,12 +1,12 @@
using Elsa.Workflows.Pipelines.WorkflowExecution;
using Elsa.Workflows.Runtime.Entities;
namespace Elsa.Workflows.Runtime.Middleware.Workflows;
/// <summary>
/// Takes care of persisting workflow execution log entries.
/// </summary>
public class PersistWorkflowExecutionLogMiddleware(WorkflowMiddlewareDelegate next, ILogRecordSink<WorkflowExecutionLogRecord> 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)
{
/// <inheritdoc />
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);
}
}

View file

@ -5,28 +5,17 @@ namespace Elsa.Workflows.Runtime.Middleware.Workflows;
/// <summary>
/// Takes care of loading and persisting workflow variables.
/// </summary>
public class PersistentVariablesMiddleware : WorkflowExecutionMiddleware
public class PersistentVariablesMiddleware(WorkflowMiddlewareDelegate next, IVariablePersistenceManager variablePersistenceManager) : WorkflowExecutionMiddleware(next)
{
private readonly IVariablePersistenceManager _variablePersistenceManager;
/// <summary>
/// Constructor.
/// </summary>
public PersistentVariablesMiddleware(WorkflowMiddlewareDelegate next, IVariablePersistenceManager variablePersistenceManager) : base(next)
{
_variablePersistenceManager = variablePersistenceManager;
}
/// <inheritdoc />
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.
}
}

View file

@ -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<Bookmark> bookmarks, CancellationToken cancellationToken)
private async Task RemoveBookmarksAsync(string workflowInstanceId, ICollection<Bookmark> 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<Bookmark> bookmarks, CancellationToken cancellationToken)
private async Task StoreBookmarksAsync(WorkflowExecutionContext context, ICollection<Bookmark> bookmarks, CancellationToken cancellationToken)
{
if (bookmarks.Count == 0)
return;
foreach (var bookmark in bookmarks)
{
var storedBookmark = context.MapBookmark(bookmark);

View file

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

View file

@ -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<ActivityExecutionRecord> activityExecutionLogRecordSink,
ILogRecordSink<WorkflowExecutionLogRecord> 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();
}
}

View file

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