mirror of https://github.com/abpframework/abp.git
21 changed files with 1361 additions and 6 deletions
@ -0,0 +1,48 @@ |
|||||
|
using System; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.ExceptionHandling; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers.Hangfire; |
||||
|
|
||||
|
public class HangfireDynamicBackgroundWorkerAdapter : ITransientDependency |
||||
|
{ |
||||
|
protected IDynamicBackgroundWorkerHandlerRegistry HandlerRegistry { get; } |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
public ILogger<HangfireDynamicBackgroundWorkerAdapter> Logger { get; set; } |
||||
|
|
||||
|
public HangfireDynamicBackgroundWorkerAdapter( |
||||
|
IDynamicBackgroundWorkerHandlerRegistry handlerRegistry, |
||||
|
IServiceProvider serviceProvider) |
||||
|
{ |
||||
|
HandlerRegistry = handlerRegistry; |
||||
|
ServiceProvider = serviceProvider; |
||||
|
Logger = NullLogger<HangfireDynamicBackgroundWorkerAdapter>.Instance; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task DoWorkAsync(string workerName, CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
var handler = HandlerRegistry.Get(workerName); |
||||
|
if (handler == null) |
||||
|
{ |
||||
|
Logger.LogWarning("No handler registered for dynamic worker: {WorkerName}", workerName); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
try |
||||
|
{ |
||||
|
await handler(new DynamicBackgroundWorkerExecutionContext(workerName, ServiceProvider), cancellationToken); |
||||
|
} |
||||
|
catch (Exception ex) |
||||
|
{ |
||||
|
await ServiceProvider.GetRequiredService<IExceptionNotifier>() |
||||
|
.NotifyAsync(new ExceptionNotificationContext(ex)); |
||||
|
|
||||
|
throw; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,170 @@ |
|||||
|
using System; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Hangfire; |
||||
|
using Hangfire.Storage; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.Hangfire; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers.Hangfire; |
||||
|
|
||||
|
[Dependency(ReplaceServices = true)] |
||||
|
public class HangfireDynamicBackgroundWorkerManager : IDynamicBackgroundWorkerManager, ISingletonDependency |
||||
|
{ |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
protected IDynamicBackgroundWorkerHandlerRegistry HandlerRegistry { get; } |
||||
|
public ILogger<HangfireDynamicBackgroundWorkerManager> Logger { get; set; } |
||||
|
|
||||
|
public HangfireDynamicBackgroundWorkerManager( |
||||
|
IServiceProvider serviceProvider, |
||||
|
IDynamicBackgroundWorkerHandlerRegistry handlerRegistry) |
||||
|
{ |
||||
|
ServiceProvider = serviceProvider; |
||||
|
HandlerRegistry = handlerRegistry; |
||||
|
Logger = NullLogger<HangfireDynamicBackgroundWorkerManager>.Instance; |
||||
|
} |
||||
|
|
||||
|
public virtual Task AddAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
Check.NotNull(handler, nameof(handler)); |
||||
|
|
||||
|
schedule.Validate(); |
||||
|
|
||||
|
var cronExpression = schedule.CronExpression; |
||||
|
if (cronExpression.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
var period = schedule.Period ?? DynamicBackgroundWorkerSchedule.DefaultPeriod; |
||||
|
cronExpression = GetCron(period); |
||||
|
} |
||||
|
|
||||
|
ScheduleRecurringJob(workerName, cronExpression, cancellationToken); |
||||
|
HandlerRegistry.Register(workerName, handler); |
||||
|
|
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
|
||||
|
public virtual Task<bool> RemoveAsync(string workerName, CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
|
||||
|
if (!HandlerRegistry.IsRegistered(workerName)) |
||||
|
{ |
||||
|
return Task.FromResult(false); |
||||
|
} |
||||
|
|
||||
|
var recurringJobId = $"DynamicWorker:{workerName}"; |
||||
|
RecurringJob.RemoveIfExists(recurringJobId); |
||||
|
HandlerRegistry.Unregister(workerName); |
||||
|
|
||||
|
return Task.FromResult(true); |
||||
|
} |
||||
|
|
||||
|
public virtual Task<bool> UpdateScheduleAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
|
||||
|
schedule.Validate(); |
||||
|
|
||||
|
if (!HandlerRegistry.IsRegistered(workerName)) |
||||
|
{ |
||||
|
return Task.FromResult(false); |
||||
|
} |
||||
|
|
||||
|
var cronExpression = schedule.CronExpression; |
||||
|
if (cronExpression.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
var period = schedule.Period ?? DynamicBackgroundWorkerSchedule.DefaultPeriod; |
||||
|
cronExpression = GetCron(period); |
||||
|
} |
||||
|
|
||||
|
ScheduleRecurringJob(workerName, cronExpression, cancellationToken); |
||||
|
|
||||
|
return Task.FromResult(true); |
||||
|
} |
||||
|
|
||||
|
public virtual bool IsRegistered(string workerName) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
return HandlerRegistry.IsRegistered(workerName); |
||||
|
} |
||||
|
|
||||
|
protected virtual void ScheduleRecurringJob(string workerName, string cronExpression, CancellationToken cancellationToken) |
||||
|
{ |
||||
|
var abpHangfireOptions = ServiceProvider.GetRequiredService<IOptions<AbpHangfireOptions>>().Value; |
||||
|
var queueName = abpHangfireOptions.DefaultQueue; |
||||
|
var recurringJobId = $"DynamicWorker:{workerName}"; |
||||
|
|
||||
|
if (!JobStorage.Current.HasFeature(JobStorageFeatures.JobQueueProperty)) |
||||
|
{ |
||||
|
Logger.LogWarning( |
||||
|
"Current storage doesn't support specifying queues ({QueueName}) directly for a specific job. Please use the QueueAttribute instead.", |
||||
|
queueName); |
||||
|
|
||||
|
RecurringJob.AddOrUpdate<HangfireDynamicBackgroundWorkerAdapter>( |
||||
|
recurringJobId, |
||||
|
adapter => adapter.DoWorkAsync(workerName, CancellationToken.None), |
||||
|
cronExpression, |
||||
|
new RecurringJobOptions |
||||
|
{ |
||||
|
TimeZone = TimeZoneInfo.Utc |
||||
|
}); |
||||
|
} |
||||
|
else |
||||
|
{ |
||||
|
RecurringJob.AddOrUpdate<HangfireDynamicBackgroundWorkerAdapter>( |
||||
|
recurringJobId, |
||||
|
queueName, |
||||
|
adapter => adapter.DoWorkAsync(workerName, CancellationToken.None), |
||||
|
cronExpression, |
||||
|
new RecurringJobOptions |
||||
|
{ |
||||
|
TimeZone = TimeZoneInfo.Utc |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected virtual string GetCron(int period) |
||||
|
{ |
||||
|
var time = TimeSpan.FromMilliseconds(period); |
||||
|
|
||||
|
if (time.TotalSeconds <= 59) |
||||
|
{ |
||||
|
var seconds = (int)Math.Round(time.TotalSeconds); |
||||
|
return $"*/{seconds} * * * * *"; |
||||
|
} |
||||
|
|
||||
|
if (time.TotalMinutes <= 59) |
||||
|
{ |
||||
|
var minutes = (int)Math.Round(time.TotalMinutes); |
||||
|
return $"*/{minutes} * * * *"; |
||||
|
} |
||||
|
|
||||
|
if (time.TotalHours <= 23) |
||||
|
{ |
||||
|
var hours = (int)Math.Round(time.TotalHours); |
||||
|
return $"0 */{hours} * * *"; |
||||
|
} |
||||
|
|
||||
|
if (time.TotalDays <= 31) |
||||
|
{ |
||||
|
var days = (int)Math.Round(time.TotalDays); |
||||
|
return $"0 0 */{days} * *"; |
||||
|
} |
||||
|
|
||||
|
throw new AbpException($"Cannot convert period: {period} to cron expression."); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,56 @@ |
|||||
|
using System; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Quartz; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.ExceptionHandling; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers.Quartz; |
||||
|
|
||||
|
public class QuartzDynamicBackgroundWorkerAdapter : IJob, ITransientDependency |
||||
|
{ |
||||
|
protected IDynamicBackgroundWorkerHandlerRegistry HandlerRegistry { get; } |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
public ILogger<QuartzDynamicBackgroundWorkerAdapter> Logger { get; set; } |
||||
|
|
||||
|
public QuartzDynamicBackgroundWorkerAdapter( |
||||
|
IDynamicBackgroundWorkerHandlerRegistry handlerRegistry, |
||||
|
IServiceProvider serviceProvider) |
||||
|
{ |
||||
|
HandlerRegistry = handlerRegistry; |
||||
|
ServiceProvider = serviceProvider; |
||||
|
Logger = NullLogger<QuartzDynamicBackgroundWorkerAdapter>.Instance; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task Execute(IJobExecutionContext context) |
||||
|
{ |
||||
|
var workerName = context.MergedJobDataMap.GetString(QuartzDynamicBackgroundWorkerManager.DynamicWorkerNameKey); |
||||
|
if (string.IsNullOrWhiteSpace(workerName)) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
var handler = HandlerRegistry.Get(workerName!); |
||||
|
if (handler == null) |
||||
|
{ |
||||
|
Logger.LogWarning("No handler registered for dynamic worker: {WorkerName}", workerName); |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
try |
||||
|
{ |
||||
|
await handler( |
||||
|
new DynamicBackgroundWorkerExecutionContext(workerName!, ServiceProvider), |
||||
|
context.CancellationToken); |
||||
|
} |
||||
|
catch (Exception ex) |
||||
|
{ |
||||
|
await ServiceProvider.GetRequiredService<IExceptionNotifier>() |
||||
|
.NotifyAsync(new ExceptionNotificationContext(ex)); |
||||
|
|
||||
|
throw; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,140 @@ |
|||||
|
using System; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Quartz; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers.Quartz; |
||||
|
|
||||
|
[Dependency(ReplaceServices = true)] |
||||
|
public class QuartzDynamicBackgroundWorkerManager : IDynamicBackgroundWorkerManager, ISingletonDependency |
||||
|
{ |
||||
|
public const string DynamicWorkerNameKey = "AbpDynamicWorkerName"; |
||||
|
|
||||
|
protected IScheduler Scheduler { get; } |
||||
|
protected IDynamicBackgroundWorkerHandlerRegistry HandlerRegistry { get; } |
||||
|
public ILogger<QuartzDynamicBackgroundWorkerManager> Logger { get; set; } |
||||
|
|
||||
|
public QuartzDynamicBackgroundWorkerManager( |
||||
|
IScheduler scheduler, |
||||
|
IDynamicBackgroundWorkerHandlerRegistry handlerRegistry) |
||||
|
{ |
||||
|
Scheduler = scheduler; |
||||
|
HandlerRegistry = handlerRegistry; |
||||
|
Logger = NullLogger<QuartzDynamicBackgroundWorkerManager>.Instance; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task AddAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
Check.NotNull(handler, nameof(handler)); |
||||
|
|
||||
|
schedule.Validate(); |
||||
|
|
||||
|
var jobKey = new JobKey($"DynamicWorker:{workerName}"); |
||||
|
var triggerKey = new TriggerKey($"DynamicWorker:{workerName}"); |
||||
|
var jobDetail = JobBuilder.Create<QuartzDynamicBackgroundWorkerAdapter>() |
||||
|
.WithIdentity(jobKey) |
||||
|
.UsingJobData(DynamicWorkerNameKey, workerName) |
||||
|
.Build(); |
||||
|
|
||||
|
var trigger = BuildTrigger(schedule, jobDetail, triggerKey); |
||||
|
|
||||
|
if (await Scheduler.CheckExists(jobDetail.Key, cancellationToken)) |
||||
|
{ |
||||
|
await Scheduler.AddJob(jobDetail, true, true, cancellationToken); |
||||
|
await Scheduler.ResumeJob(jobDetail.Key, cancellationToken); |
||||
|
await Scheduler.RescheduleJob(trigger.Key, trigger, cancellationToken); |
||||
|
} |
||||
|
else |
||||
|
{ |
||||
|
await Scheduler.ScheduleJob(jobDetail, trigger, cancellationToken); |
||||
|
} |
||||
|
|
||||
|
HandlerRegistry.Register(workerName, handler); |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<bool> RemoveAsync(string workerName, CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
|
||||
|
if (!HandlerRegistry.IsRegistered(workerName)) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
var jobKey = new JobKey($"DynamicWorker:{workerName}"); |
||||
|
await Scheduler.DeleteJob(jobKey, cancellationToken); |
||||
|
HandlerRegistry.Unregister(workerName); |
||||
|
|
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<bool> UpdateScheduleAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
|
||||
|
schedule.Validate(); |
||||
|
|
||||
|
if (!HandlerRegistry.IsRegistered(workerName)) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
var triggerKey = new TriggerKey($"DynamicWorker:{workerName}"); |
||||
|
var jobKey = new JobKey($"DynamicWorker:{workerName}"); |
||||
|
|
||||
|
var triggerBuilder = TriggerBuilder.Create() |
||||
|
.WithIdentity(triggerKey) |
||||
|
.ForJob(jobKey); |
||||
|
|
||||
|
if (!schedule.CronExpression.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
triggerBuilder.WithCronSchedule(schedule.CronExpression); |
||||
|
} |
||||
|
else |
||||
|
{ |
||||
|
triggerBuilder.WithSimpleSchedule(builder => |
||||
|
builder.WithInterval(TimeSpan.FromMilliseconds(schedule.Period!.Value)).RepeatForever()); |
||||
|
} |
||||
|
|
||||
|
var result = await Scheduler.RescheduleJob(triggerKey, triggerBuilder.Build(), cancellationToken); |
||||
|
return result != null; |
||||
|
} |
||||
|
|
||||
|
public virtual bool IsRegistered(string workerName) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
return HandlerRegistry.IsRegistered(workerName); |
||||
|
} |
||||
|
|
||||
|
protected virtual ITrigger BuildTrigger(DynamicBackgroundWorkerSchedule schedule, IJobDetail jobDetail, TriggerKey triggerKey) |
||||
|
{ |
||||
|
var triggerBuilder = TriggerBuilder.Create() |
||||
|
.ForJob(jobDetail) |
||||
|
.WithIdentity(triggerKey); |
||||
|
|
||||
|
if (!schedule.CronExpression.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
triggerBuilder.WithCronSchedule(schedule.CronExpression); |
||||
|
} |
||||
|
else |
||||
|
{ |
||||
|
triggerBuilder.WithSimpleSchedule(builder => |
||||
|
builder.WithInterval(TimeSpan.FromMilliseconds(schedule.Period!.Value)).RepeatForever()); |
||||
|
} |
||||
|
|
||||
|
return triggerBuilder.Build(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,43 @@ |
|||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers.TickerQ; |
||||
|
|
||||
|
[Dependency(ReplaceServices = true)] |
||||
|
public class TickerQDynamicBackgroundWorkerManager : IDynamicBackgroundWorkerManager, ISingletonDependency |
||||
|
{ |
||||
|
public virtual Task AddAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
"TickerQ does not support dynamic background worker registration at runtime. " + |
||||
|
"TickerQ uses FrozenDictionary for function registration, which requires all functions to be registered before the application starts. " + |
||||
|
"Please use Hangfire or Quartz provider for dynamic background workers."); |
||||
|
} |
||||
|
|
||||
|
public virtual Task<bool> RemoveAsync(string workerName, CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
"TickerQ does not support dynamic background worker registration at runtime. " + |
||||
|
"Please use Hangfire or Quartz provider for dynamic background workers."); |
||||
|
} |
||||
|
|
||||
|
public virtual Task<bool> UpdateScheduleAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
"TickerQ does not support dynamic background worker registration at runtime. " + |
||||
|
"Please use Hangfire or Quartz provider for dynamic background workers."); |
||||
|
} |
||||
|
|
||||
|
public virtual bool IsRegistered(string workerName) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,171 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Concurrent; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Logging; |
||||
|
using Microsoft.Extensions.Logging.Abstractions; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.Threading; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public class DefaultDynamicBackgroundWorkerManager : IDynamicBackgroundWorkerManager, ISingletonDependency, IDisposable |
||||
|
{ |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
public ILogger<DefaultDynamicBackgroundWorkerManager> Logger { get; set; } |
||||
|
|
||||
|
private readonly ConcurrentDictionary<string, InMemoryDynamicBackgroundWorker> _dynamicWorkers; |
||||
|
private readonly SemaphoreSlim _semaphore; |
||||
|
private bool _isDisposed; |
||||
|
|
||||
|
public DefaultDynamicBackgroundWorkerManager(IServiceProvider serviceProvider) |
||||
|
{ |
||||
|
ServiceProvider = serviceProvider; |
||||
|
Logger = NullLogger<DefaultDynamicBackgroundWorkerManager>.Instance; |
||||
|
_dynamicWorkers = new ConcurrentDictionary<string, InMemoryDynamicBackgroundWorker>(); |
||||
|
_semaphore = new SemaphoreSlim(1, 1); |
||||
|
} |
||||
|
|
||||
|
public virtual async Task AddAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
Check.NotNull(handler, nameof(handler)); |
||||
|
|
||||
|
schedule.Validate(); |
||||
|
|
||||
|
if (schedule.Period == null) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
$"The default in-memory background worker manager does not support CronExpression without Period for dynamic worker '{workerName}'. " + |
||||
|
"Please set Period, or use a scheduler-backed provider (Hangfire, Quartz, TickerQ)."); |
||||
|
} |
||||
|
|
||||
|
await _semaphore.WaitAsync(cancellationToken); |
||||
|
try |
||||
|
{ |
||||
|
if (_dynamicWorkers.TryRemove(workerName, out var existingWorker)) |
||||
|
{ |
||||
|
await existingWorker.StopAsync(cancellationToken); |
||||
|
Logger.LogInformation("Replaced existing dynamic worker: {WorkerName}", workerName); |
||||
|
} |
||||
|
|
||||
|
var worker = CreateDynamicWorker(workerName, schedule, handler); |
||||
|
_dynamicWorkers[workerName] = worker; |
||||
|
|
||||
|
await worker.StartAsync(cancellationToken); |
||||
|
} |
||||
|
finally |
||||
|
{ |
||||
|
_semaphore.Release(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<bool> RemoveAsync(string workerName, CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
|
||||
|
await _semaphore.WaitAsync(cancellationToken); |
||||
|
try |
||||
|
{ |
||||
|
if (!_dynamicWorkers.TryRemove(workerName, out var worker)) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
await worker.StopAsync(cancellationToken); |
||||
|
return true; |
||||
|
} |
||||
|
finally |
||||
|
{ |
||||
|
_semaphore.Release(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<bool> UpdateScheduleAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
|
||||
|
schedule.Validate(); |
||||
|
|
||||
|
if (schedule.Period == null) |
||||
|
{ |
||||
|
throw new AbpException( |
||||
|
$"The default in-memory background worker manager does not support CronExpression without Period for dynamic worker '{workerName}'. " + |
||||
|
"Please set Period, or use a scheduler-backed provider (Hangfire, Quartz, TickerQ)."); |
||||
|
} |
||||
|
|
||||
|
await _semaphore.WaitAsync(cancellationToken); |
||||
|
try |
||||
|
{ |
||||
|
if (!_dynamicWorkers.TryGetValue(workerName, out var worker)) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
worker.UpdateSchedule(schedule); |
||||
|
return true; |
||||
|
} |
||||
|
finally |
||||
|
{ |
||||
|
_semaphore.Release(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public virtual bool IsRegistered(string workerName) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
return _dynamicWorkers.ContainsKey(workerName); |
||||
|
} |
||||
|
|
||||
|
public virtual void Dispose() |
||||
|
{ |
||||
|
if (_isDisposed) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
_isDisposed = true; |
||||
|
|
||||
|
foreach (var kvp in _dynamicWorkers) |
||||
|
{ |
||||
|
try |
||||
|
{ |
||||
|
kvp.Value.StopAsync(CancellationToken.None).GetAwaiter().GetResult(); |
||||
|
} |
||||
|
catch (Exception ex) |
||||
|
{ |
||||
|
Logger.LogException(ex); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
_dynamicWorkers.Clear(); |
||||
|
_semaphore.Dispose(); |
||||
|
} |
||||
|
|
||||
|
protected virtual InMemoryDynamicBackgroundWorker CreateDynamicWorker( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler) |
||||
|
{ |
||||
|
var timer = ServiceProvider.GetRequiredService<AbpAsyncTimer>(); |
||||
|
var serviceScopeFactory = ServiceProvider.GetRequiredService<IServiceScopeFactory>(); |
||||
|
|
||||
|
var worker = new InMemoryDynamicBackgroundWorker( |
||||
|
workerName, schedule, handler, timer, serviceScopeFactory); |
||||
|
|
||||
|
worker.ServiceProvider = ServiceProvider; |
||||
|
worker.LazyServiceProvider = ServiceProvider.GetRequiredService<IAbpLazyServiceProvider>(); |
||||
|
|
||||
|
return worker; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,16 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public class DynamicBackgroundWorkerExecutionContext |
||||
|
{ |
||||
|
public string WorkerName { get; } |
||||
|
|
||||
|
public IServiceProvider ServiceProvider { get; } |
||||
|
|
||||
|
public DynamicBackgroundWorkerExecutionContext(string workerName, IServiceProvider serviceProvider) |
||||
|
{ |
||||
|
WorkerName = Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
ServiceProvider = Check.NotNull(serviceProvider, nameof(serviceProvider)); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,6 @@ |
|||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public delegate Task DynamicBackgroundWorkerHandler(DynamicBackgroundWorkerExecutionContext context, CancellationToken cancellationToken); |
||||
@ -0,0 +1,40 @@ |
|||||
|
using System.Collections.Concurrent; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public class DynamicBackgroundWorkerHandlerRegistry : IDynamicBackgroundWorkerHandlerRegistry, ISingletonDependency |
||||
|
{ |
||||
|
protected ConcurrentDictionary<string, DynamicBackgroundWorkerHandler> Handlers { get; } |
||||
|
|
||||
|
public DynamicBackgroundWorkerHandlerRegistry() |
||||
|
{ |
||||
|
Handlers = new ConcurrentDictionary<string, DynamicBackgroundWorkerHandler>(); |
||||
|
} |
||||
|
|
||||
|
public virtual void Register(string workerName, DynamicBackgroundWorkerHandler handler) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(handler, nameof(handler)); |
||||
|
|
||||
|
Handlers[workerName] = handler; |
||||
|
} |
||||
|
|
||||
|
public virtual bool Unregister(string workerName) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
return Handlers.TryRemove(workerName, out _); |
||||
|
} |
||||
|
|
||||
|
public virtual bool IsRegistered(string workerName) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
return Handlers.ContainsKey(workerName); |
||||
|
} |
||||
|
|
||||
|
public virtual DynamicBackgroundWorkerHandler? Get(string workerName) |
||||
|
{ |
||||
|
Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
return Handlers.TryGetValue(workerName, out var handler) ? handler : null; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,26 @@ |
|||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public static class DynamicBackgroundWorkerManagerExtensions |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// Adds a dynamic worker with the default schedule (<see cref="DynamicBackgroundWorkerSchedule.DefaultPeriod"/>).
|
||||
|
/// </summary>
|
||||
|
public static Task AddAsync( |
||||
|
this IDynamicBackgroundWorkerManager manager, |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
return manager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
Period = DynamicBackgroundWorkerSchedule.DefaultPeriod |
||||
|
}, |
||||
|
handler, |
||||
|
cancellationToken); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
using System; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public class DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
public const int DefaultPeriod = 60000; |
||||
|
|
||||
|
public int? Period { get; set; } |
||||
|
|
||||
|
public string? CronExpression { get; set; } |
||||
|
|
||||
|
public virtual void Validate() |
||||
|
{ |
||||
|
if (Period.HasValue && Period.Value <= 0) |
||||
|
{ |
||||
|
throw new ArgumentException( |
||||
|
$"Period must be greater than 0 when provided. Given value: {Period.Value}.", |
||||
|
nameof(Period)); |
||||
|
} |
||||
|
|
||||
|
if (Period == null && string.IsNullOrWhiteSpace(CronExpression)) |
||||
|
{ |
||||
|
throw new ArgumentException( |
||||
|
"At least one of 'Period' or 'CronExpression' must be set."); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,12 @@ |
|||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public interface IDynamicBackgroundWorkerHandlerRegistry |
||||
|
{ |
||||
|
void Register(string workerName, DynamicBackgroundWorkerHandler handler); |
||||
|
|
||||
|
bool Unregister(string workerName); |
||||
|
|
||||
|
bool IsRegistered(string workerName); |
||||
|
|
||||
|
DynamicBackgroundWorkerHandler? Get(string workerName); |
||||
|
} |
||||
@ -0,0 +1,41 @@ |
|||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Manages dynamic background workers that are registered at runtime
|
||||
|
/// without requiring a strongly-typed worker class.
|
||||
|
/// </summary>
|
||||
|
public interface IDynamicBackgroundWorkerManager |
||||
|
{ |
||||
|
/// <summary>
|
||||
|
/// Adds a dynamic worker by name, schedule and handler.
|
||||
|
/// If a worker with the same name already exists, it will be replaced.
|
||||
|
/// </summary>
|
||||
|
Task AddAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
CancellationToken cancellationToken = default); |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Removes a previously added dynamic worker by name.
|
||||
|
/// Returns true if the worker was found and removed; false otherwise.
|
||||
|
/// </summary>
|
||||
|
Task<bool> RemoveAsync(string workerName, CancellationToken cancellationToken = default); |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Updates the schedule of a previously added dynamic worker.
|
||||
|
/// Returns true if the worker was found and updated; false otherwise.
|
||||
|
/// </summary>
|
||||
|
Task<bool> UpdateScheduleAsync( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
CancellationToken cancellationToken = default); |
||||
|
|
||||
|
/// <summary>
|
||||
|
/// Checks whether a dynamic worker with the given name is registered.
|
||||
|
/// </summary>
|
||||
|
bool IsRegistered(string workerName); |
||||
|
} |
||||
@ -0,0 +1,51 @@ |
|||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Threading; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundWorkers; |
||||
|
|
||||
|
public class InMemoryDynamicBackgroundWorker : AsyncPeriodicBackgroundWorkerBase |
||||
|
{ |
||||
|
public string WorkerName { get; } |
||||
|
|
||||
|
private readonly DynamicBackgroundWorkerHandler _handler; |
||||
|
|
||||
|
public InMemoryDynamicBackgroundWorker( |
||||
|
string workerName, |
||||
|
DynamicBackgroundWorkerSchedule schedule, |
||||
|
DynamicBackgroundWorkerHandler handler, |
||||
|
AbpAsyncTimer timer, |
||||
|
IServiceScopeFactory serviceScopeFactory) |
||||
|
: base(timer, serviceScopeFactory) |
||||
|
{ |
||||
|
WorkerName = Check.NotNullOrWhiteSpace(workerName, nameof(workerName)); |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
_handler = Check.NotNull(handler, nameof(handler)); |
||||
|
|
||||
|
Timer.Period = schedule.Period ?? DynamicBackgroundWorkerSchedule.DefaultPeriod; |
||||
|
CronExpression = schedule.CronExpression; |
||||
|
} |
||||
|
|
||||
|
public virtual void UpdateSchedule(DynamicBackgroundWorkerSchedule schedule) |
||||
|
{ |
||||
|
Check.NotNull(schedule, nameof(schedule)); |
||||
|
|
||||
|
Timer.Stop(); |
||||
|
Timer.Period = schedule.Period ?? DynamicBackgroundWorkerSchedule.DefaultPeriod; |
||||
|
CronExpression = schedule.CronExpression; |
||||
|
Timer.Start(StartCancellationToken); |
||||
|
} |
||||
|
|
||||
|
protected override async Task DoWorkAsync(PeriodicBackgroundWorkerContext workerContext) |
||||
|
{ |
||||
|
await _handler( |
||||
|
new DynamicBackgroundWorkerExecutionContext(WorkerName, workerContext.ServiceProvider), |
||||
|
workerContext.CancellationToken); |
||||
|
} |
||||
|
|
||||
|
public override string ToString() |
||||
|
{ |
||||
|
return $"DynamicWorker:{WorkerName}"; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,386 @@ |
|||||
|
using System; |
||||
|
using System.Collections.Concurrent; |
||||
|
using System.Collections.Generic; |
||||
|
using System.Linq; |
||||
|
using System.Threading; |
||||
|
using System.Threading.Tasks; |
||||
|
using Shouldly; |
||||
|
using Volo.Abp.BackgroundWorkers; |
||||
|
using Xunit; |
||||
|
|
||||
|
namespace Volo.Abp.BackgroundJobs; |
||||
|
|
||||
|
public class DynamicBackgroundWorkerManager_Tests : BackgroundJobsTestBase |
||||
|
{ |
||||
|
private readonly IDynamicBackgroundWorkerManager _dynamicWorkerManager; |
||||
|
|
||||
|
public DynamicBackgroundWorkerManager_Tests() |
||||
|
{ |
||||
|
_dynamicWorkerManager = GetRequiredService<IDynamicBackgroundWorkerManager>(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Register_Dynamic_Worker() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
Period = 1000 |
||||
|
}, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Execute_Dynamic_Handler() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
var tcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
Period = 50 |
||||
|
}, |
||||
|
(context, _) => |
||||
|
{ |
||||
|
if (context.WorkerName == workerName) |
||||
|
{ |
||||
|
tcs.TrySetResult(true); |
||||
|
} |
||||
|
|
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
); |
||||
|
|
||||
|
var completedTask = await Task.WhenAny(tcs.Task, Task.Delay(5000)); |
||||
|
completedTask.ShouldBe(tcs.Task); |
||||
|
(await tcs.Task).ShouldBeTrue(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Add_Dynamic_Worker_With_Default_Schedule() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Remove_Dynamic_Worker() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
Period = 1000 |
||||
|
}, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
|
||||
|
var result = await _dynamicWorkerManager.RemoveAsync(workerName); |
||||
|
result.ShouldBeTrue(); |
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeFalse(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Return_False_When_Removing_NonExistent_Worker() |
||||
|
{ |
||||
|
var result = await _dynamicWorkerManager.RemoveAsync("non-existent-worker-" + Guid.NewGuid()); |
||||
|
result.ShouldBeFalse(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Update_Dynamic_Worker_Schedule() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
var executionCount = 0; |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
Period = 60000 |
||||
|
}, |
||||
|
(_, _) => |
||||
|
{ |
||||
|
Interlocked.Increment(ref executionCount); |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
); |
||||
|
|
||||
|
var result = await _dynamicWorkerManager.UpdateScheduleAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule |
||||
|
{ |
||||
|
Period = 50 |
||||
|
} |
||||
|
); |
||||
|
|
||||
|
result.ShouldBeTrue(); |
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
|
||||
|
var timeout = TimeSpan.FromSeconds(5); |
||||
|
var startTime = DateTime.UtcNow; |
||||
|
while (executionCount == 0 && DateTime.UtcNow - startTime < timeout) |
||||
|
{ |
||||
|
await Task.Delay(50); |
||||
|
} |
||||
|
|
||||
|
executionCount.ShouldBeGreaterThan(0); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Return_False_When_Updating_NonExistent_Worker() |
||||
|
{ |
||||
|
var result = await _dynamicWorkerManager.UpdateScheduleAsync( |
||||
|
"non-existent-worker-" + Guid.NewGuid(), |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 1000 } |
||||
|
); |
||||
|
|
||||
|
result.ShouldBeFalse(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Replace_Existing_Worker_When_Same_Name_Added() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
var secondHandlerTcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 60000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 50 }, |
||||
|
(_, _) => |
||||
|
{ |
||||
|
secondHandlerTcs.TrySetResult(true); |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
); |
||||
|
|
||||
|
var completedTask = await Task.WhenAny(secondHandlerTcs.Task, Task.Delay(5000)); |
||||
|
completedTask.ShouldBe(secondHandlerTcs.Task); |
||||
|
(await secondHandlerTcs.Task).ShouldBeTrue(); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
|
||||
|
var removed = await _dynamicWorkerManager.RemoveAsync(workerName); |
||||
|
removed.ShouldBeTrue(); |
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeFalse(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Throw_When_Period_Is_Zero() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
|
||||
|
await Assert.ThrowsAsync<ArgumentException>(async () => |
||||
|
{ |
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 0 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Throw_When_Period_Is_Negative() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
|
||||
|
await Assert.ThrowsAsync<ArgumentException>(async () => |
||||
|
{ |
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = -1000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Throw_When_No_Period_And_No_CronExpression() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
|
||||
|
await Assert.ThrowsAsync<ArgumentException>(async () => |
||||
|
{ |
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule(), |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Continue_Running_After_Handler_Throws_Exception() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
var callCount = 0; |
||||
|
var tcs = new TaskCompletionSource<bool>(TaskCreationOptions.RunContinuationsAsynchronously); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 50 }, |
||||
|
(_, _) => |
||||
|
{ |
||||
|
var count = Interlocked.Increment(ref callCount); |
||||
|
if (count == 1) |
||||
|
{ |
||||
|
throw new InvalidOperationException("Simulated failure"); |
||||
|
} |
||||
|
|
||||
|
tcs.TrySetResult(true); |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
); |
||||
|
|
||||
|
var completedTask = await Task.WhenAny(tcs.Task, Task.Delay(5000)); |
||||
|
completedTask.ShouldBe(tcs.Task); |
||||
|
callCount.ShouldBeGreaterThan(1); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Not_Be_Registered_After_Remove() |
||||
|
{ |
||||
|
var workerName = "dynamic-worker-" + Guid.NewGuid(); |
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeFalse(); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 1000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
|
||||
|
await _dynamicWorkerManager.RemoveAsync(workerName); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeFalse(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Handle_Concurrent_Add_With_Same_Name() |
||||
|
{ |
||||
|
var workerName = "concurrent-worker-" + Guid.NewGuid(); |
||||
|
var executedHandlerIds = new ConcurrentBag<int>(); |
||||
|
|
||||
|
var tasks = Enumerable.Range(0, 10).Select(i => |
||||
|
_dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 60000 }, |
||||
|
(_, _) => |
||||
|
{ |
||||
|
executedHandlerIds.Add(i); |
||||
|
return Task.CompletedTask; |
||||
|
} |
||||
|
) |
||||
|
).ToList(); |
||||
|
|
||||
|
await Task.WhenAll(tasks); |
||||
|
|
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeTrue(); |
||||
|
|
||||
|
var removed = await _dynamicWorkerManager.RemoveAsync(workerName); |
||||
|
removed.ShouldBeTrue(); |
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeFalse(); |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Handle_Concurrent_Add_And_Remove() |
||||
|
{ |
||||
|
var workerNames = Enumerable.Range(0, 10) |
||||
|
.Select(i => $"concurrent-worker-{i}-" + Guid.NewGuid()) |
||||
|
.ToList(); |
||||
|
|
||||
|
var addTasks = workerNames.Select(name => |
||||
|
_dynamicWorkerManager.AddAsync( |
||||
|
name, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 60000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
) |
||||
|
).ToList(); |
||||
|
|
||||
|
await Task.WhenAll(addTasks); |
||||
|
|
||||
|
foreach (var name in workerNames) |
||||
|
{ |
||||
|
_dynamicWorkerManager.IsRegistered(name).ShouldBeTrue(); |
||||
|
} |
||||
|
|
||||
|
var removeTasks = workerNames.Select(name => |
||||
|
_dynamicWorkerManager.RemoveAsync(name) |
||||
|
).ToList(); |
||||
|
|
||||
|
var results = await Task.WhenAll(removeTasks); |
||||
|
|
||||
|
results.ShouldAllBe(r => r); |
||||
|
|
||||
|
foreach (var name in workerNames) |
||||
|
{ |
||||
|
_dynamicWorkerManager.IsRegistered(name).ShouldBeFalse(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
[Fact] |
||||
|
public async Task Should_Handle_Concurrent_Add_Remove_Update() |
||||
|
{ |
||||
|
var workerName = "concurrent-mixed-" + Guid.NewGuid(); |
||||
|
|
||||
|
await _dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 60000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
); |
||||
|
|
||||
|
var tasks = new List<Task> |
||||
|
{ |
||||
|
_dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 30000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
), |
||||
|
_dynamicWorkerManager.UpdateScheduleAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 20000 } |
||||
|
), |
||||
|
_dynamicWorkerManager.AddAsync( |
||||
|
workerName, |
||||
|
new DynamicBackgroundWorkerSchedule { Period = 10000 }, |
||||
|
(_, _) => Task.CompletedTask |
||||
|
) |
||||
|
}; |
||||
|
|
||||
|
await Task.WhenAll(tasks); |
||||
|
|
||||
|
// After all concurrent operations, worker should still be in a consistent state
|
||||
|
var isRegistered = _dynamicWorkerManager.IsRegistered(workerName); |
||||
|
isRegistered.ShouldBeTrue(); |
||||
|
|
||||
|
var removed = await _dynamicWorkerManager.RemoveAsync(workerName); |
||||
|
removed.ShouldBeTrue(); |
||||
|
_dynamicWorkerManager.IsRegistered(workerName).ShouldBeFalse(); |
||||
|
} |
||||
|
} |
||||
Loading…
Reference in new issue