From e101f3161f26be7ae7c343b5431ed424522fafc9 Mon Sep 17 00:00:00 2001 From: maliming Date: Fri, 3 Jul 2026 14:14:53 +0800 Subject: [PATCH] Add tests and demo app for background job worker enhancements --- .../AbpAutoLockNameWorkerTestModule.cs | 27 ++ .../AbpBackgroundJobCleanupTestModule.cs | 23 ++ .../AbpBackgroundJobWorkerTestModule.cs | 25 ++ .../AbpDuplicateLockNameTestModule.cs | 27 ++ .../AbpDuplicateWorkerTestModule.cs | 27 ++ .../AbpMultiWorkerTestModule.cs | 28 ++ .../AbpSameJobNameTestModule.cs | 29 ++ .../BackgroundJobCleanupWorker_Tests.cs | 132 ++++++++ .../BackgroundJobWorkerTestBase.cs | 11 + .../BackgroundJobWorker_AutoLockName_Tests.cs | 34 ++ ...dJobWorker_DuplicateConfiguration_Tests.cs | 63 ++++ ...JobWorker_MultiWorkerRegistration_Tests.cs | 37 +++ .../BackgroundJobWorker_Tests.cs | 300 ++++++++++++++++++ .../RecordingBackgroundJobWorker.cs | 51 +++ .../TestableBackgroundJobCleanupWorker.cs | 26 ++ .../TestableBackgroundJobWorker.cs | 55 ++++ .../Volo/Abp/BackgroundJobs/WorkerTestJobs.cs | 99 ++++++ .../DemoAppModule.cs | 32 +- .../Jobs/CalculateAwsFeesJob.cs | 20 ++ .../Jobs/CalculateAzureFeesJob.cs | 20 ++ .../Jobs/SendEmailJob.cs | 23 ++ ...mpletionTime_To_BackgroundJobs.Designer.cs | 99 ++++++ ..._Added_CompletionTime_To_BackgroundJobs.cs | 47 +++ .../DemoAppDbContextModelSnapshot.cs | 7 +- .../BackgroundJobRepository_Tests.cs | 99 +++++- .../BackgroundJobs/BackgroundJobsTestData.cs | 1 + .../BackgroundJobsTestDataBuilder.cs | 16 + 27 files changed, 1347 insertions(+), 11 deletions(-) create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpAutoLockNameWorkerTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobCleanupTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateLockNameTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateWorkerTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpMultiWorkerTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpSameJobNameTestModule.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker_Tests.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorkerTestBase.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_AutoLockName_Tests.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_DuplicateConfiguration_Tests.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_MultiWorkerRegistration_Tests.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_Tests.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/RecordingBackgroundJobWorker.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobCleanupWorker.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobWorker.cs create mode 100644 framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/WorkerTestJobs.cs create mode 100644 modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAwsFeesJob.cs create mode 100644 modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAzureFeesJob.cs create mode 100644 modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/SendEmailJob.cs create mode 100644 modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.Designer.cs create mode 100644 modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.cs diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpAutoLockNameWorkerTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpAutoLockNameWorkerTestModule.cs new file mode 100644 index 0000000000..38b640b84b --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpAutoLockNameWorkerTestModule.cs @@ -0,0 +1,27 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpAutoLockNameWorkerTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + context.Services.AddSingleton(); + context.Services.Replace(ServiceDescriptor.Transient()); + + // No lock name is given: it is derived from the job argument type names. + Configure(options => + { + options.AddDedicatedWorker(); + options.AddDedicatedWorker(); + }); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobCleanupTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobCleanupTestModule.cs new file mode 100644 index 0000000000..b29c60784e --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobCleanupTestModule.cs @@ -0,0 +1,23 @@ +using System; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpBackgroundJobCleanupTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + // IsJobExecutionEnabled is true by default, which the cleanup worker requires. + Configure(options => + { + options.StoreSuccessfulJobs = true; + options.SuccessfulJobRetentionTime = TimeSpan.FromDays(1); + }); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerTestModule.cs new file mode 100644 index 0000000000..db56843a4b --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobWorkerTestModule.cs @@ -0,0 +1,25 @@ +using Microsoft.Extensions.DependencyInjection; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpBackgroundJobWorkerTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + // Drive the worker manually in tests; don't run the real periodic worker. + Configure(options => + { + options.IsJobExecutionEnabled = false; + }); + + context.Services.AddSingleton(); + context.Services.AddScoped(); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateLockNameTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateLockNameTestModule.cs new file mode 100644 index 0000000000..6e35c81dfc --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateLockNameTestModule.cs @@ -0,0 +1,27 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpDuplicateLockNameTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + context.Services.AddSingleton(); + context.Services.Replace(ServiceDescriptor.Transient()); + + // Two workers with different job types but the same lock name, which must fail at initialization. + Configure(options => + { + options.AddDedicatedWorker("dup-lock"); + options.AddDedicatedWorker("dup-lock"); + }); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateWorkerTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateWorkerTestModule.cs new file mode 100644 index 0000000000..f4d9aae208 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpDuplicateWorkerTestModule.cs @@ -0,0 +1,27 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpDuplicateWorkerTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + context.Services.AddSingleton(); + context.Services.Replace(ServiceDescriptor.Transient()); + + // The same job type is assigned to two dedicated workers, which must fail at initialization. + Configure(options => + { + options.AddDedicatedWorker("lock-a"); + options.AddDedicatedWorker("lock-b"); + }); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpMultiWorkerTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpMultiWorkerTestModule.cs new file mode 100644 index 0000000000..d437c9325f --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpMultiWorkerTestModule.cs @@ -0,0 +1,28 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpMultiWorkerTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + context.Services.AddSingleton(); + + // Replace the real worker with a recording one to assert how the manager resolves and starts workers. + context.Services.Replace(ServiceDescriptor.Transient()); + + Configure(options => + { + options.AddDedicatedWorker("lock-a"); + options.AddDedicatedWorker("lock-b"); + }); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpSameJobNameTestModule.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpSameJobNameTestModule.cs new file mode 100644 index 0000000000..f4ddbb8be8 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpSameJobNameTestModule.cs @@ -0,0 +1,29 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.DependencyInjection.Extensions; +using Volo.Abp.Autofac; +using Volo.Abp.Modularity; + +namespace Volo.Abp.BackgroundJobs; + +[DependsOn( + typeof(AbpBackgroundJobsModule), + typeof(AbpAutofacModule), + typeof(AbpTestBaseModule) +)] +public class AbpSameJobNameTestModule : AbpModule +{ + public override void ConfigureServices(ServiceConfigurationContext context) + { + context.Services.AddSingleton(); + context.Services.Replace(ServiceDescriptor.Transient()); + + // Two dedicated workers with different args types that both resolve to "shared-job-name". + // AddDedicatedWorker's eager check compares by type (both pass); only the manager's backstop + // (which resolves job names) can catch this. + Configure(options => + { + options.AddDedicatedWorker("lock-a"); + options.AddDedicatedWorker("lock-b"); + }); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker_Tests.cs new file mode 100644 index 0000000000..770d135cee --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobCleanupWorker_Tests.cs @@ -0,0 +1,132 @@ +using System; +using System.Collections.Generic; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Shouldly; +using Volo.Abp.BackgroundWorkers; +using Volo.Abp.DistributedLocking; +using Volo.Abp.Testing; +using Volo.Abp.Threading; +using Volo.Abp.Timing; +using Xunit; + +namespace Volo.Abp.BackgroundJobs; + +public class BackgroundJobCleanupWorker_Tests : AbpIntegratedTest +{ + private readonly IBackgroundJobStore _store; + private readonly IClock _clock; + private readonly AbpBackgroundJobWorkerOptions _workerOptions; + + public BackgroundJobCleanupWorker_Tests() + { + _store = GetRequiredService(); + _clock = GetRequiredService(); + _workerOptions = GetRequiredService>().Value; + } + + protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) + { + options.UseAutofac(); + } + + private TestableBackgroundJobCleanupWorker CreateWorker() + { + return new TestableBackgroundJobCleanupWorker( + GetRequiredService(), + GetRequiredService(), + GetRequiredService>(), + GetRequiredService>(), + GetRequiredService()); + } + + private Task RunCleanupAsync() + { + return CreateWorker().DoWorkPublicAsync(new PeriodicBackgroundWorkerContext(ServiceProvider)); + } + + private async Task InsertCompletedJobAsync(DateTime completionTime) + { + var id = Guid.NewGuid(); + await _store.InsertAsync(new BackgroundJobInfo + { + Id = id, + JobName = "job-a", + JobArgs = "{}", + CreationTime = _clock.Now, + NextTryTime = _clock.Now, + CompletionTime = completionTime + }); + return id; + } + + [Fact] + public async Task Should_Delete_Completed_Jobs_Older_Than_Retention() + { + // Retention is 1 day (configured in the module). + var oldJobId = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(2))); + var recentJobId = await InsertCompletedJobAsync(_clock.Now); + + await RunCleanupAsync(); + + (await _store.FindAsync(oldJobId)).ShouldBeNull(); // older than retention → deleted + (await _store.FindAsync(recentJobId)).ShouldNotBeNull(); // within retention → kept + } + + [Fact] + public async Task Should_Not_Delete_When_StoreSuccessfulJobs_Disabled() + { + _workerOptions.StoreSuccessfulJobs = false; + + var oldJobId = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(2))); + + await RunCleanupAsync(); + + (await _store.FindAsync(oldJobId)).ShouldNotBeNull(); + } + + [Fact] + public async Task Should_Not_Delete_When_Retention_Is_Null() + { + _workerOptions.SuccessfulJobRetentionTime = null; + + var oldJobId = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(2))); + + await RunCleanupAsync(); + + (await _store.FindAsync(oldJobId)).ShouldNotBeNull(); + } + + [Fact] + public async Task Should_Delete_All_Old_Jobs_In_Batches() + { + _workerOptions.MaxJobFetchCount = 2; + + var ids = new List(); + for (var i = 0; i < 5; i++) + { + ids.Add(await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(2)))); + } + + await RunCleanupAsync(); + + foreach (var id in ids) + { + (await _store.FindAsync(id)).ShouldBeNull(); + } + } + + [Fact] + public async Task Should_Not_Loop_Forever_When_MaxJobFetchCount_Is_Zero() + { + _workerOptions.MaxJobFetchCount = 0; + + var oldJobId = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(2))); + + // Must return (not hang) even though nothing can be fetched/deleted with a zero page size. + await RunCleanupAsync(); + + (await _store.FindAsync(oldJobId)).ShouldNotBeNull(); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorkerTestBase.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorkerTestBase.cs new file mode 100644 index 0000000000..2824c1a231 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorkerTestBase.cs @@ -0,0 +1,11 @@ +using Volo.Abp.Testing; + +namespace Volo.Abp.BackgroundJobs; + +public abstract class BackgroundJobWorkerTestBase : AbpIntegratedTest +{ + protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) + { + options.UseAutofac(); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_AutoLockName_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_AutoLockName_Tests.cs new file mode 100644 index 0000000000..76e4734b05 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_AutoLockName_Tests.cs @@ -0,0 +1,34 @@ +using System; +using System.Linq; +using Shouldly; +using Volo.Abp.Testing; +using Xunit; + +namespace Volo.Abp.BackgroundJobs; + +public class BackgroundJobWorker_AutoLockName_Tests : AbpIntegratedTest +{ + protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) + { + options.UseAutofac(); + } + + [Fact] + public void Should_Derive_Bounded_Lock_Name_From_Job_Args_Types() + { + var records = GetRequiredService().Records; + + var dedicated = records.Where(r => r.JobNameFilter?.Mode == BackgroundJobNameFilterMode.Include).ToList(); + dedicated.Count.ShouldBe(2); + + // Lock name is derived (prefix + MD5 of the full type name), so it is stable and length-bounded. + var expectedA = "AbpBackgroundJobDedicatedWorker:" + typeof(WorkerJobAArgs).FullName!.ToMd5(); + var expectedB = "AbpBackgroundJobDedicatedWorker:" + typeof(WorkerJobBArgs).FullName!.ToMd5(); + + dedicated.ShouldContain(r => r.DistributedLockName == expectedA); + dedicated.ShouldContain(r => r.DistributedLockName == expectedB); + + // Bounded length regardless of how long the type names are. + dedicated.ShouldAllBe(r => r.DistributedLockName!.Length <= 64); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_DuplicateConfiguration_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_DuplicateConfiguration_Tests.cs new file mode 100644 index 0000000000..bdd80efa7a --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_DuplicateConfiguration_Tests.cs @@ -0,0 +1,63 @@ +using Microsoft.Extensions.DependencyInjection; +using Shouldly; +using Volo.Abp.Autofac; +using Xunit; + +namespace Volo.Abp.BackgroundJobs; + +public class BackgroundJobWorker_DuplicateConfiguration_Tests +{ + [Fact] + public void Should_Throw_Without_Starting_Any_Worker_When_A_Job_Type_Is_Assigned_To_Multiple_Workers() + { + using var application = AbpApplicationFactory.Create(options => + { + options.UseAutofac(); + }); + + var exception = Record.Exception(() => application.Initialize()); + + exception.ShouldNotBeNull(); + exception.ToString().ShouldContain("dedicated worker"); + + // Validation must happen before any worker is started. + var recorder = application.ServiceProvider.GetRequiredService(); + recorder.Records.ShouldBeEmpty(); + } + + [Fact] + public void Should_Throw_Without_Starting_Any_Worker_When_Two_Workers_Share_A_Lock_Name() + { + using var application = AbpApplicationFactory.Create(options => + { + options.UseAutofac(); + }); + + var exception = Record.Exception(() => application.Initialize()); + + exception.ShouldNotBeNull(); + exception.ToString().ShouldContain("lock name"); + + var recorder = application.ServiceProvider.GetRequiredService(); + recorder.Records.ShouldBeEmpty(); + } + + [Fact] + public void Should_Throw_Without_Starting_Any_Worker_When_Different_Args_Types_Resolve_To_The_Same_Job_Name() + { + using var application = AbpApplicationFactory.Create(options => + { + options.UseAutofac(); + }); + + // The two args types pass the eager (by-type) check but resolve to the same job name, + // so only the manager's backstop validation can reject them. + var exception = Record.Exception(() => application.Initialize()); + + exception.ShouldNotBeNull(); + exception.ToString().ShouldContain("more than one dedicated worker"); + + var recorder = application.ServiceProvider.GetRequiredService(); + recorder.Records.ShouldBeEmpty(); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_MultiWorkerRegistration_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_MultiWorkerRegistration_Tests.cs new file mode 100644 index 0000000000..8279d2d493 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_MultiWorkerRegistration_Tests.cs @@ -0,0 +1,37 @@ +using System.Linq; +using Shouldly; +using Volo.Abp.Testing; +using Xunit; + +namespace Volo.Abp.BackgroundJobs; + +public class BackgroundJobWorker_MultiWorkerRegistration_Tests : AbpIntegratedTest +{ + protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) + { + options.UseAutofac(); + } + + [Fact] + public void Should_Start_Dedicated_Workers_And_A_Default_Worker() + { + // The workers are resolved from DI (RecordingBackgroundJobWorker replaces the real one), + // proving the manager honors the registered/replaced IBackgroundJobWorker. + var records = GetRequiredService().Records; + + records.Count.ShouldBe(3); + + var jobAName = BackgroundJobNameAttribute.GetName(); + var jobBName = BackgroundJobNameAttribute.GetName(); + + var dedicated = records.Where(r => r.JobNameFilter?.Mode == BackgroundJobNameFilterMode.Include).ToList(); + dedicated.Count.ShouldBe(2); + dedicated.ShouldContain(r => r.DistributedLockName == "lock-a" && r.JobNameFilter!.JobNames.Contains(jobAName)); + dedicated.ShouldContain(r => r.DistributedLockName == "lock-b" && r.JobNameFilter!.JobNames.Contains(jobBName)); + + var defaultWorker = records.Single(r => r.JobNameFilter?.Mode == BackgroundJobNameFilterMode.Exclude); + defaultWorker.DistributedLockName.ShouldBeNull(); + defaultWorker.JobNameFilter!.JobNames.ShouldContain(jobAName); + defaultWorker.JobNameFilter!.JobNames.ShouldContain(jobBName); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_Tests.cs new file mode 100644 index 0000000000..4d8d506f19 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobWorker_Tests.cs @@ -0,0 +1,300 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Shouldly; +using Volo.Abp.BackgroundWorkers; +using Volo.Abp.DependencyInjection; +using Volo.Abp.DistributedLocking; +using Volo.Abp.Threading; +using Volo.Abp.Timing; +using Xunit; + +// ReSharper disable PossibleMultipleEnumeration + +namespace Volo.Abp.BackgroundJobs; + +public class BackgroundJobWorker_Tests : BackgroundJobWorkerTestBase +{ + private readonly IBackgroundJobStore _store; + private readonly IClock _clock; + private readonly AbpBackgroundJobWorkerOptions _workerOptions; + + public BackgroundJobWorker_Tests() + { + _store = GetRequiredService(); + _clock = GetRequiredService(); + _workerOptions = GetRequiredService>().Value; + } + + private TestableBackgroundJobWorker CreateWorker() + { + return new TestableBackgroundJobWorker( + GetRequiredService(), + GetRequiredService>(), + GetRequiredService>(), + GetRequiredService(), + GetRequiredService()); + } + + private BackgroundJobInfo NewJob(string jobName) + { + return new BackgroundJobInfo + { + Id = Guid.NewGuid(), + JobName = jobName, + JobArgs = "{}", + CreationTime = _clock.Now, + NextTryTime = _clock.Now.AddMinutes(-1) + }; + } + + private PeriodicBackgroundWorkerContext Context() + { + return new PeriodicBackgroundWorkerContext(ServiceProvider); + } + + // Storing successful jobs + + [Fact] + public async Task Should_Keep_Successful_Job_As_History_When_Enabled() + { + _workerOptions.StoreSuccessfulJobs = true; + + var jobInfo = NewJob("job-a"); + await _store.InsertAsync(jobInfo); + + await CreateWorker().HandleJobSuccessPublicAsync(_store, jobInfo, _clock); + + // The job is kept (marked completed), not deleted. + var kept = await _store.FindAsync(jobInfo.Id); + kept.ShouldNotBeNull(); + kept.CompletionTime.ShouldNotBeNull(); + + // ...but it is excluded from the waiting jobs. + (await _store.GetWaitingJobsAsync(null, 1000)).ShouldNotContain(j => j.Id == jobInfo.Id); + } + + [Fact] + public async Task Should_Delete_Successful_Job_When_Disabled() + { + // StoreSuccessfulJobs is false by default. + var jobInfo = NewJob("job-a"); + await _store.InsertAsync(jobInfo); + + await CreateWorker().HandleJobSuccessPublicAsync(_store, jobInfo, _clock); + + (await _store.FindAsync(jobInfo.Id)).ShouldBeNull(); + } + + // Dedicated workers by job name + + [Fact] + public async Task Should_Return_Only_Included_Jobs() + { + await _store.InsertAsync(NewJob("job-a")); + await _store.InsertAsync(NewJob("job-a")); + await _store.InsertAsync(NewJob("job-b")); + + var worker = CreateWorker(); + worker.ConfigureTest(BackgroundJobNameFilter.Include(new[] { "job-a" })); + + var jobs = await worker.GetWaitingJobsPublicAsync(Context(), _store); + + jobs.Count.ShouldBe(2); + jobs.ShouldAllBe(j => j.JobName == "job-a"); + } + + [Fact] + public async Task Should_Exclude_Given_Jobs() + { + await _store.InsertAsync(NewJob("job-a")); + await _store.InsertAsync(NewJob("job-b")); + + var worker = CreateWorker(); + worker.ConfigureTest(BackgroundJobNameFilter.Exclude(new[] { "job-a" })); + + var jobs = await worker.GetWaitingJobsPublicAsync(Context(), _store); + + jobs.Count.ShouldBe(1); + jobs.Single().JobName.ShouldBe("job-b"); + } + + [Fact] + public void Should_Match_Job_Names_By_Filter() + { + BackgroundJobNameFilter.None.IsMatch("any").ShouldBeTrue(); + + var include = BackgroundJobNameFilter.Include(new[] { "job-a" }); + include.IsMatch("job-a").ShouldBeTrue(); + include.IsMatch("job-b").ShouldBeFalse(); + + var exclude = BackgroundJobNameFilter.Exclude(new[] { "job-a" }); + exclude.IsMatch("job-a").ShouldBeFalse(); + exclude.IsMatch("job-b").ShouldBeTrue(); + } + + [Fact] + public void BackgroundJobNameFilter_Should_Reject_Invalid_Mode_And_Names_Combinations() + { + Should.Throw(() => new BackgroundJobNameFilter(BackgroundJobNameFilterMode.Include)); + Should.Throw(() => new BackgroundJobNameFilter(BackgroundJobNameFilterMode.None, new[] { "job-a" })); + Should.Throw(() => new BackgroundJobNameFilter((BackgroundJobNameFilterMode)99, new[] { "job-a" })); + } + + // Parallel execution / eligibility + + [Fact] + public void Should_Evaluate_Job_Eligibility() + { + var worker = CreateWorker(); + + worker.IsJobEligiblePublic(null, _clock).ShouldBeFalse(); + + var future = NewJob("job-a"); + future.NextTryTime = _clock.Now.AddMinutes(5); + worker.IsJobEligiblePublic(future, _clock).ShouldBeFalse(); + + var abandoned = NewJob("job-a"); + abandoned.IsAbandoned = true; + worker.IsJobEligiblePublic(abandoned, _clock).ShouldBeFalse(); + + var completed = NewJob("job-a"); + completed.CompletionTime = _clock.Now; + worker.IsJobEligiblePublic(completed, _clock).ShouldBeFalse(); + + var eligible = NewJob("job-a"); + worker.IsJobEligiblePublic(eligible, _clock).ShouldBeTrue(); + } + + [Fact] + public void Should_Not_Be_Eligible_When_Job_Name_Filtered_Out() + { + var worker = CreateWorker(); + worker.ConfigureTest(BackgroundJobNameFilter.Include(new[] { "job-a" })); + + var otherJob = NewJob("job-b"); + worker.IsJobEligiblePublic(otherJob, _clock).ShouldBeFalse(); + } + + [Fact] + public void Should_Require_Job_Args_Types_For_A_Dedicated_Worker() + { + Should.Throw(() => new BackgroundJobWorkerConfiguration("lock-a")); + } + + [Fact] + public void AddDedicatedWorker_Should_Throw_At_Registration_When_A_Job_Type_Is_Added_Twice() + { + var options = new AbpBackgroundJobWorkerOptions(); + options.AddDedicatedWorker("lock-a"); + + Should.Throw(() => options.AddDedicatedWorker("lock-b")); + } + + [Fact] + public void AddDedicatedWorker_Should_Throw_At_Registration_When_A_Lock_Name_Is_Reused() + { + var options = new AbpBackgroundJobWorkerOptions(); + options.AddDedicatedWorker("dup-lock"); + + Should.Throw(() => options.AddDedicatedWorker("dup-lock")); + } + + [Fact] + public void AddDedicatedWorker_Should_Throw_At_Registration_When_Lock_Name_Equals_The_Default() + { + var options = new AbpBackgroundJobWorkerOptions(); + + Should.Throw(() => options.AddDedicatedWorker(options.DistributedLockName)); + } + + // Parallel execution + + [Fact] + public async Task Should_Execute_Multiple_Jobs_In_Parallel() + { + _workerOptions.MaxParallelJobExecutionCount = 3; + + var jobManager = GetRequiredService(); + await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "1" }); + await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "2" }); + await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "3" }); + + await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); + + var tracker = GetRequiredService(); + tracker.Executed.Count.ShouldBe(3); + tracker.Executed.ShouldContain("1"); + tracker.Executed.ShouldContain("2"); + tracker.Executed.ShouldContain("3"); + + // Each job must run in its own service scope (isolated DbContext/UOW). + tracker.ScopeIds.Distinct().Count().ShouldBe(3); + + (await _store.GetWaitingJobsAsync(null, 1000)).ShouldBeEmpty(); + } + + [Fact] + public async Task Should_Execute_At_Most_MaxParallel_Jobs_Per_Cycle() + { + _workerOptions.MaxParallelJobExecutionCount = 2; + + var jobManager = GetRequiredService(); + for (var i = 0; i < 5; i++) + { + await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = i.ToString() }); + } + + await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); + + var tracker = GetRequiredService(); + tracker.Executed.Count.ShouldBe(2); + (await _store.GetWaitingJobsAsync(null, 1000)).Count.ShouldBe(3); + } + + [Fact] + public async Task Should_Skip_Job_Already_Claimed_By_Another_Instance() + { + _workerOptions.MaxParallelJobExecutionCount = 5; + + var jobManager = GetRequiredService(); + var lockedJobId = Guid.Parse(await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "locked" })); + await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "free" }); + + var distributedLock = GetRequiredService(); + var lockName = _workerOptions.PerJobDistributedLockPrefix + lockedJobId; + + await using (await distributedLock.TryAcquireAsync(lockName)) + { + await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); + } + + var tracker = GetRequiredService(); + tracker.Executed.ShouldContain("free"); + tracker.Executed.ShouldNotContain("locked"); + (await _store.FindAsync(lockedJobId)).ShouldNotBeNull(); + } + + [Fact] + public async Task Should_Mark_Job_Completed_On_Successful_Execution_When_Storing_Enabled() + { + _workerOptions.StoreSuccessfulJobs = true; + _workerOptions.MaxParallelJobExecutionCount = 2; + + var jobManager = GetRequiredService(); + var jobId = Guid.Parse(await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "1" })); + + await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); + + GetRequiredService().Executed.ShouldContain("1"); + + // The job is kept (marked completed), not deleted, and excluded from the waiting list. + var job = await _store.FindAsync(jobId); + job.ShouldNotBeNull(); + job.CompletionTime.ShouldNotBeNull(); + (await _store.GetWaitingJobsAsync(null, 1000)).ShouldNotContain(j => j.Id == jobId); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/RecordingBackgroundJobWorker.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/RecordingBackgroundJobWorker.cs new file mode 100644 index 0000000000..6acb5ca78d --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/RecordingBackgroundJobWorker.cs @@ -0,0 +1,51 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// Records every call so multi-worker registration +/// (resolved from DI by ) can be asserted. Does not start any timer. +/// +public class WorkerStartRecorder +{ + public List Records { get; } = new List(); +} + +public class WorkerStartRecord +{ + public string? DistributedLockName { get; set; } + + public BackgroundJobNameFilter? JobNameFilter { get; set; } +} + +public class RecordingBackgroundJobWorker : IBackgroundJobWorker, ITransientDependency +{ + private readonly WorkerStartRecorder _recorder; + + public RecordingBackgroundJobWorker(WorkerStartRecorder recorder) + { + _recorder = recorder; + } + + public Task StartAsync( + string? distributedLockName = null, + BackgroundJobNameFilter? jobNameFilter = null, + CancellationToken cancellationToken = default) + { + _recorder.Records.Add(new WorkerStartRecord + { + DistributedLockName = distributedLockName, + JobNameFilter = jobNameFilter + }); + + return Task.CompletedTask; + } + + public Task StopAsync(CancellationToken cancellationToken = default) + { + return Task.CompletedTask; + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobCleanupWorker.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobCleanupWorker.cs new file mode 100644 index 0000000000..7c6cc16dec --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobCleanupWorker.cs @@ -0,0 +1,26 @@ +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Volo.Abp.BackgroundWorkers; +using Volo.Abp.DistributedLocking; +using Volo.Abp.Threading; + +namespace Volo.Abp.BackgroundJobs; + +public class TestableBackgroundJobCleanupWorker : BackgroundJobCleanupWorker +{ + public TestableBackgroundJobCleanupWorker( + AbpAsyncTimer timer, + IServiceScopeFactory serviceScopeFactory, + IOptions jobOptions, + IOptions workerOptions, + IAbpDistributedLock distributedLock) + : base(timer, serviceScopeFactory, jobOptions, workerOptions, distributedLock) + { + } + + public Task DoWorkPublicAsync(PeriodicBackgroundWorkerContext workerContext) + { + return DoWorkAsync(workerContext); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobWorker.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobWorker.cs new file mode 100644 index 0000000000..27274320fb --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/TestableBackgroundJobWorker.cs @@ -0,0 +1,55 @@ +#nullable enable +using System.Collections.Generic; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; +using Volo.Abp.BackgroundWorkers; +using Volo.Abp.DistributedLocking; +using Volo.Abp.Threading; +using Volo.Abp.Timing; + +namespace Volo.Abp.BackgroundJobs; + +/// +/// Exposes the protected members of for unit testing. +/// +public class TestableBackgroundJobWorker : BackgroundJobWorker +{ + public TestableBackgroundJobWorker( + AbpAsyncTimer timer, + IOptions jobOptions, + IOptions workerOptions, + IServiceScopeFactory serviceScopeFactory, + IAbpDistributedLock distributedLock) + : base(timer, jobOptions, workerOptions, serviceScopeFactory, distributedLock) + { + } + + public void ConfigureTest( + BackgroundJobNameFilter? jobNameFilter = null, + string? distributedLockName = null) + { + JobNameFilter = jobNameFilter ?? BackgroundJobNameFilter.None; + DistributedLockName = distributedLockName ?? WorkerOptions.DistributedLockName; + } + + public Task HandleJobSuccessPublicAsync(IBackgroundJobStore store, BackgroundJobInfo jobInfo, IClock clock) + { + return HandleJobSuccessAsync(store, jobInfo, clock); + } + + public Task ExecuteJobsInParallelPublicAsync(PeriodicBackgroundWorkerContext workerContext) + { + return ExecuteJobsInParallelAsync(workerContext); + } + + public Task> GetWaitingJobsPublicAsync(PeriodicBackgroundWorkerContext workerContext, IBackgroundJobStore store) + { + return GetWaitingJobsAsync(workerContext, store); + } + + public bool IsJobEligiblePublic(BackgroundJobInfo? jobInfo, IClock clock) + { + return IsJobEligible(jobInfo, clock); + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/WorkerTestJobs.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/WorkerTestJobs.cs new file mode 100644 index 0000000000..dd9feec657 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/WorkerTestJobs.cs @@ -0,0 +1,99 @@ +using System; +using System.Collections.Concurrent; +using System.Threading.Tasks; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.BackgroundJobs; + +public class ParallelJobTracker +{ + public ConcurrentBag Executed { get; } = new ConcurrentBag(); + + public ConcurrentBag ScopeIds { get; } = new ConcurrentBag(); +} + +/// +/// Scoped service used to verify that each parallel job runs in its own service scope. +/// +public class ScopeMarker +{ + public Guid Id { get; } = Guid.NewGuid(); +} + +public class ParallelTestJobArgs +{ + public string Value { get; set; } = default!; +} + +public class ParallelTestJob : AsyncBackgroundJob, ITransientDependency +{ + private readonly ParallelJobTracker _tracker; + private readonly ScopeMarker _scopeMarker; + + public ParallelTestJob(ParallelJobTracker tracker, ScopeMarker scopeMarker) + { + _tracker = tracker; + _scopeMarker = scopeMarker; + } + + public override Task ExecuteAsync(ParallelTestJobArgs args) + { + _tracker.Executed.Add(args.Value); + _tracker.ScopeIds.Add(_scopeMarker.Id); + return Task.CompletedTask; + } +} + +public class WorkerJobAArgs +{ + public string Value { get; set; } = default!; +} + +public class WorkerJobA : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(WorkerJobAArgs args) + { + return Task.CompletedTask; + } +} + +public class WorkerJobBArgs +{ + public string Value { get; set; } = default!; +} + +public class WorkerJobB : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(WorkerJobBArgs args) + { + return Task.CompletedTask; + } +} + +// Two different args types that resolve to the same job name, to exercise the manager's +// backstop validation (eager AddDedicatedWorker validation compares by type, not resolved name). +[BackgroundJobName("shared-job-name")] +public class SharedNameJobAArgs +{ +} + +[BackgroundJobName("shared-job-name")] +public class SharedNameJobBArgs +{ +} + +public class SharedNameJobA : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(SharedNameJobAArgs args) + { + return Task.CompletedTask; + } +} + +public class SharedNameJobB : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(SharedNameJobBArgs args) + { + return Task.CompletedTask; + } +} 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 976705b07d..5569f6800a 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/DemoAppModule.cs +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/DemoAppModule.cs @@ -1,5 +1,7 @@ -using System.Threading.Tasks; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; using Volo.Abp.Autofac; +using Volo.Abp.BackgroundJobs.DemoApp.Jobs; using Volo.Abp.BackgroundJobs.DemoApp.Shared; using Volo.Abp.BackgroundJobs.EntityFrameworkCore; using Volo.Abp.EntityFrameworkCore; @@ -32,17 +34,31 @@ public class DemoAppModule : AbpModule options.JobPollPeriod = 1000; options.DefaultFirstWaitDuration = 1; options.DefaultWaitFactor = 1; + + // Keep every successfully completed job as history (marks CompletionTime instead of deleting). + // Completed jobs are excluded from the waiting query and pruned after SuccessfulJobRetentionTime. + options.StoreSuccessfulJobs = true; + options.SuccessfulJobRetentionTime = System.TimeSpan.FromDays(1); + + // A dedicated worker (with its own distributed lock "DemoFeesWorkerLock") that only processes + // the slow fee-calculation jobs, so they don't block other jobs. A default worker is added automatically + // and processes all the remaining job types (e.g. SendEmailJob). + options.AddDedicatedWorker("DemoFeesWorkerLock"); + + // Let each worker execute up to 4 jobs in parallel (each job claimed with its own distributed lock, + // so multiple application instances can execute different jobs concurrently). + options.MaxParallelJobExecutionCount = 4; }); } - public override Task OnApplicationInitializationAsync(ApplicationInitializationContext context) + public override async Task OnApplicationInitializationAsync(ApplicationInitializationContext context) { - //TODO: Configure console logging - //context - // .ServiceProvider - // .GetRequiredService() - // .AddConsole(LogLevel.Debug); + // Enqueue a few demo jobs. The fee-calculation jobs are handled by the dedicated worker, + // while SendEmailJob is handled by the default worker. + var backgroundJobManager = context.ServiceProvider.GetRequiredService(); - return Task.CompletedTask; + await backgroundJobManager.EnqueueAsync(new CalculateAwsFeesJobArgs { AccountId = "acc-1" }); + await backgroundJobManager.EnqueueAsync(new CalculateAzureFeesJobArgs { SubscriptionId = "sub-1" }); + await backgroundJobManager.EnqueueAsync(new SendEmailJobArgs { To = "user@example.com", Subject = "Welcome" }); } } diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAwsFeesJob.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAwsFeesJob.cs new file mode 100644 index 0000000000..2262bbc68d --- /dev/null +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAwsFeesJob.cs @@ -0,0 +1,20 @@ +using System.Threading.Tasks; +using Microsoft.Extensions.Logging; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.BackgroundJobs.DemoApp.Jobs; + +public class CalculateAwsFeesJobArgs +{ + public string AccountId { get; set; } = default!; +} + +public class CalculateAwsFeesJob : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(CalculateAwsFeesJobArgs args) + { + // A slow, resource-intensive job that is isolated on a dedicated worker (see DemoAppModule). + Logger.LogInformation($"[AWS fees] Calculating fees for account '{args.AccountId}'..."); + return Task.CompletedTask; + } +} diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAzureFeesJob.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAzureFeesJob.cs new file mode 100644 index 0000000000..d0dc43cd94 --- /dev/null +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/CalculateAzureFeesJob.cs @@ -0,0 +1,20 @@ +using System.Threading.Tasks; +using Microsoft.Extensions.Logging; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.BackgroundJobs.DemoApp.Jobs; + +public class CalculateAzureFeesJobArgs +{ + public string SubscriptionId { get; set; } = default!; +} + +public class CalculateAzureFeesJob : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(CalculateAzureFeesJobArgs args) + { + // Isolated on the same dedicated fees worker as CalculateAwsFeesJob. + Logger.LogInformation($"[Azure fees] Calculating fees for subscription '{args.SubscriptionId}'..."); + return Task.CompletedTask; + } +} diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/SendEmailJob.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/SendEmailJob.cs new file mode 100644 index 0000000000..846b06f4fe --- /dev/null +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Jobs/SendEmailJob.cs @@ -0,0 +1,23 @@ +using System.Threading.Tasks; +using Microsoft.Extensions.Logging; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.BackgroundJobs.DemoApp.Jobs; + +public class SendEmailJobArgs +{ + public string To { get; set; } = default!; + + public string Subject { get; set; } = default!; +} + +public class SendEmailJob : AsyncBackgroundJob, ITransientDependency +{ + public override Task ExecuteAsync(SendEmailJobArgs args) + { + // A fast job. It is NOT configured for a dedicated worker, so the default worker processes it + // without waiting behind the slow fee-calculation jobs. + Logger.LogInformation($"[Email] Sending '{args.Subject}' to '{args.To}'..."); + return Task.CompletedTask; + } +} diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.Designer.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.Designer.cs new file mode 100644 index 0000000000..01241ef94a --- /dev/null +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.Designer.cs @@ -0,0 +1,99 @@ +// +using System; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Volo.Abp.BackgroundJobs.DemoApp.Db; +using Volo.Abp.EntityFrameworkCore; + +#nullable disable + +namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations +{ + [DbContext(typeof(DemoAppDbContext))] + [Migration("20260701082002_Added_CompletionTime_To_BackgroundJobs")] + partial class Added_CompletionTime_To_BackgroundJobs + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("_Abp_DatabaseProvider", EfCoreDatabaseProvider.SqlServer) + .HasAnnotation("ProductVersion", "10.0.9") + .HasAnnotation("Relational:MaxIdentifierLength", 128); + + SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + + modelBuilder.Entity("Volo.Abp.BackgroundJobs.BackgroundJobRecord", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uniqueidentifier"); + + b.Property("ApplicationName") + .HasMaxLength(96) + .HasColumnType("nvarchar(96)"); + + b.Property("CompletionTime") + .HasColumnType("datetime2"); + + b.Property("ConcurrencyStamp") + .IsConcurrencyToken() + .IsRequired() + .HasMaxLength(40) + .HasColumnType("nvarchar(40)") + .HasColumnName("ConcurrencyStamp"); + + b.Property("CreationTime") + .HasColumnType("datetime2") + .HasColumnName("CreationTime"); + + b.Property("ExtraProperties") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("ExtraProperties"); + + b.Property("IsAbandoned") + .ValueGeneratedOnAdd() + .HasColumnType("bit") + .HasDefaultValue(false); + + b.Property("JobArgs") + .IsRequired() + .HasMaxLength(1048576) + .HasColumnType("nvarchar(max)"); + + b.Property("JobName") + .IsRequired() + .HasMaxLength(128) + .HasColumnType("nvarchar(128)"); + + b.Property("LastTryTime") + .HasColumnType("datetime2"); + + b.Property("NextTryTime") + .HasColumnType("datetime2"); + + b.Property("Priority") + .ValueGeneratedOnAdd() + .HasColumnType("tinyint") + .HasDefaultValue((byte)15); + + b.Property("TryCount") + .ValueGeneratedOnAdd() + .HasColumnType("smallint") + .HasDefaultValue((short)0); + + b.HasKey("Id"); + + b.HasIndex("ApplicationName", "CompletionTime", "IsAbandoned", "NextTryTime"); + + b.ToTable("AbpBackgroundJobs", (string)null); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.cs new file mode 100644 index 0000000000..96f35f3379 --- /dev/null +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260701082002_Added_CompletionTime_To_BackgroundJobs.cs @@ -0,0 +1,47 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations +{ + /// + public partial class Added_CompletionTime_To_BackgroundJobs : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropIndex( + name: "IX_AbpBackgroundJobs_IsAbandoned_NextTryTime", + table: "AbpBackgroundJobs"); + + migrationBuilder.AddColumn( + name: "CompletionTime", + table: "AbpBackgroundJobs", + type: "datetime2", + nullable: true); + + migrationBuilder.CreateIndex( + name: "IX_AbpBackgroundJobs_ApplicationName_CompletionTime_IsAbandoned_NextTryTime", + table: "AbpBackgroundJobs", + columns: new[] { "ApplicationName", "CompletionTime", "IsAbandoned", "NextTryTime" }); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropIndex( + name: "IX_AbpBackgroundJobs_ApplicationName_CompletionTime_IsAbandoned_NextTryTime", + table: "AbpBackgroundJobs"); + + migrationBuilder.DropColumn( + name: "CompletionTime", + table: "AbpBackgroundJobs"); + + migrationBuilder.CreateIndex( + name: "IX_AbpBackgroundJobs_IsAbandoned_NextTryTime", + table: "AbpBackgroundJobs", + columns: new[] { "IsAbandoned", "NextTryTime" }); + } + } +} diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/DemoAppDbContextModelSnapshot.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/DemoAppDbContextModelSnapshot.cs index ab91bc354c..af43eb368f 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/DemoAppDbContextModelSnapshot.cs +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/DemoAppDbContextModelSnapshot.cs @@ -19,7 +19,7 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations #pragma warning disable 612, 618 modelBuilder .HasAnnotation("_Abp_DatabaseProvider", EfCoreDatabaseProvider.SqlServer) - .HasAnnotation("ProductVersion", "10.0.2") + .HasAnnotation("ProductVersion", "10.0.9") .HasAnnotation("Relational:MaxIdentifierLength", 128); SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); @@ -34,6 +34,9 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations .HasMaxLength(96) .HasColumnType("nvarchar(96)"); + b.Property("CompletionTime") + .HasColumnType("datetime2"); + b.Property("ConcurrencyStamp") .IsConcurrencyToken() .IsRequired() @@ -83,7 +86,7 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations b.HasKey("Id"); - b.HasIndex("IsAbandoned", "NextTryTime"); + b.HasIndex("ApplicationName", "CompletionTime", "IsAbandoned", "NextTryTime"); b.ToTable("AbpBackgroundJobs", (string)null); }); diff --git a/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobRepository_Tests.cs b/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobRepository_Tests.cs index c62cb43309..b717d28d15 100644 --- a/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobRepository_Tests.cs +++ b/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobRepository_Tests.cs @@ -1,4 +1,6 @@ -using System.Linq; +using System; +using System.Collections.Generic; +using System.Linq; using System.Threading.Tasks; using Shouldly; using Volo.Abp.Modularity; @@ -36,4 +38,99 @@ public abstract class BackgroundJobRepository_Tests : Background backgroundJobs.Any(j => j.ApplicationName == "App2").ShouldBeFalse(); backgroundJobs.Any(j => j.ApplicationName == null).ShouldBeFalse(); } + + [Fact] + public async Task Should_Filter_Waiting_List_By_Included_Job_Names() + { + // App1 waiting jobs: two "TestJobName" + one "OtherJobName". + var testJobs = await _backgroundJobRepository.GetWaitingListAsync("App1", 10, BackgroundJobNameFilter.Include(new[] { "TestJobName" })); + testJobs.Count.ShouldBe(2); + testJobs.ShouldAllBe(j => j.JobName == "TestJobName"); + + var otherJobs = await _backgroundJobRepository.GetWaitingListAsync("App1", 10, BackgroundJobNameFilter.Include(new[] { "OtherJobName" })); + otherJobs.Count.ShouldBe(1); + otherJobs.Single().JobName.ShouldBe("OtherJobName"); + + var none = await _backgroundJobRepository.GetWaitingListAsync("App1", 10, BackgroundJobNameFilter.Include(new[] { "NonExistentJob" })); + none.ShouldBeEmpty(); + } + + [Fact] + public async Task Should_Filter_Waiting_List_By_Excluded_Job_Names() + { + var withoutTestJob = await _backgroundJobRepository.GetWaitingListAsync("App1", 10, BackgroundJobNameFilter.Exclude(new[] { "TestJobName" })); + withoutTestJob.Count.ShouldBe(1); + withoutTestJob.Single().JobName.ShouldBe("OtherJobName"); + + var withoutNonExistent = await _backgroundJobRepository.GetWaitingListAsync("App1", 10, BackgroundJobNameFilter.Exclude(new[] { "NonExistentJob" })); + withoutNonExistent.Count.ShouldBe(3); + } + + [Fact] + public async Task Should_Return_Null_From_Store_When_Job_Not_Found() + { + // The parallel worker re-reads a claimed job under the lock; a job removed by another instance must return null (not throw). + var store = GetRequiredService(); + var found = await store.FindAsync(Guid.NewGuid()); + found.ShouldBeNull(); + } + + [Fact] + public async Task Should_Exclude_Completed_Jobs_From_Waiting_List() + { + var completedJobId = Guid.NewGuid(); + await _backgroundJobRepository.InsertAsync( + new BackgroundJobRecord(completedJobId) + { + ApplicationName = "App1", + JobName = "TestJobName", + JobArgs = "{ value: 1 }", + NextTryTime = _clock.Now.Subtract(TimeSpan.FromMinutes(1)), + Priority = BackgroundJobPriority.Normal, + IsAbandoned = false, + CompletionTime = _clock.Now, + CreationTime = _clock.Now.Subtract(TimeSpan.FromMinutes(2)), + TryCount = 1 + }, + autoSave: true); + + var waitingJobs = await _backgroundJobRepository.GetWaitingListAsync("App1", 100); + waitingJobs.ShouldNotContain(j => j.Id == completedJobId); + } + + [Fact] + public async Task Should_Delete_Old_Successful_Jobs_Of_The_Given_Application_Only() + { + var oldApp1Id = Guid.NewGuid(); + await _backgroundJobRepository.InsertAsync(NewCompletedJob(oldApp1Id, "App1", _clock.Now.Subtract(TimeSpan.FromDays(2))), autoSave: true); + + var recentApp1Id = Guid.NewGuid(); + await _backgroundJobRepository.InsertAsync(NewCompletedJob(recentApp1Id, "App1", _clock.Now), autoSave: true); + + var oldApp2Id = Guid.NewGuid(); + await _backgroundJobRepository.InsertAsync(NewCompletedJob(oldApp2Id, "App2", _clock.Now.Subtract(TimeSpan.FromDays(2))), autoSave: true); + + var deleted = await _backgroundJobRepository.DeleteAsync("App1", _clock.Now.Subtract(TimeSpan.FromDays(1)), 100); + deleted.ShouldBe(1); + + (await _backgroundJobRepository.FindAsync(oldApp1Id)).ShouldBeNull(); // App1 old completed → deleted + (await _backgroundJobRepository.FindAsync(recentApp1Id)).ShouldNotBeNull(); // App1 recent → kept + (await _backgroundJobRepository.FindAsync(oldApp2Id)).ShouldNotBeNull(); // App2 old → not touched (isolation) + } + + private BackgroundJobRecord NewCompletedJob(Guid id, string applicationName, DateTime completionTime) + { + return new BackgroundJobRecord(id) + { + ApplicationName = applicationName, + JobName = "TestJobName", + JobArgs = "{ value: 1 }", + NextTryTime = _clock.Now.Subtract(TimeSpan.FromMinutes(1)), + Priority = BackgroundJobPriority.Normal, + IsAbandoned = false, + CompletionTime = completionTime, + CreationTime = _clock.Now.Subtract(TimeSpan.FromMinutes(5)), + TryCount = 1 + }; + } } diff --git a/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestData.cs b/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestData.cs index 0d8b2bfd29..cd01e1b09a 100644 --- a/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestData.cs +++ b/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestData.cs @@ -8,4 +8,5 @@ public class BackgroundJobsTestData : ISingletonDependency public Guid JobId1 { get; } = Guid.NewGuid(); public Guid JobId2 { get; } = Guid.NewGuid(); public Guid JobId3 { get; } = Guid.NewGuid(); + public Guid JobId4 { get; } = Guid.NewGuid(); } diff --git a/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestDataBuilder.cs b/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestDataBuilder.cs index f6f4c37e12..47ebf103e2 100644 --- a/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestDataBuilder.cs +++ b/modules/background-jobs/test/Volo.Abp.BackgroundJobs.TestBase/Volo/Abp/BackgroundJobs/BackgroundJobsTestDataBuilder.cs @@ -67,5 +67,21 @@ public class BackgroundJobsTestDataBuilder : ITransientDependency TryCount = 2 } ); + + // App1 waiting job with a different job name, to verify job-name filtering. + await _backgroundJobRepository.InsertAsync( + new BackgroundJobRecord(_testData.JobId4) + { + ApplicationName = "App1", + JobName = "OtherJobName", + JobArgs = "{ value: 4 }", + NextTryTime = _clock.Now.Subtract(TimeSpan.FromMinutes(1)), + Priority = BackgroundJobPriority.Normal, + IsAbandoned = false, + LastTryTime = null, + CreationTime = _clock.Now.Subtract(TimeSpan.FromMinutes(3)), + TryCount = 0 + } + ); } }