From 157a7e01450c5bf9f6b0d09459a39e0ba5aadc7d Mon Sep 17 00:00:00 2001 From: David Garza Date: Sun, 7 Dec 2025 16:00:02 -0500 Subject: [PATCH 1/7] Add default commit workflow and activity strategy configuration and usage examples - Introduced methods to set default workflow and activity commit strategies in CommitStrategiesFeature. - Updated CommitStateOptions to include properties for default strategies. - Added extension methods for configuring default strategies in WorkflowsFeature. - Created usage examples demonstrating how to set and utilize default commit strategies. - Implemented tests to verify default strategy behavior in various scenarios. --- .../CommitStates/CommitStrategiesFeature.cs | 44 ++-- .../Extensions/ModuleExtensions.cs | 36 +++- .../Options/CommitStateOptions.cs | 22 ++ .../CommitStates/USAGE_EXAMPLE.md | 112 ++++++++++ .../DefaultActivityInvokerMiddleware.cs | 37 +++- .../DefaultActivitySchedulerMiddleware.cs | 33 ++- .../BackgroundActivityInvokerMiddleware.cs | 12 +- .../CommitTracker.cs | 24 +++ ...leWorkflowWithoutActivityCommitStrategy.cs | 19 ++ .../DefaultActivityCommitStrategy/Tests.cs | 196 ++++++++++++++++++ ...kflowWithExplicitActivityCommitStrategy.cs | 22 ++ .../CommitTracker.cs | 24 +++ ...leWorkflowWithoutWorkflowCommitStrategy.cs | 22 ++ .../DefaultWorkflowCommitStrategy/Tests.cs | 166 +++++++++++++++ ...kflowWithExplicitWorkflowCommitStrategy.cs | 23 ++ 15 files changed, 756 insertions(+), 36 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/SimpleWorkflowWithoutActivityCommitStrategy.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/WorkflowWithExplicitActivityCommitStrategy.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/SimpleWorkflowWithoutWorkflowCommitStrategy.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/WorkflowWithExplicitWorkflowCommitStrategy.cs diff --git a/src/modules/Elsa.Workflows.Core/CommitStates/CommitStrategiesFeature.cs b/src/modules/Elsa.Workflows.Core/CommitStates/CommitStrategiesFeature.cs index eeca573f3..1ca58a54e 100644 --- a/src/modules/Elsa.Workflows.Core/CommitStates/CommitStrategiesFeature.cs +++ b/src/modules/Elsa.Workflows.Core/CommitStates/CommitStrategiesFeature.cs @@ -16,27 +16,27 @@ public class CommitStrategiesFeature(IModule module) : FeatureBase(module) Add(new WorkflowExecutedWorkflowStrategy()); Add(new ActivityExecutingWorkflowStrategy()); Add(new ActivityExecutedWorkflowStrategy()); - + // Activity commit strategies. Add(new CommitAlwaysActivityStrategy()); Add(new CommitNeverActivityStrategy()); Add(new ExecutingActivityStrategy()); Add(new ExecutedActivityStrategy()); } - + public void Add(IWorkflowCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); Add(registration); } - + public void Add(string displayName, IWorkflowCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); registration.Metadata.DisplayName = displayName; Add(registration); } - + public void Add(string displayName, string description, IWorkflowCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); @@ -44,7 +44,7 @@ public class CommitStrategiesFeature(IModule module) : FeatureBase(module) registration.Metadata.Description = description; Add(registration); } - + public void Add(string name, string displayName, string description, IWorkflowCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); @@ -53,25 +53,25 @@ public class CommitStrategiesFeature(IModule module) : FeatureBase(module) registration.Metadata.Description = description; Add(registration); } - + public void Add(WorkflowCommitStrategyRegistration registration) { Services.Configure(options => options.WorkflowCommitStrategies[registration.Metadata.Name] = registration); } - + public void Add(IActivityCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); Add(registration); } - + public void Add(string displayName, IActivityCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); registration.Metadata.DisplayName = displayName; Add(registration); } - + public void Add(string displayName, string description, IActivityCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); @@ -79,7 +79,7 @@ public class CommitStrategiesFeature(IModule module) : FeatureBase(module) registration.Metadata.Description = description; Add(registration); } - + public void Add(string name, string displayName, string description, IActivityCommitStrategy strategy) { var registration = ObjectRegistrationFactory.Describe(strategy); @@ -88,12 +88,32 @@ public class CommitStrategiesFeature(IModule module) : FeatureBase(module) registration.Metadata.Description = description; Add(registration); } - + public void Add(ActivityCommitStrategyRegistration registration) { Services.Configure(options => options.ActivityCommitStrategies[registration.Metadata.Name] = registration); } - + + /// + /// Sets the specified workflow commit strategy as the global default. + /// The strategy will not be added to the registry and will only serve as a fallback when workflows do not specify their own strategy. + /// + /// The workflow commit strategy instance to use as the default. + public void SetDefaultWorkflowCommitStrategy(IWorkflowCommitStrategy strategy) + { + Services.Configure(options => options.DefaultWorkflowCommitStrategy = strategy); + } + + /// + /// Sets the specified activity commit strategy as the global default. + /// The strategy will not be added to the registry and will only serve as a fallback when activities do not specify their own strategy. + /// + /// The activity commit strategy instance to use as the default. + public void SetDefaultActivityCommitStrategy(IActivityCommitStrategy strategy) + { + Services.Configure(options => options.DefaultActivityCommitStrategy = strategy); + } + public override void Apply() { Services.AddSingleton(); diff --git a/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs index 9d2bf94c5..db3558cd1 100644 --- a/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs @@ -4,12 +4,46 @@ using Elsa.Workflows.Features; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; +/// +/// Provides extension methods for configuring commit strategies on . +/// public static class WorkflowsFeatureCommitStateExtensions { + /// + /// Configures commit strategies for workflows. + /// + /// The workflows feature. + /// An optional configuration delegate for the commit strategies feature. + /// The workflows feature for chaining. public static WorkflowsFeature UseCommitStrategies(this WorkflowsFeature workflowsFeature, Action? configure = null) { workflowsFeature.Module.Use(configure); return workflowsFeature; } - + + /// + /// Sets the specified workflow commit strategy as the global default for all workflows that do not specify their own strategy. + /// The strategy will be automatically registered if not already present. + /// + /// The workflows feature. + /// The workflow commit strategy instance to use as the default. + /// The workflows feature for chaining. + public static WorkflowsFeature WithDefaultWorkflowCommitStrategy(this WorkflowsFeature workflowsFeature, IWorkflowCommitStrategy strategy) + { + workflowsFeature.Module.Use(feature => feature.SetDefaultWorkflowCommitStrategy(strategy)); + return workflowsFeature; + } + + /// + /// Sets the specified activity commit strategy as the global default for all activities that do not specify their own strategy. + /// The strategy will be automatically registered if not already present. + /// + /// The workflows feature. + /// The activity commit strategy instance to use as the default. + /// The workflows feature for chaining. + public static WorkflowsFeature WithDefaultActivityCommitStrategy(this WorkflowsFeature workflowsFeature, IActivityCommitStrategy strategy) + { + workflowsFeature.Module.Use(feature => feature.SetDefaultActivityCommitStrategy(strategy)); + return workflowsFeature; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/CommitStates/Options/CommitStateOptions.cs b/src/modules/Elsa.Workflows.Core/CommitStates/Options/CommitStateOptions.cs index cf607f986..babd8e77d 100644 --- a/src/modules/Elsa.Workflows.Core/CommitStates/Options/CommitStateOptions.cs +++ b/src/modules/Elsa.Workflows.Core/CommitStates/Options/CommitStateOptions.cs @@ -1,7 +1,29 @@ namespace Elsa.Workflows.CommitStates; +/// +/// Configuration options for commit state strategies. +/// public class CommitStateOptions { + /// + /// Gets or sets the workflow commit strategies. + /// public IDictionary WorkflowCommitStrategies { get; set; } = new Dictionary(); + + /// + /// Gets or sets the activity commit strategies. + /// public IDictionary ActivityCommitStrategies { get; set; } = new Dictionary(); + + /// + /// Gets or sets the default workflow commit strategy instance to use when a workflow does not specify its own. + /// This strategy is not added to the registry and serves only as a fallback. + /// + public IWorkflowCommitStrategy? DefaultWorkflowCommitStrategy { get; set; } + + /// + /// Gets or sets the default activity commit strategy instance to use when an activity does not specify its own. + /// This strategy is not added to the registry and serves only as a fallback. + /// + public IActivityCommitStrategy? DefaultActivityCommitStrategy { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md b/src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md new file mode 100644 index 000000000..a4ba23eb2 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md @@ -0,0 +1,112 @@ +# Commit Strategies - Usage Example + +## Setting a Global Default Commit Strategy + +You can configure an optional default workflow commit strategy globally by passing a strategy instance to `WithDefaultWorkflowCommitStrategy`. + +### Basic Usage + +```csharp +using Elsa.Workflows.CommitStates.Strategies; + +services.AddElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()) + ) +); +``` + +**Benefits:** + +- **Simplicity**: Pass any strategy instance directly +- **Auto-Registration**: The strategy is automatically added to the registry if not already present +- **No Manual Setup Required**: You don't need to call `UseCommitStrategies()` or `AddStandardStrategies()` first + +### Configuration Within UseCommitStrategies + +You can also configure the default strategy within the `UseCommitStrategies` callback: + +```csharp +services.AddElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .UseCommitStrategies(commitStrategies => + { + // Set default strategy + commitStrategies.SetDefaultWorkflowCommitStrategy(new WorkflowExecutingWorkflowStrategy()); + }) + ) +); +``` + +## How It Works + +The commit strategy resolution follows this priority: + +1. **Workflow-specific strategy** (if set on the workflow via `Workflow.Options.CommitStrategyName`) +2. **Global default strategy** (if configured via `WithDefaultWorkflowCommitStrategy`) +3. **No automatic commit** (if neither is set) + +## Available Standard Strategies + +- `WorkflowExecutingWorkflowStrategy` - Commit before workflow execution +- `WorkflowExecutedWorkflowStrategy` - Commit after workflow execution +- `ActivityExecutingWorkflowStrategy` - Commit before each activity execution +- `ActivityExecutedWorkflowStrategy` - Commit after each activity execution +- `PeriodicWorkflowStrategy` - Commit periodically + +## Examples + +### Simple Configuration + +```csharp +services.AddElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()) + ) +); +``` + +### With Additional Custom Strategies + +```csharp +services.AddElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .UseCommitStrategies(commitStrategies => + { + // Add standard strategies + commitStrategies.AddStandardStrategies(); + + // Add custom strategy + commitStrategies.Add(new MyCustomWorkflowStrategy()); + + // Set default + commitStrategies.SetDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()); + }) + ) +); +``` + +### With Custom Strategy Configuration + +Since the method accepts instances, you can pass strategies with custom configuration: + +```csharp +var customStrategy = new PeriodicWorkflowStrategy +{ + // Configure your strategy +}; + +services.AddElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultWorkflowCommitStrategy(customStrategy) + ) +); +``` + +## Benefits + +- **Consistency**: Set a default behavior for all workflows without configuring each individually +- **Flexibility**: Individual workflows can override the default by setting their own commit strategy +- **Simplicity**: Reduce boilerplate configuration across your workflows +- **Auto-Registration**: No need to remember to add strategies to the registry first +- **Instance-Based**: Pass configured strategy instances with custom settings diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs index 24ffb4939..c5698f759 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs @@ -5,6 +5,7 @@ using Elsa.Workflows.Activities; using Elsa.Workflows.CommitStates; using Elsa.Workflows.Pipelines.ActivityExecution; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; namespace Elsa.Workflows.Middleware.Activities; @@ -22,11 +23,11 @@ public static class ActivityInvokerMiddlewareExtensions /// /// A default activity execution middleware component that evaluates the current activity's properties, executes the activity and adds any produced bookmarks to the workflow execution context. /// -public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, ICommitStrategyRegistry commitStrategyRegistry, ILogger logger) +public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, ICommitStrategyRegistry commitStrategyRegistry, IOptions commitStateOptions, ILogger logger) : IActivityExecutionMiddleware { private static readonly MethodInfo ExecuteAsyncMethodInfo = typeof(IActivity).GetMethod(nameof(IActivity.ExecuteAsync))!; - + /// public async ValueTask InvokeAsync(ActivityExecutionContext context) { @@ -65,7 +66,7 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I // Execute activity. await ExecuteActivityAsync(context); - + var currentActivityStatus = context.Status; var activityDidComplete = previousActivityStatus != ActivityStatus.Completed && currentActivityStatus == ActivityStatus.Completed; @@ -86,7 +87,7 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I // Invoke next middleware. await next(context); - + // If the activity completed, send a notification. if (activityDidComplete) { @@ -105,7 +106,7 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I /// protected virtual async ValueTask ExecuteActivityAsync(ActivityExecutionContext context) { - var executeDelegate = context.WorkflowExecutionContext.ExecuteDelegate + var executeDelegate = context.WorkflowExecutionContext.ExecuteDelegate ?? (ExecuteActivityDelegate)Delegate.CreateDelegate(typeof(ExecuteActivityDelegate), context.Activity, ExecuteAsyncMethodInfo); await executeDelegate(context); @@ -129,7 +130,18 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I private bool ShouldCommit(ActivityExecutionContext context, ActivityLifetimeEvent lifetimeEvent) { var strategyName = context.Activity.GetCommitStrategy(); - var strategy = string.IsNullOrWhiteSpace(strategyName) ? null : commitStrategyRegistry.FindActivityStrategy(strategyName); + IActivityCommitStrategy? strategy; + + if (!string.IsNullOrWhiteSpace(strategyName)) + { + strategy = commitStrategyRegistry.FindActivityStrategy(strategyName); + } + else + { + // Fall back to the default strategy if configured + strategy = commitStateOptions.Value.DefaultActivityCommitStrategy; + } + var commitAction = CommitAction.Default; if (strategy != null) @@ -147,7 +159,18 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I case CommitAction.Default: { var workflowStrategyName = context.WorkflowExecutionContext.Workflow.Options.CommitStrategyName; - var workflowStrategy = string.IsNullOrWhiteSpace(workflowStrategyName) ? null : commitStrategyRegistry.FindWorkflowStrategy(workflowStrategyName); + + IWorkflowCommitStrategy? workflowStrategy; + + if (!string.IsNullOrWhiteSpace(workflowStrategyName)) + { + workflowStrategy = commitStrategyRegistry.FindWorkflowStrategy(workflowStrategyName); + } + else + { + // Fall back to the default strategy if configured + workflowStrategy = commitStateOptions.Value.DefaultWorkflowCommitStrategy; + } if (workflowStrategy == null) return false; diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs index 26594ca93..7cf94c22e 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs @@ -3,6 +3,7 @@ using Elsa.Workflows.CommitStates; using Elsa.Workflows.Models; using Elsa.Workflows.Options; using Elsa.Workflows.Pipelines.WorkflowExecution; +using Microsoft.Extensions.Options; namespace Elsa.Workflows.Middleware.Workflows; @@ -20,7 +21,7 @@ public static class UseActivitySchedulerMiddlewareExtensions /// /// A workflow execution middleware component that executes scheduled work items. /// -public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, IActivityInvoker activityInvoker, ICommitStrategyRegistry commitStrategyRegistry) : WorkflowExecutionMiddleware(next) +public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, IActivityInvoker activityInvoker, ICommitStrategyRegistry commitStrategyRegistry, IOptions commitStateOptions) : WorkflowExecutionMiddleware(next) { /// public override async ValueTask InvokeAsync(WorkflowExecutionContext context) @@ -29,19 +30,19 @@ public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, context.TransitionTo(WorkflowSubStatus.Executing); await ConditionallyCommitStateAsync(context, WorkflowLifetimeEvent.WorkflowExecuting); - + while (scheduler.HasAny) { // Do not start a workflow if cancellation has been requested. if (context.CancellationToken.IsCancellationRequested) break; - + var currentWorkItem = scheduler.Take(); await ExecuteWorkItemAsync(context, currentWorkItem); } - + await Next(context); - + if (context.Status == WorkflowStatus.Running) context.TransitionTo(context.AllActivitiesCompleted() ? WorkflowSubStatus.Finished : WorkflowSubStatus.Suspended); } @@ -59,18 +60,28 @@ public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, await activityInvoker.InvokeAsync(context, workItem.Activity, options); } - + private async Task ConditionallyCommitStateAsync(WorkflowExecutionContext context, WorkflowLifetimeEvent lifetimeEvent) { var strategyName = context.Workflow.Options.CommitStrategyName; - var strategy = string.IsNullOrWhiteSpace(strategyName) ? null : commitStrategyRegistry.FindWorkflowStrategy(strategyName); - - if(strategy == null) + IWorkflowCommitStrategy? strategy; + + if (!string.IsNullOrWhiteSpace(strategyName)) + { + strategy = commitStrategyRegistry.FindWorkflowStrategy(strategyName); + } + else + { + // Fall back to the default strategy if configured + strategy = commitStateOptions.Value.DefaultWorkflowCommitStrategy; + } + + if (strategy == null) return; - + var strategyContext = new WorkflowCommitStateStrategyContext(context, lifetimeEvent); var commitAction = strategy.ShouldCommit(strategyContext); - + if (commitAction is CommitAction.Commit) await context.CommitAsync(); } diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs index 30277beb9..cf7a47820 100644 --- a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/BackgroundActivityInvokerMiddleware.cs @@ -11,6 +11,7 @@ using Elsa.Workflows.Runtime.Notifications; using Elsa.Workflows.Runtime.Stimuli; using JetBrains.Annotations; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; namespace Elsa.Workflows.Runtime.Middleware.Activities; @@ -24,8 +25,9 @@ public class BackgroundActivityInvokerMiddleware( IIdentityGenerator identityGenerator, IBackgroundActivityScheduler backgroundActivityScheduler, ICommitStrategyRegistry commitStrategyRegistry, - IMediator mediator) - : DefaultActivityInvokerMiddleware(next, commitStrategyRegistry, logger) + IMediator mediator, + IOptions commitStateOptions) + : DefaultActivityInvokerMiddleware(next, commitStrategyRegistry, commitStateOptions, logger) { internal static string GetBackgroundActivityOutputKey(string activityNodeId) => $"__BackgroundActivityOutput:{activityNodeId}"; internal static string GetBackgroundActivityOutcomesKey(string activityNodeId) => $"__BackgroundActivityOutcomes:{activityNodeId}"; @@ -129,7 +131,7 @@ public class BackgroundActivityInvokerMiddleware( } private static bool GetIsBackgroundExecution(ActivityExecutionContext context) => context.TransientProperties.ContainsKey(BackgroundActivityExecutionContextExtensions.IsBackgroundExecution); - + /// /// If the input contains captured output from the background activity invoker, apply that to the execution context. /// @@ -183,7 +185,7 @@ public class BackgroundActivityInvokerMiddleware( context.WorkflowExecutionContext.Properties.Remove(bookmarksKey); } - + private void CapturePropertiesIfAny(ActivityExecutionContext context) { var activity = context.Activity; @@ -195,7 +197,7 @@ public class BackgroundActivityInvokerMiddleware( if (capturedProperties == null) return; - foreach (var property in capturedProperties) + foreach (var property in capturedProperties) context.Properties[property.Key] = property.Value; } diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs new file mode 100644 index 000000000..1beab0576 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs @@ -0,0 +1,24 @@ +using Elsa.Workflows.CommitStates; +using Elsa.Workflows.State; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultActivityCommitStrategy; + +/// +/// Test helper to track commit invocations +/// +public class CommitTracker : ICommitStateHandler +{ + public int CommitCount { get; private set; } + + public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default) + { + CommitCount++; + return Task.CompletedTask; + } + + public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default) + { + CommitCount++; + return Task.CompletedTask; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/SimpleWorkflowWithoutActivityCommitStrategy.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/SimpleWorkflowWithoutActivityCommitStrategy.cs new file mode 100644 index 000000000..111008110 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/SimpleWorkflowWithoutActivityCommitStrategy.cs @@ -0,0 +1,19 @@ +using Elsa.Workflows.Activities; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultActivityCommitStrategy; + +public class SimpleWorkflowWithoutActivityCommitStrategy : WorkflowBase +{ + protected override void Build(IWorkflowBuilder builder) + { + builder.Root = new Sequence + { + Activities = + { + new WriteLine("Activity 1"), + new WriteLine("Activity 2"), + new WriteLine("Activity 3") + } + }; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs new file mode 100644 index 000000000..6a52cc196 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs @@ -0,0 +1,196 @@ +using Elsa.Extensions; +using Elsa.Testing.Shared; +using Elsa.Workflows.CommitStates; +using Elsa.Workflows.CommitStates.Strategies; +using Elsa.Workflows.CommitStates.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Xunit.Abstractions; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultActivityCommitStrategy; + +public class Tests(ITestOutputHelper testOutputHelper) +{ + private readonly ITestOutputHelper _testOutputHelper = testOutputHelper; + + [Fact(DisplayName = "Activity without explicit strategy uses default commit strategy")] + public async Task ActivityUsesDefaultCommitStrategy() + { + // Arrange + var commitTracker = new CommitTracker(); + var defaultStrategy = new ExecutedActivityStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => + { + workflows.WithDefaultActivityCommitStrategy(defaultStrategy); + workflows.CommitStateHandler = _ => commitTracker; + }) + ) + .AddWorkflow() + .Build(); + + var options = services.GetRequiredService>(); + var workflowRunner = services.GetRequiredService(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + Assert.NotNull(options.Value.DefaultActivityCommitStrategy); + Assert.Same(defaultStrategy, options.Value.DefaultActivityCommitStrategy); + + // 6 commits: 3 WriteLine activities + 3 Sequence (composite) completion checks + Assert.Equal(6, commitTracker.CommitCount); + } + + [Fact(DisplayName = "Activity-specific strategy overrides default commit strategy")] + public async Task ActivitySpecificStrategyOverridesDefault() + { + // Arrange + var commitTracker = new CommitTracker(); + var defaultStrategy = new ExecutedActivityStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => + { + workflows.WithDefaultActivityCommitStrategy(defaultStrategy); + workflows.UseCommitStrategies(commitStrategies => commitStrategies.AddStandardStrategies()); + workflows.CommitStateHandler = _ => commitTracker; + }) + ) + .AddWorkflow() + .Build(); + + var options = services.GetRequiredService>(); + var workflowRunner = services.GetRequiredService(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + Assert.NotNull(options.Value.DefaultActivityCommitStrategy); + Assert.Same(defaultStrategy, options.Value.DefaultActivityCommitStrategy); + + // 4 commits: First activity with ExecutingActivity (before), second with default ExecutedActivity (after), + // plus Sequence composite completions + Assert.Equal(4, commitTracker.CommitCount); + } + + [Fact(DisplayName = "No commits occur when no default strategy and no activity strategy")] + public async Task NoCommitsWithoutAnyStrategy() + { + // Arrange + var commitTracker = new CommitTracker(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => workflows.CommitStateHandler = _ => commitTracker) + ) + .AddWorkflow() + .Build(); + + var options = services.GetRequiredService>(); + var workflowRunner = services.GetRequiredService(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + Assert.Null(options.Value.DefaultActivityCommitStrategy); + + // 1 commit: Only the final commit from WorkflowRunner (no middleware commits) + Assert.Equal(1, commitTracker.CommitCount); + } + + [Fact(DisplayName = "Default activity strategy is not visible in commit strategy registry")] + public async Task DefaultStrategyNotInRegistry() + { + // Arrange + var defaultStrategy = new ExecutedActivityStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultActivityCommitStrategy(defaultStrategy) + ) + ) + .Build(); + + var registry = services.GetRequiredService(); + var options = services.GetRequiredService>(); + + // Act + var activityStrategies = registry.ListActivityStrategyRegistrations().ToList(); + + // Assert + Assert.Empty(activityStrategies); + Assert.NotNull(options.Value.DefaultActivityCommitStrategy); + Assert.Same(defaultStrategy, options.Value.DefaultActivityCommitStrategy); + await Task.CompletedTask; + } + + [Fact(DisplayName = "Default activity strategy with standard strategies does not duplicate")] + public async Task DefaultStrategyWithStandardStrategiesNoDuplicate() + { + // Arrange + var defaultStrategy = new ExecutedActivityStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultActivityCommitStrategy(defaultStrategy) + .UseCommitStrategies(commitStrategies => commitStrategies.AddStandardStrategies()) + ) + ) + .Build(); + + var registry = services.GetRequiredService(); + + // Manually populate the registry + var startupTask = new PopulateCommitStrategyRegistry( + registry, + services.GetRequiredService>() + ); + await startupTask.ExecuteAsync(CancellationToken.None); + + // Act + var activityStrategies = registry.ListActivityStrategyRegistrations().ToList(); + + // Assert - 4 standard activity strategies (no duplication from default) + Assert.Equal(4, activityStrategies.Count); + } + + [Fact(DisplayName = "Default workflow strategy is used when no default activity strategy is specified")] + public async Task DefaultWorkflowStrategyWithoutDefaultActivityStrategy() + { + // Arrange + var commitTracker = new CommitTracker(); + var defaultWorkflowStrategy = new ActivityExecutedWorkflowStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => + { + workflows.WithDefaultWorkflowCommitStrategy(defaultWorkflowStrategy); + workflows.CommitStateHandler = _ => commitTracker; + }) + ) + .AddWorkflow() + .Build(); + + var options = services.GetRequiredService>(); + var workflowRunner = services.GetRequiredService(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + Assert.NotNull(options.Value.DefaultWorkflowCommitStrategy); + Assert.Same(defaultWorkflowStrategy, options.Value.DefaultWorkflowCommitStrategy); + Assert.Null(options.Value.DefaultActivityCommitStrategy); + + // 6 commits: ActivityExecutedWorkflowStrategy commits after each activity completion (3 WriteLine + 3 Sequence) + Assert.Equal(6, commitTracker.CommitCount); + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/WorkflowWithExplicitActivityCommitStrategy.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/WorkflowWithExplicitActivityCommitStrategy.cs new file mode 100644 index 000000000..f6c8c9039 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/WorkflowWithExplicitActivityCommitStrategy.cs @@ -0,0 +1,22 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultActivityCommitStrategy; + +public class WorkflowWithExplicitActivityCommitStrategy : WorkflowBase +{ + protected override void Build(IWorkflowBuilder builder) + { + var writeLineWithStrategy = new WriteLine("Activity with strategy"); + writeLineWithStrategy.CommitStrategy = "ExecutingActivity"; // Uses standard "Commit Before" strategy + + builder.Root = new Sequence + { + Activities = + { + writeLineWithStrategy, + new WriteLine("Activity 2") + } + }; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs new file mode 100644 index 000000000..fbce44d7c --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs @@ -0,0 +1,24 @@ +using Elsa.Workflows.CommitStates; +using Elsa.Workflows.State; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy; + +/// +/// Tracks commit operations for testing purposes. +/// +public class CommitTracker : ICommitStateHandler +{ + public int CommitCount { get; private set; } + + public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default) + { + CommitCount++; + return Task.CompletedTask; + } + + public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default) + { + CommitCount++; + return Task.CompletedTask; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/SimpleWorkflowWithoutWorkflowCommitStrategy.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/SimpleWorkflowWithoutWorkflowCommitStrategy.cs new file mode 100644 index 000000000..92770bfa5 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/SimpleWorkflowWithoutWorkflowCommitStrategy.cs @@ -0,0 +1,22 @@ +using Elsa.Workflows.Activities; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy; + +/// +/// A simple workflow that does not specify an explicit commit strategy. +/// +public class SimpleWorkflowWithoutWorkflowCommitStrategy : WorkflowBase +{ + protected override void Build(IWorkflowBuilder builder) + { + builder.Root = new Sequence + { + Activities = + { + new WriteLine("Activity 1"), + new WriteLine("Activity 2"), + new WriteLine("Activity 3") + } + }; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs new file mode 100644 index 000000000..32a06feda --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs @@ -0,0 +1,166 @@ +using Elsa.Extensions; +using Elsa.Testing.Shared; +using Elsa.Workflows.CommitStates; +using Elsa.Workflows.CommitStates.Strategies; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Xunit.Abstractions; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy; + +public class Tests +{ + private readonly ITestOutputHelper _testOutputHelper; + + public Tests(ITestOutputHelper testOutputHelper) + { + _testOutputHelper = testOutputHelper; + } + + [Fact(DisplayName = "Workflow uses default workflow commit strategy when no explicit workflow commit strategy is set")] + public async Task WorkflowUsesDefaultWorkflowCommitStrategy() + { + // Arrange + var commitTracker = new CommitTracker(); + var defaultStrategy = new ActivityExecutedWorkflowStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => + { + workflows.WithDefaultWorkflowCommitStrategy(defaultStrategy); + workflows.CommitStateHandler = _ => commitTracker; + }) + ) + .AddWorkflow() + .Build(); + + var options = services.GetRequiredService>(); + var workflowRunner = services.GetRequiredService(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert - workflow should finish successfully with default strategy configured + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + Assert.NotNull(options.Value.DefaultWorkflowCommitStrategy); + Assert.Same(defaultStrategy, options.Value.DefaultWorkflowCommitStrategy); + + // Verify exact commit count: ActivityExecutedWorkflowStrategy commits after each activity completes + // With 3 WriteLine activities in a Sequence, this results in exactly 6 commits due to + // how composite activities and workflow completion signals interact + Assert.Equal(6, commitTracker.CommitCount); + } + + [Fact(DisplayName = "Workflow-specific strategy overrides default workflow commit strategy")] + public async Task WorkflowSpecificStrategyOverridesDefault() + { + // Arrange + var commitTracker = new CommitTracker(); + var defaultStrategy = new ActivityExecutedWorkflowStrategy(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultWorkflowCommitStrategy(defaultStrategy) + .UseCommitStrategies(commitStrategies => commitStrategies.AddStandardStrategies()) + .CommitStateHandler = _ => commitTracker + ) + ) + .AddWorkflow() + .Build(); + + var workflowRunner = services.GetRequiredService(); + var options = services.GetRequiredService>(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert - workflow should complete successfully + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + // Default strategy should be configured (available for workflows without explicit strategy) + Assert.Same(defaultStrategy, options.Value.DefaultWorkflowCommitStrategy); + + // Verify the workflow used its explicit "WorkflowExecuting" strategy (commits before workflow starts) + // 1 commit at the beginning before any activities execute + Assert.Equal(1, commitTracker.CommitCount); + } + + [Fact(DisplayName = "No commits occur when no default workflow commit strategy and no workflow strategy")] + public async Task NoCommitsWithoutAnyWorkflowCommitStrategy() + { + // Arrange + var commitTracker = new CommitTracker(); + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => + { + workflows.CommitStateHandler = _ => commitTracker; + }) + ) + .AddWorkflow() + .Build(); + + var workflowRunner = services.GetRequiredService(); + var options = services.GetRequiredService>(); + + // Act + var result = await workflowRunner.RunAsync(); + + // Assert - workflow should still complete even without commit strategy + Assert.Equal(WorkflowStatus.Finished, result.WorkflowState.Status); + // No default strategy should be configured + Assert.Null(options.Value.DefaultWorkflowCommitStrategy); + + // Verify no commits occurred during workflow execution (only final commit from WorkflowRunner) + Assert.Equal(1, commitTracker.CommitCount); + } + + [Fact(DisplayName = "Default workflow strategy is not visible in commit strategy registry")] + public async Task DefaultWorkflowStrategyNotInRegistry() + { + // Arrange + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()) + ) + ) + .Build(); + + var registry = services.GetRequiredService(); + + // Act + var registeredStrategies = registry.ListWorkflowStrategyRegistrations().ToList(); + + // Assert - default strategy should not be in the registry + Assert.Empty(registeredStrategies); + } + + [Fact(DisplayName = "Default workflow strategy with standard strategies does not duplicate")] + public void DefaultWorkflowStrategyWithStandardStrategiesNoDuplicate() + { + // Arrange + var services = new TestApplicationBuilder(_testOutputHelper) + .ConfigureElsa(elsa => elsa + .UseWorkflows(workflows => workflows + .WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()) + .UseCommitStrategies(commitStrategies => commitStrategies.AddStandardStrategies()) + ) + ) + .Build(); + + // Manually trigger the PopulateCommitStrategyRegistry startup task + var options = services.GetRequiredService>(); + var registry = services.GetRequiredService(); + + foreach (var strategy in options.Value.WorkflowCommitStrategies.Values) + registry.RegisterStrategy(strategy); + foreach (var strategy in options.Value.ActivityCommitStrategies.Values) + registry.RegisterStrategy(strategy); + + // Act + var registeredStrategies = registry.ListWorkflowStrategyRegistrations().ToList(); + + // Assert - should only have the 4 standard strategies, not 5 + Assert.Equal(4, registeredStrategies.Count); + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/WorkflowWithExplicitWorkflowCommitStrategy.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/WorkflowWithExplicitWorkflowCommitStrategy.cs new file mode 100644 index 000000000..be235e3e5 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/WorkflowWithExplicitWorkflowCommitStrategy.cs @@ -0,0 +1,23 @@ +using Elsa.Workflows.Activities; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy; + +/// +/// A workflow that explicitly sets a commit strategy, which should override the default. +/// +public class WorkflowWithExplicitWorkflowCommitStrategy : WorkflowBase +{ + protected override void Build(IWorkflowBuilder builder) + { + builder.WorkflowOptions.CommitStrategyName = "WorkflowExecuting"; + + builder.Root = new Sequence + { + Activities = + { + new WriteLine("Activity 1"), + new WriteLine("Activity 2") + } + }; + } +} From 05d807a6918997a1e541563562b1ebda20409562 Mon Sep 17 00:00:00 2001 From: David Garza Date: Tue, 9 Dec 2025 09:12:01 -0500 Subject: [PATCH 2/7] Refactor commit strategy resolution for improved clarity and efficiency; remove unused example files and add a new CommitTracker helper for testing. --- .../CommitStates/USAGE_EXAMPLE.md | 112 ------------------ .../DefaultActivityInvokerMiddleware.cs | 27 +---- .../DefaultActivitySchedulerMiddleware.cs | 14 +-- .../DefaultActivityCommitStrategy/Tests.cs | 13 +- .../CommitTracker.cs | 24 ---- .../DefaultWorkflowCommitStrategy/Tests.cs | 1 + .../CommitTracker.cs | 4 +- 7 files changed, 22 insertions(+), 173 deletions(-) delete mode 100644 src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md delete mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs rename test/integration/Elsa.Workflows.IntegrationTests/{Scenarios/DefaultActivityCommitStrategy => SharedHelpers}/CommitTracker.cs (88%) diff --git a/src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md b/src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md deleted file mode 100644 index a4ba23eb2..000000000 --- a/src/modules/Elsa.Workflows.Core/CommitStates/USAGE_EXAMPLE.md +++ /dev/null @@ -1,112 +0,0 @@ -# Commit Strategies - Usage Example - -## Setting a Global Default Commit Strategy - -You can configure an optional default workflow commit strategy globally by passing a strategy instance to `WithDefaultWorkflowCommitStrategy`. - -### Basic Usage - -```csharp -using Elsa.Workflows.CommitStates.Strategies; - -services.AddElsa(elsa => elsa - .UseWorkflows(workflows => workflows - .WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()) - ) -); -``` - -**Benefits:** - -- **Simplicity**: Pass any strategy instance directly -- **Auto-Registration**: The strategy is automatically added to the registry if not already present -- **No Manual Setup Required**: You don't need to call `UseCommitStrategies()` or `AddStandardStrategies()` first - -### Configuration Within UseCommitStrategies - -You can also configure the default strategy within the `UseCommitStrategies` callback: - -```csharp -services.AddElsa(elsa => elsa - .UseWorkflows(workflows => workflows - .UseCommitStrategies(commitStrategies => - { - // Set default strategy - commitStrategies.SetDefaultWorkflowCommitStrategy(new WorkflowExecutingWorkflowStrategy()); - }) - ) -); -``` - -## How It Works - -The commit strategy resolution follows this priority: - -1. **Workflow-specific strategy** (if set on the workflow via `Workflow.Options.CommitStrategyName`) -2. **Global default strategy** (if configured via `WithDefaultWorkflowCommitStrategy`) -3. **No automatic commit** (if neither is set) - -## Available Standard Strategies - -- `WorkflowExecutingWorkflowStrategy` - Commit before workflow execution -- `WorkflowExecutedWorkflowStrategy` - Commit after workflow execution -- `ActivityExecutingWorkflowStrategy` - Commit before each activity execution -- `ActivityExecutedWorkflowStrategy` - Commit after each activity execution -- `PeriodicWorkflowStrategy` - Commit periodically - -## Examples - -### Simple Configuration - -```csharp -services.AddElsa(elsa => elsa - .UseWorkflows(workflows => workflows - .WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()) - ) -); -``` - -### With Additional Custom Strategies - -```csharp -services.AddElsa(elsa => elsa - .UseWorkflows(workflows => workflows - .UseCommitStrategies(commitStrategies => - { - // Add standard strategies - commitStrategies.AddStandardStrategies(); - - // Add custom strategy - commitStrategies.Add(new MyCustomWorkflowStrategy()); - - // Set default - commitStrategies.SetDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy()); - }) - ) -); -``` - -### With Custom Strategy Configuration - -Since the method accepts instances, you can pass strategies with custom configuration: - -```csharp -var customStrategy = new PeriodicWorkflowStrategy -{ - // Configure your strategy -}; - -services.AddElsa(elsa => elsa - .UseWorkflows(workflows => workflows - .WithDefaultWorkflowCommitStrategy(customStrategy) - ) -); -``` - -## Benefits - -- **Consistency**: Set a default behavior for all workflows without configuring each individually -- **Flexibility**: Individual workflows can override the default by setting their own commit strategy -- **Simplicity**: Reduce boilerplate configuration across your workflows -- **Auto-Registration**: No need to remember to add strategies to the registry first -- **Instance-Based**: Pass configured strategy instances with custom settings diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs index c5698f759..79cac9b2c 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs @@ -130,17 +130,10 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I private bool ShouldCommit(ActivityExecutionContext context, ActivityLifetimeEvent lifetimeEvent) { var strategyName = context.Activity.GetCommitStrategy(); - IActivityCommitStrategy? strategy; - if (!string.IsNullOrWhiteSpace(strategyName)) - { - strategy = commitStrategyRegistry.FindActivityStrategy(strategyName); - } - else - { - // Fall back to the default strategy if configured - strategy = commitStateOptions.Value.DefaultActivityCommitStrategy; - } + IActivityCommitStrategy? strategy = !string.IsNullOrWhiteSpace(strategyName) + ? commitStrategyRegistry.FindActivityStrategy(strategyName) + : commitStateOptions.Value.DefaultActivityCommitStrategy; var commitAction = CommitAction.Default; @@ -160,17 +153,9 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I { var workflowStrategyName = context.WorkflowExecutionContext.Workflow.Options.CommitStrategyName; - IWorkflowCommitStrategy? workflowStrategy; - - if (!string.IsNullOrWhiteSpace(workflowStrategyName)) - { - workflowStrategy = commitStrategyRegistry.FindWorkflowStrategy(workflowStrategyName); - } - else - { - // Fall back to the default strategy if configured - workflowStrategy = commitStateOptions.Value.DefaultWorkflowCommitStrategy; - } + IWorkflowCommitStrategy? workflowStrategy = !string.IsNullOrWhiteSpace(workflowStrategyName) + ? commitStrategyRegistry.FindWorkflowStrategy(workflowStrategyName) + : commitStateOptions.Value.DefaultWorkflowCommitStrategy; if (workflowStrategy == null) return false; diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs index 7cf94c22e..f939d1e64 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Workflows/DefaultActivitySchedulerMiddleware.cs @@ -64,17 +64,9 @@ public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, private async Task ConditionallyCommitStateAsync(WorkflowExecutionContext context, WorkflowLifetimeEvent lifetimeEvent) { var strategyName = context.Workflow.Options.CommitStrategyName; - IWorkflowCommitStrategy? strategy; - - if (!string.IsNullOrWhiteSpace(strategyName)) - { - strategy = commitStrategyRegistry.FindWorkflowStrategy(strategyName); - } - else - { - // Fall back to the default strategy if configured - strategy = commitStateOptions.Value.DefaultWorkflowCommitStrategy; - } + IWorkflowCommitStrategy? strategy = !string.IsNullOrWhiteSpace(strategyName) + ? commitStrategyRegistry.FindWorkflowStrategy(strategyName) + : commitStateOptions.Value.DefaultWorkflowCommitStrategy; if (strategy == null) return; diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs index 6a52cc196..38fdb5210 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/Tests.cs @@ -3,15 +3,21 @@ using Elsa.Testing.Shared; using Elsa.Workflows.CommitStates; using Elsa.Workflows.CommitStates.Strategies; using Elsa.Workflows.CommitStates.Tasks; +using Elsa.Workflows.IntegrationTests.SharedHelpers; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using Xunit.Abstractions; namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultActivityCommitStrategy; -public class Tests(ITestOutputHelper testOutputHelper) +public class Tests { - private readonly ITestOutputHelper _testOutputHelper = testOutputHelper; + private readonly ITestOutputHelper _testOutputHelper; + + public Tests(ITestOutputHelper testOutputHelper) + { + _testOutputHelper = testOutputHelper; + } [Fact(DisplayName = "Activity without explicit strategy uses default commit strategy")] public async Task ActivityUsesDefaultCommitStrategy() @@ -106,7 +112,7 @@ public class Tests(ITestOutputHelper testOutputHelper) } [Fact(DisplayName = "Default activity strategy is not visible in commit strategy registry")] - public async Task DefaultStrategyNotInRegistry() + public void DefaultStrategyNotInRegistry() { // Arrange var defaultStrategy = new ExecutedActivityStrategy(); @@ -128,7 +134,6 @@ public class Tests(ITestOutputHelper testOutputHelper) Assert.Empty(activityStrategies); Assert.NotNull(options.Value.DefaultActivityCommitStrategy); Assert.Same(defaultStrategy, options.Value.DefaultActivityCommitStrategy); - await Task.CompletedTask; } [Fact(DisplayName = "Default activity strategy with standard strategies does not duplicate")] diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs deleted file mode 100644 index fbce44d7c..000000000 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/CommitTracker.cs +++ /dev/null @@ -1,24 +0,0 @@ -using Elsa.Workflows.CommitStates; -using Elsa.Workflows.State; - -namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy; - -/// -/// Tracks commit operations for testing purposes. -/// -public class CommitTracker : ICommitStateHandler -{ - public int CommitCount { get; private set; } - - public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default) - { - CommitCount++; - return Task.CompletedTask; - } - - public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default) - { - CommitCount++; - return Task.CompletedTask; - } -} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs index 32a06feda..bf71a9d37 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs @@ -2,6 +2,7 @@ using Elsa.Extensions; using Elsa.Testing.Shared; using Elsa.Workflows.CommitStates; using Elsa.Workflows.CommitStates.Strategies; +using Elsa.Workflows.IntegrationTests.SharedHelpers; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using Xunit.Abstractions; diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs b/test/integration/Elsa.Workflows.IntegrationTests/SharedHelpers/CommitTracker.cs similarity index 88% rename from test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs rename to test/integration/Elsa.Workflows.IntegrationTests/SharedHelpers/CommitTracker.cs index 1beab0576..fd26f8749 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultActivityCommitStrategy/CommitTracker.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/SharedHelpers/CommitTracker.cs @@ -1,7 +1,7 @@ using Elsa.Workflows.CommitStates; using Elsa.Workflows.State; -namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultActivityCommitStrategy; +namespace Elsa.Workflows.IntegrationTests.SharedHelpers; /// /// Test helper to track commit invocations @@ -13,12 +13,14 @@ public class CommitTracker : ICommitStateHandler public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default) { CommitCount++; + return Task.CompletedTask; } public Task CommitAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowState workflowState, CancellationToken cancellationToken = default) { CommitCount++; + return Task.CompletedTask; } } From c835f9a33d203807e5225044a85904e8e858f589 Mon Sep 17 00:00:00 2001 From: David Garza Date: Tue, 16 Dec 2025 07:22:25 -0500 Subject: [PATCH 3/7] Clarify commit strategy fallback behavior in documentation; update test method signature for consistency --- .../CommitStates/Extensions/ModuleExtensions.cs | 4 ++-- .../Scenarios/DefaultWorkflowCommitStrategy/Tests.cs | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs index db3558cd1..73464abee 100644 --- a/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/CommitStates/Extensions/ModuleExtensions.cs @@ -23,7 +23,7 @@ public static class WorkflowsFeatureCommitStateExtensions /// /// Sets the specified workflow commit strategy as the global default for all workflows that do not specify their own strategy. - /// The strategy will be automatically registered if not already present. + /// The strategy will not be added to the registry and serves only as a fallback. /// /// The workflows feature. /// The workflow commit strategy instance to use as the default. @@ -36,7 +36,7 @@ public static class WorkflowsFeatureCommitStateExtensions /// /// Sets the specified activity commit strategy as the global default for all activities that do not specify their own strategy. - /// The strategy will be automatically registered if not already present. + /// The strategy will not be added to the registry and serves only as a fallback. /// /// The workflows feature. /// The activity commit strategy instance to use as the default. diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs index bf71a9d37..9a99956ee 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/DefaultWorkflowCommitStrategy/Tests.cs @@ -116,7 +116,7 @@ public class Tests } [Fact(DisplayName = "Default workflow strategy is not visible in commit strategy registry")] - public async Task DefaultWorkflowStrategyNotInRegistry() + public void DefaultWorkflowStrategyNotInRegistry() { // Arrange var services = new TestApplicationBuilder(_testOutputHelper) From 6fc21c559028ee63bc9eb04e6e6f0aaae52fca38 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Sat, 20 Dec 2025 22:13:07 +0000 Subject: [PATCH 4/7] Initial plan From f834b040f9a4d18d0feea61823d178d943219025 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Sat, 20 Dec 2025 22:30:43 +0000 Subject: [PATCH 5/7] Add workflow dispatch notifications - Created WorkflowDefinitionDispatching and WorkflowDefinitionDispatched notifications - Created WorkflowInstanceDispatching and WorkflowInstanceDispatched notifications - Updated BackgroundWorkflowDispatcher to emit notifications before and after dispatch - Added integration tests to verify notifications are emitted correctly Co-authored-by: KnibbsyMan <23156317+KnibbsyMan@users.noreply.github.com> --- .../WorkflowDefinitionDispatched.cs | 10 +++ .../WorkflowDefinitionDispatching.cs | 9 ++ .../WorkflowInstanceDispatched.cs | 10 +++ .../WorkflowInstanceDispatching.cs | 9 ++ .../Services/BackgroundWorkflowDispatcher.cs | 23 ++++- .../WorkflowDispatchNotifications/Spy.cs | 16 ++++ .../TestHandler.cs | 46 ++++++++++ .../WorkflowDispatchNotifications/Tests.cs | 85 +++++++++++++++++++ .../Workflows.cs | 14 +++ 9 files changed, 219 insertions(+), 3 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatched.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatching.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatched.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatching.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs create mode 100644 test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Workflows.cs diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatched.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatched.cs new file mode 100644 index 000000000..d1258268c --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatched.cs @@ -0,0 +1,10 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Responses; + +namespace Elsa.Workflows.Runtime.Notifications; + +/// +/// A notification that is published when a workflow definition has been dispatched. +/// +public record WorkflowDefinitionDispatched(DispatchWorkflowDefinitionRequest Request, DispatchWorkflowResponse Response) : INotification; diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatching.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatching.cs new file mode 100644 index 000000000..f6cdb248e --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowDefinitionDispatching.cs @@ -0,0 +1,9 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Requests; + +namespace Elsa.Workflows.Runtime.Notifications; + +/// +/// A notification that is published when a workflow definition is being dispatched. +/// +public record WorkflowDefinitionDispatching(DispatchWorkflowDefinitionRequest Request) : INotification; diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatched.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatched.cs new file mode 100644 index 000000000..64690605c --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatched.cs @@ -0,0 +1,10 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Responses; + +namespace Elsa.Workflows.Runtime.Notifications; + +/// +/// A notification that is published when a workflow instance has been dispatched. +/// +public record WorkflowInstanceDispatched(DispatchWorkflowInstanceRequest Request, DispatchWorkflowResponse Response) : INotification; diff --git a/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatching.cs b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatching.cs new file mode 100644 index 000000000..0e4a6b3e7 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Notifications/WorkflowInstanceDispatching.cs @@ -0,0 +1,9 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Requests; + +namespace Elsa.Workflows.Runtime.Notifications; + +/// +/// A notification that is published when a workflow instance is being dispatched. +/// +public record WorkflowInstanceDispatching(DispatchWorkflowInstanceRequest Request) : INotification; diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs index 0af4a9278..b5b99007f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundWorkflowDispatcher.cs @@ -3,6 +3,7 @@ using Elsa.Mediator; using Elsa.Mediator.Contracts; using Elsa.Tenants.Mediator; using Elsa.Workflows.Runtime.Commands; +using Elsa.Workflows.Runtime.Notifications; using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.Runtime.Responses; @@ -11,11 +12,14 @@ namespace Elsa.Workflows.Runtime; /// /// A simple implementation that queues the specified request for workflow execution on a non-durable background worker. /// -public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher +public class BackgroundWorkflowDispatcher(ICommandSender commandSender, INotificationSender notificationSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher { /// public async Task DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { + // Emit dispatching notification + await notificationSender.SendAsync(new WorkflowDefinitionDispatching(request), cancellationToken); + var command = new DispatchWorkflowDefinitionCommand(request.DefinitionVersionId) { Input = request.Input, @@ -27,12 +31,20 @@ public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantA }; await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); - return DispatchWorkflowResponse.Success(); + var response = DispatchWorkflowResponse.Success(); + + // Emit dispatched notification + await notificationSender.SendAsync(new WorkflowDefinitionDispatched(request, response), cancellationToken); + + return response; } /// public async Task DispatchAsync(DispatchWorkflowInstanceRequest request, DispatchWorkflowOptions? options = null, CancellationToken cancellationToken = default) { + // Emit dispatching notification + await notificationSender.SendAsync(new WorkflowInstanceDispatching(request), cancellationToken); + var command = new DispatchWorkflowInstanceCommand(request.InstanceId){ BookmarkId = request.BookmarkId, ActivityHandle = request.ActivityHandle, @@ -41,7 +53,12 @@ public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantA CorrelationId = request.CorrelationId}; await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); - return DispatchWorkflowResponse.Success(); + var response = DispatchWorkflowResponse.Success(); + + // Emit dispatched notification + await notificationSender.SendAsync(new WorkflowInstanceDispatched(request, response), cancellationToken); + + return response; } /// diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs new file mode 100644 index 000000000..fbe7c2b99 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs @@ -0,0 +1,16 @@ +using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Responses; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications; + +public class Spy +{ + public bool WorkflowDefinitionDispatchingWasCalled { get; set; } + public bool WorkflowDefinitionDispatchedWasCalled { get; set; } + public bool WorkflowInstanceDispatchingWasCalled { get; set; } + public bool WorkflowInstanceDispatchedWasCalled { get; set; } + + public DispatchWorkflowDefinitionRequest? CapturedDefinitionRequest { get; set; } + public DispatchWorkflowInstanceRequest? CapturedInstanceRequest { get; set; } + public DispatchWorkflowResponse? CapturedResponse { get; set; } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs new file mode 100644 index 000000000..f2ede80cd --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs @@ -0,0 +1,46 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Runtime.Notifications; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications; + +public class TestHandler : + INotificationHandler, + INotificationHandler, + INotificationHandler, + INotificationHandler +{ + private readonly Spy _spy; + + public TestHandler(Spy spy) + { + _spy = spy; + } + + public Task HandleAsync(WorkflowDefinitionDispatching notification, CancellationToken cancellationToken) + { + _spy.WorkflowDefinitionDispatchingWasCalled = true; + _spy.CapturedDefinitionRequest = notification.Request; + return Task.CompletedTask; + } + + public Task HandleAsync(WorkflowDefinitionDispatched notification, CancellationToken cancellationToken) + { + _spy.WorkflowDefinitionDispatchedWasCalled = true; + _spy.CapturedResponse = notification.Response; + return Task.CompletedTask; + } + + public Task HandleAsync(WorkflowInstanceDispatching notification, CancellationToken cancellationToken) + { + _spy.WorkflowInstanceDispatchingWasCalled = true; + _spy.CapturedInstanceRequest = notification.Request; + return Task.CompletedTask; + } + + public Task HandleAsync(WorkflowInstanceDispatched notification, CancellationToken cancellationToken) + { + _spy.WorkflowInstanceDispatchedWasCalled = true; + _spy.CapturedResponse = notification.Response; + return Task.CompletedTask; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs new file mode 100644 index 000000000..d05f3b656 --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs @@ -0,0 +1,85 @@ +using Elsa.Testing.Shared; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Notifications; +using Elsa.Workflows.Runtime.Requests; +using Microsoft.Extensions.DependencyInjection; +using Xunit.Abstractions; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications; + +public class Tests +{ + private readonly IWorkflowDispatcher _workflowDispatcher; + private readonly Spy _spy; + + public Tests(ITestOutputHelper testOutputHelper) + { + var services = new TestApplicationBuilder(testOutputHelper) + .ConfigureServices(s => + { + s.AddSingleton(); + s.AddNotificationHandler(); + s.AddNotificationHandler(); + s.AddNotificationHandler(); + s.AddNotificationHandler(); + }) + .Build(); + + _workflowDispatcher = services.GetRequiredService(); + _spy = services.GetRequiredService(); + } + + [Fact(DisplayName = "Dispatching workflow definition should emit notifications")] + public async Task DispatchWorkflowDefinition_ShouldEmitNotifications() + { + // Arrange + var definitionVersionId = "test-definition-version-id"; + var request = new DispatchWorkflowDefinitionRequest(definitionVersionId) + { + CorrelationId = "test-correlation-id", + Input = new Dictionary { { "TestKey", "TestValue" } } + }; + + // Act + await _workflowDispatcher.DispatchAsync(request, null); + + // Allow async notification handlers to complete + await Task.Delay(100); + + // Assert + Assert.True(_spy.WorkflowDefinitionDispatchingWasCalled, "WorkflowDefinitionDispatching notification should be called"); + Assert.True(_spy.WorkflowDefinitionDispatchedWasCalled, "WorkflowDefinitionDispatched notification should be called"); + Assert.NotNull(_spy.CapturedDefinitionRequest); + Assert.Equal(definitionVersionId, _spy.CapturedDefinitionRequest.DefinitionVersionId); + Assert.Equal("test-correlation-id", _spy.CapturedDefinitionRequest.CorrelationId); + Assert.NotNull(_spy.CapturedResponse); + Assert.True(_spy.CapturedResponse.Succeeded); + } + + [Fact(DisplayName = "Dispatching workflow instance should emit notifications")] + public async Task DispatchWorkflowInstance_ShouldEmitNotifications() + { + // Arrange + var instanceId = "test-instance-id"; + var request = new DispatchWorkflowInstanceRequest(instanceId) + { + CorrelationId = "test-correlation-id", + Input = new Dictionary { { "TestKey", "TestValue" } } + }; + + // Act + await _workflowDispatcher.DispatchAsync(request, null); + + // Allow async notification handlers to complete + await Task.Delay(100); + + // Assert + Assert.True(_spy.WorkflowInstanceDispatchingWasCalled, "WorkflowInstanceDispatching notification should be called"); + Assert.True(_spy.WorkflowInstanceDispatchedWasCalled, "WorkflowInstanceDispatched notification should be called"); + Assert.NotNull(_spy.CapturedInstanceRequest); + Assert.Equal(instanceId, _spy.CapturedInstanceRequest.InstanceId); + Assert.Equal("test-correlation-id", _spy.CapturedInstanceRequest.CorrelationId); + Assert.NotNull(_spy.CapturedResponse); + Assert.True(_spy.CapturedResponse.Succeeded); + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Workflows.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Workflows.cs new file mode 100644 index 000000000..8d8ebddfe --- /dev/null +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Workflows.cs @@ -0,0 +1,14 @@ +using Elsa.Workflows.Activities; + +namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications; + +public class SimpleWorkflow : WorkflowBase +{ + public static string DefinitionId = "SimpleWorkflow"; + + protected override void Build(IWorkflowBuilder builder) + { + builder.DefinitionId = DefinitionId; + builder.Root = new WriteLine("Hello from SimpleWorkflow"); + } +} From 2ff1693cd5574c7e859e6e12dc3ee90ca9442039 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Wed, 24 Dec 2025 23:47:55 +0000 Subject: [PATCH 6/7] Replace flaky Task.Delay with TaskCompletionSource for deterministic test synchronization - Added TaskCompletionSource fields to Spy class for each notification type - Test handlers now signal completion through TaskCompletionSource - Tests await TaskCompletionSource instead of arbitrary delays - Eliminates timing-dependent flakiness in CI systems Co-authored-by: KnibbsyMan <23156317+KnibbsyMan@users.noreply.github.com> --- .../WorkflowDispatchNotifications/Spy.cs | 15 +++++++++++++++ .../WorkflowDispatchNotifications/TestHandler.cs | 4 ++++ .../WorkflowDispatchNotifications/Tests.cs | 10 ++++++---- 3 files changed, 25 insertions(+), 4 deletions(-) diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs index fbe7c2b99..d94592bc0 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Spy.cs @@ -5,6 +5,11 @@ namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotification public class Spy { + private readonly TaskCompletionSource _workflowDefinitionDispatchingTcs = new(); + private readonly TaskCompletionSource _workflowDefinitionDispatchedTcs = new(); + private readonly TaskCompletionSource _workflowInstanceDispatchingTcs = new(); + private readonly TaskCompletionSource _workflowInstanceDispatchedTcs = new(); + public bool WorkflowDefinitionDispatchingWasCalled { get; set; } public bool WorkflowDefinitionDispatchedWasCalled { get; set; } public bool WorkflowInstanceDispatchingWasCalled { get; set; } @@ -13,4 +18,14 @@ public class Spy public DispatchWorkflowDefinitionRequest? CapturedDefinitionRequest { get; set; } public DispatchWorkflowInstanceRequest? CapturedInstanceRequest { get; set; } public DispatchWorkflowResponse? CapturedResponse { get; set; } + + public Task WaitForWorkflowDefinitionDispatchingAsync() => _workflowDefinitionDispatchingTcs.Task; + public Task WaitForWorkflowDefinitionDispatchedAsync() => _workflowDefinitionDispatchedTcs.Task; + public Task WaitForWorkflowInstanceDispatchingAsync() => _workflowInstanceDispatchingTcs.Task; + public Task WaitForWorkflowInstanceDispatchedAsync() => _workflowInstanceDispatchedTcs.Task; + + public void SignalWorkflowDefinitionDispatching() => _workflowDefinitionDispatchingTcs.TrySetResult(true); + public void SignalWorkflowDefinitionDispatched() => _workflowDefinitionDispatchedTcs.TrySetResult(true); + public void SignalWorkflowInstanceDispatching() => _workflowInstanceDispatchingTcs.TrySetResult(true); + public void SignalWorkflowInstanceDispatched() => _workflowInstanceDispatchedTcs.TrySetResult(true); } diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs index f2ede80cd..9d90a4227 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/TestHandler.cs @@ -20,6 +20,7 @@ public class TestHandler : { _spy.WorkflowDefinitionDispatchingWasCalled = true; _spy.CapturedDefinitionRequest = notification.Request; + _spy.SignalWorkflowDefinitionDispatching(); return Task.CompletedTask; } @@ -27,6 +28,7 @@ public class TestHandler : { _spy.WorkflowDefinitionDispatchedWasCalled = true; _spy.CapturedResponse = notification.Response; + _spy.SignalWorkflowDefinitionDispatched(); return Task.CompletedTask; } @@ -34,6 +36,7 @@ public class TestHandler : { _spy.WorkflowInstanceDispatchingWasCalled = true; _spy.CapturedInstanceRequest = notification.Request; + _spy.SignalWorkflowInstanceDispatching(); return Task.CompletedTask; } @@ -41,6 +44,7 @@ public class TestHandler : { _spy.WorkflowInstanceDispatchedWasCalled = true; _spy.CapturedResponse = notification.Response; + _spy.SignalWorkflowInstanceDispatched(); return Task.CompletedTask; } } diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs index d05f3b656..f92bd1cbb 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/WorkflowDispatchNotifications/Tests.cs @@ -43,8 +43,9 @@ public class Tests // Act await _workflowDispatcher.DispatchAsync(request, null); - // Allow async notification handlers to complete - await Task.Delay(100); + // Wait for notification handlers to complete + await _spy.WaitForWorkflowDefinitionDispatchingAsync(); + await _spy.WaitForWorkflowDefinitionDispatchedAsync(); // Assert Assert.True(_spy.WorkflowDefinitionDispatchingWasCalled, "WorkflowDefinitionDispatching notification should be called"); @@ -70,8 +71,9 @@ public class Tests // Act await _workflowDispatcher.DispatchAsync(request, null); - // Allow async notification handlers to complete - await Task.Delay(100); + // Wait for notification handlers to complete + await _spy.WaitForWorkflowInstanceDispatchingAsync(); + await _spy.WaitForWorkflowInstanceDispatchedAsync(); // Assert Assert.True(_spy.WorkflowInstanceDispatchingWasCalled, "WorkflowInstanceDispatching notification should be called"); From ca268c16addf5ed381eaa06046aeb3770d7d61a0 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Mon, 29 Dec 2025 19:53:48 +0100 Subject: [PATCH 7/7] Enhances distributed lock handling with resilience (#7161) * Add retry mechanism for distributed locks with transient error handling and logging - Introduced Polly-based retry pipeline for distributed lock acquisition in `DistributedWorkflowClient` to handle transient errors such as network issues or database connection failures. - Added detailed logging for retry attempts and lock release errors. - Updated project dependencies to include Polly. * Refactor transient exception handling to shared resilience module. Migrated transient exception detection logic from scheduling module to a new shared resilience module. Updated services, jobs, and features to utilize the centralized `ITransientExceptionDetectionService`. This change improves maintainability and promotes reusability across modules. * Add unit tests for transient exception detection and resilience strategy evaluation. - Introduced comprehensive unit tests for `DefaultTransientExceptionDetector`, `ResilienceStrategyCatalog`, `ResilienceStrategyConfigEvaluator`, and `TransientExceptionDetectionService`. - Added helper classes and test data factories to facilitate reusable test patterns for resilience modules. - Updated solution to include `Elsa.Resilience.Core.UnitTests` project. * Add component tests for distributed lock resilience - Introduced new tests to verify retry behavior during transient lock acquisition and release failures. - Added `TestDistributedLockProvider` and related mocks for simulating transient failures. - Updated `WorkflowServer` test services to support the new distributed lock test scenarios. * Refactor distributed lock resilience tests - Consolidated test logic: streamlined test providers, injected services, and reusable test patterns. - Simplified `TestDistributedLockProvider` implementation with enhanced initialization and failure simulation. - Reorganized tests for transient acquisition/release failures to use parameterized `Theory` for improved maintainability. * Refactor transient exception handling: rename interfaces and classes for consistency, update references across codebase, and improve code readability. * Apply suggestion from @Copilot Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Simplify WorkflowServer setup and DistributedLockResilienceTests by replacing IDistributedLockProvider with TestDistributedLockProvider. * Remove unused `using` directives in unit tests to improve code cleanliness. * Remove `TransientExceptionTypes` helper and inline its usage in tests for improved maintainability. * Add descriptive `DisplayName` attributes to unit tests for improved test clarity. * Fix redundant exception checking in TransientExceptionDetector (#7162) * Initial plan * Fix redundant exception checking in TransientExceptionDetector Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Extract MaxRetryAttempts constant in DistributedLockResilienceTests (#7164) * Initial plan * Extract MaxRetryAttempts constant to eliminate hardcoded magic number Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Make TestDistributedLockProvider thread-safe with Interlocked operations (#7163) * Initial plan * Make TestDistributedLockProvider thread-safe using Interlocked operations Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/DistributedLockResilienceTests.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Resilience.Core/Services/DefaultTransientExceptionStrategy.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Include `CancellationToken` in distributed lock handling methods for improved cancellation support. * Initial plan (#7168) Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> * Cache detector list in TransientExceptionDetector to avoid repeated allocations (#7167) * Initial plan * Cache detector list in field to avoid repeated allocations Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Use IReadOnlyList instead of List for better intent expression Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Fix TestDistributedLockProvider registration to properly decorate IDistributedLockProvider (#7166) * Initial plan * Fix TestDistributedLockProvider registration to use Decorate pattern and fix variable reference bug Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Add runtime check for TestDistributedLockProvider registration Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Refactor `DistributedWorkflowClient` to simplify `Lazy` initialization. * Add integration tests for DistributedWorkflowClient lock resilience (#7165) * Initial plan * Fix compilation error: use correct parameter name transientExceptionDetector Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Add integration tests for DistributedWorkflowClient lock resilience - Add SimpleWorkflow for testing distributed lock scenarios - Add tests exercising RunInstanceAsync with transient lock failures - Verify retry logic works correctly with actual workflow execution - Test both acquisition and release failure scenarios - Decorate IDistributedLockProvider to use TestDistributedLockProvider Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Address code review feedback - Add explanatory comment for TestDistributedLockProvider cast - Remove unnecessary blank line for consistent formatting Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> Co-authored-by: Sipke Schoorstra * Simplify transient exception strategy by refactoring message pattern matching logic. * Enhance distributed lock mock to support per-lock failure configuration and improve resilience tests. * Refactor `TestDistributedLockProvider` to streamline failure handling logic and improve code clarity. * Remove unused methods and redundant test case from `DistributedLockResilienceTests`. * Refactor `DistributedLockResilienceTests` to simplify workflow client creation, consolidate assertion logic, and remove redundant test cases. * Format `ResilienceStrategyCatalogTests` by removing redundant line breaks in test setup. * Refactor `TransientExceptionDetectorTests` to simplify test setup, consolidate test cases, and remove redundant logic. * Handle `InvalidOperationException` in `XunitLogger` to suppress logging errors during inactive tests. * Update workflows to use .NET 10 and adjust resilience tests project configuration. * Refactor distributed runtime feature and integrations to improve resilience handling, configure services fluently, and add cancellation safeguards in background services. * Update `base_version` to `3.7.0` in GitHub workflow configuration. --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> Co-authored-by: Copilot <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --- .github/workflows/packages.yml | 4 +- Elsa.sln | 7 + .../BackgroundCommandSenderHostedService.cs | 64 ++++--- .../BackgroundEventPublisherHostedService.cs | 52 ++++-- src/common/Elsa.Testing.Shared/XunitLogger.cs | 14 +- .../Contracts/ITransientExceptionDetector.cs | 14 ++ .../Contracts/ITransientExceptionStrategy.cs | 14 ++ .../DefaultTransientExceptionStrategy.cs | 65 +++++++ .../Services/TransientExceptionDetector.cs | 42 +++++ .../Features/ResilienceFeature.cs | 5 + .../Elsa.Workflows.Runtime.Distributed.csproj | 4 +- .../Features/DistributedRuntimeFeature.cs | 9 +- .../Services/DistributedWorkflowClient.cs | 79 ++++++++- .../Helpers/Fixtures/WorkflowServer.cs | 15 ++ .../DistributedLockResilienceTests.cs | 165 ++++++++++++++++++ .../Mocks/TestDistributedLock.cs | 56 ++++++ .../Mocks/TestDistributedLockProvider.cs | 84 +++++++++ .../TestDistributedSynchronizationHandle.cs | 31 ++++ .../Workflows/SimpleWorkflow.cs | 26 +++ .../RunAsynchronousActivityOutput/Tests.cs | 12 +- .../DefaultTransientExceptionStrategyTests.cs | 87 +++++++++ .../Elsa.Resilience.Core.UnitTests.csproj | 13 ++ .../ResilienceStrategyCatalogTests.cs | 118 +++++++++++++ .../ResilienceStrategyConfigEvaluatorTests.cs | 135 ++++++++++++++ .../TestHelpers/TestDataFactory.cs | 21 +++ .../TransientExceptionDetectorTests.cs | 151 ++++++++++++++++ 26 files changed, 1214 insertions(+), 73 deletions(-) create mode 100644 src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionDetector.cs create mode 100644 src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionStrategy.cs create mode 100644 src/modules/Elsa.Resilience.Core/Services/DefaultTransientExceptionStrategy.cs create mode 100644 src/modules/Elsa.Resilience.Core/Services/TransientExceptionDetector.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/DistributedLockResilienceTests.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLock.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLockProvider.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedSynchronizationHandle.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Workflows/SimpleWorkflow.cs create mode 100644 test/unit/Elsa.Resilience.Core.UnitTests/DefaultTransientExceptionStrategyTests.cs create mode 100644 test/unit/Elsa.Resilience.Core.UnitTests/Elsa.Resilience.Core.UnitTests.csproj create mode 100644 test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyCatalogTests.cs create mode 100644 test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyConfigEvaluatorTests.cs create mode 100644 test/unit/Elsa.Resilience.Core.UnitTests/TestHelpers/TestDataFactory.cs create mode 100644 test/unit/Elsa.Resilience.Core.UnitTests/TransientExceptionDetectorTests.cs diff --git a/.github/workflows/packages.yml b/.github/workflows/packages.yml index d78513c43..7b58b55c5 100644 --- a/.github/workflows/packages.yml +++ b/.github/workflows/packages.yml @@ -17,7 +17,7 @@ on: types: [prereleased, published] env: - base_version: '3.6.0' + base_version: '3.7.0' feedz_feed_source: 'https://f.feedz.io/elsa-workflows/elsa-3/nuget/index.json' nuget_feed_source: 'https://api.nuget.org/v3/index.json' @@ -155,7 +155,7 @@ jobs: - uses: actions/setup-dotnet@v4 with: - dotnet-version: 9.x + dotnet-version: 10.x - name: Compile+Pack run: ./build.sh Compile+Pack --version ${VERSION} --analyseCode true diff --git a/Elsa.sln b/Elsa.sln index 7a77eb07c..c1e044cb0 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -323,6 +323,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "issue_templates", "issue_te EndProject Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "dsl", "dsl", "{477C2416-312D-46AE-BCD6-8FA1FAB43624}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Resilience.Core.UnitTests", "test\unit\Elsa.Resilience.Core.UnitTests\Elsa.Resilience.Core.UnitTests.csproj", "{B8006D70-1630-43DB-A043-FA89FAC70F37}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -579,6 +581,10 @@ Global {2B7FB49D-E4B6-4AD5-981B-3D85B94F6F48}.Debug|Any CPU.Build.0 = Debug|Any CPU {2B7FB49D-E4B6-4AD5-981B-3D85B94F6F48}.Release|Any CPU.ActiveCfg = Release|Any CPU {2B7FB49D-E4B6-4AD5-981B-3D85B94F6F48}.Release|Any CPU.Build.0 = Release|Any CPU + {B8006D70-1630-43DB-A043-FA89FAC70F37}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {B8006D70-1630-43DB-A043-FA89FAC70F37}.Debug|Any CPU.Build.0 = Debug|Any CPU + {B8006D70-1630-43DB-A043-FA89FAC70F37}.Release|Any CPU.ActiveCfg = Release|Any CPU + {B8006D70-1630-43DB-A043-FA89FAC70F37}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -680,6 +686,7 @@ Global {2B7FB49D-E4B6-4AD5-981B-3D85B94F6F48} = {B08B4E00-C2AB-48F3-8389-449F42AEF179} {477C2416-312D-46AE-BCD6-8FA1FAB43624} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} {874F5A44-DB06-47AB-A18C-2D13942E0147} = {477C2416-312D-46AE-BCD6-8FA1FAB43624} + {B8006D70-1630-43DB-A043-FA89FAC70F37} = {18453B51-25EB-4317-A4B3-B10518252E92} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {D4B5CEAA-7D70-4FCB-A68E-B03FBE5E0E5E} diff --git a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs index f0606a1f5..26916c831 100644 --- a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs @@ -45,12 +45,19 @@ public class BackgroundCommandSenderHostedService : BackgroundService } // Main dispatcher loop: read from the input channel and distribute to worker channels - await foreach (var commandContext in _commandsChannel.Reader.ReadAllAsync(cancellationToken)) + try { - var output = _outputs[index]; - await output.Writer.WriteAsync(commandContext, cancellationToken); - // Round-robin distribution - move to next worker - index = (index + 1) % _workerCount; + await foreach (var commandContext in _commandsChannel.Reader.ReadAllAsync(cancellationToken)) + { + var output = _outputs[index]; + await output.Writer.WriteAsync(commandContext, cancellationToken); + // Round-robin distribution - move to next worker + index = (index + 1) % _workerCount; + } + } + catch (OperationCanceledException ex) + { + _logger.LogDebug(ex, "An operation was cancelled while processing the queue"); } // If the input channel is completed, complete all worker channels @@ -61,31 +68,38 @@ public class BackgroundCommandSenderHostedService : BackgroundService private async Task ReadOutputAsync(Channel output, CancellationToken cancellationToken) { // Worker task: process commands from the worker's channel - await foreach (var commandContext in output.Reader.ReadAllAsync(cancellationToken)) + try { - try + await foreach (var commandContext in output.Reader.ReadAllAsync(cancellationToken)) { - // Create a fresh scope for each command to ensure proper service lifetime - using var scope = _scopeFactory.CreateScope(); - var commandSender = scope.ServiceProvider.GetRequiredService(); + try + { + // Create a fresh scope for each command to ensure proper service lifetime + using var scope = _scopeFactory.CreateScope(); + var commandSender = scope.ServiceProvider.GetRequiredService(); - // Link the service cancellation token with the command's token to ensure proper cancellation - using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource( - cancellationToken, - commandContext.CancellationToken); + // Link the service cancellation token with the command's token to ensure proper cancellation + using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken, + commandContext.CancellationToken); - // Process the command using the command sender service with the linked token - await commandSender.SendAsync( - commandContext.Command, - CommandStrategy.Default, - commandContext.Headers, - linkedTokenSource.Token); - } - catch (Exception e) - { - // Log errors but continue processing other commands - _logger.LogError(e, "An unhandled exception occurred while processing the queue"); + // Process the command using the command sender service with the linked token + await commandSender.SendAsync( + commandContext.Command, + CommandStrategy.Default, + commandContext.Headers, + linkedTokenSource.Token); + } + catch (Exception e) + { + // Log errors but continue processing other commands + _logger.LogError(e, "An unhandled exception occurred while processing the queue"); + } } } + catch (OperationCanceledException ex) + { + _logger.LogDebug(ex, "An operation was cancelled while processing the queue"); + } } } \ No newline at end of file diff --git a/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs b/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs index 00c29bba1..9bb516bdf 100644 --- a/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs @@ -52,12 +52,19 @@ public class BackgroundEventPublisherHostedService : BackgroundService // Continuously read notifications from the input channel and distribute them to worker channels // using round-robin distribution for load balancing - await foreach (var notification in channelReader.ReadAllAsync(cancellationToken)) + try { - var output = _outputs[index]; - await output.Writer.WriteAsync(notification, cancellationToken); - // Move to the next worker in a circular fashion - index = (index + 1) % _workerCount; + await foreach (var notification in channelReader.ReadAllAsync(cancellationToken)) + { + var output = _outputs[index]; + await output.Writer.WriteAsync(notification, cancellationToken); + // Move to the next worker in a circular fashion + index = (index + 1) % _workerCount; + } + } + catch (OperationCanceledException ex) + { + _logger.LogDebug(ex, "An operation was cancelled while processing the queue"); } // When the input channel is completed, complete all output channels @@ -75,23 +82,30 @@ public class BackgroundEventPublisherHostedService : BackgroundService /// Cancellation token from the hosted service private async Task ReadOutputAsync(Channel output, INotificationSender notificationSender, CancellationToken cancellationToken) { - await foreach (var notificationContext in output.Reader.ReadAllAsync(cancellationToken)) + try { - try + await foreach (var notificationContext in output.Reader.ReadAllAsync(cancellationToken)) { - var notification = notificationContext.Notification; - // Link the cancellation tokens so that cancellation can happen from either source - using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, notificationContext.CancellationToken); - await notificationSender.SendAsync(notification, NotificationStrategy.Sequential, linkedTokenSource.Token); - } - catch (OperationCanceledException e) - { - _logger.LogDebug(e, "An operation was cancelled while processing the queue"); - } - catch (Exception e) - { - _logger.LogError(e, "An unhandled exception occurred while processing the queue"); + try + { + var notification = notificationContext.Notification; + // Link the cancellation tokens so that cancellation can happen from either source + using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, notificationContext.CancellationToken); + await notificationSender.SendAsync(notification, NotificationStrategy.Sequential, linkedTokenSource.Token); + } + catch (OperationCanceledException e) + { + _logger.LogDebug(e, "An operation was cancelled while processing the queue"); + } + catch (Exception e) + { + _logger.LogError(e, "An unhandled exception occurred while processing the queue"); + } } } + catch (OperationCanceledException ex) + { + _logger.LogDebug(ex, "An operation was cancelled while processing the queue"); + } } } \ No newline at end of file diff --git a/src/common/Elsa.Testing.Shared/XunitLogger.cs b/src/common/Elsa.Testing.Shared/XunitLogger.cs index e02d37fb5..9d0e033f1 100644 --- a/src/common/Elsa.Testing.Shared/XunitLogger.cs +++ b/src/common/Elsa.Testing.Shared/XunitLogger.cs @@ -12,10 +12,18 @@ public class XunitLogger(ITestOutputHelper testOutputHelper, string categoryName public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) { - testOutputHelper.WriteLine($"{categoryName} [{eventId}] {formatter(state, exception)}"); + try + { + testOutputHelper.WriteLine($"{categoryName} [{eventId}] {formatter(state, exception)}"); - if (exception != null) - testOutputHelper.WriteLine(exception.ToString()); + if (exception != null) + testOutputHelper.WriteLine(exception.ToString()); + } + catch (InvalidOperationException) + { + // Suppress "no currently active test" exceptions that can occur when background tasks + // (like timers) try to log after tests have completed + } } private class NoopDisposable : IDisposable diff --git a/src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionDetector.cs b/src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionDetector.cs new file mode 100644 index 000000000..1a0623422 --- /dev/null +++ b/src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionDetector.cs @@ -0,0 +1,14 @@ +namespace Elsa.Resilience; + +/// +/// Service for detecting whether exceptions are transient and may be resolved by retrying. +/// +public interface ITransientExceptionDetector +{ + /// + /// Determines whether the specified exception is transient. + /// + /// The exception to check. + /// True if the exception is transient; otherwise, false. + bool IsTransient(Exception exception); +} diff --git a/src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionStrategy.cs b/src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionStrategy.cs new file mode 100644 index 000000000..0cbfcb039 --- /dev/null +++ b/src/modules/Elsa.Resilience.Core/Contracts/ITransientExceptionStrategy.cs @@ -0,0 +1,14 @@ +namespace Elsa.Resilience; + +/// +/// Defines a contract for detecting whether an exception is transient and may be resolved by retrying. +/// +public interface ITransientExceptionStrategy +{ + /// + /// Determines whether the specified exception is transient. + /// + /// The exception to check. + /// True if the exception is transient; otherwise, false. + bool IsTransient(Exception exception); +} diff --git a/src/modules/Elsa.Resilience.Core/Services/DefaultTransientExceptionStrategy.cs b/src/modules/Elsa.Resilience.Core/Services/DefaultTransientExceptionStrategy.cs new file mode 100644 index 000000000..79f85a487 --- /dev/null +++ b/src/modules/Elsa.Resilience.Core/Services/DefaultTransientExceptionStrategy.cs @@ -0,0 +1,65 @@ +namespace Elsa.Resilience; + +/// +/// Default implementation that detects common transient exceptions from the .NET framework and common patterns. +/// +public class DefaultTransientExceptionStrategy : ITransientExceptionStrategy +{ + private static readonly HashSet TransientExceptionTypeNames = new(StringComparer.OrdinalIgnoreCase) + { + // Common framework exceptions + "HttpRequestException", + "TimeoutException", + "TaskCanceledException", + "IOException", + "SocketException", + "EndOfStreamException", + + // Database-related transient exceptions (by name, not type reference) + "DbException", + "SqlException", + "NpgsqlException", + "MongoConnectionException", + "MongoExecutionTimeoutException", + "MongoNodeIsRecoveringException", + "MongoNotPrimaryException", + "MySqlException", + + // Network-related exceptions + "HttpIOException", + "WebException", + }; + + private static readonly HashSet TransientExceptionMessagePatterns = new(StringComparer.OrdinalIgnoreCase) + { + "timeout", + "timed out", + "connection reset", + "connection refused", + "broken pipe", + "network", + "end of stream", + "attempted to read past the end", + "the connection is closed", + "connection is not open", + "failed to connect", + "no connection could be made", + "an existing connection was forcibly closed", + }; + + /// + public bool IsTransient(Exception exception) + { + // Check if the exception type name matches any known transient exception + var exceptionTypeName = exception.GetType().Name; + if (TransientExceptionTypeNames.Contains(exceptionTypeName)) + return true; + + // Check if the exception message contains any transient patterns + var message = exception.Message; + + return !string.IsNullOrEmpty(message) && + TransientExceptionMessagePatterns + .Any(pattern => message.Contains(pattern, StringComparison.OrdinalIgnoreCase)); + } +} diff --git a/src/modules/Elsa.Resilience.Core/Services/TransientExceptionDetector.cs b/src/modules/Elsa.Resilience.Core/Services/TransientExceptionDetector.cs new file mode 100644 index 000000000..35c57ba1d --- /dev/null +++ b/src/modules/Elsa.Resilience.Core/Services/TransientExceptionDetector.cs @@ -0,0 +1,42 @@ +namespace Elsa.Resilience; + +/// +/// Default implementation of that delegates to registered detectors. +/// +public class TransientExceptionDetector(IEnumerable detectors) : ITransientExceptionDetector +{ + private readonly IReadOnlyList _detectors = detectors.ToList(); + + /// + public bool IsTransient(Exception exception) + { + var detectorsList = _detectors; + + // Handle aggregate exceptions specially to avoid redundant checks + if (exception is AggregateException aggregateException) + { + // Check the aggregate exception itself + if (detectorsList.Any(detector => detector.IsTransient(aggregateException))) + return true; + + // Recursively check each inner exception (this will walk their chains) + if (aggregateException.InnerExceptions.Any(IsTransient)) + return true; + + return false; + } + + // Walk the exception chain for non-aggregate exceptions + var currentException = exception; + while (currentException != null) + { + // Check if any detector identifies this exception as transient + if (detectorsList.Any(detector => detector.IsTransient(currentException))) + return true; + + currentException = currentException.InnerException; + } + + return false; + } +} diff --git a/src/modules/Elsa.Resilience/Features/ResilienceFeature.cs b/src/modules/Elsa.Resilience/Features/ResilienceFeature.cs index fe5f8f64e..cefc31421 100644 --- a/src/modules/Elsa.Resilience/Features/ResilienceFeature.cs +++ b/src/modules/Elsa.Resilience/Features/ResilienceFeature.cs @@ -84,5 +84,10 @@ public class ResilienceFeature(IModule module) : FeatureBase(module) .AddScoped(_retryAttemptRecorder) .AddScoped(_retryAttemptReader) .AddHandlersFrom(); + + // Register transient exception detection infrastructure + Services + .AddSingleton() + .AddSingleton(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj b/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj index 0439f66b4..02aedc10f 100644 --- a/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj @@ -10,11 +10,13 @@ - + + + diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs index d47e0d75d..e783dd43e 100644 --- a/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs @@ -2,6 +2,7 @@ using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; +using Elsa.Resilience.Features; using Elsa.Workflows.Runtime.Features; using Microsoft.Extensions.DependencyInjection; @@ -11,13 +12,9 @@ namespace Elsa.Workflows.Runtime.Distributed.Features; /// Installs and configures workflow runtime features. /// [DependsOn(typeof(WorkflowRuntimeFeature))] -public class DistributedRuntimeFeature : FeatureBase +[DependsOn(typeof(ResilienceFeature))] +public class DistributedRuntimeFeature(IModule module) : FeatureBase(module) { - /// - public DistributedRuntimeFeature(IModule module) : base(module) - { - } - public override void Configure() { Module.UseWorkflowRuntime(runtime => diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs index 80aa11d5a..131c71e41 100644 --- a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs @@ -1,21 +1,26 @@ using Elsa.Common.DistributedHosting; +using Elsa.Resilience; using Elsa.Workflows.Runtime.Messages; using Elsa.Workflows.State; using Medallion.Threading; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; +using Polly; namespace Elsa.Workflows.Runtime.Distributed; public class DistributedWorkflowClient( string workflowInstanceId, IDistributedLockProvider distributedLockProvider, + ITransientExceptionDetector transientExceptionDetector, IOptions distributedLockingOptions, - IServiceProvider serviceProvider) + IServiceProvider serviceProvider, + ILogger logger) : IWorkflowClient { private readonly LocalWorkflowClient _localWorkflowClient = ActivatorUtilities.CreateInstance(serviceProvider, workflowInstanceId); - + private readonly Lazy _retryPipeline = new(() => CreateRetryPipeline(transientExceptionDetector, logger, workflowInstanceId)); public string WorkflowInstanceId => workflowInstanceId; public async Task CreateInstanceAsync(CreateWorkflowInstanceRequest request, CancellationToken cancellationToken = default) @@ -25,7 +30,7 @@ public class DistributedWorkflowClient( public async Task RunInstanceAsync(RunWorkflowInstanceRequest request, CancellationToken cancellationToken = default) { - var result = await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(request, cancellationToken)); + var result = await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(request, cancellationToken), cancellationToken); return result; } @@ -52,7 +57,7 @@ public class DistributedWorkflowClient( TriggerActivityId = request.TriggerActivityId, ActivityHandle = request.ActivityHandle, IncludeWorkflowOutput = request.IncludeWorkflowOutput - }, cancellationToken)); + }, cancellationToken), cancellationToken); } public async Task CancelAsync(CancellationToken cancellationToken = default) @@ -78,15 +83,71 @@ public class DistributedWorkflowClient( public async Task DeleteAsync(CancellationToken cancellationToken = default) { // Use the same distributed lock as for execution to prevent concurrent DB writes - return await WithLockAsync(async () => await _localWorkflowClient.DeleteAsync(cancellationToken)); + return await WithLockAsync(async () => await _localWorkflowClient.DeleteAsync(cancellationToken), cancellationToken); } - private async Task WithLockAsync(Func> func) + private async Task WithLockAsync(Func> func, CancellationToken cancellationToken = default) { var lockKey = $"workflow-instance:{WorkflowInstanceId}"; + var lockHandle = await AcquireLockWithRetryAsync(lockKey, cancellationToken); + + try + { + return await func(); + } + finally + { + await ReleaseLockAsync(lockHandle); + } + } + + private async Task AcquireLockWithRetryAsync(string lockKey, CancellationToken cancellationToken = default) + { var lockTimeout = distributedLockingOptions.Value.LockAcquisitionTimeout; - await using var @lock = await distributedLockProvider.AcquireLockAsync(lockKey, lockTimeout); - var result = await func(); - return result; + + return await _retryPipeline.Value.ExecuteAsync(async ct => + await distributedLockProvider.AcquireLockAsync(lockKey, lockTimeout, ct), + cancellationToken); + } + + private async Task ReleaseLockAsync(IDistributedSynchronizationHandle? lockHandle) + { + if (lockHandle == null) + return; + + try + { + await lockHandle.DisposeAsync(); + } + catch (Exception ex) + { + // Log but don't throw - the work is already done, and the lock + // will be automatically released when the connection dies + logger.LogWarning(ex, "Failed to release distributed lock for workflow instance {WorkflowInstanceId}. The lock will be automatically released by the database.", WorkflowInstanceId); + } + } + + private static ResiliencePipeline CreateRetryPipeline( + ITransientExceptionDetector transientExceptionDetector, + ILogger logger, + string workflowInstanceId) + { + const int maxRetryAttempts = 3; + + return new ResiliencePipelineBuilder() + .AddRetry(new() + { + MaxRetryAttempts = maxRetryAttempts, + Delay = TimeSpan.FromMilliseconds(500), + BackoffType = DelayBackoffType.Exponential, + UseJitter = true, + ShouldHandle = new PredicateBuilder().Handle(transientExceptionDetector.IsTransient), + OnRetry = args => + { + logger.LogWarning(args.Outcome.Exception, "Transient error acquiring lock for workflow instance {WorkflowInstanceId}. Attempt {AttemptNumber} of {MaxAttempts}.", workflowInstanceId, args.AttemptNumber + 1, maxRetryAttempts); + return ValueTask.CompletedTask; + } + }) + .Build(); } } \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs index 1a10ecba1..7446dbe2e 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -12,11 +12,13 @@ using Elsa.Testing.Shared.Handlers; using Elsa.Testing.Shared.Services; using Elsa.Workflows.ComponentTests.Decorators; using Elsa.Workflows.ComponentTests.Materializers; +using Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks; using Elsa.Workflows.ComponentTests.WorkflowProviders; using Elsa.Workflows.Management; using Elsa.Workflows.Runtime.Distributed.Extensions; using FluentStorage; using JetBrains.Annotations; +using Medallion.Threading; using Microsoft.AspNetCore.Hosting; using Microsoft.AspNetCore.Mvc.Testing; using Microsoft.AspNetCore.TestHost; @@ -127,6 +129,18 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl builder.ConfigureTestServices(services => { + // Decorate IDistributedLockProvider with TestDistributedLockProvider so tests use it + services.Decorate(); + + // Also register TestDistributedLockProvider as itself so tests can access it directly for configuration + services.AddSingleton(sp => + { + var provider = sp.GetRequiredService(); + if (provider is not TestDistributedLockProvider testProvider) + throw new InvalidOperationException($"Expected IDistributedLockProvider to be decorated with TestDistributedLockProvider, but got {provider.GetType().Name}"); + return testProvider; + }); + services .AddSingleton() .AddScoped() @@ -138,6 +152,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl .AddWorkflowsProvider() .AddNotificationHandlersFrom() .Decorate() + .Decorate() ; }); } diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/DistributedLockResilienceTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/DistributedLockResilienceTests.cs new file mode 100644 index 000000000..efa063902 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/DistributedLockResilienceTests.cs @@ -0,0 +1,165 @@ +using Elsa.Common.DistributedHosting; +using Elsa.Common.Models; +using Elsa.Resilience; +using Elsa.Workflows.ComponentTests.Abstractions; +using Elsa.Workflows.ComponentTests.Fixtures; +using Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks; +using Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Workflows; +using Elsa.Workflows.Models; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Messages; +using Medallion.Threading; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; +using Polly; + +namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience; + +public class DistributedLockResilienceTests(App app) : AppComponentTest(app) +{ + private const int MaxRetryAttempts = 3; + + // The IDistributedLockProvider is decorated with TestDistributedLockProvider in WorkflowServer.ConfigureTestServices + // This cast is safe because the decorator pattern ensures TestDistributedLockProvider wraps the actual provider + private TestDistributedLockProvider MockProvider => (TestDistributedLockProvider)Scope.ServiceProvider.GetRequiredService(); + private ITransientExceptionDetector TransientExceptionDetector => Scope.ServiceProvider.GetRequiredService(); + private ILogger Logger => Scope.ServiceProvider.GetRequiredService>(); + private DistributedLockingOptions LockOptions => Scope.ServiceProvider.GetRequiredService>().Value; + private ResiliencePipeline RetryPipeline => CreateRetryPipeline(TransientExceptionDetector, Logger); + + [Theory] + [InlineData(1, 2, false)] // Single failure, succeeds on retry + [InlineData(2, 3, false)] // Two failures, succeeds on third attempt + [InlineData(4, 4, true)] // Four failures, exhausts retries (MaxRetryAttempts = 3) + public async Task AcquireLockWithRetry_AcquisitionFailures_BehavesAsExpected(int failureCount, int expectedAttemptCount, bool shouldThrow) + { + // Arrange + MockProvider.Reset(); + MockProvider.FailAcquisitionTimes(failureCount); + + // Act & Assert + if (shouldThrow) + { + await Assert.ThrowsAsync(async () => await AcquireLockWithRetryAsync($"test-lock-{failureCount}")); + } + else + { + await using var handle = await AcquireLockWithRetryAsync($"test-lock-{failureCount}"); + Assert.NotNull(handle); + } + + Assert.Equal(expectedAttemptCount, MockProvider.AcquisitionAttemptCount); + } + + [Theory] + [InlineData(1, false)] // Single failure, succeeds on retry (2 attempts) + [InlineData(2, false)] // Two failures, succeeds on third attempt (3 attempts) + [InlineData(4, true)] // Four failures, exhausts retries (MaxRetryAttempts = 3, so 4 attempts total) + public async Task RunInstanceAsync_TransientLockFailures_RetriesCorrectly(int failureCount, bool shouldThrow) + { + // Arrange + var workflowClient = await CreateWorkflowClientAsync(); + var workflowInstanceId = workflowClient.WorkflowInstanceId; + + // Reset and configure failures for this specific workflow instance's lock + MockProvider.Reset(); + MockProvider.FailAcquisitionTimesForLock($"workflow-instance:{workflowInstanceId}", failureCount); + var attemptCountBefore = MockProvider.AcquisitionAttemptCount; + + // Now run the instance with the configured lock failures + var runRequest = new RunWorkflowInstanceRequest(); + + // Act & Assert + if (shouldThrow) + { + await Assert.ThrowsAsync(async () => + await workflowClient.RunInstanceAsync(runRequest)); + } + else + { + var response = await workflowClient.RunInstanceAsync(runRequest); + Assert.NotNull(response); + } + + // Verify retries occurred - check the delta from before the operation to account for background noise + var expectedAttempts = failureCount + 1; // failures + 1 success (or final failure for shouldThrow case) + AssertMinimumAttempts(MockProvider.AcquisitionAttemptCount - attemptCountBefore, expectedAttempts, "acquisition"); + } + + [Fact] + public async Task RunInstanceAsync_TransientReleaseFailure_ShouldLogButNotThrow() + { + // Arrange + MockProvider.Reset(); + var workflowClient = await CreateWorkflowClientAsync(createInstance: false); + + // Configure failure after client creation to minimize background interference + MockProvider.FailReleaseOnce(); + var releaseCountBefore = MockProvider.ReleaseAttemptCount; + + // Act - Release failure should be caught and logged, not thrown + var response = await workflowClient.CreateAndRunInstanceAsync(CreateAndRunRequest()); + + // Assert + Assert.NotNull(response); + Assert.NotNull(response.WorkflowInstanceId); + AssertMinimumAttempts(MockProvider.ReleaseAttemptCount - releaseCountBefore, 1, "release"); + } + + private async Task AcquireLockWithRetryAsync(string lockName) => + await RetryPipeline.ExecuteAsync(async ct => + await MockProvider.AcquireLockAsync(lockName, LockOptions.LockAcquisitionTimeout, ct), + CancellationToken.None); + + /// + /// Creates a workflow client with an optional workflow instance already created. + /// + private async Task CreateWorkflowClientAsync(bool createInstance = true) + { + var workflowRuntime = Scope.ServiceProvider.GetRequiredService(); + var workflowClient = await workflowRuntime.CreateClientAsync(); + + if (createInstance) + { + var createRequest = new CreateWorkflowInstanceRequest + { + WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(SimpleWorkflow.DefinitionId, VersionOptions.Latest) + }; + await workflowClient.CreateInstanceAsync(createRequest); + } + + return workflowClient; + } + + private static CreateAndRunWorkflowInstanceRequest CreateAndRunRequest() => + new() + { + WorkflowDefinitionHandle = WorkflowDefinitionHandle.ByDefinitionId(SimpleWorkflow.DefinitionId, VersionOptions.Latest) + }; + + private static void AssertMinimumAttempts(int actualAttempts, int expectedAttempts, string attemptType) => + Assert.True(actualAttempts >= expectedAttempts, + $"Expected at least {expectedAttempts} {attemptType} attempts, but got {actualAttempts}"); + + private static ResiliencePipeline CreateRetryPipeline(ITransientExceptionDetector transientExceptionDetector, ILogger logger) => + new ResiliencePipelineBuilder() + .AddRetry(new() + { + MaxRetryAttempts = MaxRetryAttempts, + // NOTE: The test retry policy intentionally differs from the production configuration. + // - We use a short, constant delay (10ms) to keep tests fast. + // - We disable jitter and exponential backoff to make timing deterministic and assertions stable. + // The production pipeline uses a larger delay with exponential backoff and jitter for robustness. + Delay = TimeSpan.FromMilliseconds(10), + BackoffType = DelayBackoffType.Constant, + UseJitter = false, + ShouldHandle = new PredicateBuilder().Handle(transientExceptionDetector.IsTransient), + OnRetry = args => + { + logger.LogWarning(args.Outcome.Exception, "Transient error acquiring lock. Attempt {AttemptNumber} of {MaxAttempts}.", args.AttemptNumber + 1, MaxRetryAttempts); + return ValueTask.CompletedTask; + } + }) + .Build(); +} diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLock.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLock.cs new file mode 100644 index 000000000..1418243bc --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLock.cs @@ -0,0 +1,56 @@ +using Medallion.Threading; + +namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks; + +/// +/// Test implementation of IDistributedLock that delegates to an inner lock +/// but can simulate transient failures. +/// +public class TestDistributedLock(IDistributedLock innerLock, TestDistributedLockProvider provider, string lockName) : IDistributedLock +{ + public string Name => innerLock.Name; + + public IDistributedSynchronizationHandle Acquire(TimeSpan? timeout = null, CancellationToken cancellationToken = default) + { + if (provider.ShouldFailAcquisition(lockName)) + { + throw new TimeoutException("Simulated transient timeout during lock acquisition"); + } + + var handle = innerLock.Acquire(timeout, cancellationToken); + return new TestDistributedSynchronizationHandle(handle, provider); + } + + public async ValueTask AcquireAsync(TimeSpan? timeout = null, CancellationToken cancellationToken = default) + { + if (provider.ShouldFailAcquisition(lockName)) + { + throw new TimeoutException("Simulated transient timeout during lock acquisition"); + } + + var handle = await innerLock.AcquireAsync(timeout, cancellationToken); + return new TestDistributedSynchronizationHandle(handle, provider); + } + + public async ValueTask TryAcquireAsync(TimeSpan timeout = default, CancellationToken cancellationToken = default) + { + if (provider.ShouldFailAcquisition(lockName)) + { + throw new TimeoutException("Simulated transient timeout during lock acquisition"); + } + + var handle = await innerLock.TryAcquireAsync(timeout, cancellationToken); + return handle == null ? null : new TestDistributedSynchronizationHandle(handle, provider); + } + + public IDistributedSynchronizationHandle? TryAcquire(TimeSpan timeout = default, CancellationToken cancellationToken = default) + { + if (provider.ShouldFailAcquisition(lockName)) + { + throw new TimeoutException("Simulated transient timeout during lock acquisition"); + } + + var handle = innerLock.TryAcquire(timeout, cancellationToken); + return handle == null ? null : new TestDistributedSynchronizationHandle(handle, provider); + } +} diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLockProvider.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLockProvider.cs new file mode 100644 index 000000000..fa50c825a --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedLockProvider.cs @@ -0,0 +1,84 @@ +using JetBrains.Annotations; +using Medallion.Threading; + +namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks; + +/// +/// Test implementation of IDistributedLockProvider that allows simulating transient failures. +/// +[UsedImplicitly] +public class TestDistributedLockProvider(IDistributedLockProvider innerProvider) : IDistributedLockProvider +{ + private int _acquisitionFailuresRemaining; + private int _releaseFailuresRemaining; + private int _acquisitionAttemptCount; + private int _releaseAttemptCount; + private string? _targetLockPrefix; + + public int AcquisitionAttemptCount => _acquisitionAttemptCount; + public int ReleaseAttemptCount => _releaseAttemptCount; + + public void FailAcquisitionOnce() => ConfigureAcquisitionFailures(1); + public void FailAcquisitionTimes(int count) => ConfigureAcquisitionFailures(count); + public void FailReleaseOnce() => Interlocked.Exchange(ref _releaseFailuresRemaining, 1); + + /// + /// Configure failures for a specific lock name prefix. Only locks matching this prefix will fail. + /// + public void FailAcquisitionTimesForLock(string lockNamePrefix, int count) + { + _targetLockPrefix = lockNamePrefix; + ConfigureAcquisitionFailures(count); + } + + public void Reset() + { + ConfigureAcquisitionFailures(0); + Interlocked.Exchange(ref _releaseFailuresRemaining, 0); + Interlocked.Exchange(ref _acquisitionAttemptCount, 0); + Interlocked.Exchange(ref _releaseAttemptCount, 0); + _targetLockPrefix = null; + } + + public IDistributedLock CreateLock(string name) => + new TestDistributedLock(innerProvider.CreateLock(name), this, name); + + internal bool ShouldFailAcquisition(string lockName) + { + Interlocked.Increment(ref _acquisitionAttemptCount); + + // If a target lock prefix is configured, only fail locks matching that prefix + if (_targetLockPrefix != null && !lockName.StartsWith(_targetLockPrefix)) + return false; + + return TryConsumeFailure(ref _acquisitionFailuresRemaining); + } + + internal bool ShouldFailRelease() + { + Interlocked.Increment(ref _releaseAttemptCount); + return TryConsumeFailure(ref _releaseFailuresRemaining); + } + + private void ConfigureAcquisitionFailures(int count) => + Interlocked.Exchange(ref _acquisitionFailuresRemaining, count); + + /// + /// Atomically decrements the failure counter if it's greater than 0. + /// Returns true if a failure was consumed, false otherwise. + /// + private static bool TryConsumeFailure(ref int failureCounter) + { + int currentValue, newValue; + do + { + currentValue = Volatile.Read(ref failureCounter); + if (currentValue <= 0) + return false; + + newValue = currentValue - 1; + } while (Interlocked.CompareExchange(ref failureCounter, newValue, currentValue) != currentValue); + + return true; + } +} diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedSynchronizationHandle.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedSynchronizationHandle.cs new file mode 100644 index 000000000..92879c70f --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Mocks/TestDistributedSynchronizationHandle.cs @@ -0,0 +1,31 @@ +using Medallion.Threading; + +namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks; + +/// +/// Test implementation of IDistributedSynchronizationHandle that can simulate failures on disposal. +/// +public class TestDistributedSynchronizationHandle( + IDistributedSynchronizationHandle? innerHandle, + TestDistributedLockProvider provider) : IDistributedSynchronizationHandle +{ + public CancellationToken HandleLostToken => innerHandle?.HandleLostToken ?? CancellationToken.None; + + public void Dispose() + { + if (provider.ShouldFailRelease()) + { + throw new TimeoutException("Simulated transient timeout during lock release"); + } + innerHandle?.Dispose(); + } + + public ValueTask DisposeAsync() + { + if (provider.ShouldFailRelease()) + { + throw new TimeoutException("Simulated transient timeout during lock release"); + } + return innerHandle?.DisposeAsync() ?? ValueTask.CompletedTask; + } +} diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Workflows/SimpleWorkflow.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Workflows/SimpleWorkflow.cs new file mode 100644 index 000000000..bc66778c6 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/Workflows/SimpleWorkflow.cs @@ -0,0 +1,26 @@ +using Elsa.Extensions; +using Elsa.Workflows.Activities; + +namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Workflows; + +/// +/// A simple workflow for testing distributed lock resilience. +/// +public class SimpleWorkflow : WorkflowBase +{ + public static readonly string DefinitionId = Guid.NewGuid().ToString(); + + protected override void Build(IWorkflowBuilder builder) + { + builder.WithDefinitionId(DefinitionId); + + builder.Root = new Sequence + { + Activities = + { + new WriteLine("Workflow execution started"), + new WriteLine("Workflow execution completed") + } + }; + } +} diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs index a5b37596c..10c0c2889 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs @@ -3,9 +3,8 @@ using Elsa.Testing.Shared; using Elsa.Workflows.Activities; using Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput.Activities; using Elsa.Workflows.Memory; -using Elsa.Workflows.Runtime.Distributed; +using Elsa.Workflows.Runtime.Distributed.Extensions; using Elsa.Workflows.Runtime.Stores; -using Microsoft.Extensions.DependencyInjection; namespace Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput; @@ -116,16 +115,13 @@ public class Tests // Act var workflowFinishedRecord = await workflow.DispatchWorkflowAndRunToCompletion( - configureServices: services => - { - services.AddScoped(); - }, configureElsa: elsa => { + // Use the distributed runtime feature so the correct bookmark queue worker and its dependencies are registered. elsa.UseWorkflowRuntime(workflowRuntime => { - workflowRuntime.ActivityExecutionLogStore = sp => activityExecutionStore; - workflowRuntime.WorkflowRuntime = sp => sp.GetRequiredService(); + workflowRuntime.UseDistributedRuntime(); + workflowRuntime.ActivityExecutionLogStore = _ => activityExecutionStore; }); }); diff --git a/test/unit/Elsa.Resilience.Core.UnitTests/DefaultTransientExceptionStrategyTests.cs b/test/unit/Elsa.Resilience.Core.UnitTests/DefaultTransientExceptionStrategyTests.cs new file mode 100644 index 000000000..2dccf5bb9 --- /dev/null +++ b/test/unit/Elsa.Resilience.Core.UnitTests/DefaultTransientExceptionStrategyTests.cs @@ -0,0 +1,87 @@ +using System.Net.Sockets; + +namespace Elsa.Resilience.Core.UnitTests; + +public class DefaultTransientExceptionStrategyTests +{ + private readonly DefaultTransientExceptionStrategy _strategy = new(); + + public static TheoryData TransientExceptionTypes => + [ + typeof(HttpRequestException), + typeof(TimeoutException), + typeof(TaskCanceledException), + typeof(IOException), + typeof(SocketException), + typeof(EndOfStreamException) + ]; + + public static TheoryData TransientMessagePatterns => + [ + "timeout", + "timed out", + "connection reset", + "connection refused", + "broken pipe", + "network", + "end of stream", + "attempted to read past the end", + "the connection is closed", + "connection is not open", + "failed to connect", + "no connection could be made", + "an existing connection was forcibly closed", + "TIMEOUT", + "Connection Reset" + ]; + + public static TheoryData NonTransientMessagePatterns => + [ + "Some random error", + "Invalid operation", + "Null reference" + ]; + + [Theory(DisplayName = "Known transient exception types should be detected as transient")] + [MemberData(nameof(TransientExceptionTypes))] + public void IsTransient_KnownTransientExceptionType_ReturnsTrue(Type exceptionType) + { + var exception = (Exception)Activator.CreateInstance(exceptionType)!; + Assert.True(_strategy.IsTransient(exception)); + } + + [Theory(DisplayName = "Exceptions with transient message patterns should be detected as transient")] + [MemberData(nameof(TransientMessagePatterns))] + public void IsTransient_ExceptionWithTransientMessagePattern_ReturnsTrue(string message) + { + var exception = new Exception(message); + Assert.True(_strategy.IsTransient(exception)); + } + + [Theory(DisplayName = "Exceptions with non-transient messages should not be detected as transient")] + [MemberData(nameof(NonTransientMessagePatterns))] + public void IsTransient_ExceptionWithNonTransientMessage_ReturnsFalse(string message) + { + var exception = new Exception(message); + Assert.False(_strategy.IsTransient(exception)); + } + + [Theory(DisplayName = "Non-transient exception types should not be detected as transient regardless of message")] + [InlineData(typeof(InvalidOperationException), "Some error")] + [InlineData(typeof(ArgumentException), "Invalid argument")] + [InlineData(typeof(NullReferenceException), "Object reference not set")] + public void IsTransient_NonTransientExceptionType_ReturnsFalse(Type exceptionType, string message) + { + var exception = (Exception)Activator.CreateInstance(exceptionType, message)!; + Assert.False(_strategy.IsTransient(exception)); + } + + [Theory(DisplayName = "Exceptions with empty messages should be detected based on type only")] + [InlineData(typeof(TimeoutException), true)] + [InlineData(typeof(InvalidOperationException), false)] + public void IsTransient_ExceptionWithEmptyMessage_ChecksTypeOnly(Type exceptionType, bool expectedResult) + { + var exception = (Exception)Activator.CreateInstance(exceptionType, "")!; + Assert.Equal(expectedResult, _strategy.IsTransient(exception)); + } +} diff --git a/test/unit/Elsa.Resilience.Core.UnitTests/Elsa.Resilience.Core.UnitTests.csproj b/test/unit/Elsa.Resilience.Core.UnitTests/Elsa.Resilience.Core.UnitTests.csproj new file mode 100644 index 000000000..f86df2cad --- /dev/null +++ b/test/unit/Elsa.Resilience.Core.UnitTests/Elsa.Resilience.Core.UnitTests.csproj @@ -0,0 +1,13 @@ + + + + [Elsa.Resilience.Core]* + 49 + + + + + + + + diff --git a/test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyCatalogTests.cs b/test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyCatalogTests.cs new file mode 100644 index 000000000..3e8ae3462 --- /dev/null +++ b/test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyCatalogTests.cs @@ -0,0 +1,118 @@ +using Elsa.Resilience.Core.UnitTests.TestHelpers; +using NSubstitute; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Resilience.Core.UnitTests; + +public class ResilienceStrategyCatalogTests +{ + private static ResilienceStrategyCatalog CreateCatalog(params IResilienceStrategySource[] sources) => new(sources); + + [Fact(DisplayName = "Catalog with no sources should return empty list")] + public async Task ListAsync_NoProviders_ReturnsEmptyList() + { + var catalog = CreateCatalog(); + var result = await catalog.ListAsync(); + + Assert.Empty(result); + } + + [Fact(DisplayName = "Catalog should return all strategies from a single source")] + public async Task ListAsync_SingleProviderWithStrategies_ReturnsStrategies() + { + var strategies = new[] + { + TestDataFactory.CreateStrategy("strategy1", "Strategy 1"), + TestDataFactory.CreateStrategy("strategy2", "Strategy 2") + }; + var provider = TestDataFactory.CreateStrategySource(strategies); + var catalog = CreateCatalog(provider); + + var result = await catalog.ListAsync(); + + Assert.Collection(result, + s => Assert.Equal("strategy1", s.Id), + s => Assert.Equal("strategy2", s.Id)); + } + + [Fact(DisplayName = "Catalog should combine strategies from multiple sources")] + public async Task ListAsync_MultipleProviders_CombinesAllStrategies() + { + var provider1 = TestDataFactory.CreateStrategySource( + TestDataFactory.CreateStrategy("strategy1", "Strategy 1"), + TestDataFactory.CreateStrategy("strategy2", "Strategy 2")); + var provider2 = TestDataFactory.CreateStrategySource( + TestDataFactory.CreateStrategy("strategy3", "Strategy 3")); + var catalog = CreateCatalog(provider1, provider2); + + var result = await catalog.ListAsync().ToList(); + + Assert.Equal(3, result.Count); + Assert.Contains(result, s => s.Id == "strategy1"); + Assert.Contains(result, s => s.Id == "strategy2"); + Assert.Contains(result, s => s.Id == "strategy3"); + } + + [Theory(DisplayName = "Catalog should retrieve strategy by ID or return null if not found")] + [InlineData("test-id", "Test Strategy", true)] + [InlineData("non-existent", null, false)] + public async Task GetAsync_WithStrategyId_ReturnsExpectedResult(string searchId, string? expectedDisplayName, bool shouldExist) + { + var strategy = TestDataFactory.CreateStrategy("test-id", "Test Strategy"); + var provider = TestDataFactory.CreateStrategySource(strategy); + var catalog = CreateCatalog(provider); + + var result = await catalog.GetAsync(searchId); + + if (shouldExist) + { + Assert.NotNull(result); + Assert.Equal(searchId, result.Id); + Assert.Equal(expectedDisplayName, result.DisplayName); + } + else + { + Assert.Null(result); + } + } + + [Fact(DisplayName = "Catalog should search all sources to find a strategy")] + public async Task GetAsync_MultipleProvidersStrategyInSecond_ReturnsStrategy() + { + var provider1 = TestDataFactory.CreateStrategySource(TestDataFactory.CreateStrategy("strategy1", "Strategy 1")); + var provider2 = TestDataFactory.CreateStrategySource(TestDataFactory.CreateStrategy("strategy2", "Strategy 2")); + var catalog = CreateCatalog(provider1, provider2); + + var result = await catalog.GetAsync("strategy2"); + + Assert.NotNull(result); + Assert.Equal("strategy2", result.Id); + } + + [Fact(DisplayName = "Catalog should cache strategy list after first retrieval")] + public async Task ListAsync_CalledMultipleTimes_CachesResult() + { + var strategy = TestDataFactory.CreateStrategy("test", "Test"); + var provider = TestDataFactory.CreateStrategySource(strategy); + var catalog = CreateCatalog(provider); + + await catalog.ListAsync(); + await catalog.ListAsync(); + await catalog.ListAsync(); + + await provider.Received(1).GetStrategiesAsync(Arg.Any()); + } + + [Fact(DisplayName = "GetAsync should use cached list when available")] + public async Task GetAsync_CalledAfterList_UsesCachedResult() + { + var strategy = TestDataFactory.CreateStrategy("test", "Test"); + var provider = TestDataFactory.CreateStrategySource(strategy); + var catalog = CreateCatalog(provider); + + await catalog.ListAsync(); + await catalog.GetAsync("test"); + + await provider.Received(1).GetStrategiesAsync(Arg.Any()); + } +} diff --git a/test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyConfigEvaluatorTests.cs b/test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyConfigEvaluatorTests.cs new file mode 100644 index 000000000..ea061c445 --- /dev/null +++ b/test/unit/Elsa.Resilience.Core.UnitTests/ResilienceStrategyConfigEvaluatorTests.cs @@ -0,0 +1,135 @@ +using Elsa.Expressions.Contracts; +using Elsa.Expressions.Models; +using Elsa.Resilience.Core.UnitTests.TestHelpers; +using Elsa.Resilience.Models; +using Elsa.Resilience.Options; +using Elsa.Resilience.Serialization; +using NSubstitute; + +namespace Elsa.Resilience.Core.UnitTests; + +public class ResilienceStrategyConfigEvaluatorTests +{ + private readonly IResilienceStrategyCatalog _catalog = Substitute.For(); + private readonly IExpressionEvaluator _expressionEvaluator = Substitute.For(); + private readonly ResilienceStrategyConfigEvaluator _evaluator; + private readonly ExpressionExecutionContext _context = new(Substitute.For(), null!); + + public ResilienceStrategyConfigEvaluatorTests() + { + var options = Microsoft.Extensions.Options.Options.Create(new ResilienceOptions()); + var serializer = new ResilienceStrategySerializer(options); + _evaluator = new(_catalog, _expressionEvaluator, serializer); + } + + [Fact(DisplayName = "Evaluator should return null when config is null")] + public async Task EvaluateAsync_NullConfig_ReturnsNull() + { + var result = await _evaluator.EvaluateAsync(null, _context); + Assert.Null(result); + } + + [Fact(DisplayName = "Evaluator in identifier mode should resolve strategy from catalog")] + public async Task EvaluateAsync_IdentifierMode_WithValidId_ReturnsStrategyFromCatalog() + { + var strategy = TestDataFactory.CreateStrategy("test-strategy", "Test Strategy"); + SetupCatalogStrategy("test-strategy", strategy); + var config = CreateConfig(ResilienceStrategyConfigMode.Identifier, "test-strategy"); + + var result = await _evaluator.EvaluateAsync(config, _context); + + Assert.NotNull(result); + Assert.Equal("test-strategy", result.Id); + await _catalog.Received(1).GetAsync("test-strategy", Arg.Any()); + } + + [Theory(DisplayName = "Evaluator in identifier mode should return null for invalid strategy IDs")] + [InlineData("")] + [InlineData(" ")] + [InlineData(null)] + public async Task EvaluateAsync_IdentifierMode_WithInvalidId_ReturnsNull(string? strategyId) + { + var config = CreateConfig(ResilienceStrategyConfigMode.Identifier, strategyId); + + var result = await _evaluator.EvaluateAsync(config, _context); + + Assert.Null(result); + await _catalog.DidNotReceive().GetAsync(Arg.Any(), Arg.Any()); + } + + [Fact(DisplayName = "Evaluator in expression mode should return null when expression is null")] + public async Task EvaluateAsync_ExpressionMode_WithNullExpression_ReturnsNull() + { + var config = CreateConfig(ResilienceStrategyConfigMode.Expression); + + var result = await _evaluator.EvaluateAsync(config, _context); + + Assert.Null(result); + await _expressionEvaluator.DidNotReceive().EvaluateAsync(Arg.Any(), Arg.Any(), Arg.Any()); + } + + [Fact(DisplayName = "Evaluator in expression mode should resolve string IDs from catalog")] + public async Task EvaluateAsync_ExpressionMode_ReturnsStringId_ResolvesFromCatalog() + { + var expression = new Expression("C#", "\"test-strategy\""); + var strategy = TestDataFactory.CreateStrategy("test-strategy", "Test Strategy"); + SetupExpressionResult(expression, "test-strategy"); + SetupCatalogStrategy("test-strategy", strategy); + var config = CreateConfig(ResilienceStrategyConfigMode.Expression, expression: expression); + + var result = await _evaluator.EvaluateAsync(config, _context); + + Assert.NotNull(result); + Assert.Equal("test-strategy", result.Id); + await _catalog.Received(1).GetAsync("test-strategy", Arg.Any()); + } + + [Fact(DisplayName = "Evaluator in expression mode should return strategy objects directly")] + public async Task EvaluateAsync_ExpressionMode_ReturnsStrategyObject_ReturnsStrategyDirectly() + { + var expression = new Expression("C#", "strategy"); + var strategy = TestDataFactory.CreateStrategy("direct-strategy", "Direct Strategy"); + SetupExpressionResult(expression, strategy); + var config = CreateConfig(ResilienceStrategyConfigMode.Expression, expression: expression); + + var result = await _evaluator.EvaluateAsync(config, _context); + + Assert.NotNull(result); + Assert.Same(strategy, result); + await _catalog.DidNotReceive().GetAsync(Arg.Any(), Arg.Any()); + } + + [Theory(DisplayName = "Evaluator in expression mode should return null for unexpected result types")] + [InlineData("null", null)] + [InlineData("42", 42)] + public async Task EvaluateAsync_ExpressionMode_ReturnsUnexpectedType_ReturnsNull(string expressionCode, object? expressionResult) + { + var expression = new Expression("C#", expressionCode); + SetupExpressionResult(expression, expressionResult); + var config = CreateConfig(ResilienceStrategyConfigMode.Expression, expression: expression); + + var result = await _evaluator.EvaluateAsync(config, _context); + + Assert.Null(result); + } + + private ResilienceStrategyConfig CreateConfig(ResilienceStrategyConfigMode mode, string? strategyId = null, Expression? expression = null) + { + return new() + { + Mode = mode, + StrategyId = strategyId, + Expression = expression + }; + } + + private void SetupCatalogStrategy(string id, IResilienceStrategy strategy) + { + _catalog.GetAsync(id, Arg.Any()).Returns(strategy); + } + + private void SetupExpressionResult(Expression expression, object? result) + { + _expressionEvaluator.EvaluateAsync(expression, _context, Arg.Any()).Returns(result); + } +} diff --git a/test/unit/Elsa.Resilience.Core.UnitTests/TestHelpers/TestDataFactory.cs b/test/unit/Elsa.Resilience.Core.UnitTests/TestHelpers/TestDataFactory.cs new file mode 100644 index 000000000..e5c937eb9 --- /dev/null +++ b/test/unit/Elsa.Resilience.Core.UnitTests/TestHelpers/TestDataFactory.cs @@ -0,0 +1,21 @@ +using NSubstitute; + +namespace Elsa.Resilience.Core.UnitTests.TestHelpers; + +internal static class TestDataFactory +{ + public static IResilienceStrategy CreateStrategy(string id, string displayName) + { + var strategy = Substitute.For(); + strategy.Id.Returns(id); + strategy.DisplayName.Returns(displayName); + return strategy; + } + + public static IResilienceStrategySource CreateStrategySource(params IResilienceStrategy[] strategies) + { + var source = Substitute.For(); + source.GetStrategiesAsync(Arg.Any()).Returns(strategies); + return source; + } +} diff --git a/test/unit/Elsa.Resilience.Core.UnitTests/TransientExceptionDetectorTests.cs b/test/unit/Elsa.Resilience.Core.UnitTests/TransientExceptionDetectorTests.cs new file mode 100644 index 000000000..2862280b1 --- /dev/null +++ b/test/unit/Elsa.Resilience.Core.UnitTests/TransientExceptionDetectorTests.cs @@ -0,0 +1,151 @@ +using NSubstitute; + +namespace Elsa.Resilience.Core.UnitTests; + +public class TransientExceptionDetectorTests +{ + [Fact(DisplayName = "Service with no registered strategies should return false for any exception")] + public void IsTransient_NoDetectors_ReturnsFalse() + { + var detector = CreateDetector(); + var exception = new Exception("test"); + + Assert.False(detector.IsTransient(exception)); + } + + [Fact(DisplayName = "Service should return true when any strategy detects the exception as transient")] + public void IsTransient_DetectorReturnsTrue_ReturnsTrue() + { + var exception = new Exception("test"); + var strategy = CreateStrategy((exception, true)); + var detector = CreateDetector(strategy); + + Assert.True(detector.IsTransient(exception)); + } + + [Fact(DisplayName = "Service with multiple strategies should return true if any one detects as transient")] + public void IsTransient_MultipleDetectorsOneReturnsTrue_ReturnsTrue() + { + var exception = new Exception("test"); + var strategy1 = CreateStrategy((exception, false)); + var strategy2 = CreateStrategy((exception, true)); + var detector = CreateDetector(strategy1, strategy2); + + Assert.True(detector.IsTransient(exception)); + } + + [Fact(DisplayName = "Service should return false when all strategies detect the exception as non-transient")] + public void IsTransient_AllDetectorsReturnFalse_ReturnsFalse() + { + var exception = new Exception("test"); + var strategy1 = Substitute.For(); + var strategy2 = Substitute.For(); + strategy1.IsTransient(Arg.Any()).Returns(false); + strategy2.IsTransient(Arg.Any()).Returns(false); + var detector = CreateDetector(strategy1, strategy2); + + Assert.False(detector.IsTransient(exception)); + } + + [Theory(DisplayName = "Service should walk the inner exception chain to find transient exceptions")] + [MemberData(nameof(InnerExceptionChainTestCases))] + public void IsTransient_InnerExceptionChainHasTransient_ReturnsTrue(Exception exception, Exception transientException) + { + var strategy = Substitute.For(); + strategy.IsTransient(transientException).Returns(true); + strategy.IsTransient(Arg.Is(e => e != transientException)).Returns(false); + var detector = CreateDetector(strategy); + + Assert.True(detector.IsTransient(exception)); + } + + [Theory(DisplayName = "Service should inspect AggregateException inner exceptions")] + [MemberData(nameof(AggregateExceptionTestCases))] + public void IsTransient_AggregateException_ChecksInnerExceptions( + AggregateException aggregateException, + Action configureDetector, + bool expectedResult) + { + var strategy = Substitute.For(); + configureDetector(strategy); + var detector = CreateDetector(strategy); + + Assert.Equal(expectedResult, detector.IsTransient(aggregateException)); + } + + [Fact(DisplayName = "AggregateException with mixed inner exceptions should be transient if any inner is transient")] + public void IsTransient_AggregateExceptionWithMultipleInnerOneTransient_ReturnsTrue() + { + var transientException = new TimeoutException("timeout"); + var nonTransientException = new InvalidOperationException("invalid"); + var aggregateException = new AggregateException("aggregate", nonTransientException, transientException); + + var strategy = CreateStrategy( + (aggregateException, false), + (nonTransientException, false), + (transientException, true)); + var detector = CreateDetector(strategy); + + Assert.True(detector.IsTransient(aggregateException)); + } + + public static IEnumerable InnerExceptionChainTestCases + { + get + { + var transientException = new TimeoutException("timeout"); + + yield return + [ + new Exception("outer", transientException), + transientException + ]; + + yield return + [ + new Exception("outer", new("middle", transientException)), + transientException + ]; + } + } + + public static IEnumerable AggregateExceptionTestCases + { + get + { + // Scenario 1: One transient inner exception + var timeoutException = new TimeoutException("timeout"); + yield return + [ + new AggregateException("aggregate", timeoutException), + (Action)(detector => + { + detector.IsTransient(Arg.Is(e => e.Message == "timeout")).Returns(true); + detector.IsTransient(Arg.Is(e => e.GetType() != typeof(TimeoutException))).Returns(false); + }), + true + ]; + + // Scenario 2: No transient inner exceptions + yield return + [ + new AggregateException("aggregate", new InvalidOperationException("error1"), new ArgumentException("error2")), + (Action)(detector => + { + detector.IsTransient(Arg.Any()).Returns(false); + }), + false + ]; + } + } + + private static ITransientExceptionStrategy CreateStrategy(params (Exception exception, bool isTransient)[] behaviors) + { + var detector = Substitute.For(); + foreach (var (exception, isTransient) in behaviors) + detector.IsTransient(exception).Returns(isTransient); + return detector; + } + + private static TransientExceptionDetector CreateDetector(params ITransientExceptionStrategy[] detectors) => new(detectors); +}