From d37999970c0b72e193ba4e571518b4b360d38ebb Mon Sep 17 00:00:00 2001 From: thepra Date: Sat, 3 Oct 2026 11:10:42 +0200 Subject: [PATCH] M5: outbound requests, previews, the media proxy and served traffic - HttpScope (AsyncLocal) tags each outbound request with a purpose and a trigger: - purpose is set by the caller: actor, key, object, webfinger, context, nodeinfo; - trigger is set by the job kind, by "verify" during inbox verification, or defaults to "request". - FederationHttp records every JSON, media and stream fetch: status, time, bytes, hops, and an outcome of ok, refused or failed, with a reason: disallowed, remembered, bad-redirect, too-many-redirects, content-type, too-large, bad-json, private-address, timeout, network, or the status. A fetch a reader caused (trigger "request") is only counted per server per day. - Link previews record a 'preview' event: card, no-card or failed. - The media proxy counts cache hits. - TrafficMeter counts the client API per endpoint group, method and status class. It counts our served documents (actor, outbox, collection, object, activity, licence, webfinger, nodeinfo) by kind, status and whether signed, per day and never per server, and never names a circle's collections. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2 --- .../Statistics/LedgerOverHttpTests.cs | 28 +++ .../Statistics/OutboundRequestTests.cs | 130 ++++++++++++ PrivaPub.Tests/Support/Peer.cs | 5 +- .../Mastodon/Controllers/MediaController.cs | 9 +- PrivaPub/Domain/Content/LinkPreviews.cs | 18 +- .../Federation/Actors/RemoteActorService.cs | 4 + PrivaPub/Federation/Inbox/InboxReceiver.cs | 2 + PrivaPub/Federation/Inbox/RemotePosts.cs | 2 + .../Federation/Objects/InstanceDescriber.cs | 1 + .../Infrastructure/Http/FederationHttp.cs | 186 ++++++++++++++++-- PrivaPub/Infrastructure/Http/HttpScope.cs | 39 ++++ PrivaPub/Infrastructure/Jobs/JobWorker.cs | 2 + .../Infrastructure/Statistics/TrafficMeter.cs | 78 ++++++++ PrivaPub/Program.cs | 2 + 14 files changed, 487 insertions(+), 19 deletions(-) create mode 100644 PrivaPub.Tests/Statistics/OutboundRequestTests.cs create mode 100644 PrivaPub/Infrastructure/Http/HttpScope.cs create mode 100644 PrivaPub/Infrastructure/Statistics/TrafficMeter.cs diff --git a/PrivaPub.Tests/Statistics/LedgerOverHttpTests.cs b/PrivaPub.Tests/Statistics/LedgerOverHttpTests.cs index 698446c..dd10e30 100644 --- a/PrivaPub.Tests/Statistics/LedgerOverHttpTests.cs +++ b/PrivaPub.Tests/Statistics/LedgerOverHttpTests.cs @@ -81,5 +81,33 @@ namespace PrivaPub.Tests.Statistics Assert.DoesNotContain(bob.Name, everything); Assert.DoesNotContain("/notes/", everything); } + + [Fact] + public async Task Client_and_served_traffic_is_counted_per_day_never_per_server() + { + var token = TestContext.Current.CancellationToken; + var persona = await _host.Persona(await _host.SignUp(), "served"); + var reader = new RemoteActor(_peer, "reader"); + var day = DateTime.UtcNow.Date; + await _host.Get().Flush(token); + var before = await DB.Default.Find().Match(d => d.Day == day).ExecuteFirstAsync(token) ?? new ServerDay(); + using var client = _host.Client(); + + Assert.Equal(HttpStatusCode.OK, (await client.GetAsync("/api/v1/instance", token)).StatusCode); + using (var actorRequest = new HttpRequestMessage(HttpMethod.Get, $"/peasants/{persona.UserName}")) + { + actorRequest.Headers.Accept.ParseAdd("application/activity+json"); + Assert.Equal(HttpStatusCode.OK, (await client.SendAsync(actorRequest, token)).StatusCode); + } + Assert.Equal(HttpStatusCode.OK, (await client.SendAsync(reader.SignedGet($"/peasants/{persona.UserName}/anus"), token)).StatusCode); + await _host.Get().Flush(token); + + var after = await DB.Default.Find().Match(d => d.Day == day).ExecuteFirstAsync(token); + long Grew(Dictionary now, Dictionary then, string key) => now.GetValueOrDefault(key) - then.GetValueOrDefault(key); + Assert.True(Grew(after.Client, before.Client, "api/v1/instance:GET:2xx") >= 1); + Assert.True(Grew(after.Served, before.Served, "actor:200:unsigned") >= 1); + Assert.True(Grew(after.Served, before.Served, "outbox:200:signed") >= 1); + Assert.DoesNotContain(after.Served.Keys, k => k.Contains(persona.UserName)); + } } } diff --git a/PrivaPub.Tests/Statistics/OutboundRequestTests.cs b/PrivaPub.Tests/Statistics/OutboundRequestTests.cs new file mode 100644 index 0000000..2227e23 --- /dev/null +++ b/PrivaPub.Tests/Statistics/OutboundRequestTests.cs @@ -0,0 +1,130 @@ +using Microsoft.AspNetCore.Http; + +using PrivaPub.Infrastructure.Http; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Tests.Support; + +namespace PrivaPub.Tests.Statistics +{ + public class HttpScopeTests + { + [Fact] + public async Task Scopes_nest_and_restore_across_awaits() + { + Assert.Null(HttpScope.Purpose); + Assert.Equal("request", HttpScope.Trigger); + + using (HttpScope.Triggered("deliver")) + { + using (HttpScope.Default("object")) + { + Assert.Equal("object", HttpScope.Purpose); + using (HttpScope.For("actor")) + { + await Task.Yield(); + Assert.Equal("actor", HttpScope.Purpose); + using (HttpScope.Default("object")) + Assert.Equal("actor", HttpScope.Purpose); + } + Assert.Equal("object", HttpScope.Purpose); + } + Assert.Null(HttpScope.Purpose); + Assert.Equal("deliver", HttpScope.Trigger); + Assert.False(HttpScope.Crawling); + using (HttpScope.Crawl()) + Assert.True(HttpScope.Crawling); + } + Assert.Equal("request", HttpScope.Trigger); + } + } + + public class TrafficMeterTests + { + [Theory] + [InlineData("api/v1/accounts/{id}/follow", "api/v1/accounts")] + [InlineData("/api/v1/statuses", "api/v1/statuses")] + [InlineData("/oauth/token", "oauth/token")] + [InlineData("/clientapi/avatar/private/insert", "clientapi/avatar/private")] + [InlineData("api/{x}", "api")] + [InlineData(null, "unmatched")] + public void Client_routes_group_by_their_first_literal_segments(string template, string group) => + Assert.Equal(group, TrafficMeter.Group(template)); + + [Theory] + [InlineData("peasants/{actor}", "/peasants/alice", "actor")] + [InlineData("peasants/{actor}/anus", "/peasants/alice/anus", "outbox")] + [InlineData("peasants/{actor}/flock", "/peasants/c/flock", "collection")] + [InlineData("peasants/{actor}/groupies", "/peasants/alice/groupies", "collection")] + [InlineData("peasants/{actor}/scribbles/{postId}", "/peasants/alice/scribbles/1", "object")] + [InlineData("peasants/{actor}/grunts/{activityId}", "/peasants/alice/grunts/1", "activity")] + [InlineData("peasants/{actor}/parrot-licences/{licenceId}", "/peasants/alice/parrot-licences/1", "licence")] + [InlineData("users/{actor}", "/users/alice", "actor")] + [InlineData(".well-known/webfinger", "/.well-known/webfinger", "webfinger")] + [InlineData("nodeinfo/2.1", "/nodeinfo/2.1", "nodeinfo")] + [InlineData("@{user}", "/@alice", null)] + public void Served_documents_are_counted_by_kind_never_by_name(string template, string path, string kind) => + Assert.Equal(kind, TrafficMeter.Served(template, new PathString(path))); + } + + public sealed class OutboundRequestTests : IAsyncLifetime + { + Peer _peer; + + public async ValueTask InitializeAsync() => _peer = await Peer.Start(); + + public async ValueTask DisposeAsync() => await _peer.DisposeAsync(); + + [Fact] + public async Task A_federation_fetch_is_recorded_with_its_purpose_trigger_and_size() + { + var ledger = new MemoryLedger(); + var http = Peer.Http(ledger: ledger); + _peer.Serve("/users/a", "{\"id\":\"{A}/users/a\",\"type\":\"Person\"}"); + + using (HttpScope.Triggered("processinbox")) + using (HttpScope.For("actor")) + Assert.NotNull(await http.GetJson($"{_peer.A}/users/a", "application/activity+json", default, TestContext.Current.CancellationToken)); + + var e = Assert.Single(ledger.Of("http")); + Assert.Equal(("127.0.0.1", "actor", "processinbox", "ok", 200), (e.Host, e.Purpose, e.Trigger, e.Outcome, e.Status!.Value)); + Assert.True(e.Bytes > 10); + Assert.Equal(0, e.Redirects); + } + + [Fact] + public async Task Refusals_are_recorded_and_a_second_try_is_remembered() + { + var ledger = new MemoryLedger(); + var http = Peer.Http(ledger: ledger); + _peer.Answer("/gone", 410); + _peer.ServeText("/page", "", "text/html"); + + using (HttpScope.Triggered("fetchancestors")) + { + await http.GetJson($"{_peer.A}/gone", "application/activity+json", default, TestContext.Current.CancellationToken); + await http.GetJson($"{_peer.A}/gone", "application/activity+json", default, TestContext.Current.CancellationToken); + await http.GetJson($"{_peer.A}/page", "application/activity+json", default, TestContext.Current.CancellationToken); + await http.GetJson("ftp://example.org/x", "application/activity+json", default, TestContext.Current.CancellationToken); + } + + var events = ledger.Of("http"); + Assert.Equal(new[] { ("refused", "410"), ("refused", "remembered"), ("refused", "content-type"), ("refused", "disallowed") }, + events.Select(e => (e.Outcome, e.Reason))); + Assert.Equal("object", events[0].Purpose); + } + + [Fact] + public async Task A_fetch_a_reader_caused_is_only_counted() + { + var ledger = new MemoryLedger(); + var http = Peer.Http(ledger: ledger); + _peer.Serve("/users/b", "{\"id\":\"{A}/users/b\",\"type\":\"Person\"}"); + + using (HttpScope.For("webfinger")) + await http.GetJson($"{_peer.A}/users/b", "application/activity+json", default, TestContext.Current.CancellationToken); + + Assert.Empty(ledger.Events); + Assert.Equal(1, ledger.Counts[("127.0.0.1", "http:webfinger:ok")]); + } + } +} diff --git a/PrivaPub.Tests/Support/Peer.cs b/PrivaPub.Tests/Support/Peer.cs index d08b157..2ef1217 100644 --- a/PrivaPub.Tests/Support/Peer.cs +++ b/PrivaPub.Tests/Support/Peer.cs @@ -8,6 +8,7 @@ using Microsoft.Extensions.Options; using PrivaPub.Federation.Moderation; using PrivaPub.Infrastructure.Http; +using PrivaPub.Infrastructure.Statistics; using PrivaPub.Models.Federation; using System.Collections.Concurrent; @@ -79,7 +80,7 @@ namespace PrivaPub.Tests.Support public void Answer(string path, int status, TimeSpan delay = default) => _answers[path] = (status, delay); - public static FederationHttp Http(IMemoryCache cache = default, IDomainBlocks blocks = default) + public static FederationHttp Http(IMemoryCache cache = default, IDomainBlocks blocks = default, IInteractionLedger ledger = default) { var options = new FederationOptions { AllowPrivateNetworks = true, AllowPlainHttp = true }; var services = new ServiceCollection(); @@ -87,7 +88,7 @@ namespace PrivaPub.Tests.Support .ConfigurePrimaryHttpMessageHandler(() => SafeHttpHandlerFactory.Create(options)); return new FederationHttp(services.BuildServiceProvider().GetRequiredService(), cache ?? new MemoryCache(new MemoryCacheOptions()), new StaticOptions(options), - blocks ?? new NoBlocks(), NullLogger.Instance); + blocks ?? new NoBlocks(), NullLogger.Instance, ledger); } public async ValueTask DisposeAsync() => await _app.DisposeAsync(); diff --git a/PrivaPub/Api/Mastodon/Controllers/MediaController.cs b/PrivaPub/Api/Mastodon/Controllers/MediaController.cs index b445b3b..71928ff 100644 --- a/PrivaPub/Api/Mastodon/Controllers/MediaController.cs +++ b/PrivaPub/Api/Mastodon/Controllers/MediaController.cs @@ -1,3 +1,4 @@ +using PrivaPub.Infrastructure.Statistics; using Microsoft.AspNetCore.Authorization; using Microsoft.AspNetCore.Mvc; @@ -18,10 +19,13 @@ namespace PrivaPub.Api.Mastodon.Controllers readonly IMediaService _media; readonly IMediaProxy _proxy; - public MediaController(IMediaService media, IMediaProxy proxy) + readonly IInteractionLedger _ledger; + + public MediaController(IMediaService media, IMediaProxy proxy, IInteractionLedger ledger = default) { _media = media; _proxy = proxy; + _ledger = ledger; } [HttpPost("/api/v1/media"), HttpPost("/api/v2/media"), Scope("write:media"), RequestSizeLimit(UploadLimit), @@ -67,7 +71,10 @@ namespace PrivaPub.Api.Mastodon.Controllers Response.Headers["Content-Security-Policy"] = "default-src 'none'; sandbox"; Response.Headers["Cache-Control"] = "public, max-age=604800"; if (_proxy.Cached(url) is { Path: not null } cached) + { + _ledger?.Count(Interactions.HostOf(url), "media:hit"); return PhysicalFile(cached.Path, cached.ContentType, enableRangeProcessing: true); + } if (Request.Headers.Range.Count == 0) { var (path, contentType) = await _proxy.Fetch(signature, encoded, token); diff --git a/PrivaPub/Domain/Content/LinkPreviews.cs b/PrivaPub/Domain/Content/LinkPreviews.cs index 1bfb5aa..7ba765c 100644 --- a/PrivaPub/Domain/Content/LinkPreviews.cs +++ b/PrivaPub/Domain/Content/LinkPreviews.cs @@ -4,6 +4,8 @@ using Microsoft.Extensions.Options; using MongoDB.Entities; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Statistics; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Objects; using PrivaPub.Infrastructure.Http; @@ -46,14 +48,17 @@ namespace PrivaPub.Domain.Content readonly IJobQueue _queue; readonly IOptionsMonitor _options; + readonly IInteractionLedger _ledger; + public LinkPreviews(DbEntities dbEntities, ILocalActorService localActors, IFederationHttp http, IJobQueue queue, - IOptionsMonitor options) + IOptionsMonitor options, IInteractionLedger ledger = default) { _dbEntities = dbEntities; _localActors = localActors; _http = http; _queue = queue; _options = options; + _ledger = ledger; } public JobKind Kind => JobKind.FetchPreview; @@ -104,8 +109,19 @@ namespace PrivaPub.Domain.Content async Task Fetch(string url, CancellationToken token) { + var started = System.Diagnostics.Stopwatch.GetTimestamp(); var (finalUri, html) = await _http.GetPage(url, token); var preview = html == default ? new LinkPreview { Url = url, Failed = true } : Read(url, finalUri, html); + _ledger?.Record(new InteractionEvent + { + Channel = Interactions.Preview, + Host = Interactions.HostOf(url), + Outcome = html == default ? Interactions.Failed : Interactions.Ok, + Reason = html == default ? default : preview.Title == default && preview.ImageURL == default ? "no-card" : "card", + LatencyMs = (int)System.Diagnostics.Stopwatch.GetElapsedTime(started).TotalMilliseconds, + Bytes = html == default ? default : System.Text.Encoding.UTF8.GetByteCount(html), + Redirects = finalUri == default || finalUri.AbsoluteUri == url ? 0 : 1 + }); await DB.Default.Update() .Match(p => p.Url == url) .Modify(p => p.Title, preview.Title) diff --git a/PrivaPub/Federation/Actors/RemoteActorService.cs b/PrivaPub/Federation/Actors/RemoteActorService.cs index 13d00f8..21e72f2 100644 --- a/PrivaPub/Federation/Actors/RemoteActorService.cs +++ b/PrivaPub/Federation/Actors/RemoteActorService.cs @@ -48,6 +48,7 @@ namespace PrivaPub.Federation.Actors public async Task FetchObject(string uri, CancellationToken token) { + using var scope = HttpScope.Default("object"); var signer = await _localActors.GetInstanceActor(token); var fetched = await Get(uri, signer, token); if (fetched == default) @@ -81,6 +82,7 @@ namespace PrivaPub.Federation.Actors if (!MayFetch(actorUri)) return cached; + using var scope = HttpScope.For("actor"); using var fetched = await FetchObject(actorUri, token); var actor = fetched == default ? default : ActorDocument.Parse(fetched.Root); if (actor == default) @@ -101,6 +103,7 @@ namespace PrivaPub.Federation.Actors if (!MayFetch(keyId)) return cached; + using var scope = HttpScope.For("key"); using var fetched = await FetchObject(StripFragment(keyId), token); if (fetched == default) return cached; @@ -129,6 +132,7 @@ namespace PrivaPub.Federation.Actors if (parts is not { Length: 2 } || string.IsNullOrEmpty(parts[0]) || string.IsNullOrEmpty(parts[1])) return default; + using var scope = HttpScope.For("webfinger"); var query = $"/.well-known/webfinger?resource={Uri.EscapeDataString($"acct:{parts[0]}@{parts[1]}")}"; using var fetched = await _http.GetJson($"https://{parts[1]}{query}", "application/jrd+json, application/json", sign: default, token) ?? (_options?.CurrentValue.AllowPlainHttp == true diff --git a/PrivaPub/Federation/Inbox/InboxReceiver.cs b/PrivaPub/Federation/Inbox/InboxReceiver.cs index 8e6fc24..66400f9 100644 --- a/PrivaPub/Federation/Inbox/InboxReceiver.cs +++ b/PrivaPub/Federation/Inbox/InboxReceiver.cs @@ -2,6 +2,7 @@ using PrivaPub.Federation.Actors; using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Objects; using PrivaPub.Federation.Signing; +using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Jobs; using PrivaPub.Infrastructure.Statistics; using PrivaPub.Models.Federation; @@ -72,6 +73,7 @@ namespace PrivaPub.Federation.Inbox public async Task Receive(HttpRequest request, LocalActor recipient, CancellationToken token) { + using var scope = HttpScope.Triggered("verify"); var started = System.Diagnostics.Stopwatch.GetTimestamp(); var receipt = new Receipt { Inbox = recipient == default ? "shared" : "personal", Bytes = request.ContentLength }; InboxResult result = default; diff --git a/PrivaPub/Federation/Inbox/RemotePosts.cs b/PrivaPub/Federation/Inbox/RemotePosts.cs index db67008..2e6bf6d 100644 --- a/PrivaPub/Federation/Inbox/RemotePosts.cs +++ b/PrivaPub/Federation/Inbox/RemotePosts.cs @@ -1,6 +1,7 @@ using MongoDB.Driver; using MongoDB.Entities; +using PrivaPub.Infrastructure.Http; using PrivaPub.Domain.Content; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Moderation; @@ -123,6 +124,7 @@ namespace PrivaPub.Federation.Inbox if (existing != default || depth > MaxDepth || await DB.Default.Find().Match(d => d.ObjectURI == objectUri).ExecuteAnyAsync(token)) return existing; + using var scope = HttpScope.For("context"); using var fetched = await _remoteActors.FetchObject(objectUri, token); var note = fetched == default ? default : NoteParser.Parse(JsonNode.Parse(fetched.Root.GetRawText())); if (note == default || !Origin.Same(note.Id, note.AttributedTo)) diff --git a/PrivaPub/Federation/Objects/InstanceDescriber.cs b/PrivaPub/Federation/Objects/InstanceDescriber.cs index 9a70bf8..a5d1bf4 100644 --- a/PrivaPub/Federation/Objects/InstanceDescriber.cs +++ b/PrivaPub/Federation/Objects/InstanceDescriber.cs @@ -29,6 +29,7 @@ namespace PrivaPub.Federation.Objects public async Task Handle(Job job, CancellationToken token) { var host = job.Payload; + using var scope = HttpScope.For("nodeinfo"); using var links = await _http.GetJson($"https://{host}/.well-known/nodeinfo", "application/json", sign: default, token); var href = links == default ? default : Schemas.Select(schema => LinkFor(links.Root, schema)).FirstOrDefault(link => link != default); if (href == default) diff --git a/PrivaPub/Infrastructure/Http/FederationHttp.cs b/PrivaPub/Infrastructure/Http/FederationHttp.cs index 3d0baae..e49a24b 100644 --- a/PrivaPub/Infrastructure/Http/FederationHttp.cs +++ b/PrivaPub/Infrastructure/Http/FederationHttp.cs @@ -2,6 +2,8 @@ using Microsoft.Extensions.Caching.Memory; using Microsoft.Extensions.Options; using PrivaPub.Federation.Moderation; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Models.Statistics; using System.Net; using System.Text.Json; @@ -49,17 +51,71 @@ namespace PrivaPub.Infrastructure.Http readonly IOptionsMonitor _options; readonly IDomainBlocks _domainBlocks; readonly ILogger _logger; + readonly IInteractionLedger _ledger; public FederationHttp(IHttpClientFactory httpClientFactory, IMemoryCache cache, IOptionsMonitor options, - IDomainBlocks domainBlocks, ILogger logger) + IDomainBlocks domainBlocks, ILogger logger, IInteractionLedger ledger = default) { _httpClientFactory = httpClientFactory; _cache = cache; _options = options; _domainBlocks = domainBlocks; _logger = logger; + _ledger = ledger; } + sealed class Exchange + { + public Exchange(string url, string purpose) + { + Host = Interactions.HostOf(url); + Purpose = purpose; + } + + public readonly string Host; + public readonly string Purpose; + public readonly long Started = System.Diagnostics.Stopwatch.GetTimestamp(); + public int? Status; + public long? Bytes; + public int Hops; + public string Outcome = Interactions.Ok; + public string Reason; + + public void Refused(string reason) => (Outcome, Reason) = (Interactions.Refused, reason); + + public void Failed(string reason) => (Outcome, Reason) = (Interactions.Failed, reason); + + public void Answered(HttpResponseMessage response) => Status = (int)response.StatusCode; + } + + void Record(Exchange exchange) + { + if (_ledger == default) + return; + var trigger = HttpScope.Trigger; + if (trigger == HttpScope.Request) + { + _ledger.Count(exchange.Host, $"http:{exchange.Purpose}:{exchange.Outcome}", exchange.Bytes ?? 0); + return; + } + _ledger.Record(new InteractionEvent + { + Channel = Interactions.Http, + Host = exchange.Host, + Purpose = exchange.Purpose, + Trigger = trigger, + Outcome = exchange.Outcome, + Reason = exchange.Reason, + Status = exchange.Status, + LatencyMs = (int)System.Diagnostics.Stopwatch.GetElapsedTime(exchange.Started).TotalMilliseconds, + Bytes = exchange.Bytes, + Redirects = exchange.Hops, + Crawl = HttpScope.Crawling + }); + } + + static string StatusReason(HttpResponseMessage response) => ((int)response.StatusCode).ToString(System.Globalization.CultureInfo.InvariantCulture); + public bool IsAllowed(Uri target) { if (target is not { IsAbsoluteUri: true } || !string.IsNullOrEmpty(target.UserInfo)) @@ -79,12 +135,31 @@ namespace PrivaPub.Infrastructure.Http } public async Task GetJson(string url, string accept, Action sign, CancellationToken token) + { + var exchange = new Exchange(url, HttpScope.Purpose ?? "object"); + try + { + return await GetJson(url, accept, sign, exchange, token); + } + finally + { + Record(exchange); + } + } + + async Task GetJson(string url, string accept, Action sign, Exchange exchange, CancellationToken token) { if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) + { + exchange.Refused("disallowed"); return default; + } var negativeKey = NegativeKey(target); if (_cache.TryGetValue(negativeKey, out _)) + { + exchange.Refused("remembered"); return default; + } using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token); timeout.CancelAfter(RequestTimeout); @@ -98,46 +173,49 @@ namespace PrivaPub.Infrastructure.Http using var response = await _httpClientFactory.CreateClient(ClientName) .SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token); + exchange.Answered(response); if (IsRedirect(response.StatusCode)) { var location = response.Headers.Location; var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location); if (!IsAllowed(next)) - return Refuse(negativeKey, url, "a redirect to a disallowed location"); + return Refuse(negativeKey, url, "a redirect to a disallowed location", exchange, "bad-redirect"); target = next; + exchange.Hops++; continue; } if (!response.IsSuccessStatusCode) - return Refuse(negativeKey, url, $"status {(int)response.StatusCode}", + return Refuse(negativeKey, url, $"status {(int)response.StatusCode}", exchange, StatusReason(response), transient: (int)response.StatusCode is >= 500 or 429 or 408); var mediaType = response.Content.Headers.ContentType?.MediaType; if (mediaType == default || !JsonMediaTypes.Contains(mediaType, StringComparer.OrdinalIgnoreCase)) - return Refuse(negativeKey, url, $"content type '{mediaType}'"); + return Refuse(negativeKey, url, $"content type '{mediaType}'", exchange, "content-type"); if (response.Content.Headers.ContentLength > MaxResponseBytes) - return Refuse(negativeKey, url, "a body over the size limit"); + return Refuse(negativeKey, url, "a body over the size limit", exchange, "too-large"); var body = await ReadBounded(response.Content, MaxResponseBytes, timeout.Token); if (body == default) - return Refuse(negativeKey, url, "a body over the size limit"); + return Refuse(negativeKey, url, "a body over the size limit", exchange, "too-large"); + exchange.Bytes = body.Length; return new FetchedJson { FinalUri = target, Document = JsonDocument.Parse(body) }; } - return Refuse(negativeKey, url, "too many redirects"); + return Refuse(negativeKey, url, "too many redirects", exchange, "too-many-redirects"); } catch (OperationCanceledException) when (!token.IsCancellationRequested) { - return Refuse(negativeKey, url, "a timeout", transient: true); + return Refuse(negativeKey, url, "a timeout", exchange, "timeout", transient: true); } catch (HttpRequestException ex) { - return Refuse(negativeKey, url, ex.Message, transient: true); + return Refuse(negativeKey, url, ex.Message, exchange, "network", transient: true); } catch (Exception ex) when (ex is JsonException or BlockedDestinationException) { - return Refuse(negativeKey, url, ex.Message); + return Refuse(negativeKey, url, ex.Message, exchange, ex is JsonException ? "bad-json" : "private-address"); } } @@ -145,9 +223,28 @@ namespace PrivaPub.Infrastructure.Http Uri.TryCreate(url, UriKind.Absolute, out var target) && _cache.TryGetValue(NegativeKey(target), out bool transient) && transient; public async Task<(byte[] Bytes, string ContentType)> GetMedia(string url, long maxBytes, CancellationToken token) + { + var exchange = new Exchange(url, "media"); + try + { + var media = await GetMedia(url, maxBytes, exchange, token); + if (media.Bytes == default && exchange.Outcome == Interactions.Ok) + exchange.Refused("unusable"); + return media; + } + finally + { + Record(exchange); + } + } + + async Task<(byte[] Bytes, string ContentType)> GetMedia(string url, long maxBytes, Exchange exchange, CancellationToken token) { if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) + { + exchange.Refused("disallowed"); return default; + } using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token); timeout.CancelAfter(TimeSpan.FromSeconds(60)); try @@ -157,27 +254,52 @@ namespace PrivaPub.Infrastructure.Http using var request = new HttpRequestMessage(HttpMethod.Get, target); request.Headers.Accept.ParseAdd("image/*, video/*, audio/*"); using var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token); + exchange.Answered(response); if (IsRedirect(response.StatusCode)) { var location = response.Headers.Location; var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location); if (!IsAllowed(next)) + { + exchange.Refused("bad-redirect"); return default; + } target = next; + exchange.Hops++; continue; } var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant(); - if (!response.IsSuccessStatusCode || mediaType == default - || !(mediaType.StartsWith("image/") || mediaType.StartsWith("video/") || mediaType.StartsWith("audio/")) - || mediaType.Contains("svg") || response.Content.Headers.ContentLength > maxBytes) + if (!response.IsSuccessStatusCode) + { + exchange.Failed(StatusReason(response)); return default; + } + if (mediaType == default || !(mediaType.StartsWith("image/") || mediaType.StartsWith("video/") || mediaType.StartsWith("audio/")) + || mediaType.Contains("svg")) + { + exchange.Refused("content-type"); + return default; + } + if (response.Content.Headers.ContentLength > maxBytes) + { + exchange.Refused("too-large"); + return default; + } var bytes = await ReadBounded(response.Content, (int)Math.Min(maxBytes, int.MaxValue), timeout.Token); - return bytes == default ? default : (bytes, mediaType); + if (bytes == default) + { + exchange.Refused("too-large"); + return default; + } + exchange.Bytes = bytes.Length; + return (bytes, mediaType); } + exchange.Refused("too-many-redirects"); return default; } catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or OperationCanceledException && !token.IsCancellationRequested) { + exchange.Failed(ex switch { BlockedDestinationException => "private-address", OperationCanceledException => "timeout", _ => "network" }); _logger.LogInformation("Media {Url} refused: {Reason}", url, ex.Message); return default; } @@ -231,9 +353,29 @@ namespace PrivaPub.Infrastructure.Http } public async Task OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token) + { + var exchange = new Exchange(url, "stream"); + try + { + var response = await OpenMedia(url, range, exchange, token); + if (response == default && exchange.Outcome == Interactions.Ok) + exchange.Refused("unusable"); + exchange.Bytes = response?.Content.Headers.ContentLength; + return response; + } + finally + { + Record(exchange); + } + } + + async Task OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, Exchange exchange, CancellationToken token) { if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) + { + exchange.Refused("disallowed"); return default; + } try { for (var hop = 0; hop <= MaxRedirects; hop++) @@ -242,20 +384,29 @@ namespace PrivaPub.Infrastructure.Http request.Headers.Accept.ParseAdd("video/*, audio/*, image/*"); request.Headers.Range = range; var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, token); + exchange.Answered(response); if (IsRedirect(response.StatusCode)) { var location = response.Headers.Location; response.Dispose(); var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location); if (!IsAllowed(next)) + { + exchange.Refused("bad-redirect"); return default; + } target = next; + exchange.Hops++; continue; } var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant(); if (!response.IsSuccessStatusCode || mediaType == default || mediaType.Contains("svg") || !(mediaType.StartsWith("video/") || mediaType.StartsWith("audio/") || mediaType.StartsWith("image/") || mediaType == "application/octet-stream")) { + if (response.IsSuccessStatusCode) + exchange.Refused("content-type"); + else + exchange.Failed(StatusReason(response)); response.Dispose(); return default; } @@ -265,6 +416,7 @@ namespace PrivaPub.Infrastructure.Http } catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or OperationCanceledException && !token.IsCancellationRequested) { + exchange.Failed(ex switch { BlockedDestinationException => "private-address", OperationCanceledException => "timeout", _ => "network" }); _logger.LogInformation("Media stream {Url} refused: {Reason}", url, ex.Message); return default; } @@ -311,8 +463,12 @@ namespace PrivaPub.Infrastructure.Http static string NegativeKey(Uri target) => "federation-http:refused:" + target.AbsoluteUri; - FetchedJson Refuse(string negativeKey, string url, string reason, bool transient = false) + FetchedJson Refuse(string negativeKey, string url, string reason, Exchange exchange, string code, bool transient = false) { + if (transient) + exchange.Failed(code); + else + exchange.Refused(code); _cache.Set(negativeKey, transient, NegativeCacheLifetime); _logger.LogInformation("GET {Url} refused: {Reason}", url, reason); return default; diff --git a/PrivaPub/Infrastructure/Http/HttpScope.cs b/PrivaPub/Infrastructure/Http/HttpScope.cs new file mode 100644 index 0000000..ba5d6cd --- /dev/null +++ b/PrivaPub/Infrastructure/Http/HttpScope.cs @@ -0,0 +1,39 @@ +namespace PrivaPub.Infrastructure.Http +{ + public static class HttpScope + { + public const string Request = "request"; + + sealed record Tags(string Purpose, string Trigger, bool Crawl); + + static readonly AsyncLocal current = new(); + + public static string Purpose => current.Value?.Purpose; + public static string Trigger => current.Value?.Trigger ?? Request; + public static bool Crawling => current.Value?.Crawl == true; + + public static IDisposable For(string purpose) => Push(tags => tags with { Purpose = purpose }); + + public static IDisposable Default(string purpose) => Push(tags => tags with { Purpose = tags.Purpose ?? purpose }); + + public static IDisposable Triggered(string trigger) => Push(tags => tags with { Trigger = trigger }); + + public static IDisposable Crawl() => Push(tags => tags with { Crawl = true }); + + static IDisposable Push(Func change) + { + var previous = current.Value; + current.Value = change(previous ?? new Tags(default, default, false)); + return new Restore(previous); + } + + sealed class Restore : IDisposable + { + readonly Tags _previous; + + public Restore(Tags previous) => _previous = previous; + + public void Dispose() => current.Value = _previous; + } + } +} diff --git a/PrivaPub/Infrastructure/Jobs/JobWorker.cs b/PrivaPub/Infrastructure/Jobs/JobWorker.cs index da6dd15..b362307 100644 --- a/PrivaPub/Infrastructure/Jobs/JobWorker.cs +++ b/PrivaPub/Infrastructure/Jobs/JobWorker.cs @@ -1,3 +1,4 @@ +using PrivaPub.Infrastructure.Http; using PrivaPub.Models.Jobs; using System.Collections.Concurrent; @@ -73,6 +74,7 @@ namespace PrivaPub.Infrastructure.Jobs JobOutcome outcome; try { + using var scope = HttpScope.Triggered(handler.Kind.ToString().ToLowerInvariant()); outcome = await handler.Handle(job, stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) diff --git a/PrivaPub/Infrastructure/Statistics/TrafficMeter.cs b/PrivaPub/Infrastructure/Statistics/TrafficMeter.cs new file mode 100644 index 0000000..793b0c2 --- /dev/null +++ b/PrivaPub/Infrastructure/Statistics/TrafficMeter.cs @@ -0,0 +1,78 @@ +using System.Diagnostics; + +namespace PrivaPub.Infrastructure.Statistics +{ + public static class TrafficMeter + { + public static IApplicationBuilder UseTrafficMeter(this IApplicationBuilder app) => + app.Use(async (context, next) => + { + var started = Stopwatch.GetTimestamp(); + try + { + await next(context); + } + finally + { + Count(context, (int)Stopwatch.GetElapsedTime(started).TotalMilliseconds); + } + }); + + static void Count(HttpContext context, int latencyMs) + { + var ledger = context.RequestServices?.GetService(); + if (ledger == default) + return; + var path = context.Request.Path; + var template = (context.GetEndpoint() as RouteEndpoint)?.RoutePattern.RawText; + var method = context.Request.Method; + var status = context.Response.StatusCode; + if (path.StartsWithSegments("/api") || path.StartsWithSegments("/oauth") || path.StartsWithSegments("/clientapi")) + { + ledger.CountServer(ServerSections.Client, $"{Group(template)}:{method}:{status / 100}xx", latencyMs); + return; + } + if (method != HttpMethods.Get && method != HttpMethods.Head) + return; + var served = Served(template, path); + if (served != default) + ledger.CountServer(ServerSections.Served, + $"{served}:{status}:{(context.Request.Headers.ContainsKey("Signature") || context.Request.Headers.ContainsKey("Signature-Input") ? "signed" : "unsigned")}"); + } + + public static string Group(string template) + { + if (template == default) + return "unmatched"; + var literals = template.Trim('/').Split('/', StringSplitOptions.RemoveEmptyEntries) + .TakeWhile(segment => !segment.Contains('{')) + .Take(3); + var group = string.Join('/', literals); + return group.Length == 0 ? "unmatched" : group; + } + + public static string Served(string template, PathString path) + { + if (path.StartsWithSegments("/.well-known/webfinger")) + return "webfinger"; + if (path.StartsWithSegments("/.well-known/nodeinfo") || path.StartsWithSegments("/nodeinfo")) + return "nodeinfo"; + if (!path.StartsWithSegments("/peasants") && !path.StartsWithSegments("/users")) + return default; + var segments = (template ?? path.Value ?? string.Empty).Trim('/').Split('/'); + return segments.Length switch + { + <= 2 => "actor", + _ => segments[2] switch + { + "anus" => "outbox", + "scribbles" or "whispers" => "object", + "grunts" => "activity", + "parrot-licences" => "licence", + "mouth" or "human-centipede" => default, + _ => "collection" + } + }; + } + } +} diff --git a/PrivaPub/Program.cs b/PrivaPub/Program.cs index dd9ca6a..0a0853a 100644 --- a/PrivaPub/Program.cs +++ b/PrivaPub/Program.cs @@ -17,6 +17,7 @@ using PrivaPub.Infrastructure; using PrivaPub.Infrastructure.Cli; using PrivaPub.Infrastructure.Data; using PrivaPub.Infrastructure.Http; +using PrivaPub.Infrastructure.Statistics; using PrivaPub.Middleware; using PrivaPub.Models; using PrivaPub.Services; @@ -150,6 +151,7 @@ try app.UseRequestLocalization(await localizationService.Get()); app.UseRouting(); + app.UseTrafficMeter(); app.UseRateLimiter(); app.UseAuthentication();