diff --git a/Elsa.sln b/Elsa.sln index 24cb8186c..5508e42b2 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -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} diff --git a/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs b/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs deleted file mode 100644 index de9a16628..000000000 --- a/src/locking/Elsa.DistributedLocking.AzureBlob/AzureBlobLockProvider.cs +++ /dev/null @@ -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 _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 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 AcquireLockAsync(string name, Duration? timeout = default, CancellationToken cancellationToken = default) => CreateLockAsync(name, timeout, cancellationToken); - - private async Task 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"); - } - } - } -} \ No newline at end of file diff --git a/src/locking/Elsa.DistributedLocking.AzureBlob/DistributedLockingOptionsExtensions.cs b/src/locking/Elsa.DistributedLocking.AzureBlob/DistributedLockingOptionsExtensions.cs new file mode 100644 index 000000000..a1684e9d6 --- /dev/null +++ b/src/locking/Elsa.DistributedLocking.AzureBlob/DistributedLockingOptionsExtensions.cs @@ -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 CreateAzureDistributedLockFactory(IServiceProvider services, Uri blobContainerUrl) + { + var container = new BlobContainerClient(blobContainerUrl); + return name => new AzureBlobLeaseDistributedLock(container, name); + } + } +} \ No newline at end of file diff --git a/src/locking/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj b/src/locking/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj index 5bcd0dee9..3423ff4c0 100644 --- a/src/locking/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj +++ b/src/locking/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj @@ -14,6 +14,7 @@ + diff --git a/src/locking/Elsa.DistributedLocking.AzureBlob/ElsaOptionsExtensions.cs b/src/locking/Elsa.DistributedLocking.AzureBlob/ElsaOptionsExtensions.cs deleted file mode 100644 index 43bd615a5..000000000 --- a/src/locking/Elsa.DistributedLocking.AzureBlob/ElsaOptionsExtensions.cs +++ /dev/null @@ -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>())); - return options; - } - } -} \ No newline at end of file diff --git a/src/locking/Elsa.DistributedLocking.AzureBlob/LockedBlob.cs b/src/locking/Elsa.DistributedLocking.AzureBlob/LockedBlob.cs deleted file mode 100644 index fa1552d42..000000000 --- a/src/locking/Elsa.DistributedLocking.AzureBlob/LockedBlob.cs +++ /dev/null @@ -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!; - } -} \ No newline at end of file