diff --git a/docs/en/framework/infrastructure/background-workers/index.md b/docs/en/framework/infrastructure/background-workers/index.md index 6204857d8c..ab5f42b207 100644 --- a/docs/en/framework/infrastructure/background-workers/index.md +++ b/docs/en/framework/infrastructure/background-workers/index.md @@ -120,6 +120,64 @@ So, it resolves the given background worker and adds to the `IBackgroundWorkerMa While we generally add workers in `OnApplicationInitializationAsync`, there are no restrictions on that. You can inject `IBackgroundWorkerManager` anywhere and add workers at runtime. Background worker manager will stop and release all the registered workers when your application is being shut down. +### Dynamic Workers (Runtime Registration) + +You can add a runtime worker without pre-defining a dedicated worker class. Inject `IDynamicBackgroundWorkerManager` and pass a handler directly: + +````csharp +public class MyModule : AbpModule +{ + public override async Task OnApplicationInitializationAsync( + ApplicationInitializationContext context) + { + var dynamicWorkerManager = context.ServiceProvider + .GetRequiredService(); + + await dynamicWorkerManager.AddAsync( + "InventorySyncWorker", + new DynamicBackgroundWorkerSchedule + { + Period = 30000 //30 seconds + //CronExpression = "*/30 * * * *" //Every 30 minutes. Only for Hangfire, Quartz or TickerQ integration. + }, + async (workerContext, cancellationToken) => + { + var inventorySyncAppService = workerContext + .ServiceProvider + .GetRequiredService(); + + await inventorySyncAppService.SyncAsync(cancellationToken); + } + ); + } +} +```` + +You can also **remove** a dynamic worker or **update its schedule** at runtime: + +````csharp +//Remove a dynamic worker +var removed = await dynamicWorkerManager.RemoveAsync("InventorySyncWorker"); + +//Update the schedule of a dynamic worker +var updated = await dynamicWorkerManager.UpdateScheduleAsync( + "InventorySyncWorker", + new DynamicBackgroundWorkerSchedule + { + Period = 60000 //Change to 60 seconds + } +); +```` + +* `IDynamicBackgroundWorkerManager` is a **separate interface** from `IBackgroundWorkerManager`, dedicated to runtime (non-type-safe) worker management. +* `workerName` is the runtime identifier of the dynamic worker. If a worker with the same name already exists, it will be **replaced**. +* The `handler` receives a `DynamicBackgroundWorkerExecutionContext` containing the worker name and a scoped `IServiceProvider`. It is a good practice to **resolve dependencies** from the `workerContext.ServiceProvider` instead of constructor injection. +* At least one of `Period` or `CronExpression` must be set in `DynamicBackgroundWorkerSchedule`. +* **`CronExpression` is only supported by scheduler-backed providers ([Hangfire](./hangfire.md), [Quartz](./quartz.md)).** The default in-memory provider requires `Period` and does not support `CronExpression` alone. +* **[TickerQ](./tickerq.md) does not support dynamic background workers** because it uses `FrozenDictionary` for function registration, which requires all functions to be registered before the application starts. +* `RemoveAsync` stops and removes a dynamic worker. Returns `true` if the worker was found and removed. +* `UpdateScheduleAsync` changes the schedule of an existing dynamic worker. Returns `true` if the worker was found and updated. The handler itself is not changed. + ## Options `AbpBackgroundWorkerOptions` class is used to [set options](../../fundamentals/options.md) for the background workers. Currently, there is only one option: diff --git a/framework/src/Volo.Abp.BackgroundWorkers.Hangfire/Volo/Abp/BackgroundWorkers/Hangfire/HangfireDynamicBackgroundWorkerAdapter.cs b/framework/src/Volo.Abp.BackgroundWorkers.Hangfire/Volo/Abp/BackgroundWorkers/Hangfire/HangfireDynamicBackgroundWorkerAdapter.cs new file mode 100644 index 0000000000..9e6d5f2f0c --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers.Hangfire/Volo/Abp/BackgroundWorkers/Hangfire/HangfireDynamicBackgroundWorkerAdapter.cs @@ -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 Logger { get; set; } + + public HangfireDynamicBackgroundWorkerAdapter( + IDynamicBackgroundWorkerHandlerRegistry handlerRegistry, + IServiceProvider serviceProvider) + { + HandlerRegistry = handlerRegistry; + ServiceProvider = serviceProvider; + Logger = NullLogger.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() + .NotifyAsync(new ExceptionNotificationContext(ex)); + + throw; + } + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers.Hangfire/Volo/Abp/BackgroundWorkers/Hangfire/HangfireDynamicBackgroundWorkerManager.cs b/framework/src/Volo.Abp.BackgroundWorkers.Hangfire/Volo/Abp/BackgroundWorkers/Hangfire/HangfireDynamicBackgroundWorkerManager.cs new file mode 100644 index 0000000000..77e7d37a97 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers.Hangfire/Volo/Abp/BackgroundWorkers/Hangfire/HangfireDynamicBackgroundWorkerManager.cs @@ -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 Logger { get; set; } + + public HangfireDynamicBackgroundWorkerManager( + IServiceProvider serviceProvider, + IDynamicBackgroundWorkerHandlerRegistry handlerRegistry) + { + ServiceProvider = serviceProvider; + HandlerRegistry = handlerRegistry; + Logger = NullLogger.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 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 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>().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( + recurringJobId, + adapter => adapter.DoWorkAsync(workerName, CancellationToken.None), + cronExpression, + new RecurringJobOptions + { + TimeZone = TimeZoneInfo.Utc + }); + } + else + { + RecurringJob.AddOrUpdate( + 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."); + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers.Quartz/Volo/Abp/BackgroundWorkers/Quartz/QuartzDynamicBackgroundWorkerAdapter.cs b/framework/src/Volo.Abp.BackgroundWorkers.Quartz/Volo/Abp/BackgroundWorkers/Quartz/QuartzDynamicBackgroundWorkerAdapter.cs new file mode 100644 index 0000000000..ba8d98091f --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers.Quartz/Volo/Abp/BackgroundWorkers/Quartz/QuartzDynamicBackgroundWorkerAdapter.cs @@ -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 Logger { get; set; } + + public QuartzDynamicBackgroundWorkerAdapter( + IDynamicBackgroundWorkerHandlerRegistry handlerRegistry, + IServiceProvider serviceProvider) + { + HandlerRegistry = handlerRegistry; + ServiceProvider = serviceProvider; + Logger = NullLogger.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() + .NotifyAsync(new ExceptionNotificationContext(ex)); + + throw; + } + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers.Quartz/Volo/Abp/BackgroundWorkers/Quartz/QuartzDynamicBackgroundWorkerManager.cs b/framework/src/Volo.Abp.BackgroundWorkers.Quartz/Volo/Abp/BackgroundWorkers/Quartz/QuartzDynamicBackgroundWorkerManager.cs new file mode 100644 index 0000000000..8f409e011d --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers.Quartz/Volo/Abp/BackgroundWorkers/Quartz/QuartzDynamicBackgroundWorkerManager.cs @@ -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 Logger { get; set; } + + public QuartzDynamicBackgroundWorkerManager( + IScheduler scheduler, + IDynamicBackgroundWorkerHandlerRegistry handlerRegistry) + { + Scheduler = scheduler; + HandlerRegistry = handlerRegistry; + Logger = NullLogger.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() + .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 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 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(); + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers.TickerQ/Volo/Abp/BackgroundWorkers/TickerQ/TickerQDynamicBackgroundWorkerManager.cs b/framework/src/Volo.Abp.BackgroundWorkers.TickerQ/Volo/Abp/BackgroundWorkers/TickerQ/TickerQDynamicBackgroundWorkerManager.cs new file mode 100644 index 0000000000..b5f1cfcbe4 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers.TickerQ/Volo/Abp/BackgroundWorkers/TickerQ/TickerQDynamicBackgroundWorkerManager.cs @@ -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 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 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; + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DefaultDynamicBackgroundWorkerManager.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DefaultDynamicBackgroundWorkerManager.cs new file mode 100644 index 0000000000..a3073d0e53 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DefaultDynamicBackgroundWorkerManager.cs @@ -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 Logger { get; set; } + + private readonly ConcurrentDictionary _dynamicWorkers; + private readonly SemaphoreSlim _semaphore; + private bool _isDisposed; + + public DefaultDynamicBackgroundWorkerManager(IServiceProvider serviceProvider) + { + ServiceProvider = serviceProvider; + Logger = NullLogger.Instance; + _dynamicWorkers = new ConcurrentDictionary(); + _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 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 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(); + var serviceScopeFactory = ServiceProvider.GetRequiredService(); + + var worker = new InMemoryDynamicBackgroundWorker( + workerName, schedule, handler, timer, serviceScopeFactory); + + worker.ServiceProvider = ServiceProvider; + worker.LazyServiceProvider = ServiceProvider.GetRequiredService(); + + return worker; + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerExecutionContext.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerExecutionContext.cs new file mode 100644 index 0000000000..edb810d105 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerExecutionContext.cs @@ -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)); + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerHandler.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerHandler.cs new file mode 100644 index 0000000000..28f5c933d8 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerHandler.cs @@ -0,0 +1,6 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Volo.Abp.BackgroundWorkers; + +public delegate Task DynamicBackgroundWorkerHandler(DynamicBackgroundWorkerExecutionContext context, CancellationToken cancellationToken); \ No newline at end of file diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerHandlerRegistry.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerHandlerRegistry.cs new file mode 100644 index 0000000000..a30ea3e95a --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerHandlerRegistry.cs @@ -0,0 +1,40 @@ +using System.Collections.Concurrent; +using Volo.Abp.DependencyInjection; + +namespace Volo.Abp.BackgroundWorkers; + +public class DynamicBackgroundWorkerHandlerRegistry : IDynamicBackgroundWorkerHandlerRegistry, ISingletonDependency +{ + protected ConcurrentDictionary Handlers { get; } + + public DynamicBackgroundWorkerHandlerRegistry() + { + Handlers = new ConcurrentDictionary(); + } + + 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; + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerManagerExtensions.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerManagerExtensions.cs new file mode 100644 index 0000000000..5fedbf1360 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerManagerExtensions.cs @@ -0,0 +1,26 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Volo.Abp.BackgroundWorkers; + +public static class DynamicBackgroundWorkerManagerExtensions +{ + /// + /// Adds a dynamic worker with the default schedule (). + /// + 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); + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerSchedule.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerSchedule.cs new file mode 100644 index 0000000000..6505492e44 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/DynamicBackgroundWorkerSchedule.cs @@ -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."); + } + } +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/IDynamicBackgroundWorkerHandlerRegistry.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/IDynamicBackgroundWorkerHandlerRegistry.cs new file mode 100644 index 0000000000..7915f108db --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/IDynamicBackgroundWorkerHandlerRegistry.cs @@ -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); +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/IDynamicBackgroundWorkerManager.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/IDynamicBackgroundWorkerManager.cs new file mode 100644 index 0000000000..5054f67897 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/IDynamicBackgroundWorkerManager.cs @@ -0,0 +1,41 @@ +using System.Threading; +using System.Threading.Tasks; + +namespace Volo.Abp.BackgroundWorkers; + +/// +/// Manages dynamic background workers that are registered at runtime +/// without requiring a strongly-typed worker class. +/// +public interface IDynamicBackgroundWorkerManager +{ + /// + /// Adds a dynamic worker by name, schedule and handler. + /// If a worker with the same name already exists, it will be replaced. + /// + Task AddAsync( + string workerName, + DynamicBackgroundWorkerSchedule schedule, + DynamicBackgroundWorkerHandler handler, + CancellationToken cancellationToken = default); + + /// + /// Removes a previously added dynamic worker by name. + /// Returns true if the worker was found and removed; false otherwise. + /// + Task RemoveAsync(string workerName, CancellationToken cancellationToken = default); + + /// + /// Updates the schedule of a previously added dynamic worker. + /// Returns true if the worker was found and updated; false otherwise. + /// + Task UpdateScheduleAsync( + string workerName, + DynamicBackgroundWorkerSchedule schedule, + CancellationToken cancellationToken = default); + + /// + /// Checks whether a dynamic worker with the given name is registered. + /// + bool IsRegistered(string workerName); +} diff --git a/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/InMemoryDynamicBackgroundWorker.cs b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/InMemoryDynamicBackgroundWorker.cs new file mode 100644 index 0000000000..bab468a655 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundWorkers/Volo/Abp/BackgroundWorkers/InMemoryDynamicBackgroundWorker.cs @@ -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}"; + } +} diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/DynamicBackgroundWorkerManager_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/DynamicBackgroundWorkerManager_Tests.cs new file mode 100644 index 0000000000..be99167d85 --- /dev/null +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/DynamicBackgroundWorkerManager_Tests.cs @@ -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(); + } + + [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(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(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(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(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(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(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(); + + 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 + { + _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(); + } +} diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/DemoAppQuartzModule.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/DemoAppQuartzModule.cs index 0d521baf3d..b8aaf6e519 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/DemoAppQuartzModule.cs +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/DemoAppQuartzModule.cs @@ -1,6 +1,7 @@ using Volo.Abp.Autofac; using Volo.Abp.BackgroundJobs.DemoApp.Shared; using Volo.Abp.BackgroundJobs.Quartz; +using Volo.Abp.BackgroundWorkers.Quartz; using Volo.Abp.Modularity; namespace Volo.Abp.BackgroundJobs.DemoApp.Quartz; @@ -8,7 +9,8 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Quartz; [DependsOn( typeof(DemoAppSharedModule), typeof(AbpAutofacModule), - typeof(AbpBackgroundJobsQuartzModule) + typeof(AbpBackgroundJobsQuartzModule), + typeof(AbpBackgroundWorkersQuartzModule) )] public class DemoAppQuartzModule : AbpModule { diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/Volo.Abp.BackgroundJobs.DemoApp.Quartz.csproj b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/Volo.Abp.BackgroundJobs.DemoApp.Quartz.csproj index 8faef0c291..5f2be2b49b 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/Volo.Abp.BackgroundJobs.DemoApp.Quartz.csproj +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Quartz/Volo.Abp.BackgroundJobs.DemoApp.Quartz.csproj @@ -12,6 +12,7 @@ + diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/DemoAppSharedModule.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/DemoAppSharedModule.cs index 3c5e91958a..c69b92c7ed 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/DemoAppSharedModule.cs +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp.Shared/DemoAppSharedModule.cs @@ -2,7 +2,11 @@ using System; using System.Text.Json; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Volo.Abp.BackgroundJobs.DemoApp.Shared.Jobs; +using Volo.Abp.BackgroundWorkers; using Volo.Abp.Modularity; using Volo.Abp.MultiTenancy; @@ -11,7 +15,7 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Shared [DependsOn(typeof(AbpMultiTenancyModule))] public class DemoAppSharedModule : AbpModule { - public override void OnApplicationInitialization(ApplicationInitializationContext context) + public override async Task OnApplicationInitializationAsync(ApplicationInitializationContext context) { var dynamicJobManager = context.ServiceProvider.GetRequiredService(); @@ -28,13 +32,69 @@ namespace Volo.Abp.BackgroundJobs.DemoApp.Shared return Task.CompletedTask; } }); - } - public override void OnPostApplicationInitialization(ApplicationInitializationContext context) - { context.ServiceProvider .GetRequiredService() .CreateJobs(); + + await DynamicBackgroundWorkerDemoAsync(context); + } + + private async Task DynamicBackgroundWorkerDemoAsync(ApplicationInitializationContext context) + { + var dynamicWorkerManager = context.ServiceProvider + .GetService(); + + if (dynamicWorkerManager == null) + { + return; + } + + // AddAsync: Register a dynamic worker with a schedule and handler + await dynamicWorkerManager.AddAsync( + "DemoHeartbeatWorker", + new DynamicBackgroundWorkerSchedule + { + Period = 5000 //5 seconds + }, + async (workerContext, cancellationToken) => + { + Console.WriteLine($"[{DateTime.Now}] DemoHeartbeatWorker executed."); + await Task.CompletedTask; + } + ); + + // IsRegistered: Check if a dynamic worker is registered + var isRegistered = dynamicWorkerManager.IsRegistered("DemoHeartbeatWorker"); + Console.WriteLine($"DemoHeartbeatWorker is registered: {isRegistered}"); + + // UpdateScheduleAsync: Update the schedule of an existing dynamic worker + var updated = await dynamicWorkerManager.UpdateScheduleAsync( + "DemoHeartbeatWorker", + new DynamicBackgroundWorkerSchedule + { + Period = 10000 //Change to 10 seconds + } + ); + Console.WriteLine($"DemoHeartbeatWorker schedule updated: {updated}"); + + // RemoveAsync: Remove a dynamic worker + var removed = await dynamicWorkerManager.RemoveAsync("DemoHeartbeatWorker"); + Console.WriteLine($"DemoHeartbeatWorker removed: {removed}"); + + // Re-add the worker to keep it running for demo purposes + await dynamicWorkerManager.AddAsync( + "DemoHeartbeatWorker", + new DynamicBackgroundWorkerSchedule + { + Period = 10000 //10 seconds + }, + async (workerContext, cancellationToken) => + { + Console.WriteLine($"[{DateTime.Now}] DemoHeartbeatWorker executed."); + await Task.CompletedTask; + } + ); } } } diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260119064307_Initial.Designer.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260320082618_Initial.Designer.cs similarity index 98% rename from modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260119064307_Initial.Designer.cs rename to modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260320082618_Initial.Designer.cs index 3225815926..fd21dc852f 100644 --- a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260119064307_Initial.Designer.cs +++ b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260320082618_Initial.Designer.cs @@ -13,7 +13,7 @@ using Volo.Abp.EntityFrameworkCore; namespace Volo.Abp.BackgroundJobs.DemoApp.Migrations { [DbContext(typeof(DemoAppDbContext))] - [Migration("20260119064307_Initial")] + [Migration("20260320082618_Initial")] partial class Initial { /// diff --git a/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260119064307_Initial.cs b/modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260320082618_Initial.cs similarity index 100% rename from modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260119064307_Initial.cs rename to modules/background-jobs/app/Volo.Abp.BackgroundJobs.DemoApp/Migrations/20260320082618_Initial.cs