Merge remote-tracking branch 'origin/main' into release/3.6.0

This commit is contained in:
Sipke Schoorstra 2025-12-29 19:57:50 +01:00
commit 8f6ce7eb8c
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
48 changed files with 2059 additions and 112 deletions

View file

@ -18,7 +18,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'
@ -156,7 +156,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

View file

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

View file

@ -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<CommandContext> 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<ICommandSender>();
try
{
// Create a fresh scope for each command to ensure proper service lifetime
using var scope = _scopeFactory.CreateScope();
var commandSender = scope.ServiceProvider.GetRequiredService<ICommandSender>();
// 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");
}
}
}

View file

@ -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
/// <param name="cancellationToken">Cancellation token from the hosted service</param>
private async Task ReadOutputAsync(Channel<NotificationContext> 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");
}
}
}

View file

@ -12,10 +12,18 @@ public class XunitLogger(ITestOutputHelper testOutputHelper, string categoryName
public void Log<TState>(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func<TState, Exception?, string> 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

View file

@ -0,0 +1,14 @@
namespace Elsa.Resilience;
/// <summary>
/// Service for detecting whether exceptions are transient and may be resolved by retrying.
/// </summary>
public interface ITransientExceptionDetector
{
/// <summary>
/// Determines whether the specified exception is transient.
/// </summary>
/// <param name="exception">The exception to check.</param>
/// <returns>True if the exception is transient; otherwise, false.</returns>
bool IsTransient(Exception exception);
}

View file

@ -0,0 +1,14 @@
namespace Elsa.Resilience;
/// <summary>
/// Defines a contract for detecting whether an exception is transient and may be resolved by retrying.
/// </summary>
public interface ITransientExceptionStrategy
{
/// <summary>
/// Determines whether the specified exception is transient.
/// </summary>
/// <param name="exception">The exception to check.</param>
/// <returns>True if the exception is transient; otherwise, false.</returns>
bool IsTransient(Exception exception);
}

View file

@ -0,0 +1,65 @@
namespace Elsa.Resilience;
/// <summary>
/// Default implementation that detects common transient exceptions from the .NET framework and common patterns.
/// </summary>
public class DefaultTransientExceptionStrategy : ITransientExceptionStrategy
{
private static readonly HashSet<string> 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<string> 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",
};
/// <inheritdoc />
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));
}
}

View file

@ -0,0 +1,42 @@
namespace Elsa.Resilience;
/// <summary>
/// Default implementation of <see cref="ITransientExceptionDetector"/> that delegates to registered detectors.
/// </summary>
public class TransientExceptionDetector(IEnumerable<ITransientExceptionStrategy> detectors) : ITransientExceptionDetector
{
private readonly IReadOnlyList<ITransientExceptionStrategy> _detectors = detectors.ToList();
/// <inheritdoc />
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;
}
}

View file

@ -84,5 +84,10 @@ public class ResilienceFeature(IModule module) : FeatureBase(module)
.AddScoped(_retryAttemptRecorder)
.AddScoped(_retryAttemptReader)
.AddHandlersFrom<ResilienceFeature>();
// Register transient exception detection infrastructure
Services
.AddSingleton<ITransientExceptionStrategy, DefaultTransientExceptionStrategy>()
.AddSingleton<ITransientExceptionDetector, TransientExceptionDetector>();
}
}

View file

