using PrivaPub.Models.Jobs; using System.Collections.Concurrent; namespace PrivaPub.Infrastructure.Jobs { public interface IJobHandler { JobKind Kind { get; } int Concurrency { get; } int MaxAttempts { get; } int PerHostLimit { get; } Task Handle(Job job, CancellationToken token); } public class JobWorker : BackgroundService { static readonly TimeSpan Poll = TimeSpan.FromSeconds(5); static readonly TimeSpan ReapInterval = TimeSpan.FromSeconds(30); readonly IJobQueue _queue; readonly IEnumerable _handlers; readonly ILogger _logger; public JobWorker(IJobQueue queue, IEnumerable handlers, ILogger logger) { _queue = queue; _handlers = handlers; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { var loops = _handlers.SelectMany(handler => { var inFlight = new ConcurrentDictionary(StringComparer.OrdinalIgnoreCase); return Enumerable.Range(0, handler.Concurrency).Select(_ => Run(handler, inFlight, stoppingToken)); }).Append(Reap(stoppingToken)); await Task.WhenAll(loops); } async Task Run(IJobHandler handler, ConcurrentDictionary inFlight, CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { Job job; try { var busy = inFlight.Where(h => h.Value >= handler.PerHostLimit).Select(h => h.Key).ToList(); job = await _queue.Lease(handler.Kind, busy, stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (Exception ex) { _logger.LogError(ex, "{Worker} could not lease a {Kind} job", nameof(JobWorker), handler.Kind); await Delay(Poll, stoppingToken); continue; } if (job == default) { await _queue.WaitForWork(handler.Kind, Poll, stoppingToken); continue; } var host = job.Host ?? string.Empty; inFlight.AddOrUpdate(host, 1, (_, count) => count + 1); try { JobOutcome outcome; try { outcome = await handler.Handle(job, stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (Exception ex) { _logger.LogError(ex, "{Kind} job {Id} threw", handler.Kind, job.ID); outcome = JobOutcome.Retry(ex.GetType().Name); } if (outcome.Result != JobResult.Done && job.Attempts >= handler.MaxAttempts && outcome.Result == JobResult.Retry) _logger.LogWarning("{Kind} job {Id} for {Host} is dead after {Attempts} attempts: {Error}", handler.Kind, job.ID, job.Host, job.Attempts, outcome.Error); await _queue.Finish(job, outcome, handler.MaxAttempts, CancellationToken.None); } catch (Exception ex) { _logger.LogError(ex, "{Kind} job {Id} could not be finished; its lease will expire", handler.Kind, job.ID); } finally { inFlight.AddOrUpdate(host, 0, (_, count) => count - 1); } } } async Task Reap(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { var reaped = await _queue.Reap(stoppingToken); if (reaped > 0) _logger.LogWarning("{Worker} returned {Count} expired leases to the queue", nameof(JobWorker), reaped); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (Exception ex) { _logger.LogError(ex, "{Worker} reaper failed", nameof(JobWorker)); } await Delay(ReapInterval, stoppingToken); } } static async Task Delay(TimeSpan delay, CancellationToken token) { try { await Task.Delay(delay, token); } catch (OperationCanceledException) { } } } }