From 3beedb9ec0e8969f7d13bb3e63ceb57e0c3db99f Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 22 Feb 2025 19:45:44 +0100 Subject: [PATCH] Refactor activity execution state tracking Introduce `IsExecuting` flag to explicitly track activity execution state, improving clarity and control over workflow activity handling. Adjust scheduling intervals and add concurrency handling for database updates to enhance reliability and performance in interrupted workflows. --- src/apps/Elsa.Server.Web/Program.cs | 4 +-- .../Management/WorkflowInstanceStore.cs | 30 ++++++++++++++++++- .../Contexts/ActivityExecutionContext.cs | 12 +++++++- .../Enums/ActivityStatus.cs | 2 +- .../DefaultActivityInvokerMiddleware.cs | 6 ++++ .../Services/WorkflowRunner.cs | 2 +- .../Services/WorkflowStateExtractor.cs | 2 ++ .../State/ActivityExecutionContextState.cs | 5 ++++ 8 files changed, 57 insertions(+), 6 deletions(-) diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index 6ce024b3f..434b73cfd 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -681,10 +681,10 @@ services.Configure(options => options.Schedule.ConfigureTask(TimeSpan.FromSeconds(300)); options.Schedule.ConfigureTask(TimeSpan.FromSeconds(300)); options.Schedule.ConfigureTask(TimeSpan.FromHours(4)); - options.Schedule.ConfigureTask(TimeSpan.FromMinutes(1)); + options.Schedule.ConfigureTask(TimeSpan.FromSeconds(15)); }); -services.Configure(options => { options.InactivityThreshold = TimeSpan.FromMinutes(1); }); +services.Configure(options => { options.InactivityThreshold = TimeSpan.FromSeconds(15); }); services.Configure(options => options.Ttl = TimeSpan.FromSeconds(10)); services.Configure(options => options.CacheDuration = TimeSpan.FromDays(1)); diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs index f8295372d..9aebbfaa4 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs @@ -157,7 +157,35 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore Id = workflowInstanceId, UpdatedAt = value }; - await _store.UpdatePartialAsync(entity, [x => x.UpdatedAt], cancellationToken); + + await using var dbContext = await _store.CreateDbContextAsync(cancellationToken); + dbContext.Attach(entity); + dbContext.Entry(entity).Property(x => x.UpdatedAt).IsModified = true; + + try + { + await dbContext.SaveChangesAsync(cancellationToken); + } + catch (DbUpdateConcurrencyException e) + { + foreach (var entry in e.Entries) + { + var proposedValues = entry.CurrentValues; + var databaseValues = await entry.GetDatabaseValuesAsync(cancellationToken); + + if(databaseValues == null) + continue; + + var updatedAtProperty = entry.Metadata.GetProperty(nameof(WorkflowInstance.UpdatedAt)); + var proposedValue = (DateTimeOffset)proposedValues[updatedAtProperty]!; + var databaseValue = (DateTimeOffset)databaseValues[updatedAtProperty]!; + + if (proposedValue > databaseValue) + proposedValues[updatedAtProperty] = proposedValue; + + entry.OriginalValues.SetValues(databaseValues); + } + } } /// diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs index 4865d4abd..6ecc474d1 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs @@ -48,7 +48,6 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable var expressionExecutionContextProps = ExpressionExecutionContextExtensions.CreateActivityExecutionContextPropertiesFrom(workflowExecutionContext, workflowExecutionContext.Input); expressionExecutionContextProps[ExpressionExecutionContextExtensions.ActivityKey] = activity; ExpressionExecutionContext = new(workflowExecutionContext.ServiceProvider, new(), parentActivityExecutionContext?.ExpressionExecutionContext ?? workflowExecutionContext.ExpressionExecutionContext, expressionExecutionContextProps, Taint, CancellationToken); - ; Activity = activity; ActivityDescriptor = activityDescriptor; StartedAt = startedAt; @@ -84,6 +83,17 @@ public partial class ActivityExecutionContext : IExecutionContext, IDisposable /// public bool IsCompleted => Status is ActivityStatus.Completed or ActivityStatus.Canceled; + /// + /// Gets or sets a value indicating whether the activity is actively executing. + /// + /// + /// This flag is set to true immediately before the activity begins execution + /// and is set to false once the execution is completed. + /// It can be used to determine if an activity was in-progress in case of unexpected + /// application termination, allowing the system to retry execution upon restarting. + /// + public bool IsExecuting { get; set; } + /// /// The workflow execution context. /// diff --git a/src/modules/Elsa.Workflows.Core/Enums/ActivityStatus.cs b/src/modules/Elsa.Workflows.Core/Enums/ActivityStatus.cs index 7c34c9521..6fde9c385 100644 --- a/src/modules/Elsa.Workflows.Core/Enums/ActivityStatus.cs +++ b/src/modules/Elsa.Workflows.Core/Enums/ActivityStatus.cs @@ -11,7 +11,7 @@ public enum ActivityStatus Pending, /// - /// The activity is in the Running state. Note that event if an activity is running, it may not be executing. + /// The activity is in the Running state. While in this state, the activity is not necessarily being actively executed. This state represents a logical status rather than a physical action. /// Running, diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs index 9b5ef9cd9..b9c42bb40 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs @@ -48,6 +48,9 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I context.AddExecutionLogEntry("Precondition Failed", "Cannot execute at this time"); return; } + + // Mark activity as executing. + context.IsExecuting = true; // Conditionally commit the workflow state. if (ShouldCommit(context, ActivityLifetimeEvent.ActivityExecuting)) @@ -83,6 +86,9 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I workflowExecutionContext.Bookmarks.AddRange(context.Bookmarks); logger.LogDebug("Added {BookmarkCount} bookmarks to the workflow execution context", context.Bookmarks.Count); } + + // Mark activity as executed. + context.IsExecuting = false; // Conditionally commit the workflow state. if (ShouldCommit(context, ActivityLifetimeEvent.ActivityExecuted)) diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs index a6a0edc11..3019f1a46 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs @@ -158,7 +158,7 @@ public class WorkflowRunner( else { // Check if there are any leaf nodes in the Pending state. - var pendingActivityExecutionContexts = workflowExecutionContext.ActivityExecutionContexts.Where(x => x.Status == ActivityStatus.Pending).ToList(); + var pendingActivityExecutionContexts = workflowExecutionContext.ActivityExecutionContexts.Where(x => x.IsExecuting).ToList(); if( pendingActivityExecutionContexts.Count > 0) { diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs index 99a11a214..a39cedf73 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs @@ -145,6 +145,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor activityExecutionContext.ActivityState.Merge(activityExecutionContextState.ActivityState); activityExecutionContext.TransitionTo(activityExecutionContextState.Status); + activityExecutionContext.IsExecuting = activityExecutionContextState.IsExecuting; activityExecutionContext.StartedAt = activityExecutionContextState.StartedAt; activityExecutionContext.CompletedAt = activityExecutionContextState.CompletedAt; activityExecutionContext.Tag = activityExecutionContextState.Tag; @@ -232,6 +233,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor Properties = activityExecutionContext.Properties, ActivityState = activityExecutionContext.ActivityState, Status = activityExecutionContext.Status, + IsExecuting = activityExecutionContext.IsExecuting, StartedAt = activityExecutionContext.StartedAt, CompletedAt = activityExecutionContext.CompletedAt, Tag = activityExecutionContext.Tag, diff --git a/src/modules/Elsa.Workflows.Core/State/ActivityExecutionContextState.cs b/src/modules/Elsa.Workflows.Core/State/ActivityExecutionContextState.cs index 4f685f8b6..19a73d6dc 100644 --- a/src/modules/Elsa.Workflows.Core/State/ActivityExecutionContextState.cs +++ b/src/modules/Elsa.Workflows.Core/State/ActivityExecutionContextState.cs @@ -55,6 +55,11 @@ public class ActivityExecutionContextState /// The status of the activity. /// public ActivityStatus Status { get; set; } + + /// + /// Gets or sets a value indicating whether the activity is actively executing. + /// + public bool IsExecuting { get; set; } /// /// The time at which the activity execution began.