@ -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<CommitStateOptions>(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<CommitStateOptions>(options => options.ActivityCommitStrategies[registration.Metadata.Name] = registration);
}
/// <summary>
/// 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.
/// </summary>
/// <param name="strategy">The workflow commit strategy instance to use as the default.</param>
public void SetDefaultWorkflowCommitStrategy(IWorkflowCommitStrategy strategy)
{
Services.Configure<CommitStateOptions>(options => options.DefaultWorkflowCommitStrategy = strategy);
}
/// <summary>
/// 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.
/// </summary>
/// <param name="strategy">The activity commit strategy instance to use as the default.</param>
public void SetDefaultActivityCommitStrategy(IActivityCommitStrategy strategy)
{
Services.Configure<CommitStateOptions>(options => options.DefaultActivityCommitStrategy = strategy);
}
public override void Apply()
{
Services.AddSingleton<ICommitStrategyRegistry, DefaultCommitStrategyRegistry>();

View file

@ -4,12 +4,46 @@ using Elsa.Workflows.Features;
// ReSharper disable once CheckNamespace
namespace Elsa.Extensions;
/// <summary>
/// Provides extension methods for configuring commit strategies on <see cref="WorkflowsFeature"/>.
/// </summary>
public static class WorkflowsFeatureCommitStateExtensions
{
/// <summary>
/// Configures commit strategies for workflows.
/// </summary>
/// <param name="workflowsFeature">The workflows feature.</param>
/// <param name="configure">An optional configuration delegate for the commit strategies feature.</param>
/// <returns>The workflows feature for chaining.</returns>
public static WorkflowsFeature UseCommitStrategies(this WorkflowsFeature workflowsFeature, Action<CommitStrategiesFeature>? configure = null)
{
workflowsFeature.Module.Use(configure);
return workflowsFeature;
}
/// <summary>
/// Sets the specified workflow commit strategy as the global default for all workflows that do not specify their own strategy.
/// The strategy will not be added to the registry and serves only as a fallback.
/// </summary>
/// <param name="workflowsFeature">The workflows feature.</param>
/// <param name="strategy">The workflow commit strategy instance to use as the default.</param>
/// <returns>The workflows feature for chaining.</returns>
public static WorkflowsFeature WithDefaultWorkflowCommitStrategy(this WorkflowsFeature workflowsFeature, IWorkflowCommitStrategy strategy)
{
workflowsFeature.Module.Use<CommitStrategiesFeature>(feature => feature.SetDefaultWorkflowCommitStrategy(strategy));
return workflowsFeature;
}
/// <summary>
/// Sets the specified activity commit strategy as the global default for all activities that do not specify their own strategy.
/// The strategy will not be added to the registry and serves only as a fallback.
/// </summary>
/// <param name="workflowsFeature">The workflows feature.</param>
/// <param name="strategy">The activity commit strategy instance to use as the default.</param>
/// <returns>The workflows feature for chaining.</returns>
public static WorkflowsFeature WithDefaultActivityCommitStrategy(this WorkflowsFeature workflowsFeature, IActivityCommitStrategy strategy)
{
workflowsFeature.Module.Use<CommitStrategiesFeature>(feature => feature.SetDefaultActivityCommitStrategy(strategy));
return workflowsFeature;
}
}

View file

@ -1,7 +1,29 @@
namespace Elsa.Workflows.CommitStates;
/// <summary>
/// Configuration options for commit state strategies.
/// </summary>
public class CommitStateOptions
{
/// <summary>
/// Gets or sets the workflow commit strategies.
/// </summary>
public IDictionary<string, WorkflowCommitStrategyRegistration> WorkflowCommitStrategies { get; set; } = new Dictionary<string, WorkflowCommitStrategyRegistration>();
/// <summary>
/// Gets or sets the activity commit strategies.
/// </summary>
public IDictionary<string, ActivityCommitStrategyRegistration> ActivityCommitStrategies { get; set; } = new Dictionary<string, ActivityCommitStrategyRegistration>();
/// <summary>
/// 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.
/// </summary>
public IWorkflowCommitStrategy? DefaultWorkflowCommitStrategy { get; set; }
/// <summary>
/// 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.
/// </summary>
public IActivityCommitStrategy? DefaultActivityCommitStrategy { get; set; }
}

View file

@ -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
/// <summary>
/// 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.
/// </summary>
public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, ICommitStrategyRegistry commitStrategyRegistry, ILogger<DefaultActivityInvokerMiddleware> logger)
public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, ICommitStrategyRegistry commitStrategyRegistry, IOptions<CommitStateOptions> commitStateOptions, ILogger<DefaultActivityInvokerMiddleware> logger)
: IActivityExecutionMiddleware
{
private static readonly MethodInfo ExecuteAsyncMethodInfo = typeof(IActivity).GetMethod(nameof(IActivity.ExecuteAsync))!;
/// <inheritdoc />
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
/// </summary>
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,11 @@ 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 = !string.IsNullOrWhiteSpace(strategyName)
? commitStrategyRegistry.FindActivityStrategy(strategyName)
: commitStateOptions.Value.DefaultActivityCommitStrategy;
var commitAction = CommitAction.Default;
if (strategy != null)
@ -147,7 +152,10 @@ 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 = !string.IsNullOrWhiteSpace(workflowStrategyName)
? commitStrategyRegistry.FindWorkflowStrategy(workflowStrategyName)
: commitStateOptions.Value.DefaultWorkflowCommitStrategy;
if (workflowStrategy == null)
return false;

View file

@ -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
/// <summary>
/// A workflow execution middleware component that executes scheduled work items.
/// </summary>
public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, IActivityInvoker activityInvoker, ICommitStrategyRegistry commitStrategyRegistry) : WorkflowExecutionMiddleware(next)
public class DefaultActivitySchedulerMiddleware(WorkflowMiddlewareDelegate next, IActivityInvoker activityInvoker, ICommitStrategyRegistry commitStrategyRegistry, IOptions<CommitStateOptions> commitStateOptions) : WorkflowExecutionMiddleware(next)
{
/// <inheritdoc />
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,20 @@ 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 = !string.IsNullOrWhiteSpace(strategyName)
? commitStrategyRegistry.FindWorkflowStrategy(strategyName)
: 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();
}

