250 lines
8.1 KiB
C#
250 lines
8.1 KiB
C#
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('$', '_');
|
||
|
|
}
|
||
|
|
}
|