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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
This commit is contained in:
thepraandClaude Opus 5.5 committed 2026-10-03 10:56:13 +02:00
1 parent 125a49a1a0
commit 15cd034b29
16 files changed
+735 -1

No files matched your search

+11
View File
@@ -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 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. something about an avatar gets the same treatment: ask, then record it there.
- **Location-ranged posts never federate.** - **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 ## Data
@@ -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<StatisticsOptions>(new StatisticsOptions()), NullLogger<InteractionLedger>.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<StatisticsOptions>(new StatisticsOptions { Enabled = false }), NullLogger<InteractionLedger>.Instance);
for (var i = 0; i < 10_005; i++)
ledger.Record(new InteractionEvent { Channel = Interactions.In });
Assert.Equal(0, ledger.Dropped);
}
}
}
@@ -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<StatisticsOptions>(new StatisticsOptions()), NullLogger<InteractionLedger>.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<InteractionEvent>().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<InteractionEvent>().Match(e => e.Host == _stranger).ExecuteAnyAsync(token));
var raw = await DB.Default.Database().GetCollection<BsonDocument>(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<InstanceDay>().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<ServerDay>().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);
}
}
}
+27
View File
@@ -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<InteractionEvent> Of(string channel) => Events.Select(e => e.Event).Where(e => e.Channel == channel).ToList();
}
}
+9
View File
@@ -7,6 +7,7 @@ using PrivaPub.Models.Group;
using PrivaPub.Models.Jobs; using PrivaPub.Models.Jobs;
using PrivaPub.Models.Post; using PrivaPub.Models.Post;
using PrivaPub.Models.Social; using PrivaPub.Models.Social;
using PrivaPub.Models.Statistics;
using PrivaPub.Models.User; using PrivaPub.Models.User;
namespace PrivaPub.Infrastructure.Data namespace PrivaPub.Infrastructure.Data
@@ -118,6 +119,14 @@ namespace PrivaPub.Infrastructure.Data
await Plain<Reaction>(token, r => r.ActivityURI); await Plain<Reaction>(token, r => r.ActivityURI);
await DB.Default.Index<QuoteLicence>().Key(l => l.PostId, KeyType.Ascending).Key(l => l.QuotingObjectURI, KeyType.Ascending).CreateAsync(token); await DB.Default.Index<QuoteLicence>().Key(l => l.PostId, KeyType.Ascending).Key(l => l.QuotingObjectURI, KeyType.Ascending).CreateAsync(token);
await Unique<Domain.Content.LinkPreview>(p => p.Url, Builders<Domain.Content.LinkPreview>.Filter.Type(p => p.Url, BsonType.String), token); await Unique<Domain.Content.LinkPreview>(p => p.Url, Builders<Domain.Content.LinkPreview>.Filter.Type(p => p.Url, BsonType.String), token);
await DB.Default.Index<InteractionEvent>().Key(e => e.At, KeyType.Ascending).Option(o => o.ExpireAfter = InteractionEvent.Retention).CreateAsync(token);
await Plain<InteractionEvent>(token, e => e.Host, e => e.At);
await DB.Default.Index<InteractionSalt>().Key(s => s.Day, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token);
await DB.Default.Index<InteractionSalt>().Key(s => s.ExpiresAt, KeyType.Ascending).Option(o => o.ExpireAfter = TimeSpan.Zero).CreateAsync(token);
await DB.Default.Index<InstanceDay>().Key(d => d.Day, KeyType.Ascending).Key(d => d.Host, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token);
await Plain<InstanceDay>(token, d => d.Host, d => d.Day);
await DB.Default.Index<ServerDay>().Key(d => d.Day, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token);
} }
static async Task Unique<T>(System.Linq.Expressions.Expression<Func<T, object>> key, FilterDefinition<T> partial, static async Task Unique<T>(System.Linq.Expressions.Expression<Func<T, object>> key, FilterDefinition<T> partial,
@@ -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<Pending> _channel = Channel.CreateBounded<Pending>(new BoundedChannelOptions(Capacity)
{
FullMode = BoundedChannelFullMode.Wait,
SingleWriter = false,
SingleReader = false
});
readonly InteractionSalts _salts;
readonly IOptionsMonitor<StatisticsOptions> _options;
readonly ILogger<InteractionLedger> _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<string> _knownHosts = new(StringComparer.Ordinal);
DateTime _knownHostsAt = DateTime.MinValue;
long _dropped;
long _droppedReported;
long _written;
long _failed;
public InteractionLedger(InteractionSalts salts, IOptionsMonitor<StatisticsOptions> options, ILogger<InteractionLedger> 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<InteractionEvent>(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<InteractionEvent> 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<HashSet<string>> KnownHosts(CancellationToken token)
{
if (DateTime.UtcNow - _knownHostsAt < KnownHostsAge)
return _knownHosts;
var hosts = await DB.Default.Find<RemoteInstance, string>().Project(i => i.Host).ExecuteAsync(token);
_knownHosts = new HashSet<string>(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<InstanceDay>()
.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<ServerDay>()
.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<ServerDay>()
.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('$', '_');
}
}
@@ -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<DateTime, byte[]> _known = new();
public async Task<byte[]> 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<InteractionSalt>()
.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<InteractionSalt>().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('/', '_');
}
}
@@ -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<string> Channels = new(StringComparer.Ordinal) { Receive, In, Out, Http, Preview, Crawl };
static readonly HashSet<string> Audiences = new(StringComparer.Ordinal) { Public, Unlisted, Private, None };
static readonly HashSet<string> LocalKinds = new(StringComparer.Ordinal) { "person", "group", "application" };
static readonly HashSet<string> 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();
}
}
@@ -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;
}
}
@@ -27,6 +27,7 @@ using PrivaPub.Domain.Statuses;
using PrivaPub.Domain.Timelines; using PrivaPub.Domain.Timelines;
using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Http;
using PrivaPub.Infrastructure.Jobs; using PrivaPub.Infrastructure.Jobs;
using PrivaPub.Infrastructure.Statistics;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
namespace PrivaPub.Middleware namespace PrivaPub.Middleware
@@ -103,6 +104,14 @@ namespace PrivaPub.Middleware
.AddSingleton<IJobHandler, DeliveryJobHandler>() .AddSingleton<IJobHandler, DeliveryJobHandler>()
.AddHostedService<JobWorker>(); .AddHostedService<JobWorker>();
} }
public static IServiceCollection PrivaPubStatisticsConfiguration(this IServiceCollection service, IConfiguration configuration) =>
service
.Configure<StatisticsOptions>(configuration.GetSection("Statistics"))
.AddSingleton<InteractionSalts>()
.AddSingleton<InteractionLedger>()
.AddSingleton<IInteractionLedger>(services => services.GetRequiredService<InteractionLedger>())
.AddHostedService(services => services.GetRequiredService<InteractionLedger>());
public static IServiceCollection PrivaPubAuthServicesConfiguration(this IServiceCollection service, IConfiguration configuration) public static IServiceCollection PrivaPubAuthServicesConfiguration(this IServiceCollection service, IConfiguration configuration)
{ {
return service return service
+18
View File
@@ -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<string, long> Counters { get; set; } = new();//admin: every key
public Dictionary<string, long> PublicCounters { get; set; } = new();//all a public page may ever read
public Dictionary<string, long> Latency { get; set; } = new();
public Dictionary<string, long> Waits { get; set; } = new();
public Dictionary<string, long> Bytes { get; set; } = new();
public Dictionary<string, long> 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; }
}
}
@@ -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<string> Features { get; set; }
public long? AgeSeconds { get; set; }//Update and Delete: how old the object was
public bool Crawl { get; set; }
}
}
@@ -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
}
}
+17
View File
@@ -0,0 +1,17 @@
using MongoDB.Entities;
namespace PrivaPub.Models.Statistics
{
public class ServerDay : Entity
{
public DateTime Day { get; set; }
public Dictionary<string, long> Client { get; set; } = new();//client API per endpoint group, method and status class
public Dictionary<string, long> ClientLatency { get; set; } = new();
public Dictionary<string, long> 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; }
}
}
+1
View File
@@ -56,6 +56,7 @@ try
.PrivaPubDataBaseConfiguration() .PrivaPubDataBaseConfiguration()
.PrivaPubServicesConfiguration() .PrivaPubServicesConfiguration()
.PrivaPubFederationConfiguration(builder.Configuration) .PrivaPubFederationConfiguration(builder.Configuration)
.PrivaPubStatisticsConfiguration(builder.Configuration)
.PrivaPubCORSConfiguration() .PrivaPubCORSConfiguration()
.PrivaPubRateLimiting() .PrivaPubRateLimiting()
.PrivaPubOAuth(builder.Environment) .PrivaPubOAuth(builder.Environment)
+1 -1
View File
@@ -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] 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] 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 - [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 - [ ] P7 Threads, communities, moderation, the social graph
- [ ] P8 Signatures, discovery, the long tail - [ ] P8 Signatures, discovery, the long tail