View file

@ -10,11 +10,13 @@
<ItemGroup>
<PackageReference Include="DistributedLock.FileSystem" />
<PackageReference Include="Microsoft.Extensions.DependencyInjection" />
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions"/>
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions"/>
<PackageReference Include="Open.Linq.AsyncExtensions" />
<PackageReference Include="Polly" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Resilience\Elsa.Resilience.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>

View file

@ -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.
/// </summary>
[DependsOn(typeof(WorkflowRuntimeFeature))]
public class DistributedRuntimeFeature : FeatureBase
[DependsOn(typeof(ResilienceFeature))]
public class DistributedRuntimeFeature(IModule module) : FeatureBase(module)
{
/// <inheritdoc />
public DistributedRuntimeFeature(IModule module) : base(module)
{
}
public override void Configure()
{
Module.UseWorkflowRuntime(runtime =>

View file

@ -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> distributedLockingOptions,
IServiceProvider serviceProvider)
IServiceProvider serviceProvider,
ILogger<DistributedWorkflowClient> logger)
: IWorkflowClient
{
private readonly LocalWorkflowClient _localWorkflowClient = ActivatorUtilities.CreateInstance<LocalWorkflowClient>(serviceProvider, workflowInstanceId);
private readonly Lazy<ResiliencePipeline> _retryPipeline = new(() => CreateRetryPipeline(transientExceptionDetector, logger, workflowInstanceId));
public string WorkflowInstanceId => workflowInstanceId;
public async Task<CreateWorkflowInstanceResponse> CreateInstanceAsync(CreateWorkflowInstanceRequest request, CancellationToken cancellationToken = default)
@ -25,7 +30,7 @@ public class DistributedWorkflowClient(
public async Task<RunWorkflowInstanceResponse> 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<bool> 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<R> WithLockAsync<R>(Func<Task<R>> func)
private async Task<TReturn> WithLockAsync<TReturn>(Func<Task<TReturn>> 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<IDistributedSynchronizationHandle?> 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<DistributedWorkflowClient> 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<Exception>(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();
}
}

View file

@ -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> commitStateOptions)
: DefaultActivityInvokerMiddleware(next, commitStrategyRegistry, commitStateOptions, logger)
{
internal static string GetBackgroundActivityOutputKey(string activityNodeId) => $"__BackgroundActivityOutput:{activityNodeId}";
internal static string GetBackgroundActivityOutcomesKey(string activityNodeId) => $"__BackgroundActivityOutcomes:{activityNodeId}";
@ -121,7 +123,7 @@ public class BackgroundActivityInvokerMiddleware(
}
private static bool GetIsBackgroundExecution(ActivityExecutionContext context) => context.TransientProperties.ContainsKey(BackgroundActivityExecutionContextExtensions.IsBackgroundExecution);
/// <summary>
/// If the input contains captured output from the background activity invoker, apply that to the execution context.
/// </summary>
@ -175,7 +177,7 @@ public class BackgroundActivityInvokerMiddleware(
context.WorkflowExecutionContext.Properties.Remove(bookmarksKey);
}
private void CapturePropertiesIfAny(ActivityExecutionContext context)
{
var activity = context.Activity;
@ -187,7 +189,7 @@ public class BackgroundActivityInvokerMiddleware(
if (capturedProperties == null)
return;
foreach (var property in capturedProperties)
foreach (var property in capturedProperties)
context.Properties[property.Key] = property.Value;
}

View file

@ -0,0 +1,10 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Runtime.Requests;
using Elsa.Workflows.Runtime.Responses;
namespace Elsa.Workflows.Runtime.Notifications;
/// <summary>
/// A notification that is published when a workflow definition has been dispatched.
/// </summary>
public record WorkflowDefinitionDispatched(DispatchWorkflowDefinitionRequest Request, DispatchWorkflowResponse Response) : INotification;

View file

@ -0,0 +1,9 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Runtime.Requests;
namespace Elsa.Workflows.Runtime.Notifications;
/// <summary>
/// A notification that is published when a workflow definition is being dispatched.
/// </summary>
public record WorkflowDefinitionDispatching(DispatchWorkflowDefinitionRequest Request) : INotification;

View file

@ -0,0 +1,10 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Runtime.Requests;
using Elsa.Workflows.Runtime.Responses;
namespace Elsa.Workflows.Runtime.Notifications;
/// <summary>
/// A notification that is published when a workflow instance has been dispatched.
/// </summary>
public record WorkflowInstanceDispatched(DispatchWorkflowInstanceRequest Request, DispatchWorkflowResponse Response) : INotification;

View file

@ -0,0 +1,9 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Runtime.Requests;
namespace Elsa.Workflows.Runtime.Notifications;
/// <summary>
/// A notification that is published when a workflow instance is being dispatched.
/// </summary>
public record WorkflowInstanceDispatching(DispatchWorkflowInstanceRequest Request) : INotification;

View file

@ -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;
/// <summary>
/// A simple implementation that queues the specified request for workflow execution on a non-durable background worker.
/// </summary>
public class BackgroundWorkflowDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher
public class BackgroundWorkflowDispatcher(ICommandSender commandSender, INotificationSender notificationSender, ITenantAccessor tenantAccessor) : IWorkflowDispatcher
{
/// <inheritdoc />
public async Task<DispatchWorkflowResponse> 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;
}
/// <inheritdoc />
public async Task<DispatchWorkflowResponse> 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;
}
/// <inheritdoc />

View file

@ -12,12 +12,14 @@ 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.Scenarios.HostMethodActivities;
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;
@ -129,6 +131,18 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
builder.ConfigureTestServices(services =>
{
// Decorate IDistributedLockProvider with TestDistributedLockProvider so tests use it
services.Decorate<IDistributedLockProvider, TestDistributedLockProvider>();
// Also register TestDistributedLockProvider as itself so tests can access it directly for configuration
services.AddSingleton(sp =>
{
var provider = sp.GetRequiredService<IDistributedLockProvider>();
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<SignalManager>()
.AddScoped<AsyncWorkflowRunner>()
@ -140,6 +154,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
.AddWorkflowsProvider<TestWorkflowProvider>()
.AddNotificationHandlersFrom<WorkflowEventHandlers>()
.Decorate<IChangeTokenSignaler, EventPublishingChangeTokenSignaler>()
.Decorate<IDistributedLockProvider, TestDistributedLockProvider>()
;
});
}

View file

@ -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<IDistributedLockProvider>();
private ITransientExceptionDetector TransientExceptionDetector => Scope.ServiceProvider.GetRequiredService<ITransientExceptionDetector>();
private ILogger<DistributedLockResilienceTests> Logger => Scope.ServiceProvider.GetRequiredService<ILogger<DistributedLockResilienceTests>>();
private DistributedLockingOptions LockOptions => Scope.ServiceProvider.GetRequiredService<IOptions<DistributedLockingOptions>>().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<TimeoutException>(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<TimeoutException>(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<IDistributedSynchronizationHandle?> AcquireLockWithRetryAsync(string lockName) =>
await RetryPipeline.ExecuteAsync(async ct =>
await MockProvider.AcquireLockAsync(lockName, LockOptions.LockAcquisitionTimeout, ct),
CancellationToken.None);
/// <summary>
/// Creates a workflow client with an optional workflow instance already created.
/// </summary>
private async Task<IWorkflowClient> CreateWorkflowClientAsync(bool createInstance = true)
{
var workflowRuntime = Scope.ServiceProvider.GetRequiredService<IWorkflowRuntime>();
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<Exception>(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();
}

View file

@ -0,0 +1,56 @@
using Medallion.Threading;
namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks;
/// <summary>
/// Test implementation of IDistributedLock that delegates to an inner lock
/// but can simulate transient failures.
/// </summary>
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<IDistributedSynchronizationHandle> 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<IDistributedSynchronizationHandle?> 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);
}
}

View file

@ -0,0 +1,84 @@
using JetBrains.Annotations;
using Medallion.Threading;
namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks;
/// <summary>
/// Test implementation of IDistributedLockProvider that allows simulating transient failures.
/// </summary>
[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);
/// <summary>
/// Configure failures for a specific lock name prefix. Only locks matching this prefix will fail.
/// </summary>
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);
/// <summary>
/// Atomically decrements the failure counter if it's greater than 0.
/// Returns true if a failure was consumed, false otherwise.
/// </summary>
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;
}
}

View file

@ -0,0 +1,31 @@
using Medallion.Threading;
namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Mocks;
/// <summary>
/// Test implementation of IDistributedSynchronizationHandle that can simulate failures on disposal.
/// </summary>
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;
}
}

