Update extensions for Azure distributed lock

This commit is contained in:
Sipke Schoorstra 2021-04-07 12:51:01 +02:00
parent df6ad74eec
commit 0a02dc5cc2
6 changed files with 30 additions and 232 deletions

View file

@ -249,6 +249,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.DistributedLocking.Sql
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.DistributedLocking.Redis", "src\locking\Elsa.DistributedLocking.Redis\Elsa.DistributedLocking.Redis.csproj", "{7E5514AC-3633-4237-9490-0992417F2273}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.DistributedLocking.AzureBlob", "src\locking\Elsa.DistributedLocking.AzureBlob\Elsa.DistributedLocking.AzureBlob.csproj", "{495BE954-E6EA-41A9-8054-A5CA5DD0924E}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -603,6 +605,10 @@ Global
{7E5514AC-3633-4237-9490-0992417F2273}.Debug|Any CPU.Build.0 = Debug|Any CPU
{7E5514AC-3633-4237-9490-0992417F2273}.Release|Any CPU.ActiveCfg = Release|Any CPU
{7E5514AC-3633-4237-9490-0992417F2273}.Release|Any CPU.Build.0 = Release|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Debug|Any CPU.Build.0 = Debug|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Release|Any CPU.ActiveCfg = Release|Any CPU
{495BE954-E6EA-41A9-8054-A5CA5DD0924E}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -722,6 +728,7 @@ Global
{DBBD242E-4437-4BDB-919F-A70839BE75FA} = {DA71CDAA-8DD3-4D5F-9FBD-8E4B37A2D925}
{B09A6E42-EF42-4693-BA0E-5EE1D7F81FFD} = {DBBD242E-4437-4BDB-919F-A70839BE75FA}
{7E5514AC-3633-4237-9490-0992417F2273} = {DBBD242E-4437-4BDB-919F-A70839BE75FA}
{495BE954-E6EA-41A9-8054-A5CA5DD0924E} = {DBBD242E-4437-4BDB-919F-A70839BE75FA}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158}

View file

