mirror of https://github.com/abpframework/abp.git
Browse Source
Add dedicated workers, parallel execution and successful job retentionpull/25754/head
committed by
GitHub
49 changed files with 2499 additions and 104 deletions
@ -0,0 +1,66 @@ |
|||||
|
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; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Periodically deletes retained successfully completed jobs older than
|
||||
|
/// <see cref="AbpBackgroundJobWorkerOptions.SuccessfulJobRetentionTime"/>.
|
||||
|
/// Only relevant when <see cref="AbpBackgroundJobWorkerOptions.StoreSuccessfulJobs"/> is enabled.
|
||||
|
/// </summary>
|
||||
|
public class BackgroundJobCleanupWorker : AsyncPeriodicBackgroundWorkerBase |
||||
|
{ |
||||
|
protected AbpBackgroundJobOptions JobOptions { get; } |
||||
|
|
||||
|
protected AbpBackgroundJobWorkerOptions WorkerOptions { get; } |
||||
|
|
||||
|
protected IAbpDistributedLock DistributedLock { get; } |
||||
|
|
||||
|
public BackgroundJobCleanupWorker( |
||||
|
AbpAsyncTimer timer, |
||||
|
IServiceScopeFactory serviceScopeFactory, |
||||
|
IOptions<AbpBackgroundJobOptions> jobOptions, |
||||
|
IOptions<AbpBackgroundJobWorkerOptions> workerOptions, |
||||
|
IAbpDistributedLock distributedLock) |
||||
|
: base(timer, serviceScopeFactory) |
||||
|
{ |
||||
|
JobOptions = jobOptions.Value; |
||||
|
WorkerOptions = workerOptions.Value; |
||||
|
DistributedLock = distributedLock; |
||||
|
Timer.Period = WorkerOptions.CleanSuccessfulJobsPeriod; |
||||
|
} |
||||
|
|
||||
|
protected override async Task DoWorkAsync(PeriodicBackgroundWorkerContext workerContext) |
||||
|
{ |
||||
|
if (!JobOptions.IsJobExecutionEnabled || |
||||
|
!WorkerOptions.StoreSuccessfulJobs || |
||||
|
WorkerOptions.SuccessfulJobRetentionTime == null) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
var store = workerContext.ServiceProvider.GetRequiredService<IBackgroundJobStore>(); |
||||
|
var clock = workerContext.ServiceProvider.GetRequiredService<IClock>(); |
||||
|
var completedBefore = clock.Now.Subtract(WorkerOptions.SuccessfulJobRetentionTime.Value); |
||||
|
|
||||
|
await using (var handle = await DistributedLock.TryAcquireAsync(WorkerOptions.CleanupDistributedLockName, cancellationToken: StoppingToken)) |
||||
|
{ |
||||
|
if (handle == null) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
int deletedCount; |
||||
|
do |
||||
|
{ |
||||
|
deletedCount = await store.DeleteAsync(WorkerOptions.ApplicationName, completedBefore, WorkerOptions.MaxJobFetchCount, StoppingToken); |
||||
|
} |
||||
|
while (deletedCount > 0 && deletedCount >= WorkerOptions.MaxJobFetchCount && !StoppingToken.IsCancellationRequested); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,71 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Linq; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Filters the waiting jobs of a background job worker by job name.
|
||||
|
/// A worker is exactly one of: no filter (<see cref="None"/>), include-only (a dedicated worker) or
|
||||
|
/// exclude-only (the default worker in a multi-worker setup) — the two can never be combined.
|
||||
|
/// </summary>
|
||||
|
public class BackgroundJobNameFilter |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// A filter that matches every job name.
|
||||
|
/// </summary>
|
||||
|
public static BackgroundJobNameFilter None { get; } = new(BackgroundJobNameFilterMode.None); |
||||
|
|
||||
|
public BackgroundJobNameFilterMode Mode { get; } |
||||
|
|
||||
|
public IReadOnlyList<string> JobNames { get; } |
||||
|
|
||||
|
public BackgroundJobNameFilter(BackgroundJobNameFilterMode mode, IReadOnlyList<string>? jobNames = null) |
||||
|
{ |
||||
|
if (!Enum.IsDefined(typeof(BackgroundJobNameFilterMode), mode)) |
||||
|
{ |
||||
|
throw new ArgumentException($"Invalid background job name filter mode: {mode}", nameof(mode)); |
||||
|
} |
||||
|
|
||||
|
var names = jobNames?.Where(x => !x.IsNullOrWhiteSpace()).Distinct(StringComparer.Ordinal).ToList() ?? new List<string>(); |
||||
|
|
||||
|
if (mode == BackgroundJobNameFilterMode.None && names.Count > 0) |
||||
|
{ |
||||
|
throw new ArgumentException("Job names must be empty when the filter mode is None.", nameof(jobNames)); |
||||
|
} |
||||
|
|
||||
|
if (mode != BackgroundJobNameFilterMode.None && names.Count == 0) |
||||
|
{ |
||||
|
throw new ArgumentException("Job names cannot be empty when the filter mode is Include or Exclude.", nameof(jobNames)); |
||||
|
} |
||||
|
|
||||
|
Mode = mode; |
||||
|
JobNames = names.AsReadOnly(); |
||||
|
} |
||||
|
|
||||
|
public static BackgroundJobNameFilter Include(IReadOnlyList<string> jobNames) |
||||
|
{ |
||||
|
return new BackgroundJobNameFilter(BackgroundJobNameFilterMode.Include, jobNames); |
||||
|
} |
||||
|
|
||||
|
public static BackgroundJobNameFilter Exclude(IReadOnlyList<string> jobNames) |
||||
|
{ |
||||
|
return new BackgroundJobNameFilter(BackgroundJobNameFilterMode.Exclude, jobNames); |
||||
|
} |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Whether the given job name passes this filter, using an ordinal (case-sensitive) comparison for the
|
||||
|
/// in-memory eligibility re-check. The persistent stores translate <see cref="Mode"/> and
|
||||
|
/// <see cref="JobNames"/> into a database query instead, so their filtering follows the database collation.
|
||||
|
/// Job names are expected to be unique beyond case (they are derived from the type name by default).
|
||||
|
/// </summary>
|
||||
|
public virtual bool IsMatch(string jobName) |
||||
|
{ |
||||
|
return Mode switch |
||||
|
{ |
||||
|
BackgroundJobNameFilterMode.Include => JobNames.Contains(jobName, StringComparer.Ordinal), |
||||
|
BackgroundJobNameFilterMode.Exclude => !JobNames.Contains(jobName, StringComparer.Ordinal), |
||||
|
_ => true |
||||
|
}; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,19 @@ |
|||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
public enum BackgroundJobNameFilterMode : byte |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// No filter; all job names match.
|
||||
|
/// </summary>
|
||||
|
None = 0, |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Only the job names in the filter match.
|
||||
|
/// </summary>
|
||||
|
Include = 1, |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// All job names except those in the filter match.
|
||||
|
/// </summary>
|
||||
|
Exclude = 2 |
||||
|
} |
||||
@ -0,0 +1,37 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Linq; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Configuration of a dedicated <see cref="BackgroundJobWorker"/> that processes only specific job types.
|
||||
|
/// </summary>
|
||||
|
public class BackgroundJobWorkerConfiguration |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// A unique distributed lock name for this worker. It must be different from the names used by other workers.
|
||||
|
/// It is used to serialize the worker across application instances when
|
||||
|
/// <see cref="AbpBackgroundJobWorkerOptions.MaxParallelJobExecutionCount"/> is 1; in parallel mode
|
||||
|
/// (greater than 1) jobs are claimed with per-job locks instead and this lock is not acquired.
|
||||
|
/// </summary>
|
||||
|
public string LockName { get; } |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// The job argument types that are processed exclusively by this worker.
|
||||
|
/// </summary>
|
||||
|
public IReadOnlyList<Type> JobArgsTypes { get; } |
||||
|
|
||||
|
public BackgroundJobWorkerConfiguration(string lockName, params Type[] jobArgsTypes) |
||||
|
{ |
||||
|
LockName = Check.NotNullOrWhiteSpace(lockName, nameof(lockName)); |
||||
|
Check.NotNullOrEmpty(jobArgsTypes, nameof(jobArgsTypes)); |
||||
|
|
||||
|
if (jobArgsTypes.Any(t => t == null)) |
||||
|
{ |
||||
|
throw new ArgumentException("Job args types cannot contain null.", nameof(jobArgsTypes)); |
||||
|
} |
||||
|
|
||||
|
JobArgsTypes = jobArgsTypes.ToList(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,131 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Linq; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Owns and controls the background job workers.
|
||||
|
/// When no <see cref="AbpBackgroundJobWorkerOptions.WorkerConfigurations"/> is configured, a single
|
||||
|
/// default worker processes all jobs. Otherwise, one dedicated worker is started per configuration
|
||||
|
/// (each with its own distributed lock and job-type filter) plus a default worker for the remaining jobs.
|
||||
|
/// The workers are resolved from DI, so a replaced <see cref="IBackgroundJobWorker"/> is respected.
|
||||
|
/// </summary>
|
||||
|
public class BackgroundJobWorkerManager : IBackgroundWorker |
||||
|
{ |
||||
|
protected AbpBackgroundJobOptions JobOptions { get; } |
||||
|
|
||||
|
protected AbpBackgroundJobWorkerOptions WorkerOptions { get; } |
||||
|
|
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
|
||||
|
protected List<IBackgroundJobWorker> Workers { get; } |
||||
|
|
||||
|
public BackgroundJobWorkerManager( |
||||
|
IOptions<AbpBackgroundJobOptions> jobOptions, |
||||
|
IOptions<AbpBackgroundJobWorkerOptions> workerOptions, |
||||
|
IServiceProvider serviceProvider) |
||||
|
{ |
||||
|
JobOptions = jobOptions.Value; |
||||
|
WorkerOptions = workerOptions.Value; |
||||
|
ServiceProvider = serviceProvider; |
||||
|
Workers = new List<IBackgroundJobWorker>(); |
||||
|
} |
||||
|
|
||||
|
public virtual async Task StartAsync(CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
if (!JobOptions.IsJobExecutionEnabled) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
if (!WorkerOptions.WorkerConfigurations.Any()) |
||||
|
{ |
||||
|
await StartWorkerAsync(cancellationToken: cancellationToken); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
// AddDedicatedWorker already rejects duplicate job types and lock names eagerly. This is the backstop
|
||||
|
// for what can only be known here: two different args types that resolve to the same job name.
|
||||
|
// Validate all configurations first, so a misconfiguration does not leave already-started workers running.
|
||||
|
var dedicatedWorkers = new List<DedicatedWorkerDefinition>(); |
||||
|
var allDedicatedJobNames = new List<string>(); |
||||
|
|
||||
|
// The default worker uses WorkerOptions.DistributedLockName, so dedicated workers must not reuse it.
|
||||
|
var lockNames = new List<string> { WorkerOptions.DistributedLockName }; |
||||
|
|
||||
|
foreach (var configuration in WorkerOptions.WorkerConfigurations) |
||||
|
{ |
||||
|
var jobNames = configuration.JobArgsTypes |
||||
|
.Select(GetJobName) |
||||
|
.Distinct() |
||||
|
.ToList(); |
||||
|
|
||||
|
var alreadyConfigured = jobNames.Intersect(allDedicatedJobNames).ToList(); |
||||
|
if (alreadyConfigured.Any()) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
$"The following background job(s) are configured for more than one dedicated worker: {string.Join(", ", alreadyConfigured)}. " + |
||||
|
$"Each job type can be handled by only one dedicated worker."); |
||||
|
} |
||||
|
|
||||
|
if (lockNames.Contains(configuration.LockName)) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
$"The distributed lock name '{configuration.LockName}' is used by more than one background job worker " + |
||||
|
$"(the default worker uses '{WorkerOptions.DistributedLockName}'). Each worker must have a unique lock name to run independently."); |
||||
|
} |
||||
|
|
||||
|
lockNames.Add(configuration.LockName); |
||||
|
allDedicatedJobNames.AddRange(jobNames); |
||||
|
dedicatedWorkers.Add(new DedicatedWorkerDefinition(configuration.LockName, jobNames)); |
||||
|
} |
||||
|
|
||||
|
foreach (var dedicatedWorker in dedicatedWorkers) |
||||
|
{ |
||||
|
await StartWorkerAsync(dedicatedWorker.LockName, BackgroundJobNameFilter.Include(dedicatedWorker.JobNames), cancellationToken); |
||||
|
} |
||||
|
|
||||
|
// Default worker processes every job that is not handled by a dedicated worker.
|
||||
|
await StartWorkerAsync(null, BackgroundJobNameFilter.Exclude(allDedicatedJobNames), cancellationToken); |
||||
|
} |
||||
|
|
||||
|
protected virtual string GetJobName(Type argsType) |
||||
|
{ |
||||
|
try |
||||
|
{ |
||||
|
return JobOptions.GetJob(argsType).JobName; |
||||
|
} |
||||
|
catch (AbpException ex) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
$"No background job is registered for the args type '{argsType.FullName}' configured via AddDedicatedWorker. " + |
||||
|
$"Register the job before configuring a dedicated worker for it.", ex); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected virtual async Task StartWorkerAsync( |
||||
|
string? distributedLockName = null, |
||||
|
BackgroundJobNameFilter? jobNameFilter = null, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
var worker = ServiceProvider.GetRequiredService<IBackgroundJobWorker>(); |
||||
|
await worker.StartAsync(distributedLockName, jobNameFilter, cancellationToken); |
||||
|
Workers.Add(worker); |
||||
|
} |
||||
|
|
||||
|
public virtual async Task StopAsync(CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
foreach (var worker in Workers) |
||||
|
{ |
||||
|
await worker.StopAsync(cancellationToken); |
||||
|
} |
||||
|
|
||||
|
Workers.Clear(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,27 @@ |
|||||
|
using System.Collections.Generic; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// A validated, ready-to-start dedicated worker: the distributed lock name it runs under and the
|
||||
|
/// resolved job names it is responsible for. Built by <see cref="BackgroundJobWorkerManager"/> from a
|
||||
|
/// <see cref="BackgroundJobWorkerConfiguration"/> after all configurations have been validated.
|
||||
|
/// </summary>
|
||||
|
public class DedicatedWorkerDefinition |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// The distributed lock name this worker runs under.
|
||||
|
/// </summary>
|
||||
|
public string LockName { get; } |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// The resolved job names this worker is responsible for.
|
||||
|
/// </summary>
|
||||
|
public IReadOnlyList<string> JobNames { get; } |
||||
|
|
||||
|
public DedicatedWorkerDefinition(string lockName, IReadOnlyList<string> jobNames) |
||||
|
{ |
||||
|
LockName = lockName; |
||||
|
JobNames = jobNames; |
||||
|
} |
||||
|
} |
||||
@ -1,8 +1,27 @@ |
|||||
using Volo.Abp.BackgroundWorkers; |
using System.Collections.Generic; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
|
||||
namespace Volo.Abp.BackgroundJobs; |
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
public interface IBackgroundJobWorker : IBackgroundWorker |
/// <summary>
|
||||
|
/// A background job worker that polls and executes waiting jobs.
|
||||
|
/// Instances are created, configured and started by <see cref="BackgroundJobWorkerManager"/>.
|
||||
|
/// </summary>
|
||||
|
public interface IBackgroundJobWorker |
||||
{ |
{ |
||||
|
/// <summary>
|
||||
|
/// Starts this worker.
|
||||
|
/// </summary>
|
||||
|
/// <param name="distributedLockName">
|
||||
|
/// Distributed lock name for this worker. When null, <see cref="AbpBackgroundJobWorkerOptions.DistributedLockName"/> is used.
|
||||
|
/// </param>
|
||||
|
/// <param name="jobNameFilter">Filters the jobs this worker processes by name. When null, all jobs are processed.</param>
|
||||
|
/// <param name="cancellationToken">Cancellation token.</param>
|
||||
|
Task StartAsync( |
||||
|
string? distributedLockName = null, |
||||
|
BackgroundJobNameFilter? jobNameFilter = null, |
||||
|
CancellationToken cancellationToken = default); |
||||
|
|
||||
|
Task StopAsync(CancellationToken cancellationToken = default); |
||||
} |
} |
||||
|
|||||
@ -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<WorkerStartRecorder>(); |
||||
|
context.Services.Replace(ServiceDescriptor.Transient<IBackgroundJobWorker, RecordingBackgroundJobWorker>()); |
||||
|
|
||||
|
// No lock name is given: it is derived from the job argument type names.
|
||||
|
Configure<AbpBackgroundJobWorkerOptions>(options => |
||||
|
{ |
||||
|
options.AddDedicatedWorker<WorkerJobAArgs>(); |
||||
|
options.AddDedicatedWorker<WorkerJobBArgs>(); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -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<AbpBackgroundJobWorkerOptions>(options => |
||||
|
{ |
||||
|
options.StoreSuccessfulJobs = true; |
||||
|
options.SuccessfulJobRetentionTime = TimeSpan.FromDays(1); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -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<AbpBackgroundJobOptions>(options => |
||||
|
{ |
||||
|
options.IsJobExecutionEnabled = false; |
||||
|
}); |
||||
|
|
||||
|
context.Services.AddSingleton<ParallelJobTracker>(); |
||||
|
context.Services.AddScoped<ScopeMarker>(); |
||||
|
} |
||||
|
} |
||||
@ -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<WorkerStartRecorder>(); |
||||
|
context.Services.Replace(ServiceDescriptor.Transient<IBackgroundJobWorker, RecordingBackgroundJobWorker>()); |
||||
|
|
||||
|
// Two workers with different job types but the same lock name, which must fail at initialization.
|
||||
|
Configure<AbpBackgroundJobWorkerOptions>(options => |
||||
|
{ |
||||
|
options.AddDedicatedWorker<WorkerJobAArgs>("dup-lock"); |
||||
|
options.AddDedicatedWorker<WorkerJobBArgs>("dup-lock"); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -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<WorkerStartRecorder>(); |
||||
|
context.Services.Replace(ServiceDescriptor.Transient<IBackgroundJobWorker, RecordingBackgroundJobWorker>()); |
||||
|
|
||||
|
// The same job type is assigned to two dedicated workers, which must fail at initialization.
|
||||
|
Configure<AbpBackgroundJobWorkerOptions>(options => |
||||
|
{ |
||||
|
options.AddDedicatedWorker<WorkerJobAArgs>("lock-a"); |
||||
|
options.AddDedicatedWorker<WorkerJobAArgs>("lock-b"); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -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<WorkerStartRecorder>(); |
||||
|
|
||||
|
// Replace the real worker with a recording one to assert how the manager resolves and starts workers.
|
||||
|
context.Services.Replace(ServiceDescriptor.Transient<IBackgroundJobWorker, RecordingBackgroundJobWorker>()); |
||||
|
|
||||
|
Configure<AbpBackgroundJobWorkerOptions>(options => |
||||
|
{ |
||||
|
options.AddDedicatedWorker<WorkerJobAArgs>("lock-a"); |
||||
|
options.AddDedicatedWorker<WorkerJobBArgs>("lock-b"); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -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<WorkerStartRecorder>(); |
||||
|
context.Services.Replace(ServiceDescriptor.Transient<IBackgroundJobWorker, RecordingBackgroundJobWorker>()); |
||||
|
|
||||
|
// 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<AbpBackgroundJobWorkerOptions>(options => |
||||
|
{ |
||||
|
options.AddDedicatedWorker<SharedNameJobAArgs>("lock-a"); |
||||
|
options.AddDedicatedWorker<SharedNameJobBArgs>("lock-b"); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,148 @@ |
|||||
|
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<AbpBackgroundJobCleanupTestModule> |
||||
|
{ |
||||
|
private readonly IBackgroundJobStore _store; |
||||
|
private readonly IClock _clock; |
||||
|
private readonly AbpBackgroundJobWorkerOptions _workerOptions; |
||||
|
|
||||
|
public BackgroundJobCleanupWorker_Tests() |
||||
|
{ |
||||
|
_store = GetRequiredService<IBackgroundJobStore>(); |
||||
|
_clock = GetRequiredService<IClock>(); |
||||
|
_workerOptions = GetRequiredService<IOptions<AbpBackgroundJobWorkerOptions>>().Value; |
||||
|
} |
||||
|
|
||||
|
protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) |
||||
|
{ |
||||
|
options.UseAutofac(); |
||||
|
} |
||||
|
|
||||
|
private TestableBackgroundJobCleanupWorker CreateWorker() |
||||
|
{ |
||||
|
return new TestableBackgroundJobCleanupWorker( |
||||
|
GetRequiredService<AbpAsyncTimer>(), |
||||
|
GetRequiredService<IServiceScopeFactory>(), |
||||
|
GetRequiredService<IOptions<AbpBackgroundJobOptions>>(), |
||||
|
GetRequiredService<IOptions<AbpBackgroundJobWorkerOptions>>(), |
||||
|
GetRequiredService<IAbpDistributedLock>()); |
||||
|
} |
||||
|
|
||||
|
private Task RunCleanupAsync() |
||||
|
{ |
||||
|
return CreateWorker().DoWorkPublicAsync(new PeriodicBackgroundWorkerContext(ServiceProvider)); |
||||
|
} |
||||
|
|
||||
|
private async Task<Guid> 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<Guid>(); |
||||
|
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_Delete_Oldest_Completed_Jobs_First_When_Limited_By_MaxResultCount() |
||||
|
{ |
||||
|
var oldest = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(5))); |
||||
|
var middle = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(3))); |
||||
|
var newest = await InsertCompletedJobAsync(_clock.Now.Subtract(TimeSpan.FromDays(1))); |
||||
|
|
||||
|
// All three are completed before now, but a single call may delete only two.
|
||||
|
var deletedCount = await _store.DeleteAsync(null, _clock.Now, maxResultCount: 2); |
||||
|
|
||||
|
deletedCount.ShouldBe(2); |
||||
|
(await _store.FindAsync(oldest)).ShouldBeNull(); |
||||
|
(await _store.FindAsync(middle)).ShouldBeNull(); |
||||
|
(await _store.FindAsync(newest)).ShouldNotBeNull(); // newest survives, proving oldest-first deletion
|
||||
|
} |
||||
|
|
||||
|
[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(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
using Volo.Abp.Testing; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
public abstract class BackgroundJobWorkerTestBase : AbpIntegratedTest<AbpBackgroundJobWorkerTestModule> |
||||
|
{ |
||||
|
protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) |
||||
|
{ |
||||
|
options.UseAutofac(); |
||||
|
} |
||||
|
} |
||||
@ -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<AbpAutoLockNameWorkerTestModule> |
||||
|
{ |
||||
|
protected override void SetAbpApplicationCreationOptions(AbpApplicationCreationOptions options) |
||||
|
{ |
||||
|
options.UseAutofac(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public void Should_Derive_Bounded_Lock_Name_From_Job_Args_Types() |
||||
|
{ |
||||
|
var records = GetRequiredService<WorkerStartRecorder>().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); |
||||
|
} |
||||
|
} |
||||
@ -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<AbpDuplicateWorkerTestModule>(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<WorkerStartRecorder>(); |
||||
|
recorder.Records.ShouldBeEmpty(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public void Should_Throw_Without_Starting_Any_Worker_When_Two_Workers_Share_A_Lock_Name() |
||||
|
{ |
||||
|
using var application = AbpApplicationFactory.Create<AbpDuplicateLockNameTestModule>(options => |
||||
|
{ |
||||
|
options.UseAutofac(); |
||||
|
}); |
||||
|
|
||||
|
var exception = Record.Exception(() => application.Initialize()); |
||||
|
|
||||
|
exception.ShouldNotBeNull(); |
||||
|
exception.ToString().ShouldContain("lock name"); |
||||
|
|
||||
|
var recorder = application.ServiceProvider.GetRequiredService<WorkerStartRecorder>(); |
||||
|
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<AbpSameJobNameTestModule>(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<WorkerStartRecorder>(); |
||||
|
recorder.Records.ShouldBeEmpty(); |
||||
|
} |
||||
|
} |
||||
@ -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<AbpMultiWorkerTestModule> |
||||
|
{ |
||||
|
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<WorkerStartRecorder>().Records; |
||||
|
|
||||
|
records.Count.ShouldBe(3); |
||||
|
|
||||
|
var jobAName = BackgroundJobNameAttribute.GetName<WorkerJobAArgs>(); |
||||
|
var jobBName = BackgroundJobNameAttribute.GetName<WorkerJobBArgs>(); |
||||
|
|
||||
|
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); |
||||
|
} |
||||
|
} |
||||
@ -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<IBackgroundJobStore>(); |
||||
|
_clock = GetRequiredService<IClock>(); |
||||
|
_workerOptions = GetRequiredService<IOptions<AbpBackgroundJobWorkerOptions>>().Value; |
||||
|
} |
||||
|
|
||||
|
private TestableBackgroundJobWorker CreateWorker() |
||||
|
{ |
||||
|
return new TestableBackgroundJobWorker( |
||||
|
GetRequiredService<AbpAsyncTimer>(), |
||||
|
GetRequiredService<IOptions<AbpBackgroundJobOptions>>(), |
||||
|
GetRequiredService<IOptions<AbpBackgroundJobWorkerOptions>>(), |
||||
|
GetRequiredService<IServiceScopeFactory>(), |
||||
|
GetRequiredService<IAbpDistributedLock>()); |
||||
|
} |
||||
|
|
||||
|
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<ArgumentException>(() => new BackgroundJobNameFilter(BackgroundJobNameFilterMode.Include)); |
||||
|
Should.Throw<ArgumentException>(() => new BackgroundJobNameFilter(BackgroundJobNameFilterMode.None, new[] { "job-a" })); |
||||
|
Should.Throw<ArgumentException>(() => 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<ArgumentException>(() => 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<ParallelTestJobArgs>("lock-a"); |
||||
|
|
||||
|
Should.Throw<AbpException>(() => options.AddDedicatedWorker<ParallelTestJobArgs>("lock-b")); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public void AddDedicatedWorker_Should_Throw_At_Registration_When_A_Lock_Name_Is_Reused() |
||||
|
{ |
||||
|
var options = new AbpBackgroundJobWorkerOptions(); |
||||
|
options.AddDedicatedWorker<WorkerJobAArgs>("dup-lock"); |
||||
|
|
||||
|
Should.Throw<AbpException>(() => options.AddDedicatedWorker<WorkerJobBArgs>("dup-lock")); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public void AddDedicatedWorker_Should_Throw_At_Registration_When_Lock_Name_Equals_The_Default() |
||||
|
{ |
||||
|
var options = new AbpBackgroundJobWorkerOptions(); |
||||
|
|
||||
|
Should.Throw<AbpException>(() => options.AddDedicatedWorker<WorkerJobAArgs>(options.DistributedLockName)); |
||||
|
} |
||||
|
|
||||
|
// Parallel execution
|
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Execute_Multiple_Jobs_In_Parallel() |
||||
|
{ |
||||
|
_workerOptions.MaxParallelJobExecutionCount = 3; |
||||
|
|
||||
|
var jobManager = GetRequiredService<IBackgroundJobManager>(); |
||||
|
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<ParallelJobTracker>(); |
||||
|
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<IBackgroundJobManager>(); |
||||
|
for (var i = 0; i < 5; i++) |
||||
|
{ |
||||
|
await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = i.ToString() }); |
||||
|
} |
||||
|
|
||||
|
await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); |
||||
|
|
||||
|
var tracker = GetRequiredService<ParallelJobTracker>(); |
||||
|
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<IBackgroundJobManager>(); |
||||
|
var lockedJobId = Guid.Parse(await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "locked" })); |
||||
|
await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "free" }); |
||||
|
|
||||
|
var distributedLock = GetRequiredService<IAbpDistributedLock>(); |
||||
|
var lockName = _workerOptions.PerJobDistributedLockPrefix + lockedJobId; |
||||
|
|
||||
|
await using (await distributedLock.TryAcquireAsync(lockName)) |
||||
|
{ |
||||
|
await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); |
||||
|
} |
||||
|
|
||||
|
var tracker = GetRequiredService<ParallelJobTracker>(); |
||||
|
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<IBackgroundJobManager>(); |
||||
|
var jobId = Guid.Parse(await jobManager.EnqueueAsync(new ParallelTestJobArgs { Value = "1" })); |
||||
|
|
||||
|
await CreateWorker().ExecuteJobsInParallelPublicAsync(Context()); |
||||
|
|
||||
|
GetRequiredService<ParallelJobTracker>().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); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,51 @@ |
|||||
|
using System.Collections.Generic; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Records every <see cref="IBackgroundJobWorker.StartAsync"/> call so multi-worker registration
|
||||
|
/// (resolved from DI by <see cref="BackgroundJobWorkerManager"/>) can be asserted. Does not start any timer.
|
||||
|
/// </summary>
|
||||
|
public class WorkerStartRecorder |
||||
|
{ |
||||
|
public List<WorkerStartRecord> Records { get; } = new List<WorkerStartRecord>(); |
||||
|
} |
||||
|
|
||||
|
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; |
||||
|
} |
||||
|
} |
||||
@ -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<AbpBackgroundJobOptions> jobOptions, |
||||
|
IOptions<AbpBackgroundJobWorkerOptions> workerOptions, |
||||
|
IAbpDistributedLock distributedLock) |
||||
|
: base(timer, serviceScopeFactory, jobOptions, workerOptions, distributedLock) |
||||
|
{ |
||||
|
} |
||||
|
|
||||
|
public Task DoWorkPublicAsync(PeriodicBackgroundWorkerContext workerContext) |
||||
|
{ |
||||
|
return DoWorkAsync(workerContext); |
||||
|
} |
||||
|
} |
||||
@ -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; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Exposes the protected members of <see cref="BackgroundJobWorker"/> for unit testing.
|
||||
|
/// </summary>
|
||||
|
public class TestableBackgroundJobWorker : BackgroundJobWorker |
||||
|
{ |
||||
|
public TestableBackgroundJobWorker( |
||||
|
AbpAsyncTimer timer, |
||||
|
IOptions<AbpBackgroundJobOptions> jobOptions, |
||||
|
IOptions<AbpBackgroundJobWorkerOptions> 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<List<BackgroundJobInfo>> GetWaitingJobsPublicAsync(PeriodicBackgroundWorkerContext workerContext, IBackgroundJobStore store) |
||||
|
{ |
||||
|
return GetWaitingJobsAsync(workerContext, store); |
||||
|
} |
||||
|
|
||||
|
public bool IsJobEligiblePublic(BackgroundJobInfo? jobInfo, IClock clock) |
||||
|
{ |
||||
|
return IsJobEligible(jobInfo, clock); |
||||
|
} |
||||
|
} |
||||
@ -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<string> Executed { get; } = new ConcurrentBag<string>(); |
||||
|
|
||||
|
public ConcurrentBag<Guid> ScopeIds { get; } = new ConcurrentBag<Guid>(); |
||||
|
} |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Scoped service used to verify that each parallel job runs in its own service scope.
|
||||
|
/// </summary>
|
||||
|
public class ScopeMarker |
||||
|
{ |
||||
|
public Guid Id { get; } = Guid.NewGuid(); |
||||
|
} |
||||
|
|
||||
|
public class ParallelTestJobArgs |
||||
|
{ |
||||
|
public string Value { get; set; } = default!; |
||||
|
} |
||||
|
|
||||
|
public class ParallelTestJob : AsyncBackgroundJob<ParallelTestJobArgs>, 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<WorkerJobAArgs>, ITransientDependency |
||||
|
{ |
||||
|
public override Task ExecuteAsync(WorkerJobAArgs args) |
||||
|
{ |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public class WorkerJobBArgs |
||||
|
{ |
||||
|
public string Value { get; set; } = default!; |
||||
|
} |
||||
|
|
||||
|
public class WorkerJobB : AsyncBackgroundJob<WorkerJobBArgs>, 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<SharedNameJobAArgs>, ITransientDependency |
||||
|
{ |
||||
|
public override Task ExecuteAsync(SharedNameJobAArgs args) |
||||
|
{ |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public class SharedNameJobB : AsyncBackgroundJob<SharedNameJobBArgs>, ITransientDependency |
||||
|
{ |
||||
|
public override Task ExecuteAsync(SharedNameJobBArgs args) |
||||
|
{ |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
} |
||||
@ -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<CalculateAwsFeesJobArgs>, 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; |
||||
|
} |
||||
|
} |
||||
@ -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<CalculateAzureFeesJobArgs>, 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; |
||||
|
} |
||||
|
} |
||||
@ -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<SendEmailJobArgs>, 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; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,99 @@ |
|||||
|
// <auto-generated />
|
||||
|
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 |
||||
|
{ |
||||
|
/// <inheritdoc />
|
||||
|
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<Guid>("Id") |
||||
|
.ValueGeneratedOnAdd() |
||||
|
.HasColumnType("uniqueidentifier"); |
||||
|
|
||||
|
b.Property<string>("ApplicationName") |
||||
|
.HasMaxLength(96) |
||||
|
.HasColumnType("nvarchar(96)"); |
||||
|
|
||||
|
b.Property<DateTime?>("CompletionTime") |
||||
|
.HasColumnType("datetime2"); |
||||
|
|
||||
|
b.Property<string>("ConcurrencyStamp") |
||||
|
.IsConcurrencyToken() |
||||
|
.IsRequired() |
||||
|
.HasMaxLength(40) |
||||
|
.HasColumnType("nvarchar(40)") |
||||
|
.HasColumnName("ConcurrencyStamp"); |
||||
|
|
||||
|
b.Property<DateTime>("CreationTime") |
||||
|
.HasColumnType("datetime2") |
||||
|
.HasColumnName("CreationTime"); |
||||
|
|
||||
|
b.Property<string>("ExtraProperties") |
||||
|
.IsRequired() |
||||
|
.HasColumnType("nvarchar(max)") |
||||
|
.HasColumnName("ExtraProperties"); |
||||
|
|
||||
|
b.Property<bool>("IsAbandoned") |
||||
|
.ValueGeneratedOnAdd() |
||||
|
.HasColumnType("bit") |
||||
|
.HasDefaultValue(false); |
||||
|
|
||||
|
b.Property<string>("JobArgs") |
||||
|
.IsRequired() |
||||
|
.HasMaxLength(1048576) |
||||
|
.HasColumnType("nvarchar(max)"); |
||||
|
|
||||
|
b.Property<string>("JobName") |
||||
|
.IsRequired() |
||||
|
.HasMaxLength(128) |
||||
|
.HasColumnType("nvarchar(128)"); |
||||
|
|
||||
|
b.Property<DateTime?>("LastTryTime") |
||||
|
.HasColumnType("datetime2"); |
||||
|
|
||||
|
b.Property<DateTime>("NextTryTime") |
||||
|
.HasColumnType("datetime2"); |
||||
|
|
||||
|
b.Property<byte>("Priority") |
||||
|
.ValueGeneratedOnAdd() |
||||
|
.HasColumnType("tinyint") |
||||
|
.HasDefaultValue((byte)15); |
||||
|
|
||||
|
b.Property<short>("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
|
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,47 @@ |
|||||
|
using System; |
||||
|
using Microsoft.EntityFrameworkCore.Migrations; |
||||
|
|
||||
|
#nullable disable |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations |
||||
|
{ |
||||
|
/// <inheritdoc />
|
||||
|
public partial class Added_CompletionTime_To_BackgroundJobs : Migration |
||||
|
{ |
||||
|
/// <inheritdoc />
|
||||
|
protected override void Up(MigrationBuilder migrationBuilder) |
||||
|
{ |
||||
|
migrationBuilder.DropIndex( |
||||
|
name: "IX_AbpBackgroundJobs_IsAbandoned_NextTryTime", |
||||
|
table: "AbpBackgroundJobs"); |
||||
|
|
||||
|
migrationBuilder.AddColumn<DateTime>( |
||||
|
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" }); |
||||
|
} |
||||
|
|
||||
|
/// <inheritdoc />
|
||||
|
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" }); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
Loading…
Reference in new issue