Fix deadlocks

This commit is contained in:
Sipke Schoorstra 2021-04-06 21:45:50 +02:00
parent 023c971834
commit 4511237a35
8 changed files with 103 additions and 105 deletions

View file

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

View file

@ -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<ExecuteWorkflowDefinition> 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;
}

View file

@ -25,7 +25,7 @@ namespace Elsa.Dispatch.Handlers
public async Task<Unit> 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;

View file

@ -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<TriggerWorkflowsRequest>
public class TriggerWorkflows : IRequestHandler<TriggerWorkflowsRequest, int>
{
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<TriggerWorkflows> _logger;
public TriggerWorkflows(
IWorkflowInstanceStore workflowInstanceStore,
IBookmarkFinder bookmarkFinder,
ITriggerFinder triggerFinder,
IDistributedLockProvider distributedLockProvider,
IWorkflowDefinitionDispatcher workflowDefinitionDispatcher,
IWorkflowInstanceDispatcher workflowInstanceDispatcher,
ElsaOptions elsaOptions,
ILogger<TriggerWorkflows> logger)
{
_workflowInstanceStore = workflowInstanceStore;
_bookmarkFinder = bookmarkFinder;
_triggerFinder = triggerFinder;
_distributedLockProvider = distributedLockProvider;
_workflowDefinitionDispatcher = workflowDefinitionDispatcher;
_workflowInstanceDispatcher = workflowInstanceDispatcher;
_elsaOptions = elsaOptions;
_logger = logger;
}
public async Task<Unit> Handle(TriggerWorkflowsRequest request, CancellationToken cancellationToken)
public async Task<int> Handle(TriggerWorkflowsRequest request, CancellationToken cancellationToken)
{
var correlationId = request.CorrelationId;
var correlatedWorkflowInstanceCount = await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(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<WorkflowInstance>(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<int> 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<BookmarkFinderResult> results, object? input, CancellationToken cancellationToken)

View file

@ -106,7 +106,6 @@ namespace Microsoft.Extensions.DependencyInjection
.AddScoped<IWorkflowRunner, WorkflowRunner>()
.AddScoped<WorkflowStarter>()
.AddScoped<WorkflowResumer>()
.AddScoped<ITriggersWorkflows, TriggersWorkflows>()
.AddScoped<IStartsWorkflow>(sp => sp.GetRequiredService<WorkflowStarter>())
.AddScoped<IStartsWorkflows>(sp => sp.GetRequiredService<WorkflowStarter>())
.AddScoped<IFindsAndStartsWorkflows>(sp => sp.GetRequiredService<WorkflowStarter>())

View file

@ -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;

View file

@ -19,6 +19,7 @@
</ItemGroup>
<ItemGroup>
<PackageReference Include="DistributedLock.SqlServer" Version="1.0.0" />
<PackageReference Include="System.Data.SqlClient" Version="4.8.2" />
</ItemGroup>

View file

@ -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<string, SqlConnection> _locks = new();
private readonly ConcurrentDictionary<string, SqlDistributedLockHandle> _locks = new();
public SqlLockProvider(string connectionString, ILogger<SqlLockProvider> logger)
{
@ -34,94 +34,41 @@ namespace Elsa
public async Task<bool> 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);
}
}
}