* Add retry mechanism for distributed locks with transient error handling and logging - Introduced Polly-based retry pipeline for distributed lock acquisition in `DistributedWorkflowClient` to handle transient errors such as network issues or database connection failures. - Added detailed logging for retry attempts and lock release errors. - Updated project dependencies to include Polly. * Refactor transient exception handling to shared resilience module. Migrated transient exception detection logic from scheduling module to a new shared resilience module. Updated services, jobs, and features to utilize the centralized `ITransientExceptionDetectionService`. This change improves maintainability and promotes reusability across modules. * Add unit tests for transient exception detection and resilience strategy evaluation. - Introduced comprehensive unit tests for `DefaultTransientExceptionDetector`, `ResilienceStrategyCatalog`, `ResilienceStrategyConfigEvaluator`, and `TransientExceptionDetectionService`. - Added helper classes and test data factories to facilitate reusable test patterns for resilience modules. - Updated solution to include `Elsa.Resilience.Core.UnitTests` project. * Add component tests for distributed lock resilience - Introduced new tests to verify retry behavior during transient lock acquisition and release failures. - Added `TestDistributedLockProvider` and related mocks for simulating transient failures. - Updated `WorkflowServer` test services to support the new distributed lock test scenarios. * Refactor distributed lock resilience tests - Consolidated test logic: streamlined test providers, injected services, and reusable test patterns. - Simplified `TestDistributedLockProvider` implementation with enhanced initialization and failure simulation. - Reorganized tests for transient acquisition/release failures to use parameterized `Theory` for improved maintainability. * Refactor transient exception handling: rename interfaces and classes for consistency, update references across codebase, and improve code readability. * Apply suggestion from @Copilot Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Simplify WorkflowServer setup and DistributedLockResilienceTests by replacing IDistributedLockProvider with TestDistributedLockProvider. * Remove unused `using` directives in unit tests to improve code cleanliness. * Remove `TransientExceptionTypes` helper and inline its usage in tests for improved maintainability. * Add descriptive `DisplayName` attributes to unit tests for improved test clarity. * Fix redundant exception checking in TransientExceptionDetector (#7162) * Initial plan * Fix redundant exception checking in TransientExceptionDetector Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Extract MaxRetryAttempts constant in DistributedLockResilienceTests (#7164) * Initial plan * Extract MaxRetryAttempts constant to eliminate hardcoded magic number Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Make TestDistributedLockProvider thread-safe with Interlocked operations (#7163) * Initial plan * Make TestDistributedLockProvider thread-safe using Interlocked operations Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update test/component/Elsa.Workflows.ComponentTests/Scenarios/DistributedLockResilience/DistributedLockResilienceTests.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Resilience.Core/Services/DefaultTransientExceptionStrategy.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Include `CancellationToken` in distributed lock handling methods for improved cancellation support. * Initial plan (#7168) Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> * Cache detector list in TransientExceptionDetector to avoid repeated allocations (#7167) * Initial plan * Cache detector list in field to avoid repeated allocations Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Use IReadOnlyList instead of List for better intent expression Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Fix TestDistributedLockProvider registration to properly decorate IDistributedLockProvider (#7166) * Initial plan * Fix TestDistributedLockProvider registration to use Decorate pattern and fix variable reference bug Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Add runtime check for TestDistributedLockProvider registration Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Refactor `DistributedWorkflowClient` to simplify `Lazy<ResiliencePipeline>` initialization. * Add integration tests for DistributedWorkflowClient lock resilience (#7165) * Initial plan * Fix compilation error: use correct parameter name transientExceptionDetector Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Add integration tests for DistributedWorkflowClient lock resilience - Add SimpleWorkflow for testing distributed lock scenarios - Add tests exercising RunInstanceAsync with transient lock failures - Verify retry logic works correctly with actual workflow execution - Test both acquisition and release failure scenarios - Decorate IDistributedLockProvider to use TestDistributedLockProvider Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> * Address code review feedback - Add explanatory comment for TestDistributedLockProvider cast - Remove unnecessary blank line for consistent formatting Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com> Co-authored-by: Sipke Schoorstra <sipkeschoorstra@outlook.com> * Simplify transient exception strategy by refactoring message pattern matching logic. * Enhance distributed lock mock to support per-lock failure configuration and improve resilience tests. * Refactor `TestDistributedLockProvider` to streamline failure handling logic and improve code clarity. * Remove unused methods and redundant test case from `DistributedLockResilienceTests`. * Refactor `DistributedLockResilienceTests` to simplify workflow client creation, consolidate assertion logic, and remove redundant test cases. * Format `ResilienceStrategyCatalogTests` by removing redundant line breaks in test setup. * Refactor `TransientExceptionDetectorTests` to simplify test setup, consolidate test cases, and remove redundant logic. * Handle `InvalidOperationException` in `XunitLogger` to suppress logging errors during inactive tests. * Update workflows to use .NET 10 and adjust resilience tests project configuration. * Refactor distributed runtime feature and integrations to improve resilience handling, configure services fluently, and add cancellation safeguards in background services. * Update `base_version` to `3.7.0` in GitHub workflow configuration. --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> Co-authored-by: Copilot <198982749+Copilot@users.noreply.github.com> Co-authored-by: sfmskywalker <938393+sfmskywalker@users.noreply.github.com>
105 lines
4.4 KiB
C#
105 lines
4.4 KiB
C#
using System.Threading.Channels;
|
|
using Elsa.Mediator.Contracts;
|
|
using Elsa.Mediator.Middleware.Command;
|
|
using Elsa.Mediator.Options;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Microsoft.Extensions.Hosting;
|
|
using Microsoft.Extensions.Logging;
|
|
using Microsoft.Extensions.Options;
|
|
|
|
namespace Elsa.Mediator.HostedServices;
|
|
|
|
/// <summary>
|
|
/// Continuously reads from a channel to which commands can be sent, executing each received command.
|
|
/// </summary>
|
|
public class BackgroundCommandSenderHostedService : BackgroundService
|
|
{
|
|
private readonly int _workerCount;
|
|
private readonly ICommandsChannel _commandsChannel;
|
|
private readonly IServiceScopeFactory _scopeFactory;
|
|
private readonly List<Channel<CommandContext>> _outputs;
|
|
private readonly ILogger _logger;
|
|
|
|
/// <inheritdoc />
|
|
public BackgroundCommandSenderHostedService(IOptions<MediatorOptions> options, ICommandsChannel commandsChannel, IServiceScopeFactory scopeFactory, ILogger<BackgroundCommandSenderHostedService> logger)
|
|
{
|
|
_workerCount = options.Value.CommandWorkerCount;
|
|
_commandsChannel = commandsChannel; // The shared input channel for all commands
|
|
_scopeFactory = scopeFactory;
|
|
_logger = logger;
|
|
_outputs = new(_workerCount); // Prepare a list to hold worker-specific channels
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
protected override async Task ExecuteAsync(CancellationToken cancellationToken)
|
|
{
|
|
var index = 0; // Used for round-robin distribution of work
|
|
|
|
// Set up worker channels and start background tasks for each worker
|
|
for (var i = 0; i < _workerCount; i++)
|
|
{
|
|
var output = Channel.CreateUnbounded<CommandContext>();
|
|
_outputs.Add(output);
|
|
// Start a background task that processes commands from this worker's channel
|
|
_ = ReadOutputAsync(output, cancellationToken);
|
|
}
|
|
|
|
// Main dispatcher loop: read from the input channel and distribute to worker channels
|
|
try
|
|
{
|
|
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
|
|
foreach (var output in _outputs)
|
|
output.Writer.Complete();
|
|
}
|
|
|
|
private async Task ReadOutputAsync(Channel<CommandContext> output, CancellationToken cancellationToken)
|
|
{
|
|
// Worker task: process commands from the worker's channel
|
|
try
|
|
{
|
|
await foreach (var commandContext in output.Reader.ReadAllAsync(cancellationToken))
|
|
{
|
|
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);
|
|
|
|
// 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");
|
|
}
|
|
}
|
|
} |