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