Files
SocialPub/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs
T
thepraandClaude Opus 5.5 b2570eab38 Statistics counters survive two flushes at once
Two ledger flushes that upsert the same new day (or touch the same new
instance) at once made one of them fail on the unique key, losing its
counts; MongoDB retries that only for single-document updates. The loser
now tries again and adds to the document the winner made. Seen as a rare
failure of LedgerOverHttpTests in full runs.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
2026-10-05 03:35:10 +02:00

330 lines
11 KiB
C#

using Microsoft.Extensions.Options;
using MongoDB.Driver;
using MongoDB.Entities;
using PrivaPub.Federation.Moderation;
using PrivaPub.Infrastructure.Jobs;
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
});
static readonly TimeSpan TouchInterval = TimeSpan.FromHours(1);
static readonly HashSet<string> TouchingPurposes = new(StringComparer.Ordinal) { "actor", "key", "object", "webfinger", "context" };
readonly InteractionSalts _salts;
readonly IOptionsMonitor<StatisticsOptions> _options;
readonly ILogger<InteractionLedger> _logger;
readonly IJobQueue _queue;
readonly IDomainBlocks _blocks;
readonly ConcurrentDictionary<string, byte> _touches = new(StringComparer.Ordinal);
readonly ConcurrentDictionary<string, DateTime> _touchedAt = new(StringComparer.Ordinal);
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,
IJobQueue queue = default, IDomainBlocks blocks = default)
{
_salts = salts;
_options = options;
_logger = logger;
_queue = queue;
_blocks = blocks;
}
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);
if (key.StartsWith("http:", StringComparison.Ordinal) && key.EndsWith(":ok", StringComparison.Ordinal)
&& TouchingPurposes.Contains(key.Split(':')[1]))
_touches[host] = 0;
_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);
await FlushTouches(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 FlushTouches(token);
await FlushCounters(token);
}
static bool Touches(InteractionEvent e, bool verified) => e.Channel switch
{
Interactions.Receive => verified && e.Outcome == Interactions.Queued,
Interactions.In or Interactions.Out => true,
Interactions.Http => e.Outcome == Interactions.Ok && TouchingPurposes.Contains(e.Purpose ?? string.Empty),
_ => false
};
//a server we exchanged something with is marked as seen, and described once a week
async Task FlushTouches(CancellationToken token)
{
await _flushing.WaitAsync(token);
try
{
await Touch(token);
}
finally
{
_flushing.Release();
}
}
async Task Touch(CancellationToken token)
{
foreach (var host in _touches.Keys.ToList())
{
_touches.TryRemove(host, out _);
var now = DateTime.UtcNow;
if (host == Interactions.Unknown || _touchedAt.TryGetValue(host, out var at) && now - at < TouchInterval)
continue;
_touchedAt[host] = now;
await Upsert(() => DB.Default.Update<RemoteInstance>()
.Match(i => i.Host == host)
.Modify(i => i.Seen, "touched")
.Modify(i => i.LastSeenAt, now)
.Modify(b => b.SetOnInsert(i => i.FirstSeenAt, now))
.Option(o => o.IsUpsert = true)
.ExecuteAsync(token));
_knownHosts.Add(host);
if (_queue != default && _blocks?.IsSuspended(host) != true)
await _queue.Enqueue(JobKind.DescribeInstance, host, host, Federation.Objects.InstanceDescriber.DedupeKey(host, now), token);
}
if (_touchedAt.Count > 100_000)
_touchedAt.Clear();
}
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);
if (interaction.Host != Interactions.Unknown && Touches(interaction, pending.ActorUri != default))
_touches[interaction.Host] = 0;
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;
}
// two flushes upserting a new day's document at once: one inserts it and the other's insert fails on the unique
// key; the second try finds the document and adds to it (MongoDB retries this only for single-document updates)
static async Task Upsert(Func<Task> write)
{
try
{
await write();
}
catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey)
{
await write();
}
}
async Task FlushCounters(CancellationToken token)
{
foreach (var key in _reads.Keys.ToList())
{
if (!_reads.TryRemove(key, out var count) || count == 0)
continue;
await Upsert(() => 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;
await Upsert(() => 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)
.ExecuteAsync(token));
foreach (var key in server)
{
if (!_server.TryRemove(key, out var count) || count == 0)
continue;
await Upsert(() => 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('$', '_');
}
}