using MongoDB.Driver; using MongoDB.Entities; using PrivaPub.Models.Jobs; using System.Collections.Concurrent; using System.Linq.Expressions; namespace PrivaPub.Infrastructure.Jobs { public sealed record JobOutcome(JobResult Result, string Error = default, DateTime? RetryAt = default) { public static readonly JobOutcome Done = new(JobResult.Done); public static JobOutcome Retry(string error) => new(JobResult.Retry, error); public static JobOutcome Dead(string error) => new(JobResult.Dead, error); public static JobOutcome Defer(DateTime until, string error) => new(JobResult.Defer, error, until); } public enum JobResult { Done, Retry, Dead, Defer } public interface IJobQueue { Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token); Task EnqueueMany(IEnumerable jobs, CancellationToken token); Task Lease(JobKind kind, IReadOnlyCollection busyHosts, CancellationToken token); Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token); Task Reap(CancellationToken token); Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token); } public class JobQueue : IJobQueue { public static readonly TimeSpan LeaseTime = TimeSpan.FromMinutes(2); readonly string _owner = $"{Environment.MachineName}:{Environment.ProcessId}"; readonly ConcurrentDictionary _signals = new(); readonly FilterDefinition _scope = Builders.Filter.Empty; public JobQueue() { } public JobQueue(Expression> scope) => _scope = Builders.Filter.Where(scope); public async Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token) => await EnqueueMany(new[] { new Job { Kind = kind, Payload = payload, Host = host, DedupeKey = dedupeKey } }, token) == 1; public async Task EnqueueMany(IEnumerable jobs, CancellationToken token) { var inserted = 0; foreach (var job in jobs) { try { await DB.Default.SaveAsync(job, token); inserted++; Signal(job.Kind); } catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey) { } } return inserted; } public async Task Lease(JobKind kind, IReadOnlyCollection busyHosts, CancellationToken token) { var now = DateTime.UtcNow; var busy = busyHosts?.ToList() ?? new List(); return await DB.Default.UpdateAndGet() .Match(f => f.Where(j => j.Kind == kind && j.State == JobState.Pending && j.RunAt <= now && !busy.Contains(j.Host)) & _scope) .Modify(j => j.State, JobState.Running) .Modify(j => j.LeasedUntil, now + LeaseTime) .Modify(j => j.LeaseOwner, _owner) .Modify(b => b.Inc(j => j.Attempts, 1)) .Option(o => o.Sort = Builders.Sort.Ascending(j => j.RunAt)) .ExecuteAsync(token); } public async Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token) { var now = DateTime.UtcNow; var update = DB.Default.Update().MatchID(job.ID) .Modify(j => j.LeasedUntil, null) .Modify(j => j.LeaseOwner, null) .Modify(j => j.LastError, outcome.Error); switch (outcome.Result) { case JobResult.Done: update.Modify(j => j.State, JobState.Done).Modify(j => j.FinishedAt, now); break; case JobResult.Retry when job.Attempts < maxAttempts: update.Modify(j => j.State, JobState.Pending).Modify(j => j.RunAt, now + Backoff.After(job.Attempts)); break; case JobResult.Defer: update.Modify(j => j.State, JobState.Pending) .Modify(j => j.RunAt, outcome.RetryAt ?? now + Backoff.After(job.Attempts)) .Modify(b => b.Inc(j => j.Attempts, -1)); break; default: update.Modify(j => j.State, JobState.Dead).Modify(j => j.FinishedAt, now); break; } await update.ExecuteAsync(token); } public async Task Reap(CancellationToken token) { var now = DateTime.UtcNow; var result = await DB.Default.Update() .Match(f => f.Where(j => j.State == JobState.Running && j.LeasedUntil < now) & _scope) .Modify(j => j.State, JobState.Pending) .Modify(j => j.RunAt, now) .Modify(j => j.LeasedUntil, null) .Modify(j => j.LeaseOwner, null) .ExecuteAsync(token); return result.ModifiedCount; } public async Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token) { try { await SignalFor(kind).WaitAsync(poll, token); } catch (OperationCanceledException) when (token.IsCancellationRequested) { } } void Signal(JobKind kind) { try { SignalFor(kind).Release(); } catch (SemaphoreFullException) { } } SemaphoreSlim SignalFor(JobKind kind) => _signals.GetOrAdd(kind, _ => new SemaphoreSlim(0, 64)); } }