From 4511237a357dff3b9451364e33bd77d69e86c4fe Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 6 Apr 2021 21:45:50 +0200 Subject: [PATCH] Fix deadlocks --- .../ExecuteWorkflowInstanceRequestConsumer.cs | 2 +- .../Handlers/ExecuteWorkflowDefinition.cs | 27 ++++- .../Handlers/ExecuteWorkflowInstance.cs | 2 +- .../Dispatch/Handlers/TriggerWorkflows.cs | 68 ++++++++---- .../ElsaServiceCollectionExtensions.cs | 1 - .../ActivityPropertyUIHintResolver.cs | 2 +- .../Elsa.DistributedLocking.SqlServer.csproj | 1 + .../SqlLockProvider.cs | 105 +++++------------- 8 files changed, 103 insertions(+), 105 deletions(-) diff --git a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowInstanceRequestConsumer.cs b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowInstanceRequestConsumer.cs index 25a871617..13d850bb6 100644 --- a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowInstanceRequestConsumer.cs +++ b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteWorkflowInstanceRequestConsumer.cs @@ -32,7 +32,7 @@ namespace Elsa.Dispatch.Consumers public async Task Handle(ExecuteWorkflowInstanceRequest message) { var workflowInstanceId = message.WorkflowInstanceId; - var lockKey = $"execute-workflow-instance:{workflowInstanceId}"; + var lockKey = $"execute-workflow-instance-consumer:{workflowInstanceId}"; _logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}", workflowInstanceId); _stopwatch.Restart(); diff --git a/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs b/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs index b46d7dde8..ca42ad5b6 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowDefinition.cs @@ -1,5 +1,8 @@ +using System; using System.Threading; using System.Threading.Tasks; +using Elsa.DistributedLocking; +using Elsa.Exceptions; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; @@ -16,17 +19,23 @@ namespace Elsa.Dispatch.Handlers private readonly IStartsWorkflow _startsWorkflow; private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly IDistributedLockProvider _distributedLockProvider; + private readonly ElsaOptions _elsaOptions; private readonly ILogger _logger; public ExecuteWorkflowDefinition( IStartsWorkflow startsWorkflow, IWorkflowRegistry workflowRegistry, IWorkflowInstanceStore workflowInstanceStore, + IDistributedLockProvider distributedLockProvider, + ElsaOptions elsaOptions, ILogger logger) { _startsWorkflow = startsWorkflow; _workflowRegistry = workflowRegistry; _workflowInstanceStore = workflowInstanceStore; + _distributedLockProvider = distributedLockProvider; + _elsaOptions = elsaOptions; _logger = logger; } @@ -38,9 +47,21 @@ namespace Elsa.Dispatch.Handlers if (!ValidatePreconditions(workflowDefinitionId, workflowBlueprint)) return Unit.Value; - - if (!workflowBlueprint!.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false) - await _startsWorkflow.StartWorkflowAsync(workflowBlueprint, request.ActivityId, request.Input, request.CorrelationId, request.ContextId, cancellationToken); + + var lockKey = $"execute-workflow-definition:tenant:{tenantId}:workflow-definition:{workflowDefinitionId}"; + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken)) + throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}"); + + try + { + if (!workflowBlueprint!.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false) + await _startsWorkflow.StartWorkflowAsync(workflowBlueprint, request.ActivityId, request.Input, request.CorrelationId, request.ContextId, cancellationToken); + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken); + } return Unit.Value; } diff --git a/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowInstance.cs b/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowInstance.cs index 53204f905..a7f89f6db 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowInstance.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/ExecuteWorkflowInstance.cs @@ -25,7 +25,7 @@ namespace Elsa.Dispatch.Handlers public async Task Handle(ExecuteWorkflowInstanceRequest request, CancellationToken cancellationToken) { var workflowInstanceId = request.WorkflowInstanceId; - var workflowInstance = await _workflowInstanceStore.FindByIdAsync(request.WorkflowInstanceId, cancellationToken: cancellationToken); + var workflowInstance = await _workflowInstanceStore.FindByIdAsync(request.WorkflowInstanceId, cancellationToken); if (!ValidatePreconditions(workflowInstanceId, workflowInstance, request.ActivityId)) return Unit.Value; diff --git a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs index e81795f45..d5936437f 100644 --- a/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs +++ b/src/core/Elsa.Core/Dispatch/Handlers/TriggerWorkflows.cs @@ -1,7 +1,11 @@ -using System.Collections.Generic; +using System; +using System.Collections.Generic; +using System.Linq; using System.Threading; using System.Threading.Tasks; using Elsa.Bookmarks; +using Elsa.DistributedLocking; +using Elsa.Exceptions; using Elsa.Models; using Elsa.Persistence; using Elsa.Persistence.Specifications; @@ -12,61 +16,87 @@ using Open.Linq.AsyncExtensions; namespace Elsa.Dispatch.Handlers { - public class TriggerWorkflows : IRequestHandler + public class TriggerWorkflows : IRequestHandler { private readonly IWorkflowInstanceStore _workflowInstanceStore; private readonly IBookmarkFinder _bookmarkFinder; private readonly ITriggerFinder _triggerFinder; + private readonly IDistributedLockProvider _distributedLockProvider; private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher; private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher; + private readonly IMediator _mediator; + private readonly ElsaOptions _elsaOptions; private readonly ILogger _logger; public TriggerWorkflows( IWorkflowInstanceStore workflowInstanceStore, IBookmarkFinder bookmarkFinder, ITriggerFinder triggerFinder, + IDistributedLockProvider distributedLockProvider, IWorkflowDefinitionDispatcher workflowDefinitionDispatcher, IWorkflowInstanceDispatcher workflowInstanceDispatcher, + ElsaOptions elsaOptions, ILogger logger) { _workflowInstanceStore = workflowInstanceStore; _bookmarkFinder = bookmarkFinder; _triggerFinder = triggerFinder; + _distributedLockProvider = distributedLockProvider; _workflowDefinitionDispatcher = workflowDefinitionDispatcher; _workflowInstanceDispatcher = workflowInstanceDispatcher; + _elsaOptions = elsaOptions; _logger = logger; } - - public async Task Handle(TriggerWorkflowsRequest request, CancellationToken cancellationToken) + + public async Task Handle(TriggerWorkflowsRequest request, CancellationToken cancellationToken) { var correlationId = request.CorrelationId; - var correlatedWorkflowInstanceCount = await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId), cancellationToken); - if (correlatedWorkflowInstanceCount > 0) + if (!string.IsNullOrWhiteSpace(correlationId)) { - _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId); - var existingWorkflows = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken).ToList(); - await ResumeWorkflowsAsync(existingWorkflows, request.Input, cancellationToken); + var lockKey = $"trigger-workflows:correlation:{correlationId}"; + + if (!await _distributedLockProvider.AcquireLockAsync(lockKey, _elsaOptions.DistributedLockTimeout, cancellationToken)) + throw new LockAcquisitionException($"Failed to acquire a lock on {lockKey}"); + + int correlatedWorkflowInstanceCount; + try + { + correlatedWorkflowInstanceCount = !string.IsNullOrWhiteSpace(correlationId) ? await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId), cancellationToken) : 0; + } + finally + { + await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken); + } + + if (correlatedWorkflowInstanceCount > 0) + { + _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId); + var bookmarkResults = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken).ToList(); + await ResumeWorkflowsAsync(bookmarkResults, request.Input, cancellationToken); + return correlatedWorkflowInstanceCount; + } } - else - { - _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); - await StartWorkflowsAsync(request, cancellationToken); - } - - return Unit.Value; + + return await StartWorkflowsAsync(request, cancellationToken); } - private async Task StartWorkflowsAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken) + private async Task StartWorkflowsAsync(TriggerWorkflowsRequest request, CancellationToken cancellationToken) { + _logger.LogDebug("Triggering workflows using {ActivityType}", request.ActivityType); + var filter = request.Trigger; - var triggers = await _triggerFinder.FindTriggersAsync(request.ActivityType, filter, request.TenantId, cancellationToken); + var triggers = (await _triggerFinder.FindTriggersAsync(request.ActivityType, filter, request.TenantId, cancellationToken)).ToList(); foreach (var trigger in triggers) { var workflowBlueprint = trigger.WorkflowBlueprint; - await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowBlueprint.Id, trigger.ActivityId, request.Input, request.CorrelationId, request.ContextId, workflowBlueprint.TenantId), cancellationToken); + + await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowBlueprint.Id, trigger.ActivityId, request.Input, request.CorrelationId, request.ContextId, workflowBlueprint.TenantId), + cancellationToken); } + + return triggers.Count; } private async Task ResumeWorkflowsAsync(IEnumerable results, object? input, CancellationToken cancellationToken) diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 4c4e67546..4dc65b09e 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -106,7 +106,6 @@ namespace Microsoft.Extensions.DependencyInjection .AddScoped() .AddScoped() .AddScoped() - .AddScoped() .AddScoped(sp => sp.GetRequiredService()) .AddScoped(sp => sp.GetRequiredService()) .AddScoped(sp => sp.GetRequiredService()) diff --git a/src/core/Elsa.Core/Metadata/ActivityPropertyUIHintResolver.cs b/src/core/Elsa.Core/Metadata/ActivityPropertyUIHintResolver.cs index 67b018bb8..705354948 100644 --- a/src/core/Elsa.Core/Metadata/ActivityPropertyUIHintResolver.cs +++ b/src/core/Elsa.Core/Metadata/ActivityPropertyUIHintResolver.cs @@ -25,7 +25,7 @@ namespace Elsa.Metadata if (typeof(IEnumerable).IsAssignableFrom(type)) return ActivityPropertyUIHints.Dropdown; - + if (type.IsEnum || type.IsNullableType() && type.GetTypeOfNullable().IsEnum) return ActivityPropertyUIHints.Dropdown; diff --git a/src/locking/Elsa.DistributedLocking.SqlServer/Elsa.DistributedLocking.SqlServer.csproj b/src/locking/Elsa.DistributedLocking.SqlServer/Elsa.DistributedLocking.SqlServer.csproj index 0a2590fe1..e161594b0 100644 --- a/src/locking/Elsa.DistributedLocking.SqlServer/Elsa.DistributedLocking.SqlServer.csproj +++ b/src/locking/Elsa.DistributedLocking.SqlServer/Elsa.DistributedLocking.SqlServer.csproj @@ -19,6 +19,7 @@ + diff --git a/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs b/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs index 3b69cac5e..c5eb7f338 100644 --- a/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs +++ b/src/locking/Elsa.DistributedLocking.SqlServer/SqlLockProvider.cs @@ -6,6 +6,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.DistributedLock; using Elsa.DistributedLocking; +using Medallion.Threading.SqlServer; using Microsoft.Extensions.Logging; using NodaTime; @@ -16,9 +17,8 @@ namespace Elsa public class SqlLockProvider : IDistributedLockProvider { private readonly ILogger _logger; - private const string Prefix = "elsa"; private readonly string _connectionString; - private readonly ConcurrentDictionary _locks = new(); + private readonly ConcurrentDictionary _locks = new(); public SqlLockProvider(string connectionString, ILogger logger) { @@ -34,94 +34,41 @@ namespace Elsa public async Task AcquireLockAsync(string name, Duration? timeout = default, CancellationToken cancellationToken = default) { - var connection = new SqlConnection(_connectionString); - await connection.OpenAsync(cancellationToken); - var timeoutMils = timeout?.TotalMilliseconds ?? 0d; + var distributedLock = new SqlDistributedLock(name, _connectionString); + var timeoutTimeSpan = timeout?.ToTimeSpan() ?? TimeSpan.Zero; - try - { - var command = connection.CreateCommand(); - command.CommandText = "sp_getapplock"; - command.CommandType = CommandType.StoredProcedure; - command.Parameters.AddWithValue("@Resource", $"{Prefix}:{name}"); - command.Parameters.AddWithValue("@LockOwner", $"Session"); - command.Parameters.AddWithValue("@LockMode", $"Exclusive"); - command.Parameters.AddWithValue("@LockTimeout", timeoutMils); - command.CommandTimeout = Math.Max((int?)timeout?.TotalSeconds ?? 30, 30); + _logger.LogDebug("Acquiring a lock on {LockName}", name); - var returnParameter = command.Parameters.Add("RetVal", SqlDbType.Int); - returnParameter.Direction = ParameterDirection.ReturnValue; + if (_locks.ContainsKey(name)) + _logger.LogDebug("Waiting for existing lock {LockName} to be released", name); - await command.ExecuteNonQueryAsync(cancellationToken); - var result = Convert.ToInt32(returnParameter.Value); - - switch (result) - { - case -1: - _logger.LogDebug("The lock request timed out for {LockName}", name); - break; - - case -2: - _logger.LogDebug("The lock request was canceled for {LockName}", name); - break; - - case -3: - _logger.LogDebug("The lock request was chosen as a deadlock victim for {LockName}", name); - break; - - case -999: - _logger.LogError("Lock provider error for {LockName}", name); - break; - } - - if (result >= 0) - { - _locks[name] = connection; - return true; - } - - connection.Close(); + await using var handle = await distributedLock.AcquireAsync(timeoutTimeSpan, cancellationToken); + + if (handle == null!) return false; - } - catch (Exception) - { - connection.Close(); - throw; - } + + _locks[name] = handle; + _logger.LogDebug("Lock acquired on {LockName}", name); + + return true; } public async Task ReleaseLockAsync(string name, CancellationToken cancellationToken) { if (!_locks.ContainsKey(name)) + { + _logger.LogDebug("Failed to release lock that wasn't captured"); + return; + } + + var handle = _locks[name]; + + if (handle == null) return; - var connection = _locks[name]; - - if (connection == null) - return; - - try - { - var command = connection.CreateCommand(); - command.CommandText = "sp_releaseapplock"; - command.CommandType = CommandType.StoredProcedure; - command.Parameters.AddWithValue("@Resource", $"{Prefix}:{name}"); - command.Parameters.AddWithValue("@LockOwner", $"Session"); - - var returnParameter = command.Parameters.Add("RetVal", SqlDbType.Int); - returnParameter.Direction = ParameterDirection.ReturnValue; - - await command.ExecuteNonQueryAsync(cancellationToken); - var result = Convert.ToInt32(returnParameter.Value); - - if (result < 0) - _logger.LogError("Unable to release lock for {LockName}", name); - } - finally - { - connection.Close(); - _locks.TryRemove(name, out _); - } + await handle.DisposeAsync(); + _locks.TryRemove(name, out _); + _logger.LogDebug("Released lock on {LockName}", name); } } } \ No newline at end of file