Files
SocialPub/PrivaPub/Infrastructure/Http/FederationHttp.cs
T
thepraandClaude Opus 5.5 3ca29603ed The media proxy is bounded
Anyone could mint signed proxy URLs (a remote account changes its icon, an anonymous lookup returns the URL), and each
anonymous request held up to 40 MB in memory; the cache grew without bound between hourly trims; two clients asking
for the same new file downloaded it twice and wrote over each other in place, so a reader could get half a file with a
7-day cache header; a file over the limit was downloaded twice on every request; a failed fetch, a 404, was cached by
browsers for a week; cached media of a server suspended later were still served, and RejectMedia skipped avatars,
emoji, covers, video variants, link cards and remote edits; /clientapi/group/members returned remote pictures raw.

Now a download is shared by everyone asking at once, streamed into a .part file and renamed into place
(FederationHttp.DownloadMedia copies bounded, never into memory), at most eight at a time; a file too big to cache is
remembered for an hour and only streamed, a failure for five minutes; the cache's size is counted as it grows and
trimmed as soon as it passes the cap; a cached file is opened before it is answered; browsers may cache only a
success; nothing of a suspended server, or of one whose media are rejected, is proxied (everything remote a client
sees goes through the proxy, so that covers every kind), and blocking one purges its cache; the proxy has its own rate
limit per client address; group members' pictures are proxied; the proxy's key is loaded once, the oldest if two
were made. This changes what PrivaPub serves its clients, not what it sends to other servers.

Tests: clients asking at once share one download, a failure isn't cached by browsers, an over-limit file is fetched
three times for two requests instead of four, a blocked server's media are refused and its cache purged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
2026-10-07 10:54:33 +02:00

686 lines
25 KiB
C#

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);
Task<(int Status, FetchedJson Json)> GetJsonStatus(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);
/// <summary>Copies a media file into <paramref name="destination"/>, up to maxBytes: its type, or why it was refused
/// ("too-large", "content-type", ...).</summary>
Task<(string ContentType, string Refusal)> DownloadMedia(string url, long maxBytes, Stream destination, 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 = MediaType(response.Content.Headers);
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);
}
}
// the document and the status its origin answered with (a refusal remembered from the last five minutes keeps its
// status; 0 when it never answered)
public async Task<(int Status, FetchedJson Json)> GetJsonStatus(string url, string accept, Action<HttpRequestMessage> sign, CancellationToken token)
{
if (Uri.TryCreate(url, UriKind.Absolute, out var target) && _cache.TryGetValue(NegativeKey(target), out Refusal remembered))
return (remembered.Status, default);
var exchange = new Exchange(url, HttpScope.Purpose ?? "object");
try
{
var json = await GetJson(url, accept, sign, exchange, token);
return (exchange.Status ?? 0, json);
}
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 = MediaType(response.Content.Headers);
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 Refusal refusal) && refusal.Transient;
public async Task<(byte[] Bytes, string ContentType)> GetMedia(string url, long maxBytes, CancellationToken token)
{
var bytes = default(byte[]);
var (contentType, _) = await ReadMedia(url, maxBytes, async (content, limit, t) =>
{
bytes = await ReadBounded(content, (int)Math.Min(limit, int.MaxValue), t);
return bytes?.Length ?? -1;
}, token);
return contentType == default ? default : (bytes, contentType);
}
public Task<(string ContentType, string Refusal)> DownloadMedia(string url, long maxBytes, Stream destination, CancellationToken token) =>
ReadMedia(url, maxBytes, (content, limit, t) => CopyBounded(content, destination, limit, t), token);
// a media file read through consume (which answers how many bytes it took, or -1 past the limit): its type, or why not
async Task<(string ContentType, string Refusal)> ReadMedia(string url, long maxBytes, Func<HttpContent, long, CancellationToken, Task<long>> consume, CancellationToken token)
{
var exchange = new Exchange(url, "media");
try
{
var contentType = await ReadMedia(url, maxBytes, consume, exchange, token);
if (contentType == default && exchange.Outcome == Interactions.Ok)
exchange.Refused("unusable");
return (contentType, contentType == default ? exchange.Reason ?? "unusable" : default);
}
finally
{
Record(exchange);
}
}
async Task<string> ReadMedia(string url, long maxBytes, Func<HttpContent, long, CancellationToken, Task<long>> consume, 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 taken = await consume(response.Content, maxBytes, timeout.Token);
if (taken < 0)
{
exchange.Refused("too-large");
return default;
}
exchange.Bytes = taken;
return 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);
}
// copies at most limit bytes: how many, or -1 once there are more
static async Task<long> CopyBounded(HttpContent content, Stream destination, long limit, CancellationToken token)
{
await using var stream = await content.ReadAsStreamAsync(token);
var chunk = new byte[64 * 1024];
var copied = 0L;
int read;
while ((read = await stream.ReadAsync(chunk, token)) > 0)
{
copied += read;
if (copied > limit)
return -1;
await destination.WriteAsync(chunk.AsMemory(0, read), token);
}
return copied;
}
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;
// the media type, also from a header .NET cannot parse: Mobilizon's NodeInfo leaves its profile URL unquoted
static string MediaType(System.Net.Http.Headers.HttpContentHeaders headers) =>
headers.ContentType?.MediaType
?? (headers.NonValidated.TryGetValues("Content-Type", out var raw) ? raw.ToString().Split(';')[0].Trim() : default);
sealed record Refusal(bool Transient, int Status);
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, new Refusal(transient, exchange.Status ?? 0), NegativeCacheLifetime);
_logger.LogInformation("GET {Url} refused: {Reason}", url, reason);
return default;
}
}
}