From 87132d49c5a1189e42ea34703efd0b37f80fbe98 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 19 Aug 2021 14:45:26 +0200 Subject: [PATCH] Initial fix for deadlock situation When a parent workflow triggers another workflow using signaling and using the same correlation ID, the distributed lock provider would try to acquire a lock on the same resource, causing a deadlock situation. This isn't the final fix. --- .../Services/Models/AmbientLockContext.cs | 30 ++++++ .../LockingWorkflowInstanceExecutor.cs | 69 +++++++++----- .../Services/Workflows/WorkflowLaunchpad.cs | 94 +++++++++++-------- 3 files changed, 126 insertions(+), 67 deletions(-) create mode 100644 src/core/Elsa.Abstractions/Services/Models/AmbientLockContext.cs diff --git a/src/core/Elsa.Abstractions/Services/Models/AmbientLockContext.cs b/src/core/Elsa.Abstractions/Services/Models/AmbientLockContext.cs new file mode 100644 index 000000000..a6f6ca6e0 --- /dev/null +++ b/src/core/Elsa.Abstractions/Services/Models/AmbientLockContext.cs @@ -0,0 +1,30 @@ +using System.Collections.Generic; +using System.Threading; +using Medallion.Threading; + +namespace Elsa.Services.Models +{ + public static class AmbientLockContext + { + private static readonly AsyncLocal CorrelationLock = new(); + private static readonly AsyncLocal> WorkflowInstanceLocks = new(); + + public static IDistributedSynchronizationHandle? CurrentCorrelationLock + { + get => CorrelationLock.Value; + set => CorrelationLock.Value = value; + } + + public static IDistributedSynchronizationHandle? GetCurrentWorkflowInstanceLock(string workflowInstanceId) => + WorkflowInstanceLocks.Value != null ? WorkflowInstanceLocks.Value.TryGetValue(workflowInstanceId, out var handle) ? handle : default : default; + + public static void SetCurrentWorkflowInstanceLock(string workflowInstanceId, IDistributedSynchronizationHandle handle) + { + var dictionary = WorkflowInstanceLocks.Value ?? new Dictionary(); + dictionary[workflowInstanceId] = handle; + WorkflowInstanceLocks.Value = dictionary; + } + + public static void DeleteCurrentWorkflowInstanceLock(string workflowInstanceId) => WorkflowInstanceLocks.Value.Remove(workflowInstanceId); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs b/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs index 26c4508b8..159a2ac7e 100644 --- a/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs +++ b/src/core/Elsa.Core/Decorators/LockingWorkflowInstanceExecutor.cs @@ -1,4 +1,5 @@ -using System.Threading; +using System; +using System.Threading; using System.Threading.Tasks; using Elsa.Exceptions; using Elsa.Models; @@ -34,38 +35,54 @@ namespace Elsa.Decorators public async Task ExecuteAsync(string workflowInstanceId, string? activityId, WorkflowInput? input = default, CancellationToken cancellationToken = default) { var workflowInstanceLockKey = $"workflow-instance:{workflowInstanceId}"; - await using var workflowInstanceLockHandle = await _distributedLockProvider.AcquireLockAsync(workflowInstanceLockKey, _elsaOptions.DistributedLockTimeout, cancellationToken); + var currentWorkflowInstanceLockHandle = AmbientLockContext.GetCurrentWorkflowInstanceLock(workflowInstanceId); + var workflowInstanceLockHandle = currentWorkflowInstanceLockHandle ?? await _distributedLockProvider.AcquireLockAsync(workflowInstanceLockKey, _elsaOptions.DistributedLockTimeout, cancellationToken); if (workflowInstanceLockHandle == null) throw new LockAcquisitionException("Could not acquire a lock within the configured amount of time"); - - var workflowInstance = await _workflowInstanceStore.FindByIdAsync(workflowInstanceId, cancellationToken); - if (workflowInstance == null) + try { - _logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist", workflowInstanceId); - return new RunWorkflowResult(workflowInstance, activityId, false); + AmbientLockContext.SetCurrentWorkflowInstanceLock(workflowInstanceId, workflowInstanceLockHandle); + var workflowInstance = await _workflowInstanceStore.FindByIdAsync(workflowInstanceId, cancellationToken); + + if (workflowInstance == null) + { + _logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist", workflowInstanceId); + return new RunWorkflowResult(workflowInstance, activityId, false); + } + + var correlationId = workflowInstance.CorrelationId; + + if (!string.IsNullOrWhiteSpace(correlationId)) + { + // We need to lock on correlation ID to prevent a race condition with WorkflowLaunchpad that is used to find workflows by correlation ID to execute. + // The race condition is: when a workflow instance is done executing, the BookmarkIndexer will collect bookmarks. + // But if in the meantime an event comes in that triggers correlated workflows, the bookmarks may not have been created yet. + var currentCorrelationLockHandle = AmbientLockContext.CurrentCorrelationLock; + var correlationLockHandle = currentCorrelationLockHandle ?? await _distributedLockProvider.AcquireLockAsync(correlationId, _elsaOptions.DistributedLockTimeout, cancellationToken); + + if (correlationLockHandle == null) + throw new LockAcquisitionException($"Could not acquire a lock on correlation {correlationId} within the configured amount of time"); + + try + { + AmbientLockContext.CurrentCorrelationLock = correlationLockHandle; + return await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); + } + finally + { + AmbientLockContext.CurrentCorrelationLock = null; + await correlationLockHandle.DisposeAsync(); + } + } + + return await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); } - - var correlationId = workflowInstance.CorrelationId; - - if (!string.IsNullOrWhiteSpace(correlationId)) + finally { - // We need to lock on correlation ID to prevent a race condition with WorkflowLaunchpad that is used to find workflows by correlation ID to execute. - // The race condition is: when a workflow instance is done executing, the BookmarkIndexer will collect bookmarks. - // But if in the meantime an event comes in that triggers correlated workflows, the bookmarks may not have been created yet. - await using var correlationLockHandle = await _distributedLockProvider.AcquireLockAsync(correlationId, _elsaOptions.DistributedLockTimeout, cancellationToken); - - if(correlationLockHandle == null) - throw new LockAcquisitionException($"Could not acquire a lock on correlation {correlationId} within the configured amount of time"); - - var result = await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); - return result; - } - else - { - var result = await _workflowInstanceExecutor.ExecuteAsync(workflowInstance, activityId, input, cancellationToken); - return result; + AmbientLockContext.DeleteCurrentWorkflowInstanceLock(workflowInstanceId); + await workflowInstanceLockHandle.DisposeAsync(); } } diff --git a/src/core/Elsa.Core/Services/Workflows/WorkflowLaunchpad.cs b/src/core/Elsa.Core/Services/Workflows/WorkflowLaunchpad.cs index bd2fb55b6..42b2adc64 100644 --- a/src/core/Elsa.Core/Services/Workflows/WorkflowLaunchpad.cs +++ b/src/core/Elsa.Core/Services/Workflows/WorkflowLaunchpad.cs @@ -76,16 +76,16 @@ namespace Elsa.Services.Workflows public async Task> FindStartableWorkflowsAsync(WorkflowsQuery query, CancellationToken cancellationToken = default) { var correlationId = query.CorrelationId ?? Guid.NewGuid().ToString("N"); - var updatedContext = query with {CorrelationId = correlationId}; + var updatedContext = query with { CorrelationId = correlationId }; await using var lockHandle = await AcquireLockAsync(correlationId, cancellationToken); return await CollectStartableWorkflowsInternalAsync(updatedContext, cancellationToken); } public async Task FindStartableWorkflowAsync( - string workflowDefinitionId, - string? activityId, - string? correlationId = default, - string? contextId = default, + string workflowDefinitionId, + string? activityId, + string? correlationId = default, + string? contextId = default, string? tenantId = default, CancellationToken cancellationToken = default) { @@ -98,15 +98,15 @@ namespace Elsa.Services.Workflows } public async Task FindStartableWorkflowAsync( - IWorkflowBlueprint workflowBlueprint, - string? activityId, - string? correlationId = default, - string? contextId = default, + IWorkflowBlueprint workflowBlueprint, + string? activityId, + string? correlationId = default, + string? contextId = default, string? tenantId = default, CancellationToken cancellationToken = default) { correlationId ??= Guid.NewGuid().ToString("N"); - + // Acquire a lock on correlation ID to prevent duplicate workflow instances from being created. await using var correlationLockHandle = await AcquireLockAsync(correlationId, cancellationToken); @@ -114,11 +114,11 @@ namespace Elsa.Services.Workflows } public async Task FindAndExecuteStartableWorkflowAsync( - string workflowDefinitionId, - string? activityId, - string? correlationId = default, - string? contextId = default, - WorkflowInput? input = default, + string workflowDefinitionId, + string? activityId, + string? correlationId = default, + string? contextId = default, + WorkflowInput? input = default, string? tenantId = default, CancellationToken cancellationToken = default) { @@ -134,10 +134,10 @@ namespace Elsa.Services.Workflows } public async Task FindAndExecuteStartableWorkflowAsync( - IWorkflowBlueprint workflowBlueprint, - string? activityId, - string? correlationId = default, - string? contextId = default, + IWorkflowBlueprint workflowBlueprint, + string? activityId, + string? correlationId = default, + string? contextId = default, WorkflowInput? input = default, CancellationToken cancellationToken = default) { @@ -210,7 +210,7 @@ namespace Elsa.Services.Workflows return pendingWorkflows; } - + private async Task> CollectStartableWorkflowsInternalAsync(WorkflowsQuery query, CancellationToken cancellationToken = default) { _logger.LogDebug("Triggering workflows using {ActivityType}", query.ActivityType); @@ -232,10 +232,10 @@ namespace Elsa.Services.Workflows } private async Task CollectStartableWorkflowInternalAsync( - IWorkflowBlueprint workflowBlueprint, - string? activityId, - string correlationId, - string? contextId = default, + IWorkflowBlueprint workflowBlueprint, + string? activityId, + string correlationId, + string? contextId = default, string? tenantId = default, CancellationToken cancellationToken = default) { @@ -289,25 +289,37 @@ namespace Elsa.Services.Workflows private async Task> CollectResumableOrStartableCorrelatedWorkflowsAsync(WorkflowsQuery query, CancellationToken cancellationToken) { var correlationId = query.CorrelationId!; - - await using var handle = await AcquireLockAsync(correlationId, cancellationToken); - var correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId) - ? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId).And(new WorkflowUnfinishedStatusSpecification()), cancellationToken) - : 0; - - _logger.LogDebug("Found {CorrelatedWorkflowCount} workflows with correlation ID {CorrelationId}", correlatedWorkflowInstanceCount, correlationId); - - if (correlatedWorkflowInstanceCount > 0) + var existingHandle = AmbientLockContext.CurrentCorrelationLock; + var handle = existingHandle == null ? await AcquireLockAsync(correlationId, cancellationToken) : default; + + try { - var bookmarkResults = query.Bookmark != null - ? await _bookmarkFinder.FindBookmarksAsync(query.ActivityType, query.Bookmark, correlationId, query.TenantId, cancellationToken).ToList() - : new List(); - _logger.LogDebug("Found {BookmarkCount} bookmarks for activity type {ActivityType}", bookmarkResults.Count, query.ActivityType); - return bookmarkResults.Select(x => new CollectedWorkflow(x.WorkflowInstanceId, x.ActivityId)).ToList(); - } + var correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId) + ? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId).And(new WorkflowUnfinishedStatusSpecification()), cancellationToken) + : 0; - var startableWorkflows = await CollectStartableWorkflowsInternalAsync(query, cancellationToken); - return startableWorkflows.Select(x => new CollectedWorkflow(x.WorkflowInstance.Id, x.ActivityId)).ToList(); + _logger.LogDebug("Found {CorrelatedWorkflowCount} workflows with correlation ID {CorrelationId}", correlatedWorkflowInstanceCount, correlationId); + + if (correlatedWorkflowInstanceCount > 0) + { + var bookmarkResults = query.Bookmark != null + ? await _bookmarkFinder.FindBookmarksAsync(query.ActivityType, query.Bookmark, correlationId, query.TenantId, cancellationToken).ToList() + : new List(); + _logger.LogDebug("Found {BookmarkCount} bookmarks for activity type {ActivityType}", bookmarkResults.Count, query.ActivityType); + + //// Only return if we actually found results. If we didn't find results, continue looking for startable workflows. + //if(bookmarkResults.Any()) + return bookmarkResults.Select(x => new CollectedWorkflow(x.WorkflowInstanceId, x.ActivityId)).ToList(); + } + + var startableWorkflows = await CollectStartableWorkflowsInternalAsync(query, cancellationToken); + return startableWorkflows.Select(x => new CollectedWorkflow(x.WorkflowInstance.Id, x.ActivityId)).ToList(); + } + finally + { + if (handle != null) + await handle.DisposeAsync(); + } } private async Task AcquireLockAsync(string resource, CancellationToken cancellationToken)