using MongoDB.Driver; using MongoDB.Entities; using PrivaPub.Infrastructure.Jobs; using PrivaPub.Models.Jobs; using PrivaPub.Models.Statistics; using System.Globalization; namespace PrivaPub.Infrastructure.Statistics { public class RollupJob : IJobHandler { public static readonly TimeSpan After = TimeSpan.FromMinutes(30); readonly IJobQueue _queue; public RollupJob(IJobQueue queue) => _queue = queue; public JobKind Kind => JobKind.RollupDay; public int Concurrency => 1; public int MaxAttempts => 5; public int PerHostLimit => 1; public static string DayKey(DateTime day) => day.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture); public static Job For(DateTime day, string dedupe = default) => new() { Kind = JobKind.RollupDay, Payload = DayKey(day), DedupeKey = dedupe ?? $"rollup|{DayKey(day)}", RunAt = day.Date.AddDays(1) + After }; public async Task Handle(Job job, CancellationToken token) { if (!DateTime.TryParseExact(job.Payload, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.AdjustToUniversal | DateTimeStyles.AssumeUniversal, out var day)) return JobOutcome.Dead("not a day"); day = DateTime.SpecifyKind(day.Date, DateTimeKind.Utc); var end = day.AddDays(1); if (DateTime.UtcNow < end) return JobOutcome.Defer(end + After, "the day is not over"); await Fold(day, token); await DB.Default.DeleteAsync(s => s.Day == day); var next = end; if (next < DateTime.UtcNow.Date.AddDays(1)) await _queue.EnqueueMany(new[] { For(next) }, token); return JobOutcome.Done; } public static async Task Fold(DateTime day, CancellationToken token) { var end = day.AddDays(1); var events = new List(); using (var cursor = await DB.Default.Find().Match(e => e.At >= day && e.At < end).ExecuteCursorAsync(token)) while (await cursor.MoveNextAsync(token)) events.AddRange(cursor.Current); var fold = Rollup.Fold(events); var now = DateTime.UtcNow; foreach (var (host, hostFold) in fold.Hosts) await DB.Default.Update() .Match(d => d.Day == day && d.Host == host) .Modify(d => d.Counters, hostFold.Counters) .Modify(d => d.PublicCounters, hostFold.PublicCounters) .Modify(d => d.Latency, hostFold.Latency) .Modify(d => d.Waits, hostFold.Waits) .Modify(d => d.Bytes, hostFold.Bytes) .Modify(d => d.Accounts, hostFold.Accounts.Count) .Modify(d => d.RolledUpAt, now) .Option(o => o.IsUpsert = true) .ExecuteAsync(token); var folded = fold.Hosts.Keys.ToList(); await DB.Default.Update() .Match(d => d.Day == day && !folded.Contains(d.Host)) .Modify(d => d.Counters, new Dictionary()) .Modify(d => d.PublicCounters, new Dictionary()) .Modify(d => d.Latency, new Dictionary()) .Modify(d => d.Waits, new Dictionary()) .Modify(d => d.Bytes, new Dictionary()) .Modify(d => d.Accounts, 0) .Modify(d => d.RolledUpAt, now) .ExecuteAsync(token); await DB.Default.Update() .Match(d => d.Day == day) .Modify(d => d.Accounts, fold.Accounts.Count) .Modify(d => d.Hosts, fold.Hosts.Keys.Count(h => h != Interactions.Unknown)) .Option(o => o.IsUpsert = true) .ExecuteAsync(token); } } public class StatisticsSchedule : BackgroundService { static readonly TimeSpan Interval = TimeSpan.FromHours(1); const int CatchUpDays = 7; readonly IJobQueue _queue; readonly ILogger _logger; public StatisticsSchedule(IJobQueue queue, ILogger logger) { _queue = queue; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { await Plan(stoppingToken); } catch (Exception ex) when (ex is not OperationCanceledException) { _logger.LogWarning(ex, "Statistics rollups could not be planned"); } try { await Task.Delay(Interval, stoppingToken); } catch (OperationCanceledException) { return; } } } public async Task Plan(CancellationToken token) { var today = DateTime.UtcNow.Date; var jobs = new List { RollupJob.For(today) }; for (var back = 1; back <= CatchUpDays; back++) { var day = today.AddDays(-back); var end = day.AddDays(1); var rolled = await DB.Default.Find().Match(d => d.Day == day && d.RolledUpAt != null).ExecuteAnyAsync(token); if (!rolled && await DB.Default.Find().Match(e => e.At >= day && e.At < end).ExecuteAnyAsync(token)) jobs.Add(RollupJob.For(day)); } await _queue.EnqueueMany(jobs, token); } } }