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();