elsa-core/src/common/Elsa.Mediator/HostedServices/JobRunnerHostedService.cs
Sipke Schoorstra fc5612687d
Refactors notification system for context (#6821)
* Fix `CancellationToken` usage in `BackgroundCommandSenderHostedService`

Corrected the `CancellationToken` parameter to use `commandContext.CancellationToken` instead of the method's cancellation token, ensuring proper propagation and handling within the command sender.

Fixes #6449

* Refactor notifications system to use `NotificationContext` across channels and middleware

* Update src/common/Elsa.Mediator/Extensions/HandlerExtensions.cs

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* Simplify `InvokeAsync` method call by replacing braces with brackets in argument array.

* Refactor `NotificationContext` usage in mediator strategies and handler extensions

Replaced direct references to `Notification` and `CancellationToken` with `NotificationContext` across mediator strategies to align with updated context structure. Simplified handler invocations by passing `CommandContext` and `NotificationContext` where applicable.

* Refactor notification extraction in publishing strategies

Standardize notification retrieval by introducing `notificationContext.Notification` in `SequentialProcessingStrategy` and `ParallelProcessingStrategy` for improved clarity and consistency.

* Refactor hosted services to improve worker channel management and cancellation handling

Implemented better round-robin distribution and added linked cancellation token sources in `BackgroundCommandSender`, `JobRunner`, and `BackgroundEventPublisher` hosted services. Enhanced logging and comments for clarity.

* Update src/common/Elsa.Mediator/HostedServices/BackgroundCommandSenderHostedService.cs

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

---------

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
2025-08-05 22:22:13 +02:00

71 lines
2.7 KiB
C#

using Elsa.Mediator.Contracts;
using Elsa.Mediator.Options;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace Elsa.Mediator.HostedServices;
/// <summary>
/// A hosted service that runs jobs.
/// </summary>
public class JobRunnerHostedService : BackgroundService
{
private readonly int _workerCount;
private readonly IJobsChannel _jobsChannel;
private readonly ILogger<JobRunnerHostedService> _logger;
/// <inheritdoc />
public JobRunnerHostedService(IOptions<MediatorOptions> options, IJobsChannel jobsChannel, ILogger<JobRunnerHostedService> logger)
{
_workerCount = options.Value.JobWorkerCount;
_jobsChannel = jobsChannel;
_logger = logger;
}
/// <inheritdoc />
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
// Create an array of worker tasks to process jobs in parallel
var workers = new Task[_workerCount];
// Start multiple workers (tasks) to process jobs concurrently
for (var i = 0; i < _workerCount; i++)
workers[i] = ProcessJobsAsync(stoppingToken);
// Wait for all worker tasks to complete (typically when the application is shutting down)
await Task.WhenAll(workers);
}
/// <summary>
/// Continuously processes jobs from the job channel until cancellation is requested.
/// </summary>
/// <param name="stoppingToken">Cancellation token from the hosted service</param>
private async Task ProcessJobsAsync(CancellationToken stoppingToken)
{
// Process all jobs from the channel until it's completed or cancellation is requested
await foreach (var jobItem in _jobsChannel.Reader.ReadAllAsync(stoppingToken))
{
try
{
// Link the cancellation tokens so that cancellation can happen from either source
using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken, jobItem.CancellationTokenSource.Token);
await jobItem.Action(linkedTokenSource.Token);
_logger.LogInformation("Worker {CurrentTaskId} processed job {JobId}", Task.CurrentId, jobItem.JobId);
}
catch (OperationCanceledException)
{
_logger.LogInformation("Job {JobId} was canceled", jobItem.JobId);
}
catch (Exception ex)
{
_logger.LogError(ex, "Job {JobId} failed", jobItem.JobId);
}
finally
{
// Notify that the job has completed (whether successfully or not)
jobItem.OnJobCompleted(jobItem.JobId);
}
}
}
}