View file

@ -0,0 +1,26 @@
using Elsa.Extensions;
using Elsa.Workflows.Activities;
namespace Elsa.Workflows.ComponentTests.Scenarios.DistributedLockResilience.Workflows;
/// <summary>
/// A simple workflow for testing distributed lock resilience.
/// </summary>
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")
}
};
}
}

View file

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

View file

@ -0,0 +1,201 @@
using Elsa.Extensions;
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
{
private readonly ITestOutputHelper _testOutputHelper;
public Tests(ITestOutputHelper testOutputHelper)
{
_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<SimpleWorkflowWithoutActivityCommitStrategy>()
.Build();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
// Act
var result = await workflowRunner.RunAsync<SimpleWorkflowWithoutActivityCommitStrategy>();
// 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<WorkflowWithExplicitActivityCommitStrategy>()
.Build();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
// Act
var result = await workflowRunner.RunAsync<WorkflowWithExplicitActivityCommitStrategy>();
// 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<SimpleWorkflowWithoutActivityCommitStrategy>()
.Build();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
// Act
var result = await workflowRunner.RunAsync<SimpleWorkflowWithoutActivityCommitStrategy>();
// 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 void DefaultStrategyNotInRegistry()
{
// Arrange
var defaultStrategy = new ExecutedActivityStrategy();
var services = new TestApplicationBuilder(_testOutputHelper)
.ConfigureElsa(elsa => elsa
.UseWorkflows(workflows => workflows
.WithDefaultActivityCommitStrategy(defaultStrategy)
)
)
.Build();
var registry = services.GetRequiredService<ICommitStrategyRegistry>();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
// Act
var activityStrategies = registry.ListActivityStrategyRegistrations().ToList();
// Assert
Assert.Empty(activityStrategies);
Assert.NotNull(options.Value.DefaultActivityCommitStrategy);
Assert.Same(defaultStrategy, options.Value.DefaultActivityCommitStrategy);
}
[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<ICommitStrategyRegistry>();
// Manually populate the registry
var startupTask = new PopulateCommitStrategyRegistry(
registry,
services.GetRequiredService<IOptions<CommitStateOptions>>()
);
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<SimpleWorkflowWithoutActivityCommitStrategy>()
.Build();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
// Act
var result = await workflowRunner.RunAsync<SimpleWorkflowWithoutActivityCommitStrategy>();
// 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);
}
}

View file

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

View file

@ -0,0 +1,22 @@
using Elsa.Workflows.Activities;
namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy;
/// <summary>
/// A simple workflow that does not specify an explicit commit strategy.
/// </summary>
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")
}
};
}
}

