diff --git a/framework/Volo.Abp.sln b/framework/Volo.Abp.sln
index 991d1d08d5..f3e495c814 100644
--- a/framework/Volo.Abp.sln
+++ b/framework/Volo.Abp.sln
@@ -391,14 +391,16 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.IdentityModel.Test
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.Threading.Tests", "test\Volo.Abp.Threading.Tests\Volo.Abp.Threading.Tests.csproj", "{7B2FCAD6-86E6-49C8-ADBE-A61B4F4B101B}"
EndProject
-Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.EventBus.Boxes", "src\Volo.Abp.EventBus.Boxes\Volo.Abp.EventBus.Boxes.csproj", "{6E289F31-7924-418B-9DAC-62A7CFADF916}"
-EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.DistributedLocking", "src\Volo.Abp.DistributedLocking\Volo.Abp.DistributedLocking.csproj", "{9A7EEA08-15BE-476D-8168-53039867038E}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.Auditing.Contracts", "src\Volo.Abp.Auditing.Contracts\Volo.Abp.Auditing.Contracts.csproj", "{508B6355-AD28-4E60-8549-266D21DBF2CF}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.Http.Client.Web", "src\Volo.Abp.Http.Client.Web\Volo.Abp.Http.Client.Web.csproj", "{F7407459-8AFA-45E4-83E9-9BB01412CC08}"
EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.DistributedLocking.Abstractions", "src\Volo.Abp.DistributedLocking.Abstractions\Volo.Abp.DistributedLocking.Abstractions.csproj", "{CA805B77-D50C-431F-B3CB-1111C9C6E807}"
+EndProject
+Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.DistributedLocking.Abstractions.Tests", "test\Volo.Abp.DistributedLocking.Abstractions.Tests\Volo.Abp.DistributedLocking.Abstractions.Tests.csproj", "{C4F54FB5-C828-414D-BA03-E8E7A10C784D}"
+EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.BackgroundWorkers.Hangfire", "src\Volo.Abp.BackgroundWorkers.Hangfire\Volo.Abp.BackgroundWorkers.Hangfire.csproj", "{E5FCE710-C5A3-4F94-B9C9-BD1E99252BFB}"
EndProject
Global
@@ -1175,10 +1177,6 @@ Global
{7B2FCAD6-86E6-49C8-ADBE-A61B4F4B101B}.Debug|Any CPU.Build.0 = Debug|Any CPU
{7B2FCAD6-86E6-49C8-ADBE-A61B4F4B101B}.Release|Any CPU.ActiveCfg = Release|Any CPU
{7B2FCAD6-86E6-49C8-ADBE-A61B4F4B101B}.Release|Any CPU.Build.0 = Release|Any CPU
- {6E289F31-7924-418B-9DAC-62A7CFADF916}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {6E289F31-7924-418B-9DAC-62A7CFADF916}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {6E289F31-7924-418B-9DAC-62A7CFADF916}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {6E289F31-7924-418B-9DAC-62A7CFADF916}.Release|Any CPU.Build.0 = Release|Any CPU
{9A7EEA08-15BE-476D-8168-53039867038E}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{9A7EEA08-15BE-476D-8168-53039867038E}.Debug|Any CPU.Build.0 = Debug|Any CPU
{9A7EEA08-15BE-476D-8168-53039867038E}.Release|Any CPU.ActiveCfg = Release|Any CPU
@@ -1191,6 +1189,14 @@ Global
{F7407459-8AFA-45E4-83E9-9BB01412CC08}.Debug|Any CPU.Build.0 = Debug|Any CPU
{F7407459-8AFA-45E4-83E9-9BB01412CC08}.Release|Any CPU.ActiveCfg = Release|Any CPU
{F7407459-8AFA-45E4-83E9-9BB01412CC08}.Release|Any CPU.Build.0 = Release|Any CPU
+ {CA805B77-D50C-431F-B3CB-1111C9C6E807}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {CA805B77-D50C-431F-B3CB-1111C9C6E807}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {CA805B77-D50C-431F-B3CB-1111C9C6E807}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {CA805B77-D50C-431F-B3CB-1111C9C6E807}.Release|Any CPU.Build.0 = Release|Any CPU
+ {C4F54FB5-C828-414D-BA03-E8E7A10C784D}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
+ {C4F54FB5-C828-414D-BA03-E8E7A10C784D}.Debug|Any CPU.Build.0 = Debug|Any CPU
+ {C4F54FB5-C828-414D-BA03-E8E7A10C784D}.Release|Any CPU.ActiveCfg = Release|Any CPU
+ {C4F54FB5-C828-414D-BA03-E8E7A10C784D}.Release|Any CPU.Build.0 = Release|Any CPU
{E5FCE710-C5A3-4F94-B9C9-BD1E99252BFB}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{E5FCE710-C5A3-4F94-B9C9-BD1E99252BFB}.Debug|Any CPU.Build.0 = Debug|Any CPU
{E5FCE710-C5A3-4F94-B9C9-BD1E99252BFB}.Release|Any CPU.ActiveCfg = Release|Any CPU
@@ -1392,10 +1398,11 @@ Global
{90B1866A-EF99-40B9-970E-B898E5AA523F} = {447C8A77-E5F0-4538-8687-7383196D04EA}
{40C6740E-BFCA-4D37-8344-3D84E2044BB2} = {447C8A77-E5F0-4538-8687-7383196D04EA}
{7B2FCAD6-86E6-49C8-ADBE-A61B4F4B101B} = {447C8A77-E5F0-4538-8687-7383196D04EA}
- {6E289F31-7924-418B-9DAC-62A7CFADF916} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6}
{9A7EEA08-15BE-476D-8168-53039867038E} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6}
{508B6355-AD28-4E60-8549-266D21DBF2CF} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6}
{F7407459-8AFA-45E4-83E9-9BB01412CC08} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6}
+ {CA805B77-D50C-431F-B3CB-1111C9C6E807} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6}
+ {C4F54FB5-C828-414D-BA03-E8E7A10C784D} = {447C8A77-E5F0-4538-8687-7383196D04EA}
{E5FCE710-C5A3-4F94-B9C9-BD1E99252BFB} = {5DF0E140-0513-4D0D-BE2E-3D4D85CD70E6}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo.Abp.BackgroundJobs.csproj b/framework/src/Volo.Abp.BackgroundJobs/Volo.Abp.BackgroundJobs.csproj
index 858038e5df..7d4f2d9fa2 100644
--- a/framework/src/Volo.Abp.BackgroundJobs/Volo.Abp.BackgroundJobs.csproj
+++ b/framework/src/Volo.Abp.BackgroundJobs/Volo.Abp.BackgroundJobs.csproj
@@ -17,7 +17,7 @@
-
+
diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs
index dccf2006dc..f5be42a2d6 100644
--- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs
+++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/AbpBackgroundJobsModule.cs
@@ -1,6 +1,7 @@
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Volo.Abp.BackgroundWorkers;
+using Volo.Abp.DistributedLocking;
using Volo.Abp.Guids;
using Volo.Abp.Modularity;
using Volo.Abp.Timing;
@@ -11,7 +12,8 @@ namespace Volo.Abp.BackgroundJobs
typeof(AbpBackgroundJobsAbstractionsModule),
typeof(AbpBackgroundWorkersModule),
typeof(AbpTimingModule),
- typeof(AbpGuidsModule)
+ typeof(AbpGuidsModule),
+ typeof(AbpDistributedLockingAbstractionsModule)
)]
public class AbpBackgroundJobsModule : AbpModule
{
diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs
index 8eec35d816..2e7714b113 100644
--- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs
+++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs
@@ -5,6 +5,7 @@ using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using Volo.Abp.BackgroundWorkers;
+using Volo.Abp.DistributedLocking;
using Volo.Abp.Threading;
using Volo.Abp.Timing;
@@ -12,19 +13,25 @@ namespace Volo.Abp.BackgroundJobs
{
public class BackgroundJobWorker : AsyncPeriodicBackgroundWorkerBase, IBackgroundJobWorker
{
+ protected const string DistributedLockName = "AbpBackgroundJobWorker";
+
protected AbpBackgroundJobOptions JobOptions { get; }
protected AbpBackgroundJobWorkerOptions WorkerOptions { get; }
+ protected IAbpDistributedLock DistributedLock { get; }
+
public BackgroundJobWorker(
AbpAsyncTimer timer,
IOptions jobOptions,
IOptions workerOptions,
- IServiceScopeFactory serviceScopeFactory)
+ IServiceScopeFactory serviceScopeFactory,
+ IAbpDistributedLock distributedLock)
: base(
timer,
serviceScopeFactory)
{
+ DistributedLock = distributedLock;
WorkerOptions = workerOptions.Value;
JobOptions = jobOptions.Value;
Timer.Period = WorkerOptions.JobPollPeriod;
@@ -32,57 +39,74 @@ namespace Volo.Abp.BackgroundJobs
protected override async Task DoWorkAsync(PeriodicBackgroundWorkerContext workerContext)
{
- var store = workerContext.ServiceProvider.GetRequiredService();
-
- var waitingJobs = await store.GetWaitingJobsAsync(WorkerOptions.MaxJobFetchCount);
-
- if (!waitingJobs.Any())
+ await using (var handler = await DistributedLock.TryAcquireAsync(DistributedLockName, cancellationToken: StoppingToken))
{
- return;
- }
-
- var jobExecuter = workerContext.ServiceProvider.GetRequiredService();
- var clock = workerContext.ServiceProvider.GetRequiredService();
- var serializer = workerContext.ServiceProvider.GetRequiredService();
-
- foreach (var jobInfo in waitingJobs)
- {
- jobInfo.TryCount++;
- jobInfo.LastTryTime = clock.Now;
-
- try
+ if (handler != null)
{
- var jobConfiguration = JobOptions.GetJob(jobInfo.JobName);
- var jobArgs = serializer.Deserialize(jobInfo.JobArgs, jobConfiguration.ArgsType);
- var context = new JobExecutionContext(workerContext.ServiceProvider, jobConfiguration.JobType, jobArgs);
+ var store = workerContext.ServiceProvider.GetRequiredService();
- try
- {
- await jobExecuter.ExecuteAsync(context);
+ var waitingJobs = await store.GetWaitingJobsAsync(WorkerOptions.MaxJobFetchCount);
- await store.DeleteAsync(jobInfo.Id);
+ if (!waitingJobs.Any())
+ {
+ return;
}
- catch (BackgroundJobExecutionException)
+
+ var jobExecuter = workerContext.ServiceProvider.GetRequiredService();
+ var clock = workerContext.ServiceProvider.GetRequiredService();
+ var serializer = workerContext.ServiceProvider.GetRequiredService();
+
+ foreach (var jobInfo in waitingJobs)
{
- var nextTryTime = CalculateNextTryTime(jobInfo, clock);
+ jobInfo.TryCount++;
+ jobInfo.LastTryTime = clock.Now;
- if (nextTryTime.HasValue)
+ try
{
- jobInfo.NextTryTime = nextTryTime.Value;
+ var jobConfiguration = JobOptions.GetJob(jobInfo.JobName);
+ var jobArgs = serializer.Deserialize(jobInfo.JobArgs, jobConfiguration.ArgsType);
+ var context = new JobExecutionContext(
+ workerContext.ServiceProvider,
+ jobConfiguration.JobType,
+ jobArgs);
+
+ try
+ {
+ await jobExecuter.ExecuteAsync(context);
+
+ await store.DeleteAsync(jobInfo.Id);
+ }
+ catch (BackgroundJobExecutionException)
+ {
+ var nextTryTime = CalculateNextTryTime(jobInfo, clock);
+
+ if (nextTryTime.HasValue)
+ {
+ jobInfo.NextTryTime = nextTryTime.Value;
+ }
+ else
+ {
+ jobInfo.IsAbandoned = true;
+ }
+
+ await TryUpdateAsync(store, jobInfo);
+ }
}
- else
+ catch (Exception ex)
{
+ Logger.LogException(ex);
jobInfo.IsAbandoned = true;
+ await TryUpdateAsync(store, jobInfo);
}
-
- await TryUpdateAsync(store, jobInfo);
}
}
- catch (Exception ex)
+ else
{
- Logger.LogException(ex);
- jobInfo.IsAbandoned = true;
- await TryUpdateAsync(store, jobInfo);
+ try
+ {
+ await Task.Delay(WorkerOptions.JobPollPeriod * 12, StoppingToken);
+ }
+ catch (TaskCanceledException) { }
}
}
}
@@ -101,7 +125,8 @@ namespace Volo.Abp.BackgroundJobs
protected virtual DateTime? CalculateNextTryTime(BackgroundJobInfo jobInfo, IClock clock)
{
- var nextWaitDuration = WorkerOptions.DefaultFirstWaitDuration * (Math.Pow(WorkerOptions.DefaultWaitFactor, jobInfo.TryCount - 1));
+ var nextWaitDuration = WorkerOptions.DefaultFirstWaitDuration *
+ (Math.Pow(WorkerOptions.DefaultWaitFactor, jobInfo.TryCount - 1));
var nextTryDate = jobInfo.LastTryTime?.AddSeconds(nextWaitDuration) ??
clock.Now.AddSeconds(nextWaitDuration);
@@ -113,4 +138,4 @@ namespace Volo.Abp.BackgroundJobs
return nextTryDate;
}
}
-}
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/BackgroundWorkerBase.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/BackgroundWorkerBase.cs
index 85523d7ee9..c497555be9 100644
--- a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/BackgroundWorkerBase.cs
+++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/BackgroundWorkerBase.cs
@@ -21,6 +21,15 @@ namespace Volo.Abp.BackgroundWorkers
protected ILoggerFactory LoggerFactory => LazyServiceProvider.LazyGetRequiredService();
protected ILogger Logger => LazyServiceProvider.LazyGetService(provider => LoggerFactory?.CreateLogger(GetType().FullName) ?? NullLogger.Instance);
+
+ protected CancellationTokenSource StoppingTokenSource { get; }
+ protected CancellationToken StoppingToken { get; }
+
+ public BackgroundWorkerBase()
+ {
+ StoppingTokenSource = new CancellationTokenSource();
+ StoppingToken = StoppingTokenSource.Token;
+ }
public virtual Task StartAsync(CancellationToken cancellationToken = default)
{
@@ -31,6 +40,8 @@ namespace Volo.Abp.BackgroundWorkers
public virtual Task StopAsync(CancellationToken cancellationToken = default)
{
Logger.LogDebug("Stopped background worker: " + ToString());
+ StoppingTokenSource.Cancel();
+ StoppingTokenSource.Dispose();
return Task.CompletedTask;
}
diff --git a/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj b/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj
index 68c90489ff..87879334af 100644
--- a/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj
+++ b/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj
@@ -17,7 +17,7 @@
-
+
diff --git a/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj.DotSettings b/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj.DotSettings
deleted file mode 100644
index 58ad6c8854..0000000000
--- a/framework/src/Volo.Abp.Ddd.Domain/Volo.Abp.Ddd.Domain.csproj.DotSettings
+++ /dev/null
@@ -1,2 +0,0 @@
-
- CSharp71
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/AbpDddDomainModule.cs b/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/AbpDddDomainModule.cs
index 38e2d9fc84..523a6ee4f2 100644
--- a/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/AbpDddDomainModule.cs
+++ b/framework/src/Volo.Abp.Ddd.Domain/Volo/Abp/Domain/AbpDddDomainModule.cs
@@ -3,7 +3,6 @@ using Volo.Abp.Auditing;
using Volo.Abp.Data;
using Volo.Abp.Domain.Repositories;
using Volo.Abp.EventBus;
-using Volo.Abp.EventBus.Boxes;
using Volo.Abp.ExceptionHandling;
using Volo.Abp.Guids;
using Volo.Abp.Modularity;
@@ -19,7 +18,7 @@ namespace Volo.Abp.Domain
[DependsOn(
typeof(AbpAuditingModule),
typeof(AbpDataModule),
- typeof(AbpEventBusBoxesModule),
+ typeof(AbpEventBusModule),
typeof(AbpGuidsModule),
typeof(AbpMultiTenancyModule),
typeof(AbpThreadingModule),
diff --git a/framework/src/Volo.Abp.DistributedLocking.Abstractions/FodyWeavers.xml b/framework/src/Volo.Abp.DistributedLocking.Abstractions/FodyWeavers.xml
new file mode 100644
index 0000000000..be0de3a908
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/FodyWeavers.xml
@@ -0,0 +1,3 @@
+
+
+
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/FodyWeavers.xsd b/framework/src/Volo.Abp.DistributedLocking.Abstractions/FodyWeavers.xsd
similarity index 96%
rename from framework/src/Volo.Abp.EventBus.Boxes/FodyWeavers.xsd
rename to framework/src/Volo.Abp.DistributedLocking.Abstractions/FodyWeavers.xsd
index ffa6fc4b78..3f3946e282 100644
--- a/framework/src/Volo.Abp.EventBus.Boxes/FodyWeavers.xsd
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/FodyWeavers.xsd
@@ -1,4 +1,4 @@
-
+
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo.Abp.DistributedLocking.Abstractions.csproj
similarity index 62%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj
rename to framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo.Abp.DistributedLocking.Abstractions.csproj
index 6cbaf62aed..ab828ee9d3 100644
--- a/framework/src/Volo.Abp.EventBus.Boxes/Volo.Abp.EventBus.Boxes.csproj
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo.Abp.DistributedLocking.Abstractions.csproj
@@ -5,8 +5,8 @@
netstandard2.0
- Volo.Abp.EventBus.Boxes
- Volo.Abp.EventBus.Boxes
+ Volo.Abp.DistributedLocking.Abstractions
+ Volo.Abp.DistributedLocking.Abstractions
$(AssetTargetFallback);portable-net45+win8+wp8+wpa81;
false
false
@@ -15,9 +15,7 @@
-
-
-
+
diff --git a/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsModule.cs b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsModule.cs
new file mode 100644
index 0000000000..9c3832e49b
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsModule.cs
@@ -0,0 +1,9 @@
+using Volo.Abp.Modularity;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public class AbpDistributedLockingAbstractionsModule : AbpModule
+ {
+
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/IAbpDistributedLock.cs b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/IAbpDistributedLock.cs
new file mode 100644
index 0000000000..6619ea337c
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/IAbpDistributedLock.cs
@@ -0,0 +1,26 @@
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+using JetBrains.Annotations;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public interface IAbpDistributedLock
+ {
+ ///
+ /// Tries to acquire a named lock.
+ /// Returns a disposable object to release the lock.
+ /// It is suggested to use this method within a using block.
+ /// Returns null if the lock could not be handled.
+ ///
+ /// The name of the lock
+ /// Timeout value
+ /// Cancellation token
+ [ItemCanBeNull]
+ Task TryAcquireAsync(
+ [NotNull] string name,
+ TimeSpan timeout = default,
+ CancellationToken cancellationToken = default
+ );
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/IAbpDistributedLockHandle.cs b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/IAbpDistributedLockHandle.cs
new file mode 100644
index 0000000000..baac5296c7
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/IAbpDistributedLockHandle.cs
@@ -0,0 +1,9 @@
+using System;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public interface IAbpDistributedLockHandle : IAsyncDisposable
+ {
+
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/LocalAbpDistributedLock.cs b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/LocalAbpDistributedLock.cs
new file mode 100644
index 0000000000..4b481d3549
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/LocalAbpDistributedLock.cs
@@ -0,0 +1,30 @@
+using System;
+using System.Collections.Concurrent;
+using System.Threading;
+using System.Threading.Tasks;
+using Volo.Abp.DependencyInjection;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public class LocalAbpDistributedLock : IAbpDistributedLock, ISingletonDependency
+ {
+ private readonly ConcurrentDictionary _localSyncObjects = new();
+
+ public async Task TryAcquireAsync(
+ string name,
+ TimeSpan timeout = default,
+ CancellationToken cancellationToken = default)
+ {
+ Check.NotNullOrWhiteSpace(name, nameof(name));
+
+ var semaphore = _localSyncObjects.GetOrAdd(name, _ => new SemaphoreSlim(1, 1));
+
+ if (!await semaphore.WaitAsync(timeout, cancellationToken))
+ {
+ return null;
+ }
+
+ return new LocalAbpDistributedLockHandle(semaphore);
+ }
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/LocalAbpDistributedLockHandle.cs b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/LocalAbpDistributedLockHandle.cs
new file mode 100644
index 0000000000..da12739311
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking.Abstractions/Volo/Abp/DistributedLocking/LocalAbpDistributedLockHandle.cs
@@ -0,0 +1,21 @@
+using System.Threading;
+using System.Threading.Tasks;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public class LocalAbpDistributedLockHandle : IAbpDistributedLockHandle
+ {
+ private readonly SemaphoreSlim _semaphore;
+
+ public LocalAbpDistributedLockHandle(SemaphoreSlim semaphore)
+ {
+ _semaphore = semaphore;
+ }
+
+ public ValueTask DisposeAsync()
+ {
+ _semaphore.Release();
+ return default;
+ }
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj b/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj
index c4b590412d..17ad2f02a4 100644
--- a/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj
+++ b/framework/src/Volo.Abp.DistributedLocking/Volo.Abp.DistributedLocking.csproj
@@ -15,7 +15,7 @@
-
+
diff --git a/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockHandleExtensions.cs b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockHandleExtensions.cs
new file mode 100644
index 0000000000..f0b69116f5
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockHandleExtensions.cs
@@ -0,0 +1,14 @@
+using System;
+using Medallion.Threading;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public static class AbpDistributedLockHandleExtensions
+ {
+ public static IDistributedSynchronizationHandle ToDistributedSynchronizationHandle(
+ this IAbpDistributedLockHandle handle)
+ {
+ return handle.As().Handle;
+ }
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockingModule.cs b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockingModule.cs
index 877ac54244..3bfdbbf3f3 100644
--- a/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockingModule.cs
+++ b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/AbpDistributedLockingModule.cs
@@ -2,6 +2,9 @@
namespace Volo.Abp.DistributedLocking
{
+ [DependsOn(
+ typeof(AbpDistributedLockingAbstractionsModule)
+ )]
public class AbpDistributedLockingModule : AbpModule
{
diff --git a/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/MedallionAbpDistributedLock.cs b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/MedallionAbpDistributedLock.cs
new file mode 100644
index 0000000000..4b46096da9
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/MedallionAbpDistributedLock.cs
@@ -0,0 +1,35 @@
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+using Medallion.Threading;
+using Volo.Abp.DependencyInjection;
+
+namespace Volo.Abp.DistributedLocking
+{
+ [Dependency(ReplaceServices = true)]
+ public class MedallionAbpDistributedLock : IAbpDistributedLock, ITransientDependency
+ {
+ protected IDistributedLockProvider DistributedLockProvider { get; }
+
+ public MedallionAbpDistributedLock(IDistributedLockProvider distributedLockProvider)
+ {
+ DistributedLockProvider = distributedLockProvider;
+ }
+
+ public async Task TryAcquireAsync(
+ string name,
+ TimeSpan timeout = default,
+ CancellationToken cancellationToken = default)
+ {
+ Check.NotNullOrWhiteSpace(name, nameof(name));
+
+ var handle = await DistributedLockProvider.TryAcquireLockAsync(name, timeout, cancellationToken);
+ if (handle == null)
+ {
+ return null;
+ }
+
+ return new MedallionAbpDistributedLockHandle(handle);
+ }
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/MedallionAbpDistributedLockHandle.cs b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/MedallionAbpDistributedLockHandle.cs
new file mode 100644
index 0000000000..d22371b862
--- /dev/null
+++ b/framework/src/Volo.Abp.DistributedLocking/Volo/Abp/DistributedLocking/MedallionAbpDistributedLockHandle.cs
@@ -0,0 +1,20 @@
+using System.Threading.Tasks;
+using Medallion.Threading;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public class MedallionAbpDistributedLockHandle : IAbpDistributedLockHandle
+ {
+ public IDistributedSynchronizationHandle Handle { get; }
+
+ public MedallionAbpDistributedLockHandle(IDistributedSynchronizationHandle handle)
+ {
+ Handle = handle;
+ }
+
+ public ValueTask DisposeAsync()
+ {
+ return Handle.DisposeAsync();
+ }
+ }
+}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/FodyWeavers.xml b/framework/src/Volo.Abp.EventBus.Boxes/FodyWeavers.xml
deleted file mode 100644
index 1715698ccd..0000000000
--- a/framework/src/Volo.Abp.EventBus.Boxes/FodyWeavers.xml
+++ /dev/null
@@ -1,3 +0,0 @@
-
-
-
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs
deleted file mode 100644
index 0c7a2c4247..0000000000
--- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesModule.cs
+++ /dev/null
@@ -1,20 +0,0 @@
-using Volo.Abp.BackgroundWorkers;
-using Volo.Abp.DistributedLocking;
-using Volo.Abp.Modularity;
-
-namespace Volo.Abp.EventBus.Boxes
-{
- [DependsOn(
- typeof(AbpEventBusModule),
- typeof(AbpBackgroundWorkersModule),
- typeof(AbpDistributedLockingModule)
- )]
- public class AbpEventBusBoxesModule : AbpModule
- {
- public override void OnApplicationInitialization(ApplicationInitializationContext context)
- {
- context.AddBackgroundWorker();
- context.AddBackgroundWorker();
- }
- }
-}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/TaskHelper.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/TaskHelper.cs
deleted file mode 100644
index 87abe18da8..0000000000
--- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/TaskHelper.cs
+++ /dev/null
@@ -1,20 +0,0 @@
-using System.Threading;
-using System.Threading.Tasks;
-
-namespace Volo.Abp.EventBus.Boxes
-{
- internal static class TaskDelayHelper
- {
- public static async Task DelayAsync(int milliseconds, CancellationToken cancellationToken = default)
- {
- try
- {
- await Task.Delay(milliseconds, cancellationToken);
- }
- catch (TaskCanceledException)
- {
- return;
- }
- }
- }
-}
\ No newline at end of file
diff --git a/framework/src/Volo.Abp.EventBus/Volo.Abp.EventBus.csproj b/framework/src/Volo.Abp.EventBus/Volo.Abp.EventBus.csproj
index 5de5b0893f..bee75bf7ff 100644
--- a/framework/src/Volo.Abp.EventBus/Volo.Abp.EventBus.csproj
+++ b/framework/src/Volo.Abp.EventBus/Volo.Abp.EventBus.csproj
@@ -15,6 +15,8 @@
+
+
diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs
index 3bbdf3d72b..58280c7066 100644
--- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs
+++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs
@@ -1,7 +1,10 @@
using Microsoft.Extensions.DependencyInjection;
using System;
using System.Collections.Generic;
+using Volo.Abp.BackgroundWorkers;
+using Volo.Abp.DistributedLocking;
using Volo.Abp.EventBus.Abstractions;
+using Volo.Abp.EventBus.Boxes;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.EventBus.Local;
using Volo.Abp.Guids;
@@ -16,7 +19,9 @@ namespace Volo.Abp.EventBus
typeof(AbpEventBusAbstractionsModule),
typeof(AbpMultiTenancyModule),
typeof(AbpJsonModule),
- typeof(AbpGuidsModule)
+ typeof(AbpGuidsModule),
+ typeof(AbpBackgroundWorkersModule),
+ typeof(AbpDistributedLockingAbstractionsModule)
)]
public class AbpEventBusModule : AbpModule
{
@@ -24,6 +29,12 @@ namespace Volo.Abp.EventBus
{
AddEventHandlers(context.Services);
}
+
+ public override void OnApplicationInitialization(ApplicationInitializationContext context)
+ {
+ context.AddBackgroundWorker();
+ context.AddBackgroundWorker();
+ }
private static void AddEventHandlers(IServiceCollection services)
{
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpDistributedEventBusExtensions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusExtensions.cs
similarity index 100%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpDistributedEventBusExtensions.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusExtensions.cs
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpEventBusBoxesOptions.cs
similarity index 100%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/AbpEventBusBoxesOptions.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpEventBusBoxesOptions.cs
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IInboxProcessor.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IInboxProcessor.cs
similarity index 100%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IInboxProcessor.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IInboxProcessor.cs
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IOutboxSender.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IOutboxSender.cs
similarity index 100%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/IOutboxSender.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IOutboxSender.cs
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessManager.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessManager.cs
similarity index 100%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessManager.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessManager.cs
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs
similarity index 86%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs
index e3c71771ad..6a25804e74 100644
--- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs
+++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs
@@ -1,12 +1,12 @@
using System;
using System.Threading;
using System.Threading.Tasks;
-using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using Volo.Abp.DependencyInjection;
+using Volo.Abp.DistributedLocking;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Threading;
using Volo.Abp.Timing;
@@ -19,7 +19,7 @@ namespace Volo.Abp.EventBus.Boxes
protected IServiceProvider ServiceProvider { get; }
protected AbpAsyncTimer Timer { get; }
protected IDistributedEventBus DistributedEventBus { get; }
- protected IDistributedLockProvider DistributedLockProvider { get; }
+ protected IAbpDistributedLock DistributedLock { get; }
protected IUnitOfWorkManager UnitOfWorkManager { get; }
protected IClock Clock { get; }
protected IEventInbox Inbox { get; private set; }
@@ -28,7 +28,7 @@ namespace Volo.Abp.EventBus.Boxes
protected DateTime? LastCleanTime { get; set; }
- protected string DistributedLockName => "Inbox_" + InboxConfig.Name;
+ protected string DistributedLockName => "AbpInbox_" + InboxConfig.Name;
public ILogger Logger { get; set; }
protected CancellationTokenSource StoppingTokenSource { get; }
protected CancellationToken StoppingToken { get; }
@@ -37,7 +37,7 @@ namespace Volo.Abp.EventBus.Boxes
IServiceProvider serviceProvider,
AbpAsyncTimer timer,
IDistributedEventBus distributedEventBus,
- IDistributedLockProvider distributedLockProvider,
+ IAbpDistributedLock distributedLock,
IUnitOfWorkManager unitOfWorkManager,
IClock clock,
IOptions eventBusBoxesOptions)
@@ -45,11 +45,11 @@ namespace Volo.Abp.EventBus.Boxes
ServiceProvider = serviceProvider;
Timer = timer;
DistributedEventBus = distributedEventBus;
- DistributedLockProvider = distributedLockProvider;
+ DistributedLock = distributedLock;
UnitOfWorkManager = unitOfWorkManager;
Clock = clock;
EventBusBoxesOptions = eventBusBoxesOptions.Value;
- Timer.Period = EventBusBoxesOptions.PeriodTimeSpan.Milliseconds;
+ Timer.Period = Convert.ToInt32(EventBusBoxesOptions.PeriodTimeSpan.TotalMilliseconds);
Timer.Elapsed += TimerOnElapsed;
Logger = NullLogger.Instance;
StoppingTokenSource = new CancellationTokenSource();
@@ -84,7 +84,7 @@ namespace Volo.Abp.EventBus.Boxes
return;
}
- await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName, cancellationToken: StoppingToken))
+ await using (var handle = await DistributedLock.TryAcquireAsync(DistributedLockName, cancellationToken: StoppingToken))
{
if (handle != null)
{
@@ -120,7 +120,11 @@ namespace Volo.Abp.EventBus.Boxes
else
{
Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName);
- await TaskDelayHelper.DelayAsync(EventBusBoxesOptions.DistributedLockWaitDuration.Milliseconds, StoppingToken);
+ try
+ {
+ await Task.Delay(EventBusBoxesOptions.DistributedLockWaitDuration, StoppingToken);
+ }
+ catch (TaskCanceledException) { }
}
}
}
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/OutboxSender.cs
similarity index 83%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/OutboxSender.cs
index 32545227c3..873878bca7 100644
--- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSender.cs
+++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/OutboxSender.cs
@@ -1,12 +1,12 @@
using System;
using System.Threading;
using System.Threading.Tasks;
-using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
using Volo.Abp.DependencyInjection;
+using Volo.Abp.DistributedLocking;
using Volo.Abp.EventBus.Distributed;
using Volo.Abp.Threading;
@@ -17,11 +17,11 @@ namespace Volo.Abp.EventBus.Boxes
protected IServiceProvider ServiceProvider { get; }
protected AbpAsyncTimer Timer { get; }
protected IDistributedEventBus DistributedEventBus { get; }
- protected IDistributedLockProvider DistributedLockProvider { get; }
+ protected IAbpDistributedLock DistributedLock { get; }
protected IEventOutbox Outbox { get; private set; }
protected OutboxConfig OutboxConfig { get; private set; }
protected AbpEventBusBoxesOptions EventBusBoxesOptions { get; }
- protected string DistributedLockName => "Outbox_" + OutboxConfig.Name;
+ protected string DistributedLockName => "AbpOutbox_" + OutboxConfig.Name;
public ILogger Logger { get; set; }
protected CancellationTokenSource StoppingTokenSource { get; }
@@ -31,15 +31,15 @@ namespace Volo.Abp.EventBus.Boxes
IServiceProvider serviceProvider,
AbpAsyncTimer timer,
IDistributedEventBus distributedEventBus,
- IDistributedLockProvider distributedLockProvider,
+ IAbpDistributedLock distributedLock,
IOptions eventBusBoxesOptions)
{
ServiceProvider = serviceProvider;
- Timer = timer;
DistributedEventBus = distributedEventBus;
- DistributedLockProvider = distributedLockProvider;
+ DistributedLock = distributedLock;
EventBusBoxesOptions = eventBusBoxesOptions.Value;
- Timer.Period = EventBusBoxesOptions.PeriodTimeSpan.Milliseconds;
+ Timer = timer;
+ Timer.Period = Convert.ToInt32(EventBusBoxesOptions.PeriodTimeSpan.TotalMilliseconds);
Timer.Elapsed += TimerOnElapsed;
Logger = NullLogger.Instance;
StoppingTokenSource = new CancellationTokenSource();
@@ -69,7 +69,7 @@ namespace Volo.Abp.EventBus.Boxes
protected virtual async Task RunAsync()
{
- await using (var handle = await DistributedLockProvider.TryAcquireLockAsync(DistributedLockName, cancellationToken: StoppingToken))
+ await using (var handle = await DistributedLock.TryAcquireAsync(DistributedLockName, cancellationToken: StoppingToken))
{
if (handle != null)
{
@@ -100,7 +100,11 @@ namespace Volo.Abp.EventBus.Boxes
else
{
Logger.LogDebug("Could not obtain the distributed lock: " + DistributedLockName);
- await TaskDelayHelper.DelayAsync(EventBusBoxesOptions.DistributedLockWaitDuration.Milliseconds, StoppingToken);
+ try
+ {
+ await Task.Delay(EventBusBoxesOptions.DistributedLockWaitDuration, StoppingToken);
+ }
+ catch (TaskCanceledException) { }
}
}
}
diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSenderManager.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/OutboxSenderManager.cs
similarity index 100%
rename from framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/OutboxSenderManager.cs
rename to framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/OutboxSenderManager.cs
diff --git a/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo.Abp.DistributedLocking.Abstractions.Tests.csproj b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo.Abp.DistributedLocking.Abstractions.Tests.csproj
new file mode 100644
index 0000000000..8194846f0c
--- /dev/null
+++ b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo.Abp.DistributedLocking.Abstractions.Tests.csproj
@@ -0,0 +1,18 @@
+
+
+
+
+
+ net6.0
+ true
+
+
+
+
+
+
+
+
+
+
+
diff --git a/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsTestBase.cs b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsTestBase.cs
new file mode 100644
index 0000000000..074bd342d2
--- /dev/null
+++ b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsTestBase.cs
@@ -0,0 +1,12 @@
+using Volo.Abp.Testing;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public class AbpDistributedLockingAbstractionsTestBase : AbpIntegratedTest
+ {
+ protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options)
+ {
+ options.UseAutofac();
+ }
+ }
+}
diff --git a/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsTestModule.cs b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsTestModule.cs
new file mode 100644
index 0000000000..cfc9278335
--- /dev/null
+++ b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/AbpDistributedLockingAbstractionsTestModule.cs
@@ -0,0 +1,15 @@
+using Volo.Abp.Autofac;
+using Volo.Abp.Modularity;
+
+namespace Volo.Abp.DistributedLocking
+{
+ [DependsOn(
+ typeof(AbpTestBaseModule),
+ typeof(AbpDistributedLockingAbstractionsModule),
+ typeof(AbpAutofacModule)
+ )]
+ public class AbpDistributedLockingAbstractionsTestModule : AbpModule
+ {
+
+ }
+}
\ No newline at end of file
diff --git a/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/LocalDistributedLock_Tests.cs b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/LocalDistributedLock_Tests.cs
new file mode 100644
index 0000000000..ce7fda009d
--- /dev/null
+++ b/framework/test/Volo.Abp.DistributedLocking.Abstractions.Tests/Volo/Abp/DistributedLocking/LocalDistributedLock_Tests.cs
@@ -0,0 +1,73 @@
+using System.Threading.Tasks;
+using Shouldly;
+using Xunit;
+
+namespace Volo.Abp.DistributedLocking
+{
+ public class LocalDistributedLock_Tests : AbpDistributedLockingAbstractionsTestBase
+ {
+ private readonly IAbpDistributedLock _distributedLock;
+
+ public LocalDistributedLock_Tests()
+ {
+ _distributedLock = GetRequiredService();
+ }
+
+ [Fact]
+ public void Should_Be_Instance_Of_LocalAbpDistributedLock()
+ {
+ _distributedLock.ShouldBeOfType();
+ }
+
+ [Fact]
+ public async Task Should_Lock_With_TryAcquire()
+ {
+ await using (var handle = await _distributedLock.TryAcquireAsync("lock1"))
+ {
+ handle.ShouldNotBeNull();
+ }
+ }
+
+ [Fact]
+ public async Task Should_Not_Acquire_If_Already_Locked()
+ {
+ await using (var handle = await _distributedLock.TryAcquireAsync("lock1"))
+ {
+ handle.ShouldNotBeNull();
+
+ await Task.Run(async () =>
+ {
+ await using (var handle2 = await _distributedLock.TryAcquireAsync("lock1"))
+ {
+ handle2.ShouldBeNull();
+ }
+ });
+ }
+
+ await Task.Run(async () =>
+ {
+ await using (var handle = await _distributedLock.TryAcquireAsync("lock1"))
+ {
+ handle.ShouldNotBeNull();
+ }
+ });
+ }
+
+ [Fact]
+ public async Task Should_Obtain_Multiple_Locks()
+ {
+ await using (var handle = await _distributedLock.TryAcquireAsync("lock1"))
+ {
+ handle.ShouldNotBeNull();
+
+ await Task.Run(async () =>
+ {
+ await using (var handle2 = await _distributedLock.TryAcquireAsync("lock2"))
+ {
+ handle2.ShouldNotBeNull();
+ }
+ });
+ }
+ }
+ }
+}
\ No newline at end of file
diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/DemoAppModule.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/DemoAppModule.cs
index f64191ea3f..589892c901 100644
--- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/DemoAppModule.cs
+++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/DemoAppModule.cs
@@ -1,6 +1,4 @@
-using Microsoft.Extensions.DependencyInjection;
-using Microsoft.Extensions.Logging;
-using Volo.Abp.Autofac;
+using Volo.Abp.Autofac;
using Volo.Abp.BackgroundJobs.DemoApp.Shared;
using Volo.Abp.BackgroundJobs.EntityFrameworkCore;
using Volo.Abp.EntityFrameworkCore;
diff --git a/nupkg/common.ps1 b/nupkg/common.ps1
index fd99c6a068..75a8b0c9c0 100644
--- a/nupkg/common.ps1
+++ b/nupkg/common.ps1
@@ -93,6 +93,7 @@ $projects = (
"framework/src/Volo.Abp.Ddd.Application",
"framework/src/Volo.Abp.Ddd.Application.Contracts",
"framework/src/Volo.Abp.Ddd.Domain",
+ "framework/src/Volo.Abp.DistributedLocking.Abstractions",
"framework/src/Volo.Abp.DistributedLocking",
"framework/src/Volo.Abp.Emailing",
"framework/src/Volo.Abp.EntityFrameworkCore",
@@ -104,7 +105,6 @@ $projects = (
"framework/src/Volo.Abp.EntityFrameworkCore.SqlServer",
"framework/src/Volo.Abp.EventBus.Abstractions",
"framework/src/Volo.Abp.EventBus",
- "framework/src/Volo.Abp.EventBus.Boxes",
"framework/src/Volo.Abp.EventBus.RabbitMQ",
"framework/src/Volo.Abp.EventBus.Kafka",
"framework/src/Volo.Abp.EventBus.Rebus",