@ -1,201 +0,0 @@
using System;
using System.Collections.Generic;
using System.IO;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Elsa.DistributedLock;
using Elsa.DistributedLocking;
using Microsoft.Azure.Storage;
using Microsoft.Azure.Storage.Blob;
using Microsoft.Azure.Storage.RetryPolicies;
using Microsoft.Extensions.Logging;
using NodaTime;
namespace Elsa
{
// See also:
// The lock duration can be 15 to 60 seconds, or can be infinite, Reference: https://docs.microsoft.com/en-us/rest/api/storageservices/lease-blob
public class AzureBlobLockProvider : IDistributedLockProvider
{
private const string Prefix = "elsa";
private const string ContainerName = "elsa-lock-container";
private const int MaxLeaseTime = 60;
private const int MinLeaseTime = 15;
private static readonly object SyncRoot = new();
private readonly AutoResetEvent _leaseSemaphore = new(true);
private readonly List<LockedBlob> _lockedBlobs = new();
private readonly string _connectionString;
private readonly TimeSpan _leaseTime;
private readonly ILogger _logger;
private readonly TimeSpan _renewInterval;
private CloudBlobContainer? _cloudBlobContainer;
private Timer _renewTimer = default!;
public AzureBlobLockProvider(
string connectionString,
TimeSpan leaseTime,
TimeSpan renewInterval,
ILogger<AzureBlobLockProvider> logger)
{
_logger = logger;
_connectionString = connectionString;
if (leaseTime >= TimeSpan.FromSeconds(MaxLeaseTime) || leaseTime <= TimeSpan.FromSeconds(MinLeaseTime))
{
_logger.LogInformation("Lease time must be between 15 Seconds and 60 seconds, Found {LeaseTime} seconds. Setting default value of 60 seconds", leaseTime.TotalSeconds);
_leaseTime = TimeSpan.FromSeconds(MaxLeaseTime);
}
else
{
_leaseTime = leaseTime;
}
if (renewInterval > leaseTime)
{
_logger.LogError("Renew Interval can not be greater than LeaseTime {LeaseTime}", leaseTime.TotalSeconds);
throw new InvalidDataException($"Renew Interval can not be greater than LeaseTime {leaseTime.TotalSeconds}.");
}
_renewInterval = renewInterval;
}
public Task<bool> AcquireLockAsync(string name, Duration? timeout = default, CancellationToken cancellationToken = default) => CreateLockAsync(name, timeout, cancellationToken);
private async Task<bool> CreateLockAsync(string name, Duration? duration, CancellationToken cancellationToken)
{
// TODO: Figure out how to wait for the amount of time as specified by `duration` before giving up on acquiring a lease.
var resourceName = $"{Prefix}:{name}";
var blob = CloudBlobContainer.GetBlockBlobReference(resourceName);
if (!await blob.ExistsAsync(cancellationToken))
await blob.UploadTextAsync(string.Empty, cancellationToken);
_renewTimer = new Timer(RenewLeases, null, _renewInterval, _renewInterval);
_logger.LogInformation("Lock provider will try to acquire lock for {resourceName}", resourceName);
if (_leaseSemaphore.WaitOne())
{
try
{
var leaseId = await blob.AcquireLeaseAsync(_leaseTime, null, cancellationToken);
_lockedBlobs.Add(new LockedBlob { Blob = blob, LeaseId = leaseId, Identifier = resourceName });
_logger.LogInformation("Lock provider acquired lock for {resourceName}", resourceName);
return true;
}
catch (Exception ex)
{
_logger.LogWarning("Failed to acquire lock for {resourceName}. Reason > {ex}", resourceName, ex);
return false;
}
finally
{
_leaseSemaphore.Set();
}
}
return false;
}
public Task ReleaseLockAsync(string name, CancellationToken cancellationToken = default) => ReleaseLeaseAsync(name, cancellationToken);
private async Task ReleaseLeaseAsync(string name, CancellationToken cancellationToken = default)
{
var resourceName = $"{Prefix}:{name}";
_logger.LogInformation("Lock provider will try to release lock for {resourceName}", resourceName);
if (_leaseSemaphore.WaitOne())
{
try
{
var lockedBlob = _lockedBlobs.FirstOrDefault(x => x.Identifier == resourceName);
if (lockedBlob != null)
{
try
{
await lockedBlob.Blob
.ReleaseLeaseAsync(AccessCondition.GenerateLeaseCondition(lockedBlob.LeaseId), cancellationToken)
.ConfigureAwait(false);
_logger.LogInformation("Lock provider released lock for {resourceName}", resourceName);
}
catch (Exception ex)
{
_logger.LogWarning("Failed to release lock for {resourceName}. Reason > {ex}", resourceName, ex);
}
_lockedBlobs.Remove(lockedBlob);
}
}
finally
{
_leaseSemaphore.Set();
await _renewTimer.DisposeAsync();
}
}
}
private CloudBlobContainer CloudBlobContainer
{
get
{
if (_cloudBlobContainer != null)
return _cloudBlobContainer;
lock (SyncRoot)
{
if (_cloudBlobContainer != null)
return _cloudBlobContainer;
var blobClient = CloudStorageAccount.Parse(_connectionString).CreateCloudBlobClient();
blobClient.DefaultRequestOptions.RetryPolicy = new ExponentialRetry(TimeSpan.FromSeconds(2.0), 3);
_cloudBlobContainer = blobClient.GetContainerReference(ContainerName);
if (!_cloudBlobContainer.Exists())
_cloudBlobContainer.CreateIfNotExists(BlobContainerPublicAccessType.Off);
}
return _cloudBlobContainer;
}
}
private async void RenewLeases(object state)
{
_logger.LogDebug("Renew active leases");
if (_leaseSemaphore.WaitOne())
{
try
{
foreach (var lockedBlobs in _lockedBlobs)
await RenewLock(lockedBlobs).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Error while renewing leases");
}
finally
{
_leaseSemaphore.Set();
}
}
}
private async Task RenewLock(LockedBlob lockedBlob)
{
try
{
await lockedBlob.Blob.RenewLeaseAsync(AccessCondition.GenerateLeaseCondition(lockedBlob.LeaseId));
_logger.LogInformation("Renewed active leases for {lockedBlob.Identifier}", lockedBlob.Identifier);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Error while renewing lock");
}
}
}
}

View file

@ -0,0 +1,22 @@
using System;
using Azure.Storage.Blobs;
using Medallion.Threading;
using Medallion.Threading.Azure;
namespace Elsa
{
public static class DistributedLockingOptionsExtensions
{
public static DistributedLockingOptions UseSqlServerLockProvider(this DistributedLockingOptions options, Uri blobContainerUrl)
{
options.DistributedLockProviderFactory = sp => CreateAzureDistributedLockFactory(sp, blobContainerUrl);
return options;
}
private static Func<string, IDistributedLock> CreateAzureDistributedLockFactory(IServiceProvider services, Uri blobContainerUrl)
{
var container = new BlobContainerClient(blobContainerUrl);
return name => new AzureBlobLeaseDistributedLock(container, name);
}
}
}

View file

@ -14,6 +14,7 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="DistributedLock.Azure" Version="1.0.0" />
<PackageReference Include="Microsoft.Azure.Storage.Blob" Version="11.2.2" />
</ItemGroup>

View file

@ -1,20 +0,0 @@
using System;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
namespace Elsa
{
public static class ElsaOptionsExtensions
{
public static ElsaOptions UseAzureBlobLockProvider(this ElsaOptions options, string connectionString, TimeSpan? leaseTime = null, TimeSpan? renewInterval = null)
{
options
.UseDistributedLockProvider(sp =>
new AzureBlobLockProvider(connectionString,
leaseTime ?? TimeSpan.FromSeconds(60),
renewInterval ?? TimeSpan.FromSeconds(45),
sp.GetRequiredService<ILogger<AzureBlobLockProvider>>()));
return options;
}
}
}

View file

@ -1,11 +0,0 @@
using Microsoft.Azure.Storage.Blob;
namespace Elsa
{
internal class LockedBlob
{
public string Identifier { get; set; } = default!;
public string LeaseId { get; set; } = default!;
public CloudBlockBlob Blob { get; set; } = default!;
}
}