From 689f19f2f2ffd3fc21d7ba5c4f0336d3060d68be Mon Sep 17 00:00:00 2001 From: lukhipolito-nexxbiz Date: Mon, 23 Feb 2026 20:56:15 +0100 Subject: [PATCH] Distributed locks to prevent multiple pods executing Reload and Refresh (#7311) * Distributed locks to prevent multiple pods executing Reload and Refresh * Update src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsRefresherDistributedLocking.cs Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsReloaderDistributedLocking.cs Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionsRefresherDistributedLocking.cs Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> * Moving distributed lock decorators to proper folder * Adding logging for missing lock * Fixing logger use in distributed classes * Update src/modules/Elsa.Workflows.Runtime/Distributed/WorkflowDefinitionsReloaderDistributedLocking.cs Co-authored-by: Sipke Schoorstra * Update src/modules/Elsa.Workflows.Runtime/Distributed/WorkflowDefinitionsRefresherDistributedLocking.cs Co-authored-by: Sipke Schoorstra * Update src/modules/Elsa.Workflows.Runtime/Distributed/WorkflowDefinitionsRefresherDistributedLocking.cs Co-authored-by: Sipke Schoorstra * Update src/modules/Elsa.Workflows.Runtime/Distributed/WorkflowDefinitionsReloaderDistributedLocking.cs Co-authored-by: Sipke Schoorstra * Update src/modules/Elsa.Workflows.Runtime/Distributed/WorkflowDefinitionsRefresherDistributedLocking.cs Co-authored-by: Sipke Schoorstra * Update src/modules/Elsa.Workflows.Runtime/Distributed/WorkflowDefinitionsReloaderDistributedLocking.cs Co-authored-by: Sipke Schoorstra * Moving and renaming for clearer purpose * Fixing build after file moves * Addressing workflow refresher empty query leaking distributed lock logic * Add status to WorkflowRefresherResponse model for in progress operation * Removing typo, fixing build * Update src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsRefresher.cs Co-authored-by: Sipke Schoorstra * Refactor `RefreshWorkflowDefinitionsAsync` for lock acquisition logic simplification. --------- Co-authored-by: lucas.hipolito Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> Co-authored-by: Sipke Schoorstra --- .../WorkflowDefinitions/Refresh/Endpoint.cs | 9 +++- .../Features/DistributedRuntimeFeature.cs | 5 +- ...DistributedWorkflowDefinitionsRefresher.cs | 46 +++++++++++++++++++ .../DistributedWorkflowDefinitionsReloader.cs | 36 +++++++++++++++ .../Enums/RefreshWorkflowDefinitionsStatus.cs | 7 +++ .../RefreshWorkflowDefinitionsResponse.cs | 6 ++- 6 files changed, 106 insertions(+), 3 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsRefresher.cs create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsReloader.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Enums/RefreshWorkflowDefinitionsStatus.cs diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs index b83a54931..509998b26 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowDefinitions/Refresh/Endpoint.cs @@ -20,7 +20,14 @@ internal class Refresh(IWorkflowDefinitionsRefresher workflowDefinitionsRefreshe public override async Task HandleAsync(Request request, CancellationToken cancellationToken) { var result = await RefreshWorkflowDefinitionsAsync(request.DefinitionIds, cancellationToken); - await Send.OkAsync(new Response(result.Refreshed, result.NotFound), cancellationToken); + if (result.Status == RefreshWorkflowDefinitionsStatus.Completed) + { + await Send.OkAsync(new Response(result.Refreshed, result.NotFound), cancellationToken); + } + else + { + await Send.AcceptedAtAsync(responseBody: new Response(result.Refreshed, result.NotFound), cancellation: cancellationToken); + } } private async Task RefreshWorkflowDefinitionsAsync(ICollection? definitionIds, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs index e783dd43e..a41bcf575 100644 --- a/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs @@ -29,6 +29,9 @@ public class DistributedRuntimeFeature(IModule module) : FeatureBase(module) { Services .AddScoped() - .AddScoped(); + .AddScoped() + + .Decorate() + .Decorate(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsRefresher.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsRefresher.cs new file mode 100644 index 000000000..5a7bb16ba --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsRefresher.cs @@ -0,0 +1,46 @@ +using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Responses; +using JetBrains.Annotations; +using Medallion.Threading; +using Microsoft.Extensions.Logging; + +namespace Elsa.Workflows.Runtime.Distributed; + +/// +/// Decorator class that adds distributed locking to the Workflow Definitions Refresher. +/// +[UsedImplicitly] +public class DistributedWorkflowDefinitionsRefresher(IWorkflowDefinitionsRefresher inner, + IDistributedLockProvider distributedLockProvider, + ILogger logger) : IWorkflowDefinitionsRefresher +{ + /// + /// This ensures that only one instance of the application can refresh a set of workflow definitions at a time, preventing potential conflicts and ensuring consistency across distributed environments. + /// + public async Task RefreshWorkflowDefinitionsAsync(RefreshWorkflowDefinitionsRequest request, CancellationToken cancellationToken = default) + { + var isRefreshingAll = request.DefinitionIds == null || request.DefinitionIds.Count == 0; + var lockKey = isRefreshingAll + ? "WorkflowDefinitionsRefresher:All" + : $"WorkflowDefinitionsRefresher:{string.Join(",", request.DefinitionIds!.OrderBy(x => x))}"; + + await using var distributedLock = await distributedLockProvider.TryAcquireLockAsync( + lockKey, + TimeSpan.Zero, + cancellationToken); + + if (distributedLock == null) + { + var logMessage = isRefreshingAll + ? "Could not acquire lock for refreshing all workflow definitions. Another instance is already refreshing all workflow definitions" + : "Could not acquire lock for refreshing workflow definitions. Another instance is already refreshing these workflow definitions"; + + logger.LogInformation(logMessage); + + var failedDefinitionIds = isRefreshingAll ? Array.Empty() : request.DefinitionIds!; + return new(Array.Empty(), failedDefinitionIds, RefreshWorkflowDefinitionsStatus.AlreadyInProgress); + } + + return await inner.RefreshWorkflowDefinitionsAsync(request, cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsReloader.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsReloader.cs new file mode 100644 index 000000000..40599f034 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowDefinitionsReloader.cs @@ -0,0 +1,36 @@ +using JetBrains.Annotations; +using Medallion.Threading; +using Microsoft.Extensions.Logging; + +namespace Elsa.Workflows.Runtime.Distributed; + +/// +/// Decorator class that adds distributed locking to the Workflow Definitions Reloader. +/// +[UsedImplicitly] +public class DistributedWorkflowDefinitionsReloader( + IWorkflowDefinitionsReloader inner, + IDistributedLockProvider distributedLockProvider, + ILogger logger) : IWorkflowDefinitionsReloader +{ + private const string LockKey = "WorkflowDefinitionsReloader"; + + /// + /// This ensures that only one instance of the application can reload workflow definitions at a time, preventing potential conflicts and ensuring consistency across distributed environments. + /// + public async Task ReloadWorkflowDefinitionsAsync(CancellationToken cancellationToken = default) + { + await using var distributedLock = await distributedLockProvider.TryAcquireLockAsync( + LockKey, + TimeSpan.Zero, + cancellationToken); + + if (distributedLock == null) + { + logger.LogInformation("Could not acquire lock for workflow definitions reload. Another instance is already reloading workflow definitions"); + return; + } + + await inner.ReloadWorkflowDefinitionsAsync(cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Enums/RefreshWorkflowDefinitionsStatus.cs b/src/modules/Elsa.Workflows.Runtime/Enums/RefreshWorkflowDefinitionsStatus.cs new file mode 100644 index 000000000..d7b260756 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Enums/RefreshWorkflowDefinitionsStatus.cs @@ -0,0 +1,7 @@ +namespace Elsa.Workflows.Runtime; + +public enum RefreshWorkflowDefinitionsStatus +{ + Completed, + AlreadyInProgress +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs b/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs index 6a9843f15..be2619ca1 100644 --- a/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs +++ b/src/modules/Elsa.Workflows.Runtime/Responses/RefreshWorkflowDefinitionsResponse.cs @@ -3,4 +3,8 @@ namespace Elsa.Workflows.Runtime.Responses; /// /// Represents a response to a request to refresh workflow definitions. /// -public record RefreshWorkflowDefinitionsResponse(ICollection Refreshed, ICollection NotFound); \ No newline at end of file +public record RefreshWorkflowDefinitionsResponse( + ICollection Refreshed, + ICollection NotFound, + RefreshWorkflowDefinitionsStatus Status = RefreshWorkflowDefinitionsStatus.Completed + ); \ No newline at end of file