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.
This commit is contained in:
Sipke Schoorstra 2021-08-19 14:45:26 +02:00
parent e24c9dd59d
commit 87132d49c5
3 changed files with 126 additions and 67 deletions

View file

@ -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<IDistributedSynchronizationHandle?> CorrelationLock = new();
private static readonly AsyncLocal<IDictionary<string, IDistributedSynchronizationHandle?>> 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<string, IDistributedSynchronizationHandle?>();
dictionary[workflowInstanceId] = handle;
WorkflowInstanceLocks.Value = dictionary;
}
public static void DeleteCurrentWorkflowInstanceLock(string workflowInstanceId) => WorkflowInstanceLocks.Value.Remove(workflowInstanceId);
}
}

View file

@ -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<RunWorkflowResult> 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();
}
}

View file

@ -76,16 +76,16 @@ namespace Elsa.Services.Workflows
public async Task<IEnumerable<StartableWorkflow>> 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<StartableWorkflow?> 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<StartableWorkflow?> 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<RunWorkflowResult> 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<IEnumerable<StartableWorkflow>> 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<StartableWorkflow?> 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<IEnumerable<CollectedWorkflow>> 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<WorkflowInstance>(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<BookmarkFinderResult>();
_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<WorkflowInstance>(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<BookmarkFinderResult>();
_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<IDistributedSynchronizationHandle> AcquireLockAsync(string resource, CancellationToken cancellationToken)