diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs index 061de2fda..f03bf0e17 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/WorkflowExecutionPipelineBuilderExtensions.cs @@ -13,8 +13,8 @@ public static class WorkflowExecutionPipelineBuilderExtensions public static IWorkflowExecutionPipelineBuilder UseDefaultPipeline(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder .Reset() - .UseBackgroundActivities() .UseDeferredActivityTasks() + .UseBackgroundActivities() .UseBookmarkPersistence() .UseActivityExecutionLogPersistence() .UseWorkflowExecutionLogPersistence() diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs index 6446b229e..ed84ef148 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Workflows/ScheduleBackgroundActivitiesMiddleware.cs @@ -1,5 +1,6 @@ using Elsa.Extensions; using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management; using Elsa.Workflows.Pipelines.WorkflowExecution; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Middleware.Activities; @@ -14,7 +15,8 @@ public class ScheduleBackgroundActivitiesMiddleware( WorkflowMiddlewareDelegate next, IBackgroundActivityScheduler backgroundActivityScheduler, IStimulusHasher stimulusHasher, - IBookmarkStore bookmarkStore) + IBookmarkStore bookmarkStore, + IWorkflowInstanceManager workflowInstanceManager) : WorkflowExecutionMiddleware(next) { /// @@ -29,42 +31,53 @@ public class ScheduleBackgroundActivitiesMiddleware( var scheduledBackgroundActivities = workflowExecutionContext .TransientProperties .GetOrAdd(BackgroundActivityInvokerMiddleware.BackgroundActivitySchedulesKey, () => new List()); + + if (scheduledBackgroundActivities.Count == 0) + return; - foreach (var scheduledBackgroundActivity in scheduledBackgroundActivities) + context.DeferTask(async () => { - // Schedule the background activity. - var jobId = await backgroundActivityScheduler.ScheduleAsync(scheduledBackgroundActivity, cancellationToken); - - // Select the bookmark associated with the background activity. - var bookmark = workflowExecutionContext.Bookmarks.First(x => x.Id == scheduledBackgroundActivity.BookmarkId); - var stimulus = bookmark.GetPayload(); - - // Store the created job ID. - workflowExecutionContext.Bookmarks.Remove(bookmark); - stimulus.JobId = jobId; - bookmark = bookmark with + // Commit state. + await workflowInstanceManager.SaveAsync(context, cancellationToken); + + foreach (var scheduledBackgroundActivity in scheduledBackgroundActivities) { - Payload = bookmark.Payload, - Hash = stimulusHasher.Hash(bookmark.Name, stimulus) - }; - workflowExecutionContext.Bookmarks.Add(bookmark); + // Schedule the background activity. + var jobId = await backgroundActivityScheduler.ScheduleAsync(scheduledBackgroundActivity, cancellationToken); - // Update the bookmark. - var storedBookmark = new StoredBookmark - { - Id = bookmark.Id, - TenantId = tenantId, - ActivityInstanceId = bookmark.ActivityInstanceId, - ActivityTypeName = bookmark.Name, - Hash = bookmark.Hash, - WorkflowInstanceId = workflowExecutionContext.Id, - CreatedAt = bookmark.CreatedAt, - CorrelationId = workflowExecutionContext.CorrelationId, - Payload = bookmark.Payload, - Metadata = bookmark.Metadata, - }; + // Select the bookmark associated with the background activity. + var bookmark = workflowExecutionContext.Bookmarks.First(x => x.Id == scheduledBackgroundActivity.BookmarkId); + var stimulus = bookmark.GetPayload(); + + // Store the created job ID. + workflowExecutionContext.Bookmarks.Remove(bookmark); + stimulus.JobId = jobId; + bookmark = bookmark with + { + Payload = bookmark.Payload, + Hash = stimulusHasher.Hash(bookmark.Name, stimulus) + }; + workflowExecutionContext.Bookmarks.Add(bookmark); + + // Update the bookmark. + var storedBookmark = new StoredBookmark + { + Id = bookmark.Id, + TenantId = tenantId, + ActivityInstanceId = bookmark.ActivityInstanceId, + ActivityTypeName = bookmark.Name, + Hash = bookmark.Hash, + WorkflowInstanceId = workflowExecutionContext.Id, + CreatedAt = bookmark.CreatedAt, + CorrelationId = workflowExecutionContext.CorrelationId, + Payload = bookmark.Payload, + Metadata = bookmark.Metadata, + }; - await bookmarkStore.SaveAsync(storedBookmark, cancellationToken); - } + await bookmarkStore.SaveAsync(storedBookmark, cancellationToken); + } + }); + + } } \ No newline at end of file