Files
SocialPub/PrivaPub/Infrastructure/Statistics/InteractionLedger.cs
T
thepraandClaude Opus 5.5 c8a305ae92 M8: every server we touch is described, located and snapshotted weekly
- Touches: the ledger marks a server as touched when it sends us a verified activity, when
  we exchange activities with it, or when we read its actors, keys, objects or WebFinger.
  It upserts RemoteInstance.Seen, FirstSeenAt and LastSeenAt at most hourly per server, and
  queues one DescribeInstance a week with the same dedupe key ObjectRecords uses. Suspended
  servers and pages behind link previews are never described. Migration _010 marks the
  servers already known as touched, with their dates.
- InstanceDescriber.Describe(host, crawled, allowed) reads:
  - NodeInfo 2.2/2.1/2.0, now with its published user counts, posts, comments,
    description, languages and schema version;
  - for software with a Mastodon API, /api/v2/instance falling back to v1: title,
    languages, registration mode, character limit, API version, source URL.

  It never keeps a contact as a field; the raw document is kept for the admin only. It
  locates the server from the address our connection reached (DB-IP Lite city and ASN, the
  CDN named when fronted) and writes a RemoteInstanceSnapshot per ISO week, unreachable
  weeks included. A crawled server is upserted as crawled only on insert, so it never
  downgrades a touched one, and robots.txt can deny any path.
- PublicGeo.Project is the only public form of a location: a CDN-fronted server shows its
  CDN only, a server reporting at least ten users shows its city, coordinates and network,
  any other only its country.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
2026-10-03 12:07:57 +02:00

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