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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
This commit is contained in:
thepraandClaude Opus 5.5 committed 2026-10-03 11:10:42 +02:00
1 parent e247001bdb
commit d37999970c
14 files changed
+488 -20

No files matched your search

@@ -81,5 +81,33 @@ namespace PrivaPub.Tests.Statistics
Assert.DoesNotContain(bob.Name, everything); Assert.DoesNotContain(bob.Name, everything);
Assert.DoesNotContain("/notes/", 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<InteractionLedger>().Flush(token);
var before = await DB.Default.Find<ServerDay>().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<InteractionLedger>().Flush(token);
var after = await DB.Default.Find<ServerDay>().Match(d => d.Day == day).ExecuteFirstAsync(token);
long Grew(Dictionary<string, long> now, Dictionary<string, long> 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));
}
} }
} }
@@ -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", "<html></html>", "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")]);
}
}
}
+3 -2
View File
@@ -8,6 +8,7 @@ using Microsoft.Extensions.Options;
using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Moderation;
using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Http;
using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Models.Federation; using PrivaPub.Models.Federation;
using System.Collections.Concurrent; 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 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 options = new FederationOptions { AllowPrivateNetworks = true, AllowPlainHttp = true };
var services = new ServiceCollection(); var services = new ServiceCollection();
@@ -87,7 +88,7 @@ namespace PrivaPub.Tests.Support
.ConfigurePrimaryHttpMessageHandler(() => SafeHttpHandlerFactory.Create(options)); .ConfigurePrimaryHttpMessageHandler(() => SafeHttpHandlerFactory.Create(options));
return new FederationHttp(services.BuildServiceProvider().GetRequiredService<IHttpClientFactory>(), return new FederationHttp(services.BuildServiceProvider().GetRequiredService<IHttpClientFactory>(),
cache ?? new MemoryCache(new MemoryCacheOptions()), new StaticOptions<FederationOptions>(options), cache ?? new MemoryCache(new MemoryCacheOptions()), new StaticOptions<FederationOptions>(options),
blocks ?? new NoBlocks(), NullLogger<FederationHttp>.Instance); blocks ?? new NoBlocks(), NullLogger<FederationHttp>.Instance, ledger);
} }
public async ValueTask DisposeAsync() => await _app.DisposeAsync(); public async ValueTask DisposeAsync() => await _app.DisposeAsync();
@@ -1,3 +1,4 @@
using PrivaPub.Infrastructure.Statistics;
using Microsoft.AspNetCore.Authorization; using Microsoft.AspNetCore.Authorization;
using Microsoft.AspNetCore.Mvc; using Microsoft.AspNetCore.Mvc;
@@ -18,10 +19,13 @@ namespace PrivaPub.Api.Mastodon.Controllers
readonly IMediaService _media; readonly IMediaService _media;
readonly IMediaProxy _proxy; readonly IMediaProxy _proxy;
public MediaController(IMediaService media, IMediaProxy proxy) readonly IInteractionLedger _ledger;
public MediaController(IMediaService media, IMediaProxy proxy, IInteractionLedger ledger = default)
{ {
_media = media; _media = media;
_proxy = proxy; _proxy = proxy;
_ledger = ledger;
} }
[HttpPost("/api/v1/media"), HttpPost("/api/v2/media"), Scope("write:media"), RequestSizeLimit(UploadLimit), [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["Content-Security-Policy"] = "default-src 'none'; sandbox";
Response.Headers["Cache-Control"] = "public, max-age=604800"; Response.Headers["Cache-Control"] = "public, max-age=604800";
if (_proxy.Cached(url) is { Path: not null } cached) if (_proxy.Cached(url) is { Path: not null } cached)
{
_ledger?.Count(Interactions.HostOf(url), "media:hit");
return PhysicalFile(cached.Path, cached.ContentType, enableRangeProcessing: true); return PhysicalFile(cached.Path, cached.ContentType, enableRangeProcessing: true);
}
if (Request.Headers.Range.Count == 0) if (Request.Headers.Range.Count == 0)
{ {
var (path, contentType) = await _proxy.Fetch(signature, encoded, token); var (path, contentType) = await _proxy.Fetch(signature, encoded, token);
+17 -1
View File
@@ -4,6 +4,8 @@ using Microsoft.Extensions.Options;
using MongoDB.Entities; using MongoDB.Entities;
using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Models.Statistics;
using PrivaPub.Federation.Actors; using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Objects; using PrivaPub.Federation.Objects;
using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Http;
@@ -46,14 +48,17 @@ namespace PrivaPub.Domain.Content
readonly IJobQueue _queue; readonly IJobQueue _queue;
readonly IOptionsMonitor<FederationOptions> _options; readonly IOptionsMonitor<FederationOptions> _options;
readonly IInteractionLedger _ledger;
public LinkPreviews(DbEntities dbEntities, ILocalActorService localActors, IFederationHttp http, IJobQueue queue, public LinkPreviews(DbEntities dbEntities, ILocalActorService localActors, IFederationHttp http, IJobQueue queue,
IOptionsMonitor<FederationOptions> options) IOptionsMonitor<FederationOptions> options, IInteractionLedger ledger = default)
{ {
_dbEntities = dbEntities; _dbEntities = dbEntities;
_localActors = localActors; _localActors = localActors;
_http = http; _http = http;
_queue = queue; _queue = queue;
_options = options; _options = options;
_ledger = ledger;
} }
public JobKind Kind => JobKind.FetchPreview; public JobKind Kind => JobKind.FetchPreview;
@@ -104,8 +109,19 @@ namespace PrivaPub.Domain.Content
async Task<LinkPreview> Fetch(string url, CancellationToken token) async Task<LinkPreview> Fetch(string url, CancellationToken token)
{ {
var started = System.Diagnostics.Stopwatch.GetTimestamp();
var (finalUri, html) = await _http.GetPage(url, token); var (finalUri, html) = await _http.GetPage(url, token);
var preview = html == default ? new LinkPreview { Url = url, Failed = true } : Read(url, finalUri, html); 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<LinkPreview>() await DB.Default.Update<LinkPreview>()
.Match(p => p.Url == url) .Match(p => p.Url == url)
.Modify(p => p.Title, preview.Title) .Modify(p => p.Title, preview.Title)
@@ -48,6 +48,7 @@ namespace PrivaPub.Federation.Actors
public async Task<FetchedJson> FetchObject(string uri, CancellationToken token) public async Task<FetchedJson> FetchObject(string uri, CancellationToken token)
{ {
using var scope = HttpScope.Default("object");
var signer = await _localActors.GetInstanceActor(token); var signer = await _localActors.GetInstanceActor(token);
var fetched = await Get(uri, signer, token); var fetched = await Get(uri, signer, token);
if (fetched == default) if (fetched == default)
@@ -81,6 +82,7 @@ namespace PrivaPub.Federation.Actors
if (!MayFetch(actorUri)) if (!MayFetch(actorUri))
return cached; return cached;
using var scope = HttpScope.For("actor");
using var fetched = await FetchObject(actorUri, token); using var fetched = await FetchObject(actorUri, token);
var actor = fetched == default ? default : ActorDocument.Parse(fetched.Root); var actor = fetched == default ? default : ActorDocument.Parse(fetched.Root);
if (actor == default) if (actor == default)
@@ -101,6 +103,7 @@ namespace PrivaPub.Federation.Actors
if (!MayFetch(keyId)) if (!MayFetch(keyId))
return cached; return cached;
using var scope = HttpScope.For("key");
using var fetched = await FetchObject(StripFragment(keyId), token); using var fetched = await FetchObject(StripFragment(keyId), token);
if (fetched == default) if (fetched == default)
return cached; return cached;
@@ -129,6 +132,7 @@ namespace PrivaPub.Federation.Actors
if (parts is not { Length: 2 } || string.IsNullOrEmpty(parts[0]) || string.IsNullOrEmpty(parts[1])) if (parts is not { Length: 2 } || string.IsNullOrEmpty(parts[0]) || string.IsNullOrEmpty(parts[1]))
return default; return default;
using var scope = HttpScope.For("webfinger");
var query = $"/.well-known/webfinger?resource={Uri.EscapeDataString($"acct:{parts[0]}@{parts[1]}")}"; 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) using var fetched = await _http.GetJson($"https://{parts[1]}{query}", "application/jrd+json, application/json", sign: default, token)
?? (_options?.CurrentValue.AllowPlainHttp == true ?? (_options?.CurrentValue.AllowPlainHttp == true
@@ -2,6 +2,7 @@ using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Moderation;
using PrivaPub.Federation.Objects; using PrivaPub.Federation.Objects;
using PrivaPub.Federation.Signing; using PrivaPub.Federation.Signing;
using PrivaPub.Infrastructure.Http;
using PrivaPub.Infrastructure.Jobs; using PrivaPub.Infrastructure.Jobs;
using PrivaPub.Infrastructure.Statistics; using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Models.Federation; using PrivaPub.Models.Federation;
@@ -72,6 +73,7 @@ namespace PrivaPub.Federation.Inbox
public async Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, CancellationToken token) public async Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, CancellationToken token)
{ {
using var scope = HttpScope.Triggered("verify");
var started = System.Diagnostics.Stopwatch.GetTimestamp(); var started = System.Diagnostics.Stopwatch.GetTimestamp();
var receipt = new Receipt { Inbox = recipient == default ? "shared" : "personal", Bytes = request.ContentLength }; var receipt = new Receipt { Inbox = recipient == default ? "shared" : "personal", Bytes = request.ContentLength };
InboxResult result = default; InboxResult result = default;
+2
View File
@@ -1,6 +1,7 @@
using MongoDB.Driver; using MongoDB.Driver;
using MongoDB.Entities; using MongoDB.Entities;
using PrivaPub.Infrastructure.Http;
using PrivaPub.Domain.Content; using PrivaPub.Domain.Content;
using PrivaPub.Federation.Actors; using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Moderation;
@@ -123,6 +124,7 @@ namespace PrivaPub.Federation.Inbox
if (existing != default || depth > MaxDepth || await DB.Default.Find<DeletedObject>().Match(d => d.ObjectURI == objectUri).ExecuteAnyAsync(token)) if (existing != default || depth > MaxDepth || await DB.Default.Find<DeletedObject>().Match(d => d.ObjectURI == objectUri).ExecuteAnyAsync(token))
return existing; return existing;
using var scope = HttpScope.For("context");
using var fetched = await _remoteActors.FetchObject(objectUri, token); using var fetched = await _remoteActors.FetchObject(objectUri, token);
var note = fetched == default ? default : NoteParser.Parse(JsonNode.Parse(fetched.Root.GetRawText())); var note = fetched == default ? default : NoteParser.Parse(JsonNode.Parse(fetched.Root.GetRawText()));
if (note == default || !Origin.Same(note.Id, note.AttributedTo)) if (note == default || !Origin.Same(note.Id, note.AttributedTo))
@@ -29,6 +29,7 @@ namespace PrivaPub.Federation.Objects
public async Task<JobOutcome> Handle(Job job, CancellationToken token) public async Task<JobOutcome> Handle(Job job, CancellationToken token)
{ {
var host = job.Payload; 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); 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); var href = links == default ? default : Schemas.Select(schema => LinkFor(links.Root, schema)).FirstOrDefault(link => link != default);
if (href == default) if (href == default)
+172 -16
View File
@@ -2,6 +2,8 @@ using Microsoft.Extensions.Caching.Memory;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Moderation;
using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Models.Statistics;
using System.Net; using System.Net;
using System.Text.Json; using System.Text.Json;
@@ -49,17 +51,71 @@ namespace PrivaPub.Infrastructure.Http
readonly IOptionsMonitor<FederationOptions> _options; readonly IOptionsMonitor<FederationOptions> _options;
readonly IDomainBlocks _domainBlocks; readonly IDomainBlocks _domainBlocks;
readonly ILogger<FederationHttp> _logger; readonly ILogger<FederationHttp> _logger;
readonly IInteractionLedger _ledger;
public FederationHttp(IHttpClientFactory httpClientFactory, IMemoryCache cache, IOptionsMonitor<FederationOptions> options, public FederationHttp(IHttpClientFactory httpClientFactory, IMemoryCache cache, IOptionsMonitor<FederationOptions> options,
IDomainBlocks domainBlocks, ILogger<FederationHttp> logger) IDomainBlocks domainBlocks, ILogger<FederationHttp> logger, IInteractionLedger ledger = default)
{ {
_httpClientFactory = httpClientFactory; _httpClientFactory = httpClientFactory;
_cache = cache; _cache = cache;
_options = options; _options = options;
_domainBlocks = domainBlocks; _domainBlocks = domainBlocks;
_logger = logger; _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) public bool IsAllowed(Uri target)
{ {
if (target is not { IsAbsoluteUri: true } || !string.IsNullOrEmpty(target.UserInfo)) if (target is not { IsAbsoluteUri: true } || !string.IsNullOrEmpty(target.UserInfo))
@@ -79,12 +135,31 @@ namespace PrivaPub.Infrastructure.Http
} }
public async Task<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> sign, CancellationToken token) public async Task<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> 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<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> sign, Exchange exchange, CancellationToken token)
{ {
if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target))
{
exchange.Refused("disallowed");
return default; return default;
}
var negativeKey = NegativeKey(target); var negativeKey = NegativeKey(target);
if (_cache.TryGetValue(negativeKey, out _)) if (_cache.TryGetValue(negativeKey, out _))
{
exchange.Refused("remembered");
return default; return default;
}
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token); using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token);
timeout.CancelAfter(RequestTimeout); timeout.CancelAfter(RequestTimeout);
@@ -98,46 +173,49 @@ namespace PrivaPub.Infrastructure.Http
using var response = await _httpClientFactory.CreateClient(ClientName) using var response = await _httpClientFactory.CreateClient(ClientName)
.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token); .SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token);
exchange.Answered(response);
if (IsRedirect(response.StatusCode)) if (IsRedirect(response.StatusCode))
{ {
var location = response.Headers.Location; var location = response.Headers.Location;
var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location); var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location);
if (!IsAllowed(next)) 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; target = next;
exchange.Hops++;
continue; continue;
} }
if (!response.IsSuccessStatusCode) 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); transient: (int)response.StatusCode is >= 500 or 429 or 408);
var mediaType = response.Content.Headers.ContentType?.MediaType; var mediaType = response.Content.Headers.ContentType?.MediaType;
if (mediaType == default || !JsonMediaTypes.Contains(mediaType, StringComparer.OrdinalIgnoreCase)) 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) 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); var body = await ReadBounded(response.Content, MaxResponseBytes, timeout.Token);
if (body == default) 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 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) 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) 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) 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; 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) 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)) if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target))
{
exchange.Refused("disallowed");
return default; return default;
}
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token); using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token);
timeout.CancelAfter(TimeSpan.FromSeconds(60)); timeout.CancelAfter(TimeSpan.FromSeconds(60));
try try
@@ -157,27 +254,52 @@ namespace PrivaPub.Infrastructure.Http
using var request = new HttpRequestMessage(HttpMethod.Get, target); using var request = new HttpRequestMessage(HttpMethod.Get, target);
request.Headers.Accept.ParseAdd("image/*, video/*, audio/*"); request.Headers.Accept.ParseAdd("image/*, video/*, audio/*");
using var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token); using var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token);
exchange.Answered(response);
if (IsRedirect(response.StatusCode)) if (IsRedirect(response.StatusCode))
{ {
var location = response.Headers.Location; var location = response.Headers.Location;
var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location); var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location);
if (!IsAllowed(next)) if (!IsAllowed(next))
{
exchange.Refused("bad-redirect");
return default; return default;
}
target = next; target = next;
exchange.Hops++;
continue; continue;
} }
var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant(); var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant();
if (!response.IsSuccessStatusCode || mediaType == default if (!response.IsSuccessStatusCode)
|| !(mediaType.StartsWith("image/") || mediaType.StartsWith("video/") || mediaType.StartsWith("audio/")) {
|| mediaType.Contains("svg") || response.Content.Headers.ContentLength > maxBytes) exchange.Failed(StatusReason(response));
return default; return default;
var bytes = await ReadBounded(response.Content, (int)Math.Min(maxBytes, int.MaxValue), timeout.Token);
return bytes == default ? default : (bytes, mediaType);
} }
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);
if (bytes == default)
{
exchange.Refused("too-large");
return default;
}
exchange.Bytes = bytes.Length;
return (bytes, mediaType);
}
exchange.Refused("too-many-redirects");
return default; return default;
} }
catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or OperationCanceledException && !token.IsCancellationRequested) 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); _logger.LogInformation("Media {Url} refused: {Reason}", url, ex.Message);
return default; return default;
} }
@@ -231,9 +353,29 @@ namespace PrivaPub.Infrastructure.Http
} }
public async Task<HttpResponseMessage> OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token) public async Task<HttpResponseMessage> 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<HttpResponseMessage> OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, Exchange exchange, CancellationToken token)
{ {
if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target))
{
exchange.Refused("disallowed");
return default; return default;
}
try try
{ {
for (var hop = 0; hop <= MaxRedirects; hop++) for (var hop = 0; hop <= MaxRedirects; hop++)
@@ -242,20 +384,29 @@ namespace PrivaPub.Infrastructure.Http
request.Headers.Accept.ParseAdd("video/*, audio/*, image/*"); request.Headers.Accept.ParseAdd("video/*, audio/*, image/*");
request.Headers.Range = range; request.Headers.Range = range;
var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, token); var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, token);
exchange.Answered(response);
if (IsRedirect(response.StatusCode)) if (IsRedirect(response.StatusCode))
{ {
var location = response.Headers.Location; var location = response.Headers.Location;
response.Dispose(); response.Dispose();
var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location); var next = location == default ? default : location.IsAbsoluteUri ? location : new Uri(target, location);
if (!IsAllowed(next)) if (!IsAllowed(next))
{
exchange.Refused("bad-redirect");
return default; return default;
}
target = next; target = next;
exchange.Hops++;
continue; continue;
} }
var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant(); var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant();
if (!response.IsSuccessStatusCode || mediaType == default || mediaType.Contains("svg") if (!response.IsSuccessStatusCode || mediaType == default || mediaType.Contains("svg")
|| !(mediaType.StartsWith("video/") || mediaType.StartsWith("audio/") || mediaType.StartsWith("image/") || mediaType == "application/octet-stream")) || !(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(); response.Dispose();
return default; return default;
} }
@@ -265,6 +416,7 @@ namespace PrivaPub.Infrastructure.Http
} }
catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or OperationCanceledException && !token.IsCancellationRequested) 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); _logger.LogInformation("Media stream {Url} refused: {Reason}", url, ex.Message);
return default; return default;
} }
@@ -311,8 +463,12 @@ namespace PrivaPub.Infrastructure.Http
static string NegativeKey(Uri target) => "federation-http:refused:" + target.AbsoluteUri; 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); _cache.Set(negativeKey, transient, NegativeCacheLifetime);
_logger.LogInformation("GET {Url} refused: {Reason}", url, reason); _logger.LogInformation("GET {Url} refused: {Reason}", url, reason);
return default; return default;
+39
View File
@@ -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<Tags> 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<Tags, Tags> 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;
}
}
}
@@ -1,3 +1,4 @@
using PrivaPub.Infrastructure.Http;
using PrivaPub.Models.Jobs; using PrivaPub.Models.Jobs;
using System.Collections.Concurrent; using System.Collections.Concurrent;
@@ -73,6 +74,7 @@ namespace PrivaPub.Infrastructure.Jobs
JobOutcome outcome; JobOutcome outcome;
try try
{ {
using var scope = HttpScope.Triggered(handler.Kind.ToString().ToLowerInvariant());
outcome = await handler.Handle(job, stoppingToken); outcome = await handler.Handle(job, stoppingToken);
} }
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
@@ -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<IInteractionLedger>();
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"
}
};
}
}
}
+2
View File
@@ -17,6 +17,7 @@ using PrivaPub.Infrastructure;
using PrivaPub.Infrastructure.Cli; using PrivaPub.Infrastructure.Cli;
using PrivaPub.Infrastructure.Data; using PrivaPub.Infrastructure.Data;
using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Http;
using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Middleware; using PrivaPub.Middleware;
using PrivaPub.Models; using PrivaPub.Models;
using PrivaPub.Services; using PrivaPub.Services;
@@ -150,6 +151,7 @@ try
app.UseRequestLocalization(await localizationService.Get()); app.UseRequestLocalization(await localizationService.Get());
app.UseRouting(); app.UseRouting();
app.UseTrafficMeter();
app.UseRateLimiter(); app.UseRateLimiter();
app.UseAuthentication(); app.UseAuthentication();