Files
SocialPub/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs
T

250 lines
8.1 KiB
C#
Raw Normal View History

2026-10-03 10:56:13 +02:00
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('$', '_');
}
}