From 6ef7e2891a0f9892b83a8d728eac9e4746db957d Mon Sep 17 00:00:00 2001 From: thepra Date: Sat, 3 Oct 2026 11:35:40 +0200 Subject: [PATCH] M10: the admin statistics API Under /clientapi/admin/statistics, admin only (the root JWT's IsAdmin policy, like the domain blocks; /api tokens are persona tokens): - GET overview?days: totals per channel, inbound and outbound outcomes, the delivery success rate, delivery latency p50/p95 from the buckets, active and known servers, distinct accounts, and the ledger's written/dropped/failed counts; - GET hosts?days&sort=volume|failures|latency|host&software&page&limit: per server traffic, refusals, failures, latency, accounts, software and delivery health; - GET hosts/{host}?days: the server's description, its daily series and its last 100 events; - GET events?host&channel&outcome&reason&before&limit, and GET server?days for the ServerDay rows; - POST rollups/{day} refolds a past day; POST hosts/{host}/describe describes a server again. Days not yet rolled up, today included, are folded live from the events, so the numbers are current. Answers that may carry locations carry the DB-IP attribution. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2 --- .../Statistics/AdminStatisticsTests.cs | 126 +++++++++++ .../ClientToServer/StatisticsController.cs | 72 +++++++ .../Domain/Statistics/StatisticsQueries.cs | 195 ++++++++++++++++++ .../Middleware/SocialPubConfigurations.cs | 1 + 4 files changed, 394 insertions(+) create mode 100644 PrivaPub.Tests/Statistics/AdminStatisticsTests.cs create mode 100644 PrivaPub/Controllers/ClientToServer/StatisticsController.cs create mode 100644 PrivaPub/Domain/Statistics/StatisticsQueries.cs diff --git a/PrivaPub.Tests/Statistics/AdminStatisticsTests.cs b/PrivaPub.Tests/Statistics/AdminStatisticsTests.cs new file mode 100644 index 0000000..23c15c0 --- /dev/null +++ b/PrivaPub.Tests/Statistics/AdminStatisticsTests.cs @@ -0,0 +1,126 @@ +using MongoDB.Entities; + +using PrivaPub.Domain.Statistics; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Jobs; +using PrivaPub.Tests.Support; +using PrivaPub.Tests.Support.Host; + +using System.Net; +using System.Text.Json.Nodes; + +namespace PrivaPub.Tests.Statistics +{ + public class StatisticsQueriesTests + { + [Fact] + public void Percentiles_come_from_the_latency_buckets() + { + var latency = new Dictionary { ["out:le50ms"] = 50, ["out:le250ms"] = 45, ["out:le5000ms"] = 4, ["out:inf"] = 1, ["in:le50ms"] = 9 }; + + Assert.Equal(50, StatisticsQueries.Percentile(latency, "out", 0.5)); + Assert.Equal(250, StatisticsQueries.Percentile(latency, "out", 0.95)); + Assert.Equal(5000, StatisticsQueries.Percentile(latency, "out", 0.99)); + Assert.Equal(30001, StatisticsQueries.Percentile(latency, "out", 1)); + Assert.Null(StatisticsQueries.Percentile(latency, "http:actor", 0.5)); + } + + [Fact] + public void Outcomes_add_up_per_channel() + { + var counters = new Dictionary + { + ["out:Create:Note:ok"] = 3, ["out:Create:Note:retry:503"] = 2, ["out:Like:-:ok"] = 1, ["in:Create:Note:accepted:stored"] = 4, ["recv:Create:202:queued"] = 9 + }; + + Assert.Equal(new Dictionary { ["ok"] = 4, ["retry"] = 2 }, StatisticsQueries.Outcomes(counters, "out")); + Assert.Equal(new Dictionary { ["accepted"] = 4 }, StatisticsQueries.Outcomes(counters, "in")); + } + } + + [Trait("Category", "Integration")] + public sealed class AdminStatisticsTests : IAsyncLifetime + { + PrivaPubHost _host; + Peer _peer; + + public async ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + _host = await PrivaPubHost.Shared(); + _peer = await Peer.Start(); + } + + public async ValueTask DisposeAsync() + { + if (_peer != default) + await _peer.DisposeAsync(); + } + + [Fact] + public async Task Only_the_admin_reads_the_statistics() + { + var token = TestContext.Current.CancellationToken; + var plain = await _host.SignUp(); + using var anonymous = _host.Client(); + using var user = _host.As(plain.Jwt); + + Assert.Equal(HttpStatusCode.Unauthorized, (await anonymous.GetAsync("/clientapi/admin/statistics/overview", token)).StatusCode); + Assert.Equal(HttpStatusCode.Forbidden, (await user.GetAsync("/clientapi/admin/statistics/overview", token)).StatusCode); + Assert.Equal(HttpStatusCode.Forbidden, (await user.PostAsync("/clientapi/admin/statistics/rollups/2001-01-01", default, token)).StatusCode); + } + + [Fact] + public async Task The_admin_sees_a_server_its_events_and_the_overview() + { + var token = TestContext.Current.CancellationToken; + var admin = await _host.Admin(); + var persona = await _host.Persona(await _host.SignUp(), "stats"); + var bob = new RemoteActor(_peer, "bob"); + var noteId = $"{_peer.A}/notes/{Guid.NewGuid():N}"; + var to = new JsonArray($"{PrivaPubHost.Base}/peasants/{persona.UserName}"); + using (var client = _host.Client()) + Assert.Equal(HttpStatusCode.Accepted, (await client.SendAsync(bob.SignedPost($"/peasants/{persona.UserName}/mouth", new JsonObject + { + ["id"] = noteId + "/activity", ["type"] = "Create", ["actor"] = bob.Id, ["to"] = to.DeepClone(), + ["object"] = new JsonObject { ["id"] = noteId, ["type"] = "Note", ["attributedTo"] = bob.Id, ["to"] = to.DeepClone(), ["content"] = "hi" } + }), token)).StatusCode); + await _host.RunInbox(noteId + "/activity", token); + await _host.Get().Flush(token); + using var adminClient = _host.As(admin.Jwt); + + var host = JsonNode.Parse(await adminClient.GetStringAsync("/clientapi/admin/statistics/hosts/127.0.0.1?days=1", token))!; + var overview = JsonNode.Parse(await adminClient.GetStringAsync("/clientapi/admin/statistics/overview?days=1", token))!; + var hosts = JsonNode.Parse(await adminClient.GetStringAsync("/clientapi/admin/statistics/hosts?days=1&sort=host", token))!; + var events = JsonNode.Parse(await adminClient.GetStringAsync("/clientapi/admin/statistics/events?host=127.0.0.1&channel=in&limit=5", token))!; + + var today = host["days"]!.AsArray().Last()!; + Assert.True(today["live"]!.GetValue()); + Assert.True(today["counters"]!["in:Create:Note:accepted:stored"]!.GetValue() >= 1); + Assert.Contains("DB-IP", host["attribution"]!.GetValue()); + Assert.True(overview["channels"]!["in"]!.GetValue() >= 1); + Assert.Contains(hosts["hosts"]!.AsArray(), h => h!["host"]!.GetValue() == "127.0.0.1"); + Assert.NotEmpty(events.AsArray()); + Assert.DoesNotContain(persona.UserName, host.ToJsonString()); + Assert.DoesNotContain(bob.Id, host.ToJsonString()); + Assert.Equal(HttpStatusCode.NotFound, (await adminClient.GetAsync($"/clientapi/admin/statistics/hosts/never{Guid.NewGuid():N}.example", token)).StatusCode); + } + + [Fact] + public async Task The_admin_can_refold_a_past_day_and_describe_a_server_again() + { + var token = TestContext.Current.CancellationToken; + var admin = await _host.Admin(); + using var client = _host.As(admin.Jwt); + var describe = $"described{Guid.NewGuid():N}.example"; + + Assert.Equal(HttpStatusCode.BadRequest, (await client.PostAsync($"/clientapi/admin/statistics/rollups/{DateTime.UtcNow:yyyy-MM-dd}", default, token)).StatusCode); + Assert.Equal(HttpStatusCode.BadRequest, (await client.PostAsync("/clientapi/admin/statistics/rollups/yesterday", default, token)).StatusCode); + Assert.Equal(HttpStatusCode.Accepted, (await client.PostAsync("/clientapi/admin/statistics/rollups/2001-01-05", default, token)).StatusCode); + Assert.Equal(HttpStatusCode.Accepted, (await client.PostAsync($"/clientapi/admin/statistics/hosts/{describe}/describe", default, token)).StatusCode); + + Assert.True(await DB.Default.Find().Match(j => j.Kind == JobKind.RollupDay && j.Payload == "2001-01-05" && j.DedupeKey.StartsWith("rollup|2001-01-05|rerun|")).ExecuteAnyAsync(token)); + Assert.True(await DB.Default.Find().Match(j => j.Kind == JobKind.DescribeInstance && j.Payload == describe).ExecuteAnyAsync(token)); + } + } +} diff --git a/PrivaPub/Controllers/ClientToServer/StatisticsController.cs b/PrivaPub/Controllers/ClientToServer/StatisticsController.cs new file mode 100644 index 0000000..7f40b31 --- /dev/null +++ b/PrivaPub/Controllers/ClientToServer/StatisticsController.cs @@ -0,0 +1,72 @@ +using Microsoft.AspNetCore.Authorization; +using Microsoft.AspNetCore.Mvc; + +using PrivaPub.ClientModels; +using PrivaPub.Domain.Statistics; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Jobs; + +using System.Globalization; + +namespace PrivaPub.Controllers.ClientToServer +{ + [ApiController, + Route("clientapi/admin/statistics"), + Authorize(Policy = Policies.IsAdmin)] + public class StatisticsController : ControllerBase + { + readonly StatisticsQueries _queries; + readonly IJobQueue _queue; + + public StatisticsController(StatisticsQueries queries, IJobQueue queue) + { + _queries = queries; + _queue = queue; + } + + [HttpGet, Route("/clientapi/admin/statistics/overview")] + public async Task Overview([FromQuery] int days = 30, CancellationToken token = default) => + Ok(await _queries.Overview(days, token)); + + [HttpGet, Route("/clientapi/admin/statistics/hosts")] + public async Task Hosts([FromQuery] int days = 30, [FromQuery] string sort = default, [FromQuery] string software = default, + [FromQuery] int page = 1, [FromQuery] int limit = 50, CancellationToken token = default) => + Ok(await _queries.Hosts(days, sort, software, page, limit, token)); + + [HttpGet, Route("/clientapi/admin/statistics/hosts/{host}")] + public async Task Host(string host, [FromQuery] int days = 90, CancellationToken token = default) => + await _queries.Host(host, days, token) is { } found ? Ok(found) : NotFound(); + + [HttpGet, Route("/clientapi/admin/statistics/events")] + public async Task Events([FromQuery] string host = default, [FromQuery] string channel = default, [FromQuery] string outcome = default, + [FromQuery] string reason = default, [FromQuery] DateTime? before = default, [FromQuery] int limit = 100, CancellationToken token = default) => + Ok(await _queries.Events(host, channel, outcome, reason, before, limit, token)); + + [HttpGet, Route("/clientapi/admin/statistics/server")] + public async Task Server([FromQuery] int days = 30, CancellationToken token = default) => + Ok(await _queries.Server(days, token)); + + [HttpPost, Route("/clientapi/admin/statistics/rollups/{day}")] + public async Task Rollup(string day, CancellationToken token) + { + if (!DateTime.TryParseExact(day, "yyyy-MM-dd", CultureInfo.InvariantCulture, DateTimeStyles.AdjustToUniversal | DateTimeStyles.AssumeUniversal, out var parsed) + || parsed.Date >= DateTime.UtcNow.Date) + return BadRequest(); + var job = RollupJob.For(parsed.Date, $"rollup|{day}|rerun|{DateTime.UtcNow.Ticks}"); + job.RunAt = DateTime.UtcNow; + await _queue.EnqueueMany(new[] { job }, token); + return Accepted(); + } + + [HttpPost, Route("/clientapi/admin/statistics/hosts/{host}/describe")] + public async Task Describe(string host, CancellationToken token) + { + host = Interactions.Host(host); + if (host == Interactions.Unknown) + return BadRequest(); + await _queue.Enqueue(JobKind.DescribeInstance, host, host, $"describe|{host}|admin|{DateTime.UtcNow.Ticks}", token); + return Accepted(); + } + } +} diff --git a/PrivaPub/Domain/Statistics/StatisticsQueries.cs b/PrivaPub/Domain/Statistics/StatisticsQueries.cs new file mode 100644 index 0000000..4ba30a7 --- /dev/null +++ b/PrivaPub/Domain/Statistics/StatisticsQueries.cs @@ -0,0 +1,195 @@ +using MongoDB.Entities; + +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Statistics; + +namespace PrivaPub.Domain.Statistics +{ + public sealed record HostDay(DateTime Day, string Host, IReadOnlyDictionary Counters, IReadOnlyDictionary Latency, + IReadOnlyDictionary Bytes, IReadOnlyDictionary Reads, int Accounts, bool Live); + + public class StatisticsQueries + { + public const string Attribution = "IP geolocation by DB-IP (https://db-ip.com), CC BY 4.0"; + static readonly string[] Channels = { Interactions.Receive, Interactions.In, Interactions.Out, Interactions.Http, Interactions.Preview, Interactions.Crawl }; + static readonly int[] LatencyLimits = { 50, 100, 250, 500, 1000, 2500, 5000, 10000, 30000 }; + + public async Task> Days(int days, string host, CancellationToken token) + { + var today = DateTime.UtcNow.Date; + var since = today.AddDays(1 - Math.Clamp(days, 1, 3660)); + var rows = await DB.Default.Find() + .Match(d => d.Day >= since) + .Match(d => host == null || d.Host == host) + .ExecuteAsync(token); + var result = rows.Where(r => r.RolledUpAt != null) + .Select(r => new HostDay(r.Day, r.Host, r.Counters, r.Latency, r.Bytes, r.Reads, r.Accounts, false)) + .ToList(); + var rolledDays = rows.Where(r => r.RolledUpAt != null).Select(r => r.Day).ToHashSet(); + for (var day = since; day <= today; day = day.AddDays(1)) + { + if (rolledDays.Contains(day)) + continue; + var end = day.AddDays(1); + var start = day; + var events = await DB.Default.Find() + .Match(e => e.At >= start && e.At < end) + .Match(e => host == null || e.Host == host) + .ExecuteAsync(token); + var reads = rows.Where(r => r.Day == day).ToDictionary(r => r.Host, r => r.Reads); + foreach (var (name, fold) in Rollup.Fold(events).Hosts) + result.Add(new HostDay(day, name, fold.Counters, fold.Latency, fold.Bytes, reads.GetValueOrDefault(name) ?? new(), fold.Accounts.Count, true)); + foreach (var (name, read) in reads.Where(r => result.All(h => h.Day != day || h.Host != r.Key))) + result.Add(new HostDay(day, name, new Dictionary(), new Dictionary(), new Dictionary(), read, 0, true)); + } + return result.OrderBy(d => d.Day).ThenBy(d => d.Host, StringComparer.Ordinal).ToList(); + } + + public async Task Overview(int days, CancellationToken token) + { + var hostDays = await Days(days, default, token); + var counters = Sum(hostDays.Select(d => d.Counters)); + var latency = Sum(hostDays.Select(d => d.Latency)); + var server = await DB.Default.Find().Match(d => d.Day >= DateTime.UtcNow.Date.AddDays(1 - Math.Clamp(days, 1, 3660))).ExecuteAsync(token); + var delivered = Outcomes(counters, Interactions.Out); + var attempts = delivered.Where(o => o.Key != Interactions.Deferred).Sum(o => o.Value); + return new + { + Days = days, + Channels = Channels.ToDictionary(c => c, c => counters.Where(k => k.Key.StartsWith(c + ":", StringComparison.Ordinal)).Sum(k => k.Value)), + Inbound = Outcomes(counters, Interactions.In), + Outbound = delivered, + DeliverySuccess = attempts == 0 ? (double?)null : Math.Round(delivered.GetValueOrDefault(Interactions.Ok) / (double)attempts, 4), + DeliveryLatencyMs = new { P50 = Percentile(latency, Interactions.Out, 0.5), P95 = Percentile(latency, Interactions.Out, 0.95) }, + Hosts = new + { + Active = hostDays.Where(d => d.Host != Interactions.Unknown && d.Counters.Count > 0).Select(d => d.Host).Distinct().Count(), + Known = await DB.Default.CountAsync(cancellation: token) + }, + Accounts = server.Sum(d => d.Accounts) + hostDays.Where(d => d.Live).Sum(d => d.Accounts), + Ledger = new + { + Written = server.Sum(d => d.LedgerWritten), + Dropped = server.Sum(d => d.LedgerDropped), + Failed = server.Sum(d => d.LedgerFailed) + } + }; + } + + public async Task Hosts(int days, string sort, string software, int page, int limit, CancellationToken token) + { + var hostDays = await Days(days, default, token); + var instances = (await DB.Default.Find().ExecuteAsync(token)).ToDictionary(i => i.Host, StringComparer.Ordinal); + var rows = hostDays.GroupBy(d => d.Host) + .Select(g => + { + var counters = Sum(g.Select(d => d.Counters)); + var outbound = Outcomes(counters, Interactions.Out); + var instance = instances.GetValueOrDefault(g.Key); + return new + { + Host = g.Key, + instance?.Software, + Version = instance?.SoftwareVersion, + In = counters.Where(k => k.Key.StartsWith("in:", StringComparison.Ordinal)).Sum(k => k.Value), + Out = counters.Where(k => k.Key.StartsWith("out:", StringComparison.Ordinal)).Sum(k => k.Value), + Refused = counters.Where(k => k.Key.StartsWith("recv:", StringComparison.Ordinal) && !k.Key.Contains(":202:", StringComparison.Ordinal)).Sum(k => k.Value), + Failures = outbound.Where(o => o.Key is Interactions.Retry or Interactions.Dead or Interactions.Failed).Sum(o => o.Value), + LatencyP50 = Percentile(Sum(g.Select(d => d.Latency)), Interactions.Out, 0.5), + Accounts = g.Sum(d => d.Accounts), + instance?.DescribedAt, + Delivery = instance == default ? default : new { instance.ConsecutiveFailures, instance.UnavailableUntil, instance.LastSuccessAt, instance.LastFailureAt } + }; + }) + .Where(r => software == default || string.Equals(r.Software, software, StringComparison.OrdinalIgnoreCase)); + rows = sort switch + { + "failures" => rows.OrderByDescending(r => r.Failures), + "latency" => rows.OrderByDescending(r => r.LatencyP50 ?? 0), + "host" => rows.OrderBy(r => r.Host, StringComparer.Ordinal), + _ => rows.OrderByDescending(r => r.In + r.Out) + }; + var list = rows.ToList(); + limit = Math.Clamp(limit, 1, 200); + return new { Total = list.Count, Page = Math.Max(page, 1), Hosts = list.Skip((Math.Max(page, 1) - 1) * limit).Take(limit), Attribution }; + } + + public async Task Host(string host, int days, CancellationToken token) + { + host = Interactions.Host(host); + var instance = await DB.Default.Find().Match(i => i.Host == host).ExecuteFirstAsync(token); + var hostDays = await Days(days, host, token); + var events = await DB.Default.Find().Match(e => e.Host == host).Sort(e => e.At, Order.Descending).Limit(100).ExecuteAsync(token); + if (instance == default && hostDays.Count == 0 && events.Count == 0) + return default; + return new + { + Host = host, + Instance = instance, + Days = hostDays.Select(d => new { d.Day, d.Counters, d.Latency, d.Bytes, d.Reads, d.Accounts, d.Live }), + Events = events, + Attribution + }; + } + + public async Task> Events(string host, string channel, string outcome, string reason, DateTime? before, int limit, CancellationToken token) + { + var hostName = host == default ? default : Interactions.Host(host); + return await DB.Default.Find() + .Match(e => hostName == null || e.Host == hostName) + .Match(e => channel == null || e.Channel == channel) + .Match(e => outcome == null || e.Outcome == outcome) + .Match(e => reason == null || e.Reason == reason) + .Match(e => before == null || e.At < before) + .Sort(e => e.At, Order.Descending) + .Limit(Math.Clamp(limit, 1, 200)) + .ExecuteAsync(token); + } + + public async Task> Server(int days, CancellationToken token) + { + var since = DateTime.UtcNow.Date.AddDays(1 - Math.Clamp(days, 1, 3660)); + return await DB.Default.Find().Match(d => d.Day >= since).Sort(d => d.Day, Order.Ascending).ExecuteAsync(token); + } + + public static Dictionary Outcomes(IReadOnlyDictionary counters, string channel) + { + var outcomes = new Dictionary(); + foreach (var (key, value) in counters) + { + var parts = key.Split(':'); + if (parts.Length >= 4 && parts[0] == channel) + outcomes[parts[3]] = outcomes.GetValueOrDefault(parts[3]) + value; + } + return outcomes; + } + + public static int? Percentile(IReadOnlyDictionary latency, string prefix, double share) + { + var buckets = LatencyLimits.Select(limit => (Limit: (int?)limit, Count: latency.GetValueOrDefault($"{prefix}:le{limit}ms"))) + .Append((Limit: (int?)null, Count: latency.GetValueOrDefault($"{prefix}:inf"))) + .ToList(); + var total = buckets.Sum(b => b.Count); + if (total == 0) + return default; + long seen = 0; + foreach (var (limit, count) in buckets) + { + seen += count; + if (seen >= share * total) + return limit ?? LatencyLimits[^1] + 1; + } + return default; + } + + static Dictionary Sum(IEnumerable> dictionaries) + { + var sum = new Dictionary(); + foreach (var dictionary in dictionaries) + foreach (var (key, value) in dictionary ?? new Dictionary()) + sum[key] = sum.GetValueOrDefault(key) + value; + return sum; + } + } +} diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index 478f216..19bafe1 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -115,6 +115,7 @@ namespace PrivaPub.Middleware .AddHostedService(services => services.GetRequiredService()) .AddSingleton() .AddSingleton() + .AddSingleton() .AddHostedService(); public static IServiceCollection PrivaPubAuthServicesConfiguration(this IServiceCollection service, IConfiguration configuration)