From 15cd034b295d0b5243eb262497d106e2f97c1db4 Mon Sep 17 00:00:00 2001 From: thepra Date: Sat, 3 Oct 2026 10:56:13 +0200 Subject: [PATCH] M1: the interaction ledger InteractionEvent records one interaction with a remote server: its channel (recv, in, out, http, preview, crawl), activity and object type, outcome and reason, status, latency, wait, bytes, attempt, audience, local actor kind, inbox, signature, features and the object's age. IInteractionLedger.Record never blocks and never throws: events go into a bounded channel of 10k, a full channel drops and counts, and a hosted service writes batches of up to 1000 every two seconds. Privacy, as decided by the owner: - no persona, root, group or activity id, inbox URL, actor URI or sender IP is stored; - distinct accounts are counted with an HMAC keyed by a per-day salt (InteractionSalt, upserted so restarts agree, never created for a past day); - the local actor kind survives only on public and unlisted traffic; - a host claimed by an unverified sender is kept only if it is already known. Traffic caused by reading is only counted per day (InstanceDay.Reads, ServerDay). Indexes: a 90-day TTL on events, unique day rows, and a TTL safety net on salts. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2 --- CLAUDE.md | 11 + .../Statistics/InteractionsTests.cs | 105 ++++++++ PrivaPub.Tests/Statistics/LedgerStoreTests.cs | 89 +++++++ PrivaPub.Tests/Support/MemoryLedger.cs | 27 ++ PrivaPub/Infrastructure/Data/Indexes.cs | 9 + .../Statistics/InteractionLedger.cs | 250 ++++++++++++++++++ .../Statistics/InteractionSalts.cs | 44 +++ .../Infrastructure/Statistics/Interactions.cs | 101 +++++++ .../Statistics/StatisticsOptions.cs | 9 + .../Middleware/SocialPubConfigurations.cs | 9 + PrivaPub/Models/Statistics/InstanceDay.cs | 18 ++ .../Models/Statistics/InteractionEvent.cs | 33 +++ PrivaPub/Models/Statistics/InteractionSalt.cs | 11 + PrivaPub/Models/Statistics/ServerDay.cs | 17 ++ PrivaPub/Program.cs | 1 + docs/ROADMAP.md | 2 +- 16 files changed, 735 insertions(+), 1 deletion(-) create mode 100644 PrivaPub.Tests/Statistics/InteractionsTests.cs create mode 100644 PrivaPub.Tests/Statistics/LedgerStoreTests.cs create mode 100644 PrivaPub.Tests/Support/MemoryLedger.cs create mode 100644 PrivaPub/Infrastructure/Statistics/InteractionLedger.cs create mode 100644 PrivaPub/Infrastructure/Statistics/InteractionSalts.cs create mode 100644 PrivaPub/Infrastructure/Statistics/Interactions.cs create mode 100644 PrivaPub/Infrastructure/Statistics/StatisticsOptions.cs create mode 100644 PrivaPub/Models/Statistics/InstanceDay.cs create mode 100644 PrivaPub/Models/Statistics/InteractionEvent.cs create mode 100644 PrivaPub/Models/Statistics/InteractionSalt.cs create mode 100644 PrivaPub/Models/Statistics/ServerDay.cs diff --git a/CLAUDE.md b/CLAUDE.md index af03a68..eaf6c0f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -248,6 +248,17 @@ cd /var/www/privapub.thepra.dev && sudo -u www-data ASPNETCORE_ENVIRONMENT=Produ decided in `docs/ROADMAP.md` ("Owner decisions on what PrivaPub reveals"). Anything new that tells another server something about an avatar gets the same treatment: ask, then record it there. - **Location-ranged posts never federate.** +- **Statistics name servers, never people** (owner decisions on statistics, 2026-10-03). Every interaction is recorded + through `IInteractionLedger` (`Infrastructure/Statistics`) as an `InteractionEvent` (90 days), but an event never + holds a root, persona or group id, an activity id, an inbox URL, a remote actor URI or a sender IP: + - Distinct remote accounts are counted with `ActorHash`, an HMAC keyed by that day's `InteractionSalt`, which the + day's rollup deletes. + - `LocalKind` survives `Interactions.Sanitize` only on public and unlisted traffic, so a circle can show neither by + presence nor by absence, and no reason code names one. + - Traffic a reader causes (the media proxy, lookups, the client API, fetches of our own documents) is only counted + per day (`Count`, `CountServer`), never logged per event. + - A host claimed by an unverified sender is kept only if it is already a `RemoteInstance`. + - `Record` never blocks and never throws: it writes to a bounded channel, and a full channel drops and counts. ## Data diff --git a/PrivaPub.Tests/Statistics/InteractionsTests.cs b/PrivaPub.Tests/Statistics/InteractionsTests.cs new file mode 100644 index 0000000..e0f685b --- /dev/null +++ b/PrivaPub.Tests/Statistics/InteractionsTests.cs @@ -0,0 +1,105 @@ +using Microsoft.Extensions.Logging.Abstractions; + +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Statistics; +using PrivaPub.Tests.Support; + +namespace PrivaPub.Tests.Statistics +{ + public class InteractionsTests + { + [Fact] + public void The_local_kind_is_kept_only_on_public_and_unlisted_traffic() + { + Assert.Equal("person", Interactions.Sanitize(new InteractionEvent { Audience = "public", LocalKind = "person" }).LocalKind); + Assert.Equal("group", Interactions.Sanitize(new InteractionEvent { Audience = "unlisted", LocalKind = "group" }).LocalKind); + Assert.Null(Interactions.Sanitize(new InteractionEvent { Audience = "private", LocalKind = "group" }).LocalKind); + Assert.Null(Interactions.Sanitize(new InteractionEvent { Audience = "none", LocalKind = "person" }).LocalKind); + Assert.Null(Interactions.Sanitize(new InteractionEvent { LocalKind = "person" }).LocalKind); + Assert.Null(Interactions.Sanitize(new InteractionEvent { Audience = "public", LocalKind = "circle" }).LocalKind); + } + + [Fact] + public void Junk_becomes_other_and_hosts_are_normalised() + { + var e = Interactions.Sanitize(new InteractionEvent + { + Host = "Mastodon.Social.", + Channel = "sideways", + Activity = "https://www.w3.org/ns/activitystreams#Create", + Object = "Note", + Outcome = "Dropped!", + Reason = "not-addressed", + Audience = "everyone", + Inbox = "front-door", + Signature = "cavage:rsa-sha256", + Features = new() { "fep-044f", "Bad Feature", "fep-044f" } + }); + + Assert.Equal("mastodon.social", e.Host); + Assert.Equal("other", e.Channel); + Assert.Equal("other", e.Activity); + Assert.Equal("Note", e.Object); + Assert.Equal("other", e.Outcome); + Assert.Equal("not-addressed", e.Reason); + Assert.Null(e.Audience); + Assert.Null(e.Inbox); + Assert.Equal("cavage:rsa-sha256", e.Signature); + Assert.Equal(new[] { "fep-044f", "other" }, e.Features); + Assert.Equal("-", Interactions.Host("evil host/with path")); + Assert.Equal("-", Interactions.Host(default)); + Assert.Equal("social.example", Interactions.HostOf("https://Social.Example:8443/users/x")); + } + + [Fact] + public void Latencies_and_waits_fall_into_buckets() + { + Assert.Equal("le50ms", Interactions.Latency(0)); + Assert.Equal("le250ms", Interactions.Latency(101)); + Assert.Equal("inf", Interactions.Latency(31_000)); + Assert.Equal("le1s", Interactions.Wait(900)); + Assert.Equal("le60s", Interactions.Wait(59_000)); + Assert.Equal("inf", Interactions.Wait(100_000_000)); + } + + [Fact] + public void An_actor_hash_is_stable_for_a_salt_and_unlinkable_across_salts() + { + var monday = new byte[32]; + var tuesday = Enumerable.Repeat((byte)1, 32).ToArray(); + const string actor = "https://social.example/users/alice"; + + var first = InteractionSalts.Hash(monday, actor); + + Assert.Equal(first, InteractionSalts.Hash(monday, actor)); + Assert.NotEqual(first, InteractionSalts.Hash(tuesday, actor)); + Assert.NotEqual(first, InteractionSalts.Hash(monday, actor + "2")); + Assert.Equal(16, first.Length); + Assert.DoesNotContain("alice", first); + Assert.Null(InteractionSalts.Hash(default, actor)); + } + + [Fact] + public void A_full_ledger_drops_and_counts_but_never_blocks_or_throws() + { + var ledger = new InteractionLedger(new InteractionSalts(), new StaticOptions(new StatisticsOptions()), NullLogger.Instance); + + for (var i = 0; i < 10_005; i++) + ledger.Record(new InteractionEvent { Channel = Interactions.In }); + ledger.Record(default); + + Assert.Equal(5, ledger.Dropped); + } + + [Fact] + public void A_disabled_ledger_records_nothing() + { + var ledger = new InteractionLedger(new InteractionSalts(), new StaticOptions(new StatisticsOptions { Enabled = false }), NullLogger.Instance); + + for (var i = 0; i < 10_005; i++) + ledger.Record(new InteractionEvent { Channel = Interactions.In }); + + Assert.Equal(0, ledger.Dropped); + } + } +} diff --git a/PrivaPub.Tests/Statistics/LedgerStoreTests.cs b/PrivaPub.Tests/Statistics/LedgerStoreTests.cs new file mode 100644 index 0000000..216b7b8 --- /dev/null +++ b/PrivaPub.Tests/Statistics/LedgerStoreTests.cs @@ -0,0 +1,89 @@ +using Microsoft.Extensions.Logging.Abstractions; + +using MongoDB.Bson; +using MongoDB.Driver; +using MongoDB.Entities; + +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Statistics; +using PrivaPub.Tests.Support; + +namespace PrivaPub.Tests.Statistics +{ + [Trait("Category", "Integration")] + public sealed class LedgerStoreTests : IAsyncLifetime + { + readonly string _known = $"known{Guid.NewGuid():N}.example"; + readonly string _stranger = $"stranger{Guid.NewGuid():N}.example"; + + public async ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + await DB.Default.SaveAsync(new RemoteInstance { Host = _known }); + } + + public ValueTask DisposeAsync() => ValueTask.CompletedTask; + + static InteractionLedger Ledger() => + new(new InteractionSalts(), new StaticOptions(new StatisticsOptions()), NullLogger.Instance); + + [Fact] + public async Task Events_are_stored_with_a_hash_and_never_the_actor() + { + var token = TestContext.Current.CancellationToken; + var ledger = Ledger(); + var actor = $"https://{_known}/users/someone{Guid.NewGuid():N}"; + + ledger.Record(new InteractionEvent { Host = _known, Channel = Interactions.In, Activity = "Create", Object = "Note", Outcome = Interactions.Accepted }, actor); + ledger.Record(new InteractionEvent { Host = _stranger, Channel = Interactions.Receive, Status = 401, Outcome = Interactions.Refused, Reason = "no-signature" }, hostClaimed: true); + ledger.Record(new InteractionEvent { Host = _known, Channel = Interactions.Receive, Status = 401, Outcome = Interactions.Refused, Reason = "signature-invalid" }, hostClaimed: true); + await ledger.Flush(token); + + var stored = await DB.Default.Find().Match(e => e.Host == _known).ExecuteAsync(token); + Assert.Equal(2, stored.Count); + var accepted = Assert.Single(stored, e => e.Channel == Interactions.In); + Assert.NotNull(accepted.ActorHash); + Assert.Equal("signature-invalid", Assert.Single(stored, e => e.Channel == Interactions.Receive).Reason); + Assert.False(await DB.Default.Find().Match(e => e.Host == _stranger).ExecuteAnyAsync(token)); + var raw = await DB.Default.Database().GetCollection(nameof(InteractionEvent)).Find(new BsonDocument("_id", ObjectId.Parse(accepted.ID))).FirstAsync(token); + Assert.DoesNotContain("someone", raw.ToJson()); + Assert.DoesNotContain("/users/", raw.ToJson()); + } + + [Fact] + public async Task Two_processes_agree_on_a_days_salt_and_no_salt_is_made_for_the_past() + { + var token = TestContext.Current.CancellationToken; + + var one = await new InteractionSalts().For(DateTime.UtcNow, token); + var two = await new InteractionSalts().For(DateTime.UtcNow, token); + + Assert.Equal(Convert.ToBase64String(one), Convert.ToBase64String(two)); + Assert.Null(await new InteractionSalts().For(new DateTime(2001, 1, 1, 0, 0, 0, DateTimeKind.Utc), token)); + } + + [Fact] + public async Task Counters_add_up_per_day() + { + var token = TestContext.Current.CancellationToken; + var ledger = Ledger(); + + ledger.Count(_known, "media:hit", 1000); + ledger.Count(_known, "media:hit", 500); + ledger.Count(_known, "media.miss"); + ledger.CountServer(ServerSections.Client, $"test{_known}:GET:2xx", 120); + ledger.CountServer(ServerSections.Served, $"test{_known}:200:signed"); + await ledger.Flush(token); + + var day = await DB.Default.Find().Match(d => d.Host == _known && d.Day == DateTime.UtcNow.Date).ExecuteFirstAsync(token); + Assert.Equal(2, day.Reads["media:hit"]); + Assert.Equal(1500, day.Reads["media:hit:bytes"]); + Assert.Equal(1, day.Reads["media_miss"]); + var server = await DB.Default.Find().Match(d => d.Day == DateTime.UtcNow.Date).ExecuteFirstAsync(token); + Assert.Equal(1, server.Client[$"test{_known.Replace('.', '_')}:GET:2xx"]); + Assert.Equal(1, server.Served[$"test{_known.Replace('.', '_')}:200:signed"]); + Assert.True(server.ClientLatency["le250ms"] >= 1); + } + } +} diff --git a/PrivaPub.Tests/Support/MemoryLedger.cs b/PrivaPub.Tests/Support/MemoryLedger.cs new file mode 100644 index 0000000..cdaa6c3 --- /dev/null +++ b/PrivaPub.Tests/Support/MemoryLedger.cs @@ -0,0 +1,27 @@ +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Statistics; + +using System.Collections.Concurrent; + +namespace PrivaPub.Tests.Support +{ + public sealed class MemoryLedger : IInteractionLedger + { + public ConcurrentQueue<(InteractionEvent Event, string ActorUri, bool HostClaimed)> Events { get; } = new(); + public ConcurrentDictionary<(string Host, string Key), long> Counts { get; } = new(); + public ConcurrentDictionary<(string Section, string Key), long> ServerCounts { get; } = new(); + + public long Dropped => 0; + + public void Record(InteractionEvent interaction, string actorUri = default, bool hostClaimed = false) => + Events.Enqueue((Interactions.Sanitize(interaction), actorUri, hostClaimed)); + + public void Count(string host, string key, long bytes = 0) => + Counts.AddOrUpdate((Interactions.Host(host), key), 1, (_, count) => count + 1); + + public void CountServer(string section, string key, int? latencyMs = default) => + ServerCounts.AddOrUpdate((section, key), 1, (_, count) => count + 1); + + public IReadOnlyList Of(string channel) => Events.Select(e => e.Event).Where(e => e.Channel == channel).ToList(); + } +} diff --git a/PrivaPub/Infrastructure/Data/Indexes.cs b/PrivaPub/Infrastructure/Data/Indexes.cs index 3c2ae50..a06150d 100644 --- a/PrivaPub/Infrastructure/Data/Indexes.cs +++ b/PrivaPub/Infrastructure/Data/Indexes.cs @@ -7,6 +7,7 @@ using PrivaPub.Models.Group; using PrivaPub.Models.Jobs; using PrivaPub.Models.Post; using PrivaPub.Models.Social; +using PrivaPub.Models.Statistics; using PrivaPub.Models.User; namespace PrivaPub.Infrastructure.Data @@ -118,6 +119,14 @@ namespace PrivaPub.Infrastructure.Data await Plain(token, r => r.ActivityURI); await DB.Default.Index().Key(l => l.PostId, KeyType.Ascending).Key(l => l.QuotingObjectURI, KeyType.Ascending).CreateAsync(token); await Unique(p => p.Url, Builders.Filter.Type(p => p.Url, BsonType.String), token); + + await DB.Default.Index().Key(e => e.At, KeyType.Ascending).Option(o => o.ExpireAfter = InteractionEvent.Retention).CreateAsync(token); + await Plain(token, e => e.Host, e => e.At); + await DB.Default.Index().Key(s => s.Day, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token); + await DB.Default.Index().Key(s => s.ExpiresAt, KeyType.Ascending).Option(o => o.ExpireAfter = TimeSpan.Zero).CreateAsync(token); + await DB.Default.Index().Key(d => d.Day, KeyType.Ascending).Key(d => d.Host, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token); + await Plain(token, d => d.Host, d => d.Day); + await DB.Default.Index().Key(d => d.Day, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token); } static async Task Unique(System.Linq.Expressions.Expression> key, FilterDefinition partial, diff --git a/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs b/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs new file mode 100644 index 0000000..a0b05b1 --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs @@ -0,0 +1,250 @@ +using Microsoft.Extensions.Options; + +using MongoDB.Driver; +using MongoDB.Entities; + +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Statistics; + +using System.Collections.Concurrent; +using System.Threading.Channels; + +namespace PrivaPub.Infrastructure.Statistics +{ + public interface IInteractionLedger + { + void Record(InteractionEvent interaction, string actorUri = default, bool hostClaimed = false); + void Count(string host, string key, long bytes = 0); + void CountServer(string section, string key, int? latencyMs = default); + long Dropped { get; } + } + + public static class ServerSections + { + public const string Client = "client"; + public const string Served = "served"; + } + + public class InteractionLedger : BackgroundService, IInteractionLedger + { + const int Capacity = 10_000; + const int BatchSize = 1000; + static readonly TimeSpan BatchWait = TimeSpan.FromSeconds(2); + static readonly TimeSpan CounterInterval = TimeSpan.FromSeconds(10); + static readonly TimeSpan KnownHostsAge = TimeSpan.FromMinutes(10); + + sealed record Pending(InteractionEvent Event, string ActorUri, bool HostClaimed); + + readonly Channel _channel = Channel.CreateBounded(new BoundedChannelOptions(Capacity) + { + FullMode = BoundedChannelFullMode.Wait, + SingleWriter = false, + SingleReader = false + }); + readonly InteractionSalts _salts; + readonly IOptionsMonitor _options; + readonly ILogger _logger; + readonly SemaphoreSlim _flushing = new(1, 1); + readonly ConcurrentDictionary<(DateTime Day, string Host, string Key), long> _reads = new(); + readonly ConcurrentDictionary<(DateTime Day, string Field, string Key), long> _server = new(); + HashSet _knownHosts = new(StringComparer.Ordinal); + DateTime _knownHostsAt = DateTime.MinValue; + long _dropped; + long _droppedReported; + long _written; + long _failed; + + public InteractionLedger(InteractionSalts salts, IOptionsMonitor options, ILogger logger) + { + _salts = salts; + _options = options; + _logger = logger; + } + + public long Dropped => Interlocked.Read(ref _dropped); + + public void Record(InteractionEvent interaction, string actorUri = default, bool hostClaimed = false) + { + try + { + if (interaction == default || !_options.CurrentValue.Enabled) + return; + if (!_channel.Writer.TryWrite(new Pending(interaction, actorUri, hostClaimed))) + Interlocked.Increment(ref _dropped); + } + catch (Exception ex) + { + _logger.LogDebug(ex, "An interaction could not be recorded"); + } + } + + public void Count(string host, string key, long bytes = 0) + { + if (!_options.CurrentValue.Enabled || string.IsNullOrEmpty(key)) + return; + var day = DateTime.UtcNow.Date; + host = Interactions.Host(host); + key = Key(key); + _reads.AddOrUpdate((day, host, key), 1, (_, count) => count + 1); + if (bytes > 0) + _reads.AddOrUpdate((day, host, key + ":bytes"), bytes, (_, count) => count + bytes); + } + + public void CountServer(string section, string key, int? latencyMs = default) + { + if (!_options.CurrentValue.Enabled || string.IsNullOrEmpty(key)) + return; + var day = DateTime.UtcNow.Date; + var field = section == ServerSections.Served ? nameof(ServerDay.Served) : nameof(ServerDay.Client); + _server.AddOrUpdate((day, field, Key(key)), 1, (_, count) => count + 1); + if (latencyMs is { } latency && field == nameof(ServerDay.Client)) + _server.AddOrUpdate((day, nameof(ServerDay.ClientLatency), Interactions.Latency(latency)), 1, (_, count) => count + 1); + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + var countersAt = DateTime.UtcNow; + while (!stoppingToken.IsCancellationRequested) + { + try + { + using var wait = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); + wait.CancelAfter(BatchWait); + try + { + await _channel.Reader.WaitToReadAsync(wait.Token); + } + catch (OperationCanceledException) when (!stoppingToken.IsCancellationRequested) + { + } + await FlushEvents(stoppingToken); + if (DateTime.UtcNow - countersAt >= CounterInterval) + { + countersAt = DateTime.UtcNow; + await FlushCounters(stoppingToken); + } + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + } + catch (Exception ex) + { + _logger.LogWarning(ex, "The interaction ledger could not flush"); + await Task.Delay(BatchWait, CancellationToken.None); + } + } + } + + public override async Task StopAsync(CancellationToken cancellationToken) + { + await base.StopAsync(cancellationToken); + using var budget = new CancellationTokenSource(TimeSpan.FromSeconds(5)); + try + { + await Flush(budget.Token); + } + catch (Exception ex) + { + _logger.LogWarning(ex, "The interaction ledger could not flush on shutdown"); + } + } + + public async Task Flush(CancellationToken token) + { + await FlushEvents(token); + await FlushCounters(token); + } + + async Task FlushEvents(CancellationToken token) + { + await _flushing.WaitAsync(token); + try + { + while (_channel.Reader.TryPeek(out _)) + { + var batch = new List(BatchSize); + while (batch.Count < BatchSize && _channel.Reader.TryRead(out var pending)) + batch.Add(await Prepare(pending, token)); + if (batch.Count == 0) + return; + try + { + await DB.Default.InsertAsync(batch, token); + Interlocked.Add(ref _written, batch.Count); + } + catch (Exception ex) when (ex is not OperationCanceledException) + { + Interlocked.Add(ref _failed, batch.Count); + _logger.LogWarning(ex, "{Count} interactions could not be written", batch.Count); + } + } + } + finally + { + _flushing.Release(); + } + } + + async Task Prepare(Pending pending, CancellationToken token) + { + var interaction = Interactions.Sanitize(pending.Event); + if (pending.HostClaimed && interaction.Host != Interactions.Unknown && !(await KnownHosts(token)).Contains(interaction.Host)) + interaction.Host = Interactions.Unknown; + if (pending.ActorUri != default) + interaction.ActorHash = InteractionSalts.Hash(await _salts.For(interaction.At, token), pending.ActorUri); + return interaction; + } + + async Task> KnownHosts(CancellationToken token) + { + if (DateTime.UtcNow - _knownHostsAt < KnownHostsAge) + return _knownHosts; + var hosts = await DB.Default.Find().Project(i => i.Host).ExecuteAsync(token); + _knownHosts = new HashSet(hosts.Where(h => h != default), StringComparer.Ordinal); + _knownHostsAt = DateTime.UtcNow; + return _knownHosts; + } + + async Task FlushCounters(CancellationToken token) + { + foreach (var key in _reads.Keys.ToList()) + { + if (!_reads.TryRemove(key, out var count) || count == 0) + continue; + await DB.Default.Update() + .Match(d => d.Day == key.Day && d.Host == key.Host) + .Modify(b => b.Inc($"{nameof(InstanceDay.Reads)}.{key.Key}", count)) + .Option(o => o.IsUpsert = true) + .ExecuteAsync(token); + } + + var day = DateTime.UtcNow.Date; + var written = Interlocked.Exchange(ref _written, 0); + var failed = Interlocked.Exchange(ref _failed, 0); + var dropped = Interlocked.Read(ref _dropped); + var newlyDropped = dropped - Interlocked.Exchange(ref _droppedReported, dropped); + var server = _server.Keys.ToList(); + if (server.Count == 0 && written == 0 && failed == 0 && newlyDropped == 0) + return; + var update = DB.Default.Update() + .Match(d => d.Day == day) + .Modify(b => b.Inc(d => d.LedgerWritten, written)) + .Modify(b => b.Inc(d => d.LedgerFailed, failed)) + .Modify(b => b.Inc(d => d.LedgerDropped, newlyDropped)) + .Option(o => o.IsUpsert = true); + await update.ExecuteAsync(token); + foreach (var key in server) + { + if (!_server.TryRemove(key, out var count) || count == 0) + continue; + await DB.Default.Update() + .Match(d => d.Day == key.Day) + .Modify(b => b.Inc($"{key.Field}.{key.Key}", count)) + .Option(o => o.IsUpsert = true) + .ExecuteAsync(token); + } + } + + static string Key(string key) => key.Replace('.', '_').Replace('$', '_'); + } +} diff --git a/PrivaPub/Infrastructure/Statistics/InteractionSalts.cs b/PrivaPub/Infrastructure/Statistics/InteractionSalts.cs new file mode 100644 index 0000000..47ea8bf --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/InteractionSalts.cs @@ -0,0 +1,44 @@ +using MongoDB.Entities; + +using PrivaPub.Models.Statistics; + +using System.Collections.Concurrent; +using System.Security.Cryptography; +using System.Text; + +namespace PrivaPub.Infrastructure.Statistics +{ + public class InteractionSalts + { + readonly ConcurrentDictionary _known = new(); + + public async Task For(DateTime day, CancellationToken token) + { + day = day.Date; + if (_known.TryGetValue(day, out var cached)) + return cached; + InteractionSalt salt; + if (day == DateTime.UtcNow.Date) + salt = await DB.Default.UpdateAndGet() + .Match(s => s.Day == day) + .Modify(b => b.SetOnInsert(s => s.Key, Convert.ToBase64String(RandomNumberGenerator.GetBytes(32)))) + .Modify(b => b.SetOnInsert(s => s.ExpiresAt, day.AddHours(30))) + .Option(o => o.IsUpsert = true) + .ExecuteAsync(token); + else + salt = await DB.Default.Find().Match(s => s.Day == day).ExecuteFirstAsync(token); + if (salt?.Key == default) + return default; + var key = Convert.FromBase64String(salt.Key); + foreach (var old in _known.Keys.Where(k => k < day.AddDays(-1))) + _known.TryRemove(old, out _); + _known[day] = key; + return key; + } + + public static string Hash(byte[] salt, string actorUri) => + salt == default || string.IsNullOrEmpty(actorUri) + ? default + : Convert.ToBase64String(HMACSHA256.HashData(salt, Encoding.UTF8.GetBytes(actorUri))[..12]).TrimEnd('=').Replace('+', '-').Replace('/', '_'); + } +} diff --git a/PrivaPub/Infrastructure/Statistics/Interactions.cs b/PrivaPub/Infrastructure/Statistics/Interactions.cs new file mode 100644 index 0000000..15fc0cd --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/Interactions.cs @@ -0,0 +1,101 @@ +using PrivaPub.Models.Statistics; + +using System.Text.RegularExpressions; + +namespace PrivaPub.Infrastructure.Statistics +{ + public static partial class Interactions + { + public const string Receive = "recv"; + public const string In = "in"; + public const string Out = "out"; + public const string Http = "http"; + public const string Preview = "preview"; + public const string Crawl = "crawl"; + + public const string Queued = "queued"; + public const string Refused = "refused"; + public const string Accepted = "accepted"; + public const string Dropped = "dropped"; + public const string Rejected = "rejected"; + public const string Ok = "ok"; + public const string Failed = "failed"; + public const string Deferred = "deferred"; + public const string Retry = "retry"; + public const string Dead = "dead"; + + public const string Public = "public"; + public const string Unlisted = "unlisted"; + public const string Private = "private"; + public const string None = "none"; + + public const string Unknown = "-"; + public const string Other = "other"; + + static readonly HashSet Channels = new(StringComparer.Ordinal) { Receive, In, Out, Http, Preview, Crawl }; + static readonly HashSet Audiences = new(StringComparer.Ordinal) { Public, Unlisted, Private, None }; + static readonly HashSet LocalKinds = new(StringComparer.Ordinal) { "person", "group", "application" }; + static readonly HashSet Inboxes = new(StringComparer.Ordinal) { "shared", "personal" }; + static readonly int[] LatencyBuckets = { 50, 100, 250, 500, 1000, 2500, 5000, 10000, 30000 }; + static readonly int[] WaitBuckets = { 1, 10, 60, 600, 3600, 21600, 86400 }; + + public static InteractionEvent Sanitize(InteractionEvent e) + { + e.Host = Host(e.Host); + e.Channel = Channels.Contains(e.Channel) ? e.Channel : Other; + e.Activity = TypeName(e.Activity); + e.Object = TypeName(e.Object); + e.Purpose = Word(e.Purpose); + e.Trigger = Word(e.Trigger); + e.Outcome = Word(e.Outcome); + e.Reason = Word(e.Reason); + e.Audience = e.Audience != default && Audiences.Contains(e.Audience) ? e.Audience : default; + e.LocalKind = e.Audience is Public or Unlisted && e.LocalKind != default && LocalKinds.Contains(e.LocalKind) ? e.LocalKind : default; + e.Inbox = e.Inbox != default && Inboxes.Contains(e.Inbox) ? e.Inbox : default; + e.Signature = e.Signature == default ? default : SignatureName().IsMatch(e.Signature) ? e.Signature : Other; + e.Features = e.Features?.Select(Word).Where(f => f != default).Distinct().Take(32).ToList(); + if (e.Features?.Count == 0) + e.Features = default; + return e; + } + + public static string Host(string host) + { + host = host?.Trim().TrimEnd('.').ToLowerInvariant(); + return string.IsNullOrEmpty(host) || host.Length > 253 || !HostName().IsMatch(host) ? Unknown : host; + } + + public static string HostOf(string uri) => + Uri.TryCreate(uri, UriKind.Absolute, out var parsed) ? Host(parsed.Host) : Unknown; + + public static string Latency(int milliseconds) => Bucket("le", milliseconds, LatencyBuckets, "ms"); + + public static string Wait(int milliseconds) => Bucket("le", milliseconds / 1000, WaitBuckets, "s"); + + static string Bucket(string prefix, int value, int[] limits, string unit) + { + foreach (var limit in limits) + if (value <= limit) + return $"{prefix}{limit}{unit}"; + return "inf"; + } + + static string TypeName(string value) => + value == default ? default : TypeWord().IsMatch(value) ? value : Other; + + static string Word(string value) => + value == default ? default : KebabWord().IsMatch(value) ? value : Other; + + [GeneratedRegex("^[A-Za-z][A-Za-z0-9_-]{0,31}$")] + private static partial Regex TypeWord(); + + [GeneratedRegex("^[a-z0-9][a-z0-9-]{0,39}$")] + private static partial Regex KebabWord(); + + [GeneratedRegex("^[a-z0-9][a-z0-9:-]{0,39}$")] + private static partial Regex SignatureName(); + + [GeneratedRegex("^[a-z0-9]([a-z0-9.-]*[a-z0-9])?$")] + private static partial Regex HostName(); + } +} diff --git a/PrivaPub/Infrastructure/Statistics/StatisticsOptions.cs b/PrivaPub/Infrastructure/Statistics/StatisticsOptions.cs new file mode 100644 index 0000000..fb109d7 --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/StatisticsOptions.cs @@ -0,0 +1,9 @@ +namespace PrivaPub.Infrastructure.Statistics +{ + public class StatisticsOptions + { + public bool Enabled { get; set; } = true; + public string GeoDirectory { get; set; } = "/var/lib/privapub/geo"; + public int PublicCityMinUsers { get; set; } = 10; + } +} diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index 7361712..c78d0de 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -27,6 +27,7 @@ using PrivaPub.Domain.Statuses; using PrivaPub.Domain.Timelines; using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Infrastructure.Statistics; using Microsoft.Extensions.Options; namespace PrivaPub.Middleware @@ -103,6 +104,14 @@ namespace PrivaPub.Middleware .AddSingleton() .AddHostedService(); } + public static IServiceCollection PrivaPubStatisticsConfiguration(this IServiceCollection service, IConfiguration configuration) => + service + .Configure(configuration.GetSection("Statistics")) + .AddSingleton() + .AddSingleton() + .AddSingleton(services => services.GetRequiredService()) + .AddHostedService(services => services.GetRequiredService()); + public static IServiceCollection PrivaPubAuthServicesConfiguration(this IServiceCollection service, IConfiguration configuration) { return service diff --git a/PrivaPub/Models/Statistics/InstanceDay.cs b/PrivaPub/Models/Statistics/InstanceDay.cs new file mode 100644 index 0000000..954d7e9 --- /dev/null +++ b/PrivaPub/Models/Statistics/InstanceDay.cs @@ -0,0 +1,18 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Statistics +{ + public class InstanceDay : Entity + { + public DateTime Day { get; set; } + public string Host { get; set; } + public Dictionary Counters { get; set; } = new();//admin: every key + public Dictionary PublicCounters { get; set; } = new();//all a public page may ever read + public Dictionary Latency { get; set; } = new(); + public Dictionary Waits { get; set; } = new(); + public Dictionary Bytes { get; set; } = new(); + public Dictionary Reads { get; set; } = new();//caused by reading (media proxy, lookups): counted live, never by the rollup + public int Accounts { get; set; }//distinct remote accounts that day + public DateTime? RolledUpAt { get; set; } + } +} diff --git a/PrivaPub/Models/Statistics/InteractionEvent.cs b/PrivaPub/Models/Statistics/InteractionEvent.cs new file mode 100644 index 0000000..444569d --- /dev/null +++ b/PrivaPub/Models/Statistics/InteractionEvent.cs @@ -0,0 +1,33 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Statistics +{ + public class InteractionEvent : Entity + { + public static readonly TimeSpan Retention = TimeSpan.FromDays(90); + + public DateTime At { get; set; } = DateTime.UtcNow; + public string Host { get; set; }//the remote server, lowercase; "-" when unverified and unknown + public string Channel { get; set; }//recv | in | out | http | preview | crawl + public string Activity { get; set; }//Create, Follow, ...; "other" when it is no plain type name + public string Object { get; set; }//Note, Article, Person, ... + public string Purpose { get; set; }//http and crawl: actor, key, object, context, webfinger, nodeinfo, ... + public string Trigger { get; set; }//http: what caused the request (inbox, deliver, describe, request, ...) + public string Outcome { get; set; }//queued, refused, accepted, dropped, ok, failed, deferred, retry, dead + public string Reason { get; set; } + public int? Status { get; set; } + public int? LatencyMs { get; set; } + public int? WaitMs { get; set; }//in: since it reached the inbox; out: since it was queued + public long? Bytes { get; set; } + public int? Attempt { get; set; } + public int? Redirects { get; set; } + public string Audience { get; set; }//public | unlisted | private | none + public string LocalKind { get; set; }//person | group | application, only on public and unlisted traffic + public string Inbox { get; set; }//shared | personal + public string Signature { get; set; }//"cavage:rsa-sha256", "none" + public string ActorHash { get; set; }//keyed with that day's salt, never reversible, never linkable across days + public List Features { get; set; } + public long? AgeSeconds { get; set; }//Update and Delete: how old the object was + public bool Crawl { get; set; } + } +} diff --git a/PrivaPub/Models/Statistics/InteractionSalt.cs b/PrivaPub/Models/Statistics/InteractionSalt.cs new file mode 100644 index 0000000..db83af6 --- /dev/null +++ b/PrivaPub/Models/Statistics/InteractionSalt.cs @@ -0,0 +1,11 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Statistics +{ + public class InteractionSalt : Entity + { + public DateTime Day { get; set; } + public string Key { get; set; } + public DateTime ExpiresAt { get; set; }//a safety net: the day's rollup deletes it + } +} diff --git a/PrivaPub/Models/Statistics/ServerDay.cs b/PrivaPub/Models/Statistics/ServerDay.cs new file mode 100644 index 0000000..bf295e4 --- /dev/null +++ b/PrivaPub/Models/Statistics/ServerDay.cs @@ -0,0 +1,17 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Statistics +{ + public class ServerDay : Entity + { + public DateTime Day { get; set; } + public Dictionary Client { get; set; } = new();//client API per endpoint group, method and status class + public Dictionary ClientLatency { get; set; } = new(); + public Dictionary Served { get; set; } = new();//our own documents fetched, never per server + public long LedgerWritten { get; set; } + public long LedgerDropped { get; set; } + public long LedgerFailed { get; set; } + public int Accounts { get; set; } + public int Hosts { get; set; } + } +} diff --git a/PrivaPub/Program.cs b/PrivaPub/Program.cs index f6405b5..dd9ca6a 100644 --- a/PrivaPub/Program.cs +++ b/PrivaPub/Program.cs @@ -56,6 +56,7 @@ try .PrivaPubDataBaseConfiguration() .PrivaPubServicesConfiguration() .PrivaPubFederationConfiguration(builder.Configuration) + .PrivaPubStatisticsConfiguration(builder.Configuration) .PrivaPubCORSConfiguration() .PrivaPubRateLimiting() .PrivaPubOAuth(builder.Environment) diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 999821b..05fdcc3 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -12,7 +12,7 @@ Written 2026-10-01 from the original 2023 code, the decePubClient UI, a federati - [x] P4 Groups and privacy features: v1.6.0, deployed 2026-10-01; communities, circles and local-only located posts verified by tests. v1.6.1 adds the pasture (`tools/pasture/`): live interop with GoToSocial 0.22.1 passes all 25 checks, three runs in a row. Lemmy and a live Mastodon circle member are not run yet; the pasture has GoToSocial only - [x] P5 Lose nothing: v1.7.0 to v1.9.1, deployed 2026-10-01; parsing, typed details, provenance, downvotes, tombstones, federated blocks, all checked live against GoToSocial. Book reviews and forum threads keep their raw form only (typed in P7 and P8) - [x] P6 Emoji, polls, quotes, reactions, cards, players: v1.10.0 to v1.15.0, deployed 2026-10-01; polls, link cards, ranged video streaming and quote policies checked live against GoToSocial -- [ ] T/M Test sweep, the full pasture and the interaction ledger (see "Sweep and statistics" below) +- [ ] T/M Test sweep, the full pasture and the interaction ledger (see "Sweep and statistics" below). T1–T4: v1.15.1, deployed and verified 2026-10-03 (CI runs all tests on a throwaway mongod; nothing answers 500) - [ ] P7 Threads, communities, moderation, the social graph - [ ] P8 Signatures, discovery, the long tail