View file

@ -0,0 +1,167 @@
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;
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<SimpleWorkflowWithoutWorkflowCommitStrategy>()
.Build();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
// Act
var result = await workflowRunner.RunAsync<SimpleWorkflowWithoutWorkflowCommitStrategy>();
// 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<WorkflowWithExplicitWorkflowCommitStrategy>()
.Build();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
// Act
var result = await workflowRunner.RunAsync<WorkflowWithExplicitWorkflowCommitStrategy>();
// 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<SimpleWorkflowWithoutWorkflowCommitStrategy>()
.Build();
var workflowRunner = services.GetRequiredService<IWorkflowRunner>();
var options = services.GetRequiredService<IOptions<CommitStateOptions>>();
// Act
var result = await workflowRunner.RunAsync<SimpleWorkflowWithoutWorkflowCommitStrategy>();
// 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 void DefaultWorkflowStrategyNotInRegistry()
{
// Arrange
var services = new TestApplicationBuilder(_testOutputHelper)
.ConfigureElsa(elsa => elsa
.UseWorkflows(workflows => workflows
.WithDefaultWorkflowCommitStrategy(new ActivityExecutedWorkflowStrategy())
)
)
.Build();
var registry = services.GetRequiredService<ICommitStrategyRegistry>();
// 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<IOptions<CommitStateOptions>>();
var registry = services.GetRequiredService<ICommitStrategyRegistry>();
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);
}
}

