2026-10-03 10:56:13 +02:00
|
|
|
using Microsoft.Extensions.Options;
|
|
|
|
|
|
|
|
|
|
using MongoDB.Driver;
|
|
|
|
|
using MongoDB.Entities;
|
|
|
|
|
|
2026-10-03 12:07:57 +02:00
|
|
|
using PrivaPub.Federation.Moderation;
|
|
|
|
|
using PrivaPub.Infrastructure.Jobs;
|
2026-10-03 10:56:13 +02:00
|
|
|
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
|
|
|
|
|
});
|
2026-10-03 12:07:57 +02:00
|
|
|
static readonly TimeSpan TouchInterval = TimeSpan.FromHours(1);
|
|
|
|
|
static readonly HashSet<string> TouchingPurposes = new(StringComparer.Ordinal) { "actor", "key", "object", "webfinger", "context" };
|
|
|
|
|
|
2026-10-03 10:56:13 +02:00
|
|
|
readonly InteractionSalts _salts;
|
|
|
|
|
readonly IOptionsMonitor<StatisticsOptions> _options;
|
|
|
|
|
readonly ILogger<InteractionLedger> _logger;
|
2026-10-03 12:07:57 +02:00
|
|
|
readonly IJobQueue _queue;
|
|
|
|
|
readonly IDomainBlocks _blocks;
|
|
|
|
|
readonly ConcurrentDictionary<string, byte> _touches = new(StringComparer.Ordinal);
|
|
|
|
|
readonly ConcurrentDictionary<string, DateTime> _touchedAt = new(StringComparer.Ordinal);
|
2026-10-03 10:56:13 +02:00
|
|
|
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;
|
|
|
|
|
|
2026-10-03 12:07:57 +02:00
|
|
|
public InteractionLedger(InteractionSalts salts, IOptionsMonitor<StatisticsOptions> options, ILogger<InteractionLedger> logger,
|
|
|
|
|
IJobQueue queue = default, IDomainBlocks blocks = default)
|
2026-10-03 10:56:13 +02:00
|
|
|
{
|
|
|
|
|
_salts = salts;
|
|
|
|
|
_options = options;
|
|
|
|
|
_logger = logger;
|
2026-10-03 12:07:57 +02:00
|
|
|
_queue = queue;
|
|
|
|
|
_blocks = blocks;
|
2026-10-03 10:56:13 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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);
|
2026-10-03 12:07:57 +02:00
|
|
|
if (key.StartsWith("http:", StringComparison.Ordinal) && key.EndsWith(":ok", StringComparison.Ordinal)
|
|
|
|
|
&& TouchingPurposes.Contains(key.Split(':')[1]))
|
|
|
|
|
_touches[host] = 0;
|
2026-10-03 10:56:13 +02:00
|
|
|
_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);
|
2026-10-03 12:07:57 +02:00
|
|
|
await FlushTouches(stoppingToken);
|
2026-10-03 10:56:13 +02:00
|
|
|
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);
|
2026-10-03 12:07:57 +02:00
|
|
|
await FlushTouches(token);
|
2026-10-03 10:56:13 +02:00
|
|
|
await FlushCounters(token);
|
|
|
|
|
}
|
|
|
|
|
|
2026-10-03 12:07:57 +02:00
|
|
|
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 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();
|
|
|
|
|
}
|
|
|
|
|
|
2026-10-03 10:56:13 +02:00
|
|
|
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);
|
2026-10-03 12:07:57 +02:00
|
|
|
if (interaction.Host != Interactions.Unknown && Touches(interaction, pending.ActorUri != default))
|
|
|
|
|
_touches[interaction.Host] = 0;
|
2026-10-03 10:56:13 +02:00
|
|
|
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('$', '_');
|
|
|
|
|
}
|
|
|
|
|
}
|