From fc15f6356d2e4480a1df6078e1168fa673b33327 Mon Sep 17 00:00:00 2001 From: thepra Date: Sat, 3 Oct 2026 11:30:14 +0200 Subject: [PATCH] M7: daily rollups A RollupDay job folds each finished day's InteractionEvents into one InstanceDay per server: - Counters for the admin, keyed channel:activity:object:outcome:reason, plus signature schemes, audiences, local kinds and features; - PublicCounters, everything a public page may ever read: - inbound and outbound activities from an allowlist, on public, unlisted or unaddressed traffic only, with the outcome collapsed (accepted or dropped, delivered or failed); - no reasons, no Flag or Block; - features of public objects; - health:ok or health:failed from reachability (a 4xx means the server answered); - latency, wait and byte histograms, and the number of distinct accounts from that day's hashes. Re-running a day replaces it and keeps the live Reads counters. The day's salt is then deleted, so its hashes can never be recomputed, and the next day is queued. StatisticsSchedule plans today's rollup every hour and catches up any of the last seven days that have events but no rollup. Waits round up into their bucket. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2 --- .gitignore | 1 + PrivaPub.Tests/Statistics/RollupTests.cs | 152 ++++++++++++++++++ PrivaPub.Tests/Support/Host/PrivaPubHost.cs | 3 +- .../Infrastructure/Statistics/Interactions.cs | 2 +- PrivaPub/Infrastructure/Statistics/Rollup.cs | 132 +++++++++++++++ .../Infrastructure/Statistics/RollupJob.cs | 148 +++++++++++++++++ .../Middleware/SocialPubConfigurations.cs | 4 +- PrivaPub/Models/Jobs/Job.cs | 3 +- 8 files changed, 441 insertions(+), 4 deletions(-) create mode 100644 PrivaPub.Tests/Statistics/RollupTests.cs create mode 100644 PrivaPub/Infrastructure/Statistics/Rollup.cs create mode 100644 PrivaPub/Infrastructure/Statistics/RollupJob.cs diff --git a/.gitignore b/.gitignore index fe790b0..879cc77 100644 --- a/.gitignore +++ b/.gitignore @@ -402,3 +402,4 @@ FodyWeavers.xsd PrivaPub/media-store/ PrivaPub/media-store-proxy/ tools/pasture/.publish/ +.claude/worktrees/ diff --git a/PrivaPub.Tests/Statistics/RollupTests.cs b/PrivaPub.Tests/Statistics/RollupTests.cs new file mode 100644 index 0000000..a35ff3b --- /dev/null +++ b/PrivaPub.Tests/Statistics/RollupTests.cs @@ -0,0 +1,152 @@ +using MongoDB.Entities; + +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Statistics; +using PrivaPub.Tests.Support; + +namespace PrivaPub.Tests.Statistics +{ + public class RollupFoldTests + { + static IEnumerable Day(string host) => new[] + { + new InteractionEvent { Host = host, Channel = "recv", Activity = "Create", Status = 202, Reason = "queued", Signature = "cavage:rsa-sha256", LatencyMs = 40, Bytes = 900 }, + new InteractionEvent { Host = host, Channel = "recv", Activity = "Create", Status = 401, Reason = "no-signature" }, + new InteractionEvent { Host = host, Channel = "in", Activity = "Create", Object = "Note", Outcome = "accepted", Reason = "stored", Audience = "public", LocalKind = "person", ActorHash = "a", Features = new() { "fep-044f-quote", "source:text/markdown" }, WaitMs = 1500 }, + new InteractionEvent { Host = host, Channel = "in", Activity = "Create", Object = "Note", Outcome = "accepted", Reason = "stored", Audience = "private", ActorHash = "b", Features = new() { "interaction-policy" } }, + new InteractionEvent { Host = host, Channel = "in", Activity = "Flag", Outcome = "accepted", Reason = "reported", Audience = "none", ActorHash = "a" }, + new InteractionEvent { Host = host, Channel = "in", Activity = "Like", Object = "Note", Outcome = "dropped", Reason = "unknown-object", Audience = "public" }, + new InteractionEvent { Host = host, Channel = "out", Activity = "Create", Object = "Note", Outcome = "ok", Status = 202, Audience = "public", LocalKind = "person", LatencyMs = 120, Bytes = 2000 }, + new InteractionEvent { Host = host, Channel = "out", Activity = "Block", Outcome = "ok", Status = 202, Audience = "none" }, + new InteractionEvent { Host = host, Channel = "out", Activity = "Create", Object = "Note", Outcome = "retry", Reason = "503", Audience = "private" }, + new InteractionEvent { Host = host, Channel = "http", Purpose = "actor", Trigger = "verify", Outcome = "ok", LatencyMs = 300, Bytes = 1400 }, + new InteractionEvent { Host = host, Channel = "http", Purpose = "object", Trigger = "fetchancestors", Outcome = "refused", Reason = "410" } + }; + + [Fact] + public void Every_event_counts_for_the_admin_with_its_reason() + { + var fold = Rollup.Fold(Day("social.example")).Hosts["social.example"]; + + Assert.Equal(1, fold.Counters["recv:Create:202:queued"]); + Assert.Equal(1, fold.Counters["recv:Create:401:no-signature"]); + Assert.Equal(1, fold.Counters["sig:cavage:rsa-sha256"]); + Assert.Equal(2, fold.Counters["in:Create:Note:accepted:stored"]); + Assert.Equal(1, fold.Counters["in:Like:Note:dropped:unknown-object"]); + Assert.Equal(1, fold.Counters["out:Create:Note:retry:503"]); + Assert.Equal(1, fold.Counters["out:Block:-:ok"]); + Assert.Equal(1, fold.Counters["http:object:refused:410"]); + Assert.Equal(1, fold.Counters["aud:in:private"]); + Assert.Equal(1, fold.Counters["kind:out:person"]); + Assert.Equal(1, fold.Counters["feat:source:text/markdown"]); + Assert.Equal(1, fold.Latency["out:le250ms"]); + Assert.Equal(1, fold.Latency["http:actor:le500ms"]); + Assert.Equal(1, fold.Waits["in:le10s"]); + Assert.Equal(2000, fold.Bytes["out"]); + Assert.Equal(1400, fold.Bytes["http:actor"]); + Assert.Equal(2, fold.Accounts.Count); + } + + [Fact] + public void The_public_counters_keep_only_public_kinds_of_traffic_without_reasons() + { + var fold = Rollup.Fold(Day("social.example")).Hosts["social.example"]; + + Assert.Equal(new Dictionary + { + ["in:Create:Note:accepted"] = 1, + ["in:Like:Note:dropped"] = 1, + ["out:Create:Note:delivered"] = 1, + ["feat:fep-044f-quote"] = 1, + ["feat:source:text/markdown"] = 1, + ["health:ok"] = 4, + ["health:failed"] = 1 + }, fold.PublicCounters); + Assert.DoesNotContain(fold.PublicCounters.Keys, k => k.Contains("Flag") || k.Contains("Block") || k.Contains("recv") || k.Contains("interaction-policy")); + } + + [Fact] + public void Keys_are_safe_field_names_and_folding_is_repeatable() + { + Assert.Equal("in:-:Note:accepted", Rollup.Key("in", null, "Note", "accepted", null)); + Assert.Equal("feat:content:text/x_mfm", Rollup.Key("feat", "content:text/x.mfm")); + var first = Rollup.Fold(Day("a.example").Concat(Day("b.example"))); + var second = Rollup.Fold(Day("a.example").Concat(Day("b.example"))); + + Assert.Equal(first.Hosts["a.example"].Counters, second.Hosts["a.example"].Counters); + Assert.Equal(2, first.Accounts.Count); + Assert.Equal(2, first.Hosts.Count); + } + } + + [Trait("Category", "Integration")] + public sealed class RollupJobTests : IAsyncLifetime + { + static readonly DateTime Day = new(2001, 2, 3, 0, 0, 0, DateTimeKind.Utc); + readonly string _host = $"rollup{Guid.NewGuid():N}.example"; + + public ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + return ValueTask.CompletedTask; + } + + public ValueTask DisposeAsync() => ValueTask.CompletedTask; + + [Fact] + public async Task A_past_day_folds_into_server_rows_replaces_itself_keeps_reads_and_forgets_its_salt() + { + var token = TestContext.Current.CancellationToken; + await DB.Default.SaveAsync(new[] + { + new InteractionEvent { At = Day.AddHours(3), Host = _host, Channel = "in", Activity = "Create", Object = "Note", Outcome = "accepted", Audience = "public", ActorHash = "x" }, + new InteractionEvent { At = Day.AddHours(5), Host = _host, Channel = "out", Activity = "Create", Object = "Note", Outcome = "ok", Audience = "public" }, + new InteractionEvent { At = Day.AddDays(1).AddMinutes(1), Host = _host, Channel = "in", Activity = "Like", Outcome = "accepted" } + }, token); + await DB.Default.SaveAsync(new InstanceDay { Day = Day, Host = _host, Reads = new() { ["media:hit"] = 7 } }, token); + await DB.Default.SaveAsync(new InteractionSalt { Day = Day, Key = "AAAA", ExpiresAt = DateTime.UtcNow.AddHours(1) }, token); + var job = new RollupJob(new JobQueue(j => j.DedupeKey == "never")); + + Assert.Equal(JobResult.Done, (await job.Handle(RollupJob.For(Day, $"test|{Guid.NewGuid():N}"), token)).Result); + var first = await DB.Default.Find().Match(d => d.Day == Day && d.Host == _host).ExecuteSingleAsync(token); + Assert.Equal(JobResult.Done, (await job.Handle(RollupJob.For(Day, $"test|{Guid.NewGuid():N}"), token)).Result); + var second = await DB.Default.Find().Match(d => d.Day == Day && d.Host == _host).ExecuteSingleAsync(token); + + Assert.Equal(1, first.Counters["in:Create:Note:accepted"]); + Assert.False(first.Counters.ContainsKey("in:Like:-:accepted")); + Assert.Equal(1, first.PublicCounters["out:Create:Note:delivered"]); + Assert.Equal(1, first.Accounts); + Assert.Equal(7, first.Reads["media:hit"]); + Assert.Equal(first.Counters, second.Counters); + Assert.Equal(7, second.Reads["media:hit"]); + Assert.NotNull(second.RolledUpAt); + Assert.False(await DB.Default.Find().Match(s => s.Day == Day).ExecuteAnyAsync(token)); + Assert.True(await DB.Default.Find().Match(j => j.DedupeKey == "rollup|2001-02-04" && j.Kind == JobKind.RollupDay).ExecuteAnyAsync(token)); + } + + [Fact] + public async Task Today_is_not_folded_before_it_ends() + { + var outcome = await new RollupJob(new JobQueue(j => j.DedupeKey == "never")).Handle(RollupJob.For(DateTime.UtcNow.Date, "x"), TestContext.Current.CancellationToken); + + Assert.Equal(JobResult.Defer, outcome.Result); + Assert.Equal(DateTime.UtcNow.Date.AddDays(1) + RollupJob.After, outcome.RetryAt); + } + + [Fact] + public async Task The_schedule_plans_today_and_catches_up_a_missed_day() + { + var token = TestContext.Current.CancellationToken; + var missed = DateTime.UtcNow.Date.AddDays(-3); + await DB.Default.SaveAsync(new InteractionEvent { At = missed.AddHours(12), Host = _host, Channel = "in", Activity = "Create", Outcome = "accepted" }, token); + + await new StatisticsSchedule(new JobQueue(j => j.DedupeKey == "never"), Microsoft.Extensions.Logging.Abstractions.NullLogger.Instance).Plan(token); + + var today = await DB.Default.Find().Match(j => j.DedupeKey == $"rollup|{RollupJob.DayKey(DateTime.UtcNow.Date)}").ExecuteFirstAsync(token); + Assert.Equal(DateTime.UtcNow.Date.AddDays(1) + RollupJob.After, today.RunAt); + Assert.True(await DB.Default.Find().Match(j => j.DedupeKey == $"rollup|{RollupJob.DayKey(missed)}").ExecuteAnyAsync(token)); + } + } +} diff --git a/PrivaPub.Tests/Support/Host/PrivaPubHost.cs b/PrivaPub.Tests/Support/Host/PrivaPubHost.cs index b349f07..7f42dbe 100644 --- a/PrivaPub.Tests/Support/Host/PrivaPubHost.cs +++ b/PrivaPub.Tests/Support/Host/PrivaPubHost.cs @@ -10,6 +10,7 @@ using Microsoft.Extensions.Hosting; using PrivaPub.Api.Mastodon.Auth; using PrivaPub.Domain.Media; using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Infrastructure.Statistics; using System.Net; @@ -22,7 +23,7 @@ namespace PrivaPub.Tests.Support.Host public const string ClientHeader = "X-Test-Client"; static readonly SemaphoreSlim Boot = new(1, 1); - static readonly Type[] Unwanted = { typeof(JobWorker), typeof(MediaJanitor), typeof(OAuthPruner) }; + static readonly Type[] Unwanted = { typeof(JobWorker), typeof(MediaJanitor), typeof(OAuthPruner), typeof(StatisticsSchedule) }; static PrivaPubHost _shared; readonly string _mediaRoot = Path.Combine(Path.GetTempPath(), $"privapub-tests-{Guid.NewGuid():N}"); diff --git a/PrivaPub/Infrastructure/Statistics/Interactions.cs b/PrivaPub/Infrastructure/Statistics/Interactions.cs index 534025b..c2d31db 100644 --- a/PrivaPub/Infrastructure/Statistics/Interactions.cs +++ b/PrivaPub/Infrastructure/Statistics/Interactions.cs @@ -70,7 +70,7 @@ namespace PrivaPub.Infrastructure.Statistics public static string Latency(int milliseconds) => Bucket("le", milliseconds, LatencyBuckets, "ms"); - public static string Wait(int milliseconds) => Bucket("le", milliseconds / 1000, WaitBuckets, "s"); + public static string Wait(int milliseconds) => Bucket("le", (int)Math.Ceiling(milliseconds / 1000.0), WaitBuckets, "s"); static string Bucket(string prefix, int value, int[] limits, string unit) { diff --git a/PrivaPub/Infrastructure/Statistics/Rollup.cs b/PrivaPub/Infrastructure/Statistics/Rollup.cs new file mode 100644 index 0000000..37edde9 --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/Rollup.cs @@ -0,0 +1,132 @@ +using PrivaPub.Models.Statistics; + +namespace PrivaPub.Infrastructure.Statistics +{ + public sealed class HostFold + { + public Dictionary Counters { get; } = new(); + public Dictionary PublicCounters { get; } = new(); + public Dictionary Latency { get; } = new(); + public Dictionary Waits { get; } = new(); + public Dictionary Bytes { get; } = new(); + public HashSet Accounts { get; } = new(StringComparer.Ordinal); + } + + public sealed class DayFold + { + public Dictionary Hosts { get; } = new(StringComparer.Ordinal); + public HashSet Accounts { get; } = new(StringComparer.Ordinal); + } + + public static class Rollup + { + static readonly HashSet PublicActivities = new(StringComparer.Ordinal) + { + "Create", "Update", "Delete", "Announce", "Like", "EmojiReact", "Dislike", "Follow", "Accept", "Reject", "Undo", "Move", "Add", "Remove", + "Join", "Leave", "QuoteRequest" + }; + + static readonly HashSet PublicObjects = new(StringComparer.Ordinal) + { + "Note", "Article", "Page", "Question", "Video", "Audio", "Image", "Event", "Document", "Tombstone", "Person", "Group", "Service", + "Application", "Organization", "Create", "Update", "Delete", "Announce", "Like", "Dislike", "Follow", "Undo", "EmojiReact" + }; + + static readonly HashSet HealthPurposes = new(StringComparer.Ordinal) { "actor", "key", "object" }; + + public static DayFold Fold(IEnumerable events) + { + var day = new DayFold(); + foreach (var e in events) + { + var host = e.Host ?? Interactions.Unknown; + if (!day.Hosts.TryGetValue(host, out var fold)) + day.Hosts[host] = fold = new HostFold(); + Count(fold, e); + if (e.ActorHash != default) + { + fold.Accounts.Add(e.ActorHash); + day.Accounts.Add(e.ActorHash); + } + } + return day; + } + + static void Count(HostFold fold, InteractionEvent e) + { + var channel = e.Channel ?? Interactions.Other; + switch (channel) + { + case Interactions.Receive: + Add(fold.Counters, Key(channel, e.Activity, e.Status?.ToString(System.Globalization.CultureInfo.InvariantCulture), e.Reason)); + if (e.Signature != default) + Add(fold.Counters, Key("sig", e.Signature)); + break; + case Interactions.In or Interactions.Out: + Add(fold.Counters, Key(channel, e.Activity, e.Object, e.Outcome, e.Reason)); + if (e.Audience != default) + Add(fold.Counters, Key("aud", channel, e.Audience)); + if (e.LocalKind != default) + Add(fold.Counters, Key("kind", channel, e.LocalKind)); + PublicActivity(fold, e, channel); + break; + case Interactions.Http: + Add(fold.Counters, Key(channel, e.Purpose, e.Outcome, e.Reason)); + break; + default: + Add(fold.Counters, Key(channel, e.Outcome, e.Reason)); + break; + } + if (Health(e) is { } health) + Add(fold.PublicCounters, Key("health", health)); + foreach (var feature in e.Features ?? new List()) + { + Add(fold.Counters, Key("feat", feature)); + if (e.Audience == Interactions.Public) + Add(fold.PublicCounters, Key("feat", feature)); + } + var timing = channel == Interactions.Http ? Key(channel, e.Purpose) : channel; + if (e.LatencyMs is { } latency) + Add(fold.Latency, Key(timing, Interactions.Latency(latency))); + if (e.WaitMs is { } wait) + Add(fold.Waits, Key(channel, Interactions.Wait(wait))); + if (e.Bytes is { } bytes && bytes > 0) + Add(fold.Bytes, timing, bytes); + } + + static void PublicActivity(HostFold fold, InteractionEvent e, string channel) + { + if (e.Audience is not (Interactions.Public or Interactions.Unlisted or Interactions.None) + || e.Activity == default || !PublicActivities.Contains(e.Activity)) + return; + var outcome = channel == Interactions.Out + ? e.Outcome == Interactions.Ok ? "delivered" : e.Outcome is Interactions.Retry or Interactions.Dead or Interactions.Failed ? "failed" : default + : e.Outcome == Interactions.Accepted ? "accepted" : e.Outcome is Interactions.Dropped or Interactions.Rejected ? "dropped" : default; + if (outcome == default) + return; + var type = e.Object == default ? Interactions.Unknown : PublicObjects.Contains(e.Object) ? e.Object : Interactions.Other; + Add(fold.PublicCounters, Key(channel, e.Activity, type, outcome)); + } + + static string Health(InteractionEvent e) => e.Channel switch + { + Interactions.Out when e.Outcome == Interactions.Ok || e.Outcome == Interactions.Dead && e.Status is >= 400 and < 500 => "ok", + Interactions.Out when e.Outcome is Interactions.Retry or Interactions.Failed + || e.Outcome == Interactions.Dead && (e.Status >= 500 || e.Reason is "timeout" or "network") => "failed", + Interactions.Http when HealthPurposes.Contains(e.Purpose ?? string.Empty) => + e.Outcome is Interactions.Ok or Interactions.Refused ? "ok" : e.Outcome == Interactions.Failed ? "failed" : default, + _ => default + }; + + public static string Key(params string[] segments) + { + var count = segments.Length; + while (count > 1 && segments[count - 1] == default) + count--; + return string.Join(':', segments.Take(count).Select(s => (s ?? Interactions.Unknown).Replace('.', '_').Replace('$', '_'))); + } + + static void Add(Dictionary counters, string key, long by = 1) => + counters[key] = counters.GetValueOrDefault(key) + by; + } +} diff --git a/PrivaPub/Infrastructure/Statistics/RollupJob.cs b/PrivaPub/Infrastructure/Statistics/RollupJob.cs new file mode 100644 index 0000000..c1e74c1 --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/RollupJob.cs @@ -0,0 +1,148 @@ +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); + } + } +} diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index c78d0de..a8f2018 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -110,7 +110,9 @@ namespace PrivaPub.Middleware .AddSingleton() .AddSingleton() .AddSingleton(services => services.GetRequiredService()) - .AddHostedService(services => services.GetRequiredService()); + .AddHostedService(services => services.GetRequiredService()) + .AddSingleton() + .AddHostedService(); public static IServiceCollection PrivaPubAuthServicesConfiguration(this IServiceCollection service, IConfiguration configuration) { diff --git a/PrivaPub/Models/Jobs/Job.cs b/PrivaPub/Models/Jobs/Job.cs index f1fcebd..8b9e6cc 100644 --- a/PrivaPub/Models/Jobs/Job.cs +++ b/PrivaPub/Models/Jobs/Job.cs @@ -26,7 +26,8 @@ namespace PrivaPub.Models.Jobs DescribeInstance, PollRefresh, PollClose, - FetchPreview + FetchPreview, + RollupDay } public enum JobState