View file

@ -0,0 +1,23 @@
using Elsa.Workflows.Activities;
namespace Elsa.Workflows.IntegrationTests.Scenarios.DefaultWorkflowCommitStrategy;
/// <summary>
/// A workflow that explicitly sets a commit strategy, which should override the default.
/// </summary>
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")
}
};
}
}

View file

@ -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<DistributedWorkflowRuntime>();
},
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<DistributedWorkflowRuntime>();
workflowRuntime.UseDistributedRuntime();
workflowRuntime.ActivityExecutionLogStore = _ => activityExecutionStore;
});
});

View file

@ -0,0 +1,31 @@
using Elsa.Workflows.Runtime.Requests;
using Elsa.Workflows.Runtime.Responses;
namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications;
public class Spy
{
private readonly TaskCompletionSource<bool> _workflowDefinitionDispatchingTcs = new();
private readonly TaskCompletionSource<bool> _workflowDefinitionDispatchedTcs = new();
private readonly TaskCompletionSource<bool> _workflowInstanceDispatchingTcs = new();
private readonly TaskCompletionSource<bool> _workflowInstanceDispatchedTcs = new();
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; }
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);
}

View file

@ -0,0 +1,50 @@
using Elsa.Mediator.Contracts;
using Elsa.Workflows.Runtime.Notifications;
namespace Elsa.Workflows.IntegrationTests.Scenarios.WorkflowDispatchNotifications;
public class TestHandler :
INotificationHandler<WorkflowDefinitionDispatching>,
INotificationHandler<WorkflowDefinitionDispatched>,
INotificationHandler<WorkflowInstanceDispatching>,
INotificationHandler<WorkflowInstanceDispatched>
{
private readonly Spy _spy;
public TestHandler(Spy spy)
{
_spy = spy;
}
public Task HandleAsync(WorkflowDefinitionDispatching notification, CancellationToken cancellationToken)
{
_spy.WorkflowDefinitionDispatchingWasCalled = true;
_spy.CapturedDefinitionRequest = notification.Request;
_spy.SignalWorkflowDefinitionDispatching();
return Task.CompletedTask;
}
public Task HandleAsync(WorkflowDefinitionDispatched notification, CancellationToken cancellationToken)
{
_spy.WorkflowDefinitionDispatchedWasCalled = true;
_spy.CapturedResponse = notification.Response;
_spy.SignalWorkflowDefinitionDispatched();
return Task.CompletedTask;
}
public Task HandleAsync(WorkflowInstanceDispatching notification, CancellationToken cancellationToken)
{
_spy.WorkflowInstanceDispatchingWasCalled = true;
_spy.CapturedInstanceRequest = notification.Request;
_spy.SignalWorkflowInstanceDispatching();
return Task.CompletedTask;
}
public Task HandleAsync(WorkflowInstanceDispatched notification, CancellationToken cancellationToken)
{
_spy.WorkflowInstanceDispatchedWasCalled = true;
_spy.CapturedResponse = notification.Response;
_spy.SignalWorkflowInstanceDispatched();
return Task.CompletedTask;
}
}

View file

@ -0,0 +1,87 @@
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<Spy>();
s.AddNotificationHandler<TestHandler, WorkflowDefinitionDispatching>();
s.AddNotificationHandler<TestHandler, WorkflowDefinitionDispatched>();
s.AddNotificationHandler<TestHandler, WorkflowInstanceDispatching>();
s.AddNotificationHandler<TestHandler, WorkflowInstanceDispatched>();
})
.Build();
_workflowDispatcher = services.GetRequiredService<IWorkflowDispatcher>();
_spy = services.GetRequiredService<Spy>();
}
[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<string, object> { { "TestKey", "TestValue" } }
};
// Act
await _workflowDispatcher.DispatchAsync(request, null);
// Wait for notification handlers to complete
await _spy.WaitForWorkflowDefinitionDispatchingAsync();
await _spy.WaitForWorkflowDefinitionDispatchedAsync();
// 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<string, object> { { "TestKey", "TestValue" } }
};
// Act
await _workflowDispatcher.DispatchAsync(request, null);
// Wait for notification handlers to complete
await _spy.WaitForWorkflowInstanceDispatchingAsync();
await _spy.WaitForWorkflowInstanceDispatchedAsync();
// 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);
}
}

