Files
SocialPub/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs
T
thepraandClaude Opus 5.5 15cd034b29 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
2026-10-03 10:56:13 +02:00

251 lines
8.1 KiB
C#

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('$', '_');
}
}