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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
This commit is contained in:
thepraandClaude Opus 5.5 committed 2026-10-03 11:30:19 +02:00
1 parent bc8ce63707
commit fc15f6356d
8 files changed
+441 -4

No files matched your search

@@ -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<JobOutcome> 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<InteractionSalt>(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<InteractionEvent>();
using (var cursor = await DB.Default.Find<InteractionEvent>().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<InstanceDay>()
.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<InstanceDay>()
.Match(d => d.Day == day && !folded.Contains(d.Host))
.Modify(d => d.Counters, new Dictionary<string, long>())
.Modify(d => d.PublicCounters, new Dictionary<string, long>())
.Modify(d => d.Latency, new Dictionary<string, long>())
.Modify(d => d.Waits, new Dictionary<string, long>())
.Modify(d => d.Bytes, new Dictionary<string, long>())
.Modify(d => d.Accounts, 0)
.Modify(d => d.RolledUpAt, now)
.ExecuteAsync(token);
await DB.Default.Update<ServerDay>()
.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<StatisticsSchedule> _logger;
public StatisticsSchedule(IJobQueue queue, ILogger<StatisticsSchedule> 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<Job> { 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<InstanceDay>().Match(d => d.Day == day && d.RolledUpAt != null).ExecuteAnyAsync(token);
if (!rolled && await DB.Default.Find<InteractionEvent>().Match(e => e.At >= day && e.At < end).ExecuteAnyAsync(token))
jobs.Add(RollupJob.For(day));
}
await _queue.EnqueueMany(jobs, token);
}
}
}