diff --git a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs index 5828def98..2ac94c073 100644 --- a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs +++ b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs @@ -14,7 +14,6 @@ using Elsa.Mediator.Options; using Elsa.Mediator.Services; using JetBrains.Annotations; using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; // ReSharper disable once CheckNamespace namespace Microsoft.Extensions.DependencyInjection; @@ -54,22 +53,10 @@ public static class DependencyInjectionExtensions .AddSingleton() .AddSingleton() .AddSingleton() - .AddSingleton() + .AddSingleton() .AddHostedService() - .AddHostedService(sp => - { - using var scope = sp.CreateScope(); - - var options = scope.ServiceProvider.GetRequiredService>().Value; - return ActivatorUtilities.CreateInstance(scope.ServiceProvider, options.CommandWorkerCount); - }) - .AddHostedService(sp => - { - using var scope = sp.CreateScope(); - - var options = scope.ServiceProvider.GetRequiredService>().Value; - return ActivatorUtilities.CreateInstance(scope.ServiceProvider, options.NotificationWorkerCount); - }); + .AddHostedService() + .AddHostedService(); } /// diff --git a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs index 6b5092dbc..8b4a0db0c 100644 --- a/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs @@ -1,8 +1,10 @@ using System.Threading.Channels; using Elsa.Mediator.Contracts; +using Elsa.Mediator.Options; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; namespace Elsa.Mediator.HostedServices; @@ -18,13 +20,13 @@ public class BackgroundCommandSenderHostedService : BackgroundService private readonly ILogger _logger; /// - public BackgroundCommandSenderHostedService(int workerCount, ICommandsChannel commandsChannel, IServiceScopeFactory scopeFactory, ILogger logger) + public BackgroundCommandSenderHostedService(IOptions options, ICommandsChannel commandsChannel, IServiceScopeFactory scopeFactory, ILogger logger) { - _workerCount = workerCount; + _workerCount = options.Value.CommandWorkerCount; _commandsChannel = commandsChannel; _scopeFactory = scopeFactory; _logger = logger; - _outputs = new List>(workerCount); + _outputs = new List>(_workerCount); } /// diff --git a/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs b/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs index 19d6e686d..67e43ddfb 100644 --- a/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/BackgroundEventPublisherHostedService.cs @@ -1,8 +1,10 @@ using System.Threading.Channels; using Elsa.Mediator.Contracts; +using Elsa.Mediator.Options; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; namespace Elsa.Mediator.HostedServices; @@ -18,13 +20,13 @@ public class BackgroundEventPublisherHostedService : BackgroundService private readonly ILogger _logger; /// - public BackgroundEventPublisherHostedService(int workerCount, INotificationsChannel notificationsChannel, IServiceScopeFactory scopeFactory, ILogger logger) + public BackgroundEventPublisherHostedService(IOptions options, INotificationsChannel notificationsChannel, IServiceScopeFactory scopeFactory, ILogger logger) { - _workerCount = workerCount; + _workerCount = options.Value.NotificationWorkerCount; _notificationsChannel = notificationsChannel; _scopeFactory = scopeFactory; _logger = logger; - _outputs = new List>(workerCount); + _outputs = new List>(_workerCount); } /// diff --git a/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs b/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs index 223eb35f4..4eeec9061 100644 --- a/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs @@ -1,6 +1,8 @@ using Elsa.Mediator.Contracts; +using Elsa.Mediator.Options; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; namespace Elsa.Mediator.HostedServices; @@ -9,13 +11,14 @@ namespace Elsa.Mediator.HostedServices; /// public class JobRunnerHostedService : BackgroundService { - private const int WorkerCount = 4; + private readonly int _workerCount; private readonly IJobsChannel _jobsChannel; private readonly ILogger _logger; /// - public JobRunnerHostedService(IJobsChannel jobsChannel, ILogger logger) - { + public JobRunnerHostedService(IOptions options, IJobsChannel jobsChannel, ILogger logger) + { + _workerCount = options.Value.JobWorkerCount; _jobsChannel = jobsChannel; _logger = logger; } @@ -23,9 +26,9 @@ public class JobRunnerHostedService : BackgroundService /// protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - var workers = new Task[WorkerCount]; - - for (var i = 0; i < WorkerCount; i++) + var workers = new Task[_workerCount]; + + for (var i = 0; i < _workerCount; i++) workers[i] = ProcessJobsAsync(stoppingToken); await Task.WhenAll(workers); diff --git a/src/common/Elsa.Mediator/Options/MediatorOptions.cs b/src/common/Elsa.Mediator/Options/MediatorOptions.cs index a8e646683..d42f2a4aa 100644 --- a/src/common/Elsa.Mediator/Options/MediatorOptions.cs +++ b/src/common/Elsa.Mediator/Options/MediatorOptions.cs @@ -10,12 +10,17 @@ public class MediatorOptions /// /// The max number of workers to process commands. /// - public int CommandWorkerCount { get; set; } = 4; - + public int CommandWorkerCount { get; set; } = 4; + /// /// The max number of workers to process notifications. /// - public int NotificationWorkerCount { get; set; } = 4; + public int NotificationWorkerCount { get; set; } = 4; + + /// + /// The max number of workers to process jobs. + /// + public int JobWorkerCount { get; set; } = 4; /// /// The default publishing strategy to use.