| | | 1 | | using Elsa.Common.DistributedHosting; |
| | | 2 | | using Elsa.Resilience; |
| | | 3 | | using Elsa.Workflows.Runtime.Messages; |
| | | 4 | | using Elsa.Workflows.State; |
| | | 5 | | using Medallion.Threading; |
| | | 6 | | using Microsoft.Extensions.DependencyInjection; |
| | | 7 | | using Microsoft.Extensions.Logging; |
| | | 8 | | using Microsoft.Extensions.Options; |
| | | 9 | | using Polly; |
| | | 10 | | |
| | | 11 | | namespace Elsa.Workflows.Runtime.Distributed; |
| | | 12 | | |
| | 138 | 13 | | public class DistributedWorkflowClient( |
| | 138 | 14 | | string workflowInstanceId, |
| | 138 | 15 | | IDistributedLockProvider distributedLockProvider, |
| | 138 | 16 | | ITransientExceptionDetector transientExceptionDetector, |
| | 138 | 17 | | IOptions<DistributedLockingOptions> distributedLockingOptions, |
| | 138 | 18 | | IServiceProvider serviceProvider, |
| | 138 | 19 | | ILogger<DistributedWorkflowClient> logger) |
| | | 20 | | : IWorkflowClient |
| | | 21 | | { |
| | 138 | 22 | | private readonly LocalWorkflowClient _localWorkflowClient = ActivatorUtilities.CreateInstance<LocalWorkflowClient>(s |
| | 262 | 23 | | private readonly Lazy<ResiliencePipeline> _retryPipeline = new(() => CreateRetryPipeline(transientExceptionDetector, |
| | 135 | 24 | | public string WorkflowInstanceId => workflowInstanceId; |
| | | 25 | | |
| | | 26 | | public async Task<CreateWorkflowInstanceResponse> CreateInstanceAsync(CreateWorkflowInstanceRequest request, Cancell |
| | | 27 | | { |
| | 32 | 28 | | return await _localWorkflowClient.CreateInstanceAsync(request, cancellationToken); |
| | 32 | 29 | | } |
| | | 30 | | |
| | | 31 | | public async Task<RunWorkflowInstanceResponse> RunInstanceAsync(RunWorkflowInstanceRequest request, CancellationToke |
| | | 32 | | { |
| | 133 | 33 | | var result = await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(request, cancellationTo |
| | 66 | 34 | | return result; |
| | 66 | 35 | | } |
| | | 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 |
| | 112 | 41 | | return await WithLockAsync(async () => await _localWorkflowClient.CreateAndRunInstanceAsync(request, cancellatio |
| | 56 | 42 | | } |
| | | 43 | | |
| | | 44 | | public async Task CancelAsync(CancellationToken cancellationToken = default) |
| | | 45 | | { |
| | 0 | 46 | | await _localWorkflowClient.CancelAsync(cancellationToken); |
| | 0 | 47 | | } |
| | | 48 | | |
| | | 49 | | public async Task<WorkflowState> ExportStateAsync(CancellationToken cancellationToken = default) |
| | | 50 | | { |
| | 14 | 51 | | return await _localWorkflowClient.ExportStateAsync(cancellationToken); |
| | 14 | 52 | | } |
| | | 53 | | |
| | | 54 | | public async Task ImportStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) |
| | | 55 | | { |
| | 0 | 56 | | await _localWorkflowClient.ImportStateAsync(workflowState, cancellationToken); |
| | 0 | 57 | | } |
| | | 58 | | |
| | | 59 | | public async Task<bool> InstanceExistsAsync(CancellationToken cancellationToken = default) |
| | | 60 | | { |
| | 0 | 61 | | return await _localWorkflowClient.InstanceExistsAsync(cancellationToken); |
| | 0 | 62 | | } |
| | | 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 |
| | 16 | 67 | | return await WithLockAsync(async () => await _localWorkflowClient.DeleteAsync(cancellationToken), cancellationTo |
| | 8 | 68 | | } |
| | | 69 | | |
| | | 70 | | private async Task<TReturn> WithLockAsync<TReturn>(Func<Task<TReturn>> func, CancellationToken cancellationToken = d |
| | | 71 | | { |
| | 131 | 72 | | var lockKey = $"workflow-instance:{WorkflowInstanceId}"; |
| | 131 | 73 | | var lockHandle = await AcquireLockWithRetryAsync(lockKey, cancellationToken); |
| | | 74 | | |
| | | 75 | | try |
| | | 76 | | { |
| | 130 | 77 | | return await func(); |
| | | 78 | | } |
| | | 79 | | finally |
| | | 80 | | { |
| | 130 | 81 | | await ReleaseLockAsync(lockHandle); |
| | | 82 | | } |
| | 130 | 83 | | } |
| | | 84 | | |
| | | 85 | | private async Task<IDistributedSynchronizationHandle?> AcquireLockWithRetryAsync(string lockKey, CancellationToken c |
| | | 86 | | { |
| | 131 | 87 | | var lockTimeout = distributedLockingOptions.Value.LockAcquisitionTimeout; |
| | | 88 | | |
| | 131 | 89 | | return await _retryPipeline.Value.ExecuteAsync(async ct => |
| | 137 | 90 | | await distributedLockProvider.AcquireLockAsync(lockKey, lockTimeout, ct), |
| | 131 | 91 | | cancellationToken); |
| | 130 | 92 | | } |
| | | 93 | | |
| | | 94 | | private async Task ReleaseLockAsync(IDistributedSynchronizationHandle? lockHandle) |
| | | 95 | | { |
| | 130 | 96 | | if (lockHandle == null) |
| | 0 | 97 | | return; |
| | | 98 | | |
| | | 99 | | try |
| | | 100 | | { |
| | 130 | 101 | | await lockHandle.DisposeAsync(); |
| | 129 | 102 | | } |
| | 1 | 103 | | 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 |
| | 1 | 107 | | logger.LogWarning(ex, "Failed to release distributed lock for workflow instance {WorkflowInstanceId}. The lo |
| | 1 | 108 | | } |
| | 130 | 109 | | } |
| | | 110 | | |
| | | 111 | | private static ResiliencePipeline CreateRetryPipeline( |
| | | 112 | | ITransientExceptionDetector transientExceptionDetector, |
| | | 113 | | ILogger<DistributedWorkflowClient> logger, |
| | | 114 | | string workflowInstanceId) |
| | | 115 | | { |
| | | 116 | | const int maxRetryAttempts = 3; |
| | | 117 | | |
| | 124 | 118 | | return new ResiliencePipelineBuilder() |
| | 124 | 119 | | .AddRetry(new() |
| | 124 | 120 | | { |
| | 124 | 121 | | MaxRetryAttempts = maxRetryAttempts, |
| | 124 | 122 | | Delay = TimeSpan.FromMilliseconds(500), |
| | 124 | 123 | | BackoffType = DelayBackoffType.Exponential, |
| | 124 | 124 | | UseJitter = true, |
| | 124 | 125 | | ShouldHandle = new PredicateBuilder().Handle<Exception>(transientExceptionDetector.IsTransient), |
| | 124 | 126 | | OnRetry = args => |
| | 124 | 127 | | { |
| | 6 | 128 | | logger.LogWarning(args.Outcome.Exception, "Transient error acquiring lock for workflow instance {Wor |
| | 6 | 129 | | return ValueTask.CompletedTask; |
| | 124 | 130 | | } |
| | 124 | 131 | | }) |
| | 124 | 132 | | .Build(); |
| | | 133 | | } |
| | | 134 | | } |