Files
SocialPub/PrivaPub/Infrastructure/Http/FederationHttp.cs
T

624 lines
22 KiB
C#
Raw Normal View History

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;
namespace PrivaPub.Infrastructure.Http
{
public sealed class FetchedJson : IDisposable
{
public Uri FinalUri { get; init; }
public JsonDocument Document { get; init; }
public JsonElement Root => Document.RootElement;
public void Dispose() => Document?.Dispose();
}
public interface IFederationHttp
{
bool IsAllowed(Uri target);
Task<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> sign, CancellationToken token);
bool FailedTemporarily(string url);
Task<(Uri FinalUri, string Html)> GetPage(string url, CancellationToken token);
Task<HttpResponseMessage> OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token);
Task<HttpResponseMessage> Send(HttpRequestMessage request, CancellationToken token);
Task<(byte[] Bytes, string ContentType)> GetMedia(string url, long maxBytes, CancellationToken token);
Task<(int Status, string Text)> GetText(string url, int maxBytes, CancellationToken token);
Task<IReadOnlyList<string>> GetStringArray(string url, int maxItems, CancellationToken token);
}
public class FederationHttp : IFederationHttp
{
public const string ClientName = "Federation";
public const int MaxResponseBytes = 1024 * 1024;
public static readonly TimeSpan RequestTimeout = TimeSpan.FromSeconds(15);
const int MaxRedirects = 3;
static readonly TimeSpan NegativeCacheLifetime = TimeSpan.FromMinutes(5);
static readonly string[] JsonMediaTypes =
{
"application/activity+json",
"application/ld+json",
"application/jrd+json",
"application/json"
};
readonly IHttpClientFactory _httpClientFactory;
readonly IMemoryCache _cache;
readonly IOptionsMonitor<FederationOptions> _options;
readonly IDomainBlocks _domainBlocks;
readonly ILogger<FederationHttp> _logger;
readonly IInteractionLedger _ledger;
public FederationHttp(IHttpClientFactory httpClientFactory, IMemoryCache cache, IOptionsMonitor<FederationOptions> options,
IDomainBlocks domainBlocks, ILogger<FederationHttp> 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 HttpRequestMessage Request(Uri target, string accept)
{
var request = new HttpRequestMessage(HttpMethod.Get, target);
request.Headers.Accept.ParseAdd(accept);
if (HttpScope.UserAgent is { } userAgent)
request.Headers.UserAgent.ParseAdd(userAgent);
return request;
}
//robots.txt: the status alone matters when it is not 2xx (RFC 9309), so it is returned whatever it is; 0 when unreachable
public async Task<(int Status, string Text)> GetText(string url, int maxBytes, CancellationToken token)
{
var exchange = new Exchange(url, HttpScope.Purpose ?? "text");
try
{
if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target))
{
exchange.Refused("disallowed");
return default;
}
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token);
timeout.CancelAfter(RequestTimeout);
for (var hop = 0; hop <= MaxRedirects; hop++)
{
using var request = Request(target, "text/plain");
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 ((int)response.StatusCode, default);
}
target = next;
exchange.Hops++;
continue;
}
if (!response.IsSuccessStatusCode)
{
exchange.Refused(StatusReason(response));
return ((int)response.StatusCode, default);
}
var bytes = await ReadPrefix(response.Content, maxBytes, timeout.Token);
exchange.Bytes = bytes.Length;
return ((int)response.StatusCode, System.Text.Encoding.UTF8.GetString(bytes));
}
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" });
return default;
}
finally
{
Record(exchange);
}
}
//a JSON array of strings such as /api/v1/instance/peers, read from at most the first MaxResponseBytes: a large server's
//list is cut short rather than refused
public async Task<IReadOnlyList<string>> GetStringArray(string url, int maxItems, CancellationToken token)
{
var exchange = new Exchange(url, HttpScope.Purpose ?? "list");
var items = new List<string>();
try
{
if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target))
{
exchange.Refused("disallowed");
return items;
}
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token);
timeout.CancelAfter(RequestTimeout);
for (var hop = 0; hop <= MaxRedirects; hop++)
{
using var request = Request(target, "application/json");
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 items;
}
target = next;
exchange.Hops++;
continue;
}
if (!response.IsSuccessStatusCode)
{
exchange.Refused(StatusReason(response));
return items;
}
var mediaType = response.Content.Headers.ContentType?.MediaType;
if (mediaType == default || !JsonMediaTypes.Contains(mediaType, StringComparer.OrdinalIgnoreCase))
{
exchange.Refused("content-type");
return items;
}
var bytes = await ReadPrefix(response.Content, MaxResponseBytes, timeout.Token);
exchange.Bytes = bytes.Length;
ReadStrings(bytes, maxItems, items);
return items;
}
exchange.Refused("too-many-redirects");
return items;
}
catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or OperationCanceledException && !token.IsCancellationRequested)
{
exchange.Failed(ex switch { BlockedDestinationException => "private-address", OperationCanceledException => "timeout", _ => "network" });
return items;
}
finally
{
Record(exchange);
}
}
public static void ReadStrings(byte[] json, int maxItems, List<string> items)
{
var reader = new Utf8JsonReader(json, isFinalBlock: false, state: default);
try
{
if (!reader.Read() || reader.TokenType != JsonTokenType.StartArray)
return;
while (items.Count < maxItems && reader.Read() && reader.TokenType != JsonTokenType.EndArray)
{
if (reader.TokenType == JsonTokenType.String)
items.Add(reader.GetString());
else if (!reader.TrySkip())
return;
}
}
catch (JsonException)
{
}
}
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))
return false;
if (_domainBlocks?.IsSuspended(target.Host) == true)
return false;
var options = _options.CurrentValue;
if (target.Scheme != Uri.UriSchemeHttps && !(options.AllowPlainHttp && target.Scheme == Uri.UriSchemeHttp))
return false;
if (options.AllowPrivateNetworks)
return true;
return target.HostNameType == UriHostNameType.Dns
&& target.Host.Contains('.')
&& !target.Host.EndsWith(".localhost", StringComparison.OrdinalIgnoreCase)
&& !target.Host.EndsWith(".local", StringComparison.OrdinalIgnoreCase)
&& !target.Host.EndsWith(".internal", StringComparison.OrdinalIgnoreCase);
}
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))
{
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);
try
{
for (var hop = 0; hop <= MaxRedirects; hop++)
{
using var request = Request(target, accept);
sign?.Invoke(request);
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", exchange, "bad-redirect");
target = next;
exchange.Hops++;
continue;
}
if (!response.IsSuccessStatusCode)
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}'", exchange, "content-type");
if (response.Content.Headers.ContentLength > MaxResponseBytes)
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", exchange, "too-large");
exchange.Bytes = body.Length;
return new FetchedJson { FinalUri = target, Document = JsonDocument.Parse(body) };
}
return Refuse(negativeKey, url, "too many redirects", exchange, "too-many-redirects");
}
catch (OperationCanceledException) when (!token.IsCancellationRequested)
{
return Refuse(negativeKey, url, "a timeout", exchange, "timeout", transient: true);
}
catch (HttpRequestException ex)
{
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, exchange, ex is JsonException ? "bad-json" : "private-address");
}
}
public bool FailedTemporarily(string url) =>
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
{
for (var hop = 0; hop <= MaxRedirects; hop++)
{
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)
{
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);
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;
}
}
public async Task<(Uri FinalUri, string Html)> GetPage(string url, CancellationToken token)
{
if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target))
return default;
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(token);
timeout.CancelAfter(RequestTimeout);
try
{
for (var hop = 0; hop <= MaxRedirects; hop++)
{
using var request = new HttpRequestMessage(HttpMethod.Get, target);
request.Headers.Accept.ParseAdd("text/html, application/xhtml+xml");
using var response = await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token);
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 default;
target = next;
continue;
}
var mediaType = response.Content.Headers.ContentType?.MediaType?.ToLowerInvariant();
if (!response.IsSuccessStatusCode || mediaType is not ("text/html" or "application/xhtml+xml"))
return default;
var bytes = await ReadPrefix(response.Content, MaxPageBytes, timeout.Token);
var charset = response.Content.Headers.ContentType?.CharSet?.Trim('"');
var encoding = System.Text.Encoding.UTF8;
try
{
if (!string.IsNullOrEmpty(charset))
encoding = System.Text.Encoding.GetEncoding(charset);
}
catch (ArgumentException)
{
}
return (target, encoding.GetString(bytes));
}
return default;
}
catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or OperationCanceledException && !token.IsCancellationRequested)
{
_logger.LogInformation("Page {Url} refused: {Reason}", url, ex.Message);
return default;
}
}
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))
{
exchange.Refused("disallowed");
return default;
}
try
{
for (var hop = 0; hop <= MaxRedirects; hop++)
{
using var request = new HttpRequestMessage(HttpMethod.Get, target);
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;
}
return response;
}
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 stream {Url} refused: {Reason}", url, ex.Message);
return default;
}
}
const int MaxPageBytes = 512 * 1024;
static async Task<byte[]> ReadPrefix(HttpContent content, int limit, CancellationToken token)
{
await using var stream = await content.ReadAsStreamAsync(token);
var buffer = new byte[limit];
var total = 0;
int read;
while (total < limit && (read = await stream.ReadAsync(buffer.AsMemory(total, limit - total), token)) > 0)
total += read;
return buffer[..total];
}
public async Task<HttpResponseMessage> Send(HttpRequestMessage request, CancellationToken token)
{
if (!IsAllowed(request.RequestUri))
throw new BlockedDestinationException(request.RequestUri?.Host);
return await _httpClientFactory.CreateClient(ClientName).SendAsync(request, HttpCompletionOption.ResponseHeadersRead, token);
}
public static async Task<byte[]> ReadBounded(HttpContent content, int limit, CancellationToken token)
{
await using var stream = await content.ReadAsStreamAsync(token);
using var buffer = new MemoryStream();
var chunk = new byte[16 * 1024];
int read;
while ((read = await stream.ReadAsync(chunk, token)) > 0)
{
if (buffer.Length + read > limit)
return default;
buffer.Write(chunk, 0, read);
}
return buffer.ToArray();
}
static bool IsRedirect(HttpStatusCode status) =>
status is HttpStatusCode.MovedPermanently or HttpStatusCode.Found or HttpStatusCode.SeeOther
or HttpStatusCode.TemporaryRedirect or HttpStatusCode.PermanentRedirect;
static string NegativeKey(Uri target) => "federation-http:refused:" + target.AbsoluteUri;
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;
}
}
}