< Summary

Information
Class: Elsa.Workflows.Runtime.Distributed.DistributedWorkflowClient
Assembly: Elsa.Workflows.Runtime.Distributed
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs
Line coverage
88%
Covered lines: 53
Uncovered lines: 7
Coverable lines: 60
Total lines: 134
Line coverage: 88.3%
Branch coverage
50%
Covered branches: 1
Total branches: 2
Branch coverage: 50%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
get_WorkflowInstanceId()100%11100%
CreateInstanceAsync()100%11100%
<RunInstanceAsync()100%11100%
RunInstanceAsync()100%11100%
<CreateAndRunInstanceAsync()100%11100%
CreateAndRunInstanceAsync()100%11100%
CancelAsync()100%210%
ExportStateAsync()100%11100%
ImportStateAsync()100%210%
InstanceExistsAsync()100%210%
<DeleteAsync()100%11100%
DeleteAsync()100%11100%
WithLockAsync()100%11100%
AcquireLockWithRetryAsync()100%11100%
<AcquireLockWithRetryAsync()100%11100%
ReleaseLockAsync()50%2287.5%
CreateRetryPipeline(...)100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs

#LineLine coverage
 1using Elsa.Common.DistributedHosting;
 2using Elsa.Resilience;
 3using Elsa.Workflows.Runtime.Messages;
 4using Elsa.Workflows.State;
 5using Medallion.Threading;
 6using Microsoft.Extensions.DependencyInjection;
 7using Microsoft.Extensions.Logging;
 8using Microsoft.Extensions.Options;
 9using Polly;
 10
 11namespace Elsa.Workflows.Runtime.Distributed;
 12
 13813public class DistributedWorkflowClient(
 13814    string workflowInstanceId,
 13815    IDistributedLockProvider distributedLockProvider,
 13816    ITransientExceptionDetector transientExceptionDetector,
 13817    IOptions<DistributedLockingOptions> distributedLockingOptions,
 13818    IServiceProvider serviceProvider,
 13819    ILogger<DistributedWorkflowClient> logger)
 20    : IWorkflowClient
 21{
 13822    private readonly LocalWorkflowClient _localWorkflowClient = ActivatorUtilities.CreateInstance<LocalWorkflowClient>(s
 26223    private readonly Lazy<ResiliencePipeline> _retryPipeline = new(() => CreateRetryPipeline(transientExceptionDetector,
 13524    public string WorkflowInstanceId => workflowInstanceId;
 25
 26    public async Task<CreateWorkflowInstanceResponse> CreateInstanceAsync(CreateWorkflowInstanceRequest request, Cancell
 27    {
 3228        return await _localWorkflowClient.CreateInstanceAsync(request, cancellationToken);
 3229    }
 30
 31    public async Task<RunWorkflowInstanceResponse> RunInstanceAsync(RunWorkflowInstanceRequest request, CancellationToke
 32    {
 13333        var result = await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(request, cancellationTo
 6634        return result;
 6635    }
 36
 37    public async Task<RunWorkflowInstanceResponse> CreateAndRunInstanceAsync(CreateAndRunWorkflowInstanceRequest request
 38    {
 39        // We need to lock newly created workflow instances too, because it might dispatch child workflows that attempt 
 40        // For example, when using a DispatchWorkflow activity configured to wait for the dispatched workflow to complet
 11241        return await WithLockAsync(async () => await _localWorkflowClient.CreateAndRunInstanceAsync(request, cancellatio
 5642    }
 43
 44    public async Task CancelAsync(CancellationToken cancellationToken = default)
 45    {
 046        await _localWorkflowClient.CancelAsync(cancellationToken);
 047    }
 48
 49    public async Task<WorkflowState> ExportStateAsync(CancellationToken cancellationToken = default)
 50    {
 1451        return await _localWorkflowClient.ExportStateAsync(cancellationToken);
 1452    }
 53
 54    public async Task ImportStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default)
 55    {
 056        await _localWorkflowClient.ImportStateAsync(workflowState, cancellationToken);
 057    }
 58
 59    public async Task<bool> InstanceExistsAsync(CancellationToken cancellationToken = default)
 60    {
 061        return await _localWorkflowClient.InstanceExistsAsync(cancellationToken);
 062    }
 63
 64    public async Task<bool> DeleteAsync(CancellationToken cancellationToken = default)
 65    {
 66        // Use the same distributed lock as for execution to prevent concurrent DB writes
 1667        return await WithLockAsync(async () => await _localWorkflowClient.DeleteAsync(cancellationToken), cancellationTo
 868    }
 69
 70    private async Task<TReturn> WithLockAsync<TReturn>(Func<Task<TReturn>> func, CancellationToken cancellationToken = d
 71    {
 13172        var lockKey = $"workflow-instance:{WorkflowInstanceId}";
 13173        var lockHandle = await AcquireLockWithRetryAsync(lockKey, cancellationToken);
 74
 75        try
 76        {
 13077            return await func();
 78        }
 79        finally
 80        {
 13081            await ReleaseLockAsync(lockHandle);
 82        }
 13083    }
 84
 85    private async Task<IDistributedSynchronizationHandle?> AcquireLockWithRetryAsync(string lockKey, CancellationToken c
 86    {
 13187        var lockTimeout = distributedLockingOptions.Value.LockAcquisitionTimeout;
 88
 13189        return await _retryPipeline.Value.ExecuteAsync(async ct =>
 13790            await distributedLockProvider.AcquireLockAsync(lockKey, lockTimeout, ct),
 13191            cancellationToken);
 13092    }
 93
 94    private async Task ReleaseLockAsync(IDistributedSynchronizationHandle? lockHandle)
 95    {
 13096        if (lockHandle == null)
 097            return;
 98
 99        try
 100        {
 130101            await lockHandle.DisposeAsync();
 129102        }
 1103        catch (Exception ex)
 104        {
 105            // Log but don't throw - the work is already done, and the lock
 106            // will be automatically released when the connection dies
 1107            logger.LogWarning(ex, "Failed to release distributed lock for workflow instance {WorkflowInstanceId}. The lo
 1108        }
 130109    }
 110
 111    private static ResiliencePipeline CreateRetryPipeline(
 112        ITransientExceptionDetector transientExceptionDetector,
 113        ILogger<DistributedWorkflowClient> logger,
 114        string workflowInstanceId)
 115    {
 116        const int maxRetryAttempts = 3;
 117
 124118        return new ResiliencePipelineBuilder()
 124119            .AddRetry(new()
 124120            {
 124121                MaxRetryAttempts = maxRetryAttempts,
 124122                Delay = TimeSpan.FromMilliseconds(500),
 124123                BackoffType = DelayBackoffType.Exponential,
 124124                UseJitter = true,
 124125                ShouldHandle = new PredicateBuilder().Handle<Exception>(transientExceptionDetector.IsTransient),
 124126                OnRetry = args =>
 124127                {
 6128                    logger.LogWarning(args.Outcome.Exception, "Transient error acquiring lock for workflow instance {Wor
 6129                    return ValueTask.CompletedTask;
 124130                }
 124131            })
 124132            .Build();
 133    }
 134}