using PrivaPub.Infrastructure.Http; 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); // the lease is renewed while the handler runs; once it is lost (another worker leased the job again) the // handler is cancelled and its outcome dropped using var running = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); var renewing = KeepLease(job, running); try { JobOutcome outcome; try { using var scope = HttpScope.Triggered(handler.Kind.ToString().ToLowerInvariant()); outcome = await handler.Handle(job, running.Token); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (OperationCanceledException) when (running.IsCancellationRequested) { _logger.LogWarning("{Kind} job {Id} lost its lease and was stopped", handler.Kind, job.ID); continue; } 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); if (!await _queue.Finish(job, outcome, handler.MaxAttempts, CancellationToken.None)) _logger.LogWarning("{Kind} job {Id} lost its lease before it finished; its outcome is dropped", handler.Kind, job.ID); } catch (Exception ex) { _logger.LogError(ex, "{Kind} job {Id} could not be finished; its lease will expire", handler.Kind, job.ID); } finally { await running.CancelAsync(); await renewing; inFlight.AddOrUpdate(host, 0, (_, count) => count - 1); } } } // renews the lease every third of its length until the handler ends; a lease found lost cancels the handler async Task KeepLease(Job job, CancellationTokenSource running) { var every = _queue.LeaseFor / 3; while (!running.IsCancellationRequested) { try { await Task.Delay(every, running.Token); } catch (OperationCanceledException) { return; } try { if (await _queue.Renew(job, CancellationToken.None)) continue; await running.CancelAsync(); return; } catch (Exception ex) { // a passing database error: the lease has two more thirds to go, so try again at the next turn _logger.LogWarning(ex, "{Worker} could not renew the lease of job {Id}", nameof(JobWorker), job.ID); } } } 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) { } } } }