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 _channel = Channel.CreateBounded(new BoundedChannelOptions(Capacity) { FullMode = BoundedChannelFullMode.Wait, SingleWriter = false, SingleReader = false }); readonly InteractionSalts _salts; readonly IOptionsMonitor _options; readonly ILogger _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 _knownHosts = new(StringComparer.Ordinal); DateTime _knownHostsAt = DateTime.MinValue; long _dropped; long _droppedReported; long _written; long _failed; public InteractionLedger(InteractionSalts salts, IOptionsMonitor options, ILogger 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(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 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> KnownHosts(CancellationToken token) { if (DateTime.UtcNow - _knownHostsAt < KnownHostsAge) return _knownHosts; var hosts = await DB.Default.Find().Project(i => i.Host).ExecuteAsync(token); _knownHosts = new HashSet(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() .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() .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() .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('$', '_'); } }