View file

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

View file

@ -0,0 +1,26 @@
using Elsa.Workflows.CommitStates;
using Elsa.Workflows.State;
namespace Elsa.Workflows.IntegrationTests.SharedHelpers;
/// <summary>
/// Test helper to track commit invocations
/// </summary>
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;
}
}

View file

@ -0,0 +1,87 @@
using System.Net.Sockets;
namespace Elsa.Resilience.Core.UnitTests;
public class DefaultTransientExceptionStrategyTests
{
private readonly DefaultTransientExceptionStrategy _strategy = new();
public static TheoryData<Type> TransientExceptionTypes =>
[
typeof(HttpRequestException),
typeof(TimeoutException),
typeof(TaskCanceledException),
typeof(IOException),
typeof(SocketException),
typeof(EndOfStreamException)
];
public static TheoryData<string> 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<string> 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));
}
}

View file

@ -0,0 +1,13 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<Include>[Elsa.Resilience.Core]*</Include>
<Threshold>49</Threshold>
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\..\..\src\common\Elsa.Testing.Shared\Elsa.Testing.Shared.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Resilience.Core\Elsa.Resilience.Core.csproj" />
</ItemGroup>
</Project>

View file

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

View file

@ -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<IResilienceStrategyCatalog>();
private readonly IExpressionEvaluator _expressionEvaluator = Substitute.For<IExpressionEvaluator>();
private readonly ResilienceStrategyConfigEvaluator _evaluator;
private readonly ExpressionExecutionContext _context = new(Substitute.For<IServiceProvider>(), 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<CancellationToken>());
}
[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<string>(), Arg.Any<CancellationToken>());
}
[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<object>(Arg.Any<Expression>(), Arg.Any<ExpressionExecutionContext>(), Arg.Any<ExpressionEvaluatorOptions>());
}
[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<CancellationToken>());
}
[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<string>(), Arg.Any<CancellationToken>());
}
[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<CancellationToken>()).Returns(strategy);
}
private void SetupExpressionResult(Expression expression, object? result)
{
_expressionEvaluator.EvaluateAsync<object>(expression, _context, Arg.Any<ExpressionEvaluatorOptions>()).Returns(result);
}
}

View file

@ -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<IResilienceStrategy>();
strategy.Id.Returns(id);
strategy.DisplayName.Returns(displayName);
return strategy;
}
public static IResilienceStrategySource CreateStrategySource(params IResilienceStrategy[] strategies)
{
var source = Substitute.For<IResilienceStrategySource>();
source.GetStrategiesAsync(Arg.Any<CancellationToken>()).Returns(strategies);
return source;
}
}

View file

@ -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<ITransientExceptionStrategy>();
var strategy2 = Substitute.For<ITransientExceptionStrategy>();
strategy1.IsTransient(Arg.Any<Exception>()).Returns(false);
strategy2.IsTransient(Arg.Any<Exception>()).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<ITransientExceptionStrategy>();
strategy.IsTransient(transientException).Returns(true);
strategy.IsTransient(Arg.Is<Exception>(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<ITransientExceptionStrategy> configureDetector,
bool expectedResult)
{
var strategy = Substitute.For<ITransientExceptionStrategy>();
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<object[]> 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<object[]> AggregateExceptionTestCases
{
get
{
// Scenario 1: One transient inner exception
var timeoutException = new TimeoutException("timeout");
yield return
[
new AggregateException("aggregate", timeoutException),
(Action<ITransientExceptionStrategy>)(detector =>
{
detector.IsTransient(Arg.Is<TimeoutException>(e => e.Message == "timeout")).Returns(true);
detector.IsTransient(Arg.Is<Exception>(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<ITransientExceptionStrategy>)(detector =>
{
detector.IsTransient(Arg.Any<Exception>()).Returns(false);
}),
false
];
}
}
private static ITransientExceptionStrategy CreateStrategy(params (Exception exception, bool isTransient)[] behaviors)
{
var detector = Substitute.For<ITransientExceptionStrategy>();
foreach (var (exception, isTransient) in behaviors)
detector.IsTransient(exception).Returns(isTransient);
return detector;
}
private static TransientExceptionDetector CreateDetector(params ITransientExceptionStrategy[] detectors) => new(detectors);
}