From 84be58fcebe814d31cc4492452637aedfeb5667b Mon Sep 17 00:00:00 2001 From: Thomas Trummer Date: Fri, 14 Mar 2025 10:14:47 +0100 Subject: [PATCH 1/2] Allow to configure JobRunner worker count --- .../Extensions/DependencyInjectionExtensions.cs | 10 ++++++++-- .../HostedServices/JobRunnerHostedService.cs | 13 +++++++------ src/common/Elsa.Mediator/Options/MediatorOptions.cs | 11 ++++++++--- 3 files changed, 23 insertions(+), 11 deletions(-) diff --git a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs index 5828def98..bfff39946 100644 --- a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs +++ b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs @@ -54,8 +54,14 @@ public static class DependencyInjectionExtensions .AddSingleton() .AddSingleton() .AddSingleton() - .AddSingleton() - .AddHostedService() + .AddSingleton() + .AddHostedService(sp => + { + using var scope = sp.CreateScope(); + + var options = scope.ServiceProvider.GetRequiredService>().Value; + return ActivatorUtilities.CreateInstance(scope.ServiceProvider, options.JobWorkerCount); + }) .AddHostedService(sp => { using var scope = sp.CreateScope(); diff --git a/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs b/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs index 223eb35f4..b759367b1 100644 --- a/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs +++ b/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs @@ -9,13 +9,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(int workerCount, IJobsChannel jobsChannel, ILogger logger) + { + _workerCount = workerCount; _jobsChannel = jobsChannel; _logger = logger; } @@ -23,9 +24,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. From 84191f149d5bc8062393d7341f72b5d077a09e80 Mon Sep 17 00:00:00 2001 From: Thomas Trummer Date: Mon, 17 Mar 2025 08:45:34 +0100 Subject: [PATCH 2/2] Use IOptions to get workerCount for BackgroundCommandSenderHostedService, BackgroundEventPublisherHostedService and JobRunnerHostedService --- .../DependencyInjectionExtensions.cs | 25 +++---------------- .../BackgroundCommandSenderHostedService.cs | 8 +++--- .../BackgroundEventPublisherHostedService.cs | 8 +++--- .../HostedServices/JobRunnerHostedService.cs | 6 +++-- 4 files changed, 17 insertions(+), 30 deletions(-) diff --git a/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs b/src/common/Elsa.Mediator/Extensions/DependencyInjectionExtensions.cs index bfff39946..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; @@ -55,27 +54,9 @@ public static class DependencyInjectionExtensions .AddSingleton() .AddSingleton() .AddSingleton() - .AddHostedService(sp => - { - using var scope = sp.CreateScope(); - - var options = scope.ServiceProvider.GetRequiredService>().Value; - return ActivatorUtilities.CreateInstance(scope.ServiceProvider, options.JobWorkerCount); - }) - .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() + .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 b759367b1..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; @@ -14,9 +16,9 @@ public class JobRunnerHostedService : BackgroundService private readonly ILogger _logger; /// - public JobRunnerHostedService(int workerCount, IJobsChannel jobsChannel, ILogger logger) + public JobRunnerHostedService(IOptions options, IJobsChannel jobsChannel, ILogger logger) { - _workerCount = workerCount; + _workerCount = options.Value.JobWorkerCount; _jobsChannel = jobsChannel; _logger = logger; }