diff --git a/CLAUDE.md b/CLAUDE.md index 715132b..74eb8e0 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -336,10 +336,23 @@ group www-data and reaches the private mongod; `sudo -u www-data` works too. 6. **A client never contacts a remote server for media:** every remote URL the API returns goes through `IMediaProxy.Wrap`, an HMAC-signed `/media/proxy/` URL fetched by `IFederationHttp.GetMedia`. 7. **The proxy serves three ways:** - - **Cached:** a file already cached is served from disk, ranges included. - - **Downloaded:** a request without a `Range` is downloaded whole, up to `Media:MaxProxiedBytes`, then cached. + - **Cached:** a file already cached is served from disk, ranges included. It is opened before the answer, so a trim + that deletes it meanwhile breaks nothing. + - **Downloaded:** a request without a `Range` is downloaded whole, up to `Media:MaxProxiedBytes`, then cached: + - once for everyone asking for it at the same time; + - streamed into a `.part` file beside its place and renamed there, never held in memory; + - with its host in its `.type` sidecar, so that blocking a server can purge it. - **Streamed:** a ranged request, or anything too big to cache, is streamed from the origin with the range passed on, - and never cached. That is how remote video plays. + and never cached. That is how remote video plays. A file found too big is remembered for an hour, so it isn't + fetched twice; one that failed is remembered for five minutes. + + Also: + - The cache's size is counted as it grows and trimmed as soon as it passes `ProxyCacheBytes`; the janitor + recounts hourly. + - Browsers may cache only a success. + - Nothing of a suspended server, or of one whose media are rejected, is proxied (avatars, emoji, covers and link + cards included, since they all go through it), and blocking one purges its cache. + - The `proxy` rate limit counts per client address. nginx has a `/media/proxy/` location with `proxy_buffering off` and a 600 s read timeout for those streams. 8. **A focal point is two finite numbers** within -1..1 (`FocalPoint.Parse`); anything else is ignored. A stored NaN made diff --git a/PrivaPub.Tests/Domain/MediaFlowTests.cs b/PrivaPub.Tests/Domain/MediaFlowTests.cs index 29c0a0e..8ddc669 100644 --- a/PrivaPub.Tests/Domain/MediaFlowTests.cs +++ b/PrivaPub.Tests/Domain/MediaFlowTests.cs @@ -1,3 +1,4 @@ +using PrivaPub.Models.Federation; using Microsoft.AspNetCore.Http; using Microsoft.Extensions.Caching.Memory; @@ -98,5 +99,37 @@ namespace PrivaPub.Tests.Domain Assert.StartsWith(_harness.Media.ProxyRoot, file); Assert.Null(tampered); } + + // nothing of a suspended server, or of one whose media are rejected, is proxied; blocking one purges what was cached + [Fact] + public async Task A_blocked_servers_media_are_not_proxied_and_their_cache_goes() + { + var token = TestContext.Current.CancellationToken; + var path = $"/files/{Guid.NewGuid():N}.png"; + _harness.Peer.ServeFile(path, Png(10, 10), "image/png"); + var remote = _harness.Peer.A + path; + var host = new Uri(remote).Host; + var blocks = new Blocks(); + var proxy = new MediaProxy(_harness.Local, Peer.Http(), _harness.Media, new StaticOptions(new MediaOptions()), blocks); + Assert.Equal(ProxyOutcome.Cached, (await proxy.Download(remote, token)).Outcome); + + blocks.Blocked[host] = new DomainBlock { Domain = host, Severity = DomainBlockSeverity.Silence, RejectMedia = true }; + Assert.True(proxy.Refuses(remote)); + Assert.True(proxy.Purge(host) >= 1); + Assert.Equal(default, proxy.Cached(remote)); + + blocks.Blocked[host] = new DomainBlock { Domain = host, Severity = DomainBlockSeverity.Silence }; + Assert.False(proxy.Refuses(remote)); + blocks.Blocked[host] = new DomainBlock { Domain = host, Severity = DomainBlockSeverity.Suspend }; + Assert.True(proxy.Refuses(remote)); + } + + sealed class Blocks : PrivaPub.Federation.Moderation.IDomainBlocks + { + public Dictionary Blocked { get; } = new(); + public DomainBlock Find(string host) => Blocked.GetValueOrDefault(host); + public bool IsSuspended(string host) => Find(host)?.Severity == DomainBlockSeverity.Suspend; + public Task Reload(CancellationToken token) => Task.CompletedTask; + } } } diff --git a/PrivaPub.Tests/Http/MastodonMediaTests.cs b/PrivaPub.Tests/Http/MastodonMediaTests.cs index a5d2075..a6a45f1 100644 --- a/PrivaPub.Tests/Http/MastodonMediaTests.cs +++ b/PrivaPub.Tests/Http/MastodonMediaTests.cs @@ -415,6 +415,36 @@ namespace PrivaPub.Tests.Http Assert.Single(_peer.Requests); } + // clients asking at once for a file not cached yet share one download + [Fact] + public async Task Clients_asking_at_once_share_one_download() + { + var proxy = _host.Get(); + var bytes = Bytes(200_000); + var remote = Served(bytes, "image/png"); + var path = new Uri(remote).AbsolutePath; + using var client = _host.Client(); + + var answers = await Task.WhenAll(Enumerable.Range(0, 5).Select(_ => client.GetAsync(proxy.Wrap(remote), Token))); + + foreach (var answer in answers) + Assert.Equal(bytes, await answer.Content.ReadAsByteArrayAsync(Token)); + Assert.Single(_peer.Requests, r => r.Path == path); + } + + // a remote file that can't be had is a 404 no browser keeps + [Fact] + public async Task A_remote_file_that_fails_is_not_cached_by_browsers() + { + var proxy = _host.Get(); + using var client = _host.Client(); + + using var missing = await client.GetAsync(proxy.Wrap($"{_peer.A}/media/{Guid.NewGuid():N}.png"), Token); + + Assert.Equal(HttpStatusCode.NotFound, missing.StatusCode); + Assert.Null(missing.Headers.CacheControl); + } + [Fact] public async Task Anything_over_the_proxy_limit_is_streamed_and_never_cached() { @@ -432,7 +462,8 @@ namespace PrivaPub.Tests.Http Assert.Equal(default, proxy.Cached(remote)); using var again = await client.GetAsync(proxy.Wrap(remote), Token); Assert.Equal(bytes, await again.Content.ReadAsByteArrayAsync(Token)); - Assert.Equal(4, _peer.Requests.Count); + // the first time, an attempt to cache it and the stream; then it is known to be too big, and only streamed + Assert.Equal(3, _peer.Requests.Count); var small = Served(Bytes(SmallProxyHost.Limit / 2)); using (var fits = await client.GetAsync(proxy.Wrap(small), Token)) Assert.Equal(HttpStatusCode.OK, fits.StatusCode); diff --git a/PrivaPub.Tests/Support/Host/PrivaPubHost.cs b/PrivaPub.Tests/Support/Host/PrivaPubHost.cs index 975b434..a6e55f4 100644 --- a/PrivaPub.Tests/Support/Host/PrivaPubHost.cs +++ b/PrivaPub.Tests/Support/Host/PrivaPubHost.cs @@ -79,6 +79,7 @@ namespace PrivaPub.Tests.Support.Host ["Statistics:Cdn:AutoUpdate"] = "false", ["Media:Root"] = _mediaRoot, ["RateLimits:UploadsBurst"] = "1000", + ["RateLimits:ProxyBurst"] = "100000", ["Logging:LogLevel:Default"] = "Warning", ["Serilog:MinimumLevel:Default"] = Environment.GetEnvironmentVariable("PRIVAPUB_TEST_LOGS") == "1" ? "Information" : "Fatal" }; diff --git a/PrivaPub/Api/Mastodon/Controllers/MediaController.cs b/PrivaPub/Api/Mastodon/Controllers/MediaController.cs index 1f19579..f7108ef 100644 --- a/PrivaPub/Api/Mastodon/Controllers/MediaController.cs +++ b/PrivaPub/Api/Mastodon/Controllers/MediaController.cs @@ -78,35 +78,55 @@ namespace PrivaPub.Api.Mastodon.Controllers return Json(View(attachment)); } - [HttpGet("/media/proxy/{signature}/{encoded}"), AllowAnonymous, ApiExplorerSettings(IgnoreApi = true)] + [HttpGet("/media/proxy/{signature}/{encoded}"), AllowAnonymous, EnableRateLimiting(RateLimiting.Proxy), ApiExplorerSettings(IgnoreApi = true)] public async Task Proxy(string signature, string encoded, CancellationToken token) { var url = _proxy.Verified(signature, encoded); - if (url == default) + if (url == default || _proxy.Refuses(url)) return NotFound(); Response.Headers["X-Content-Type-Options"] = "nosniff"; 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) + if (_proxy.Cached(url) is { Path: not null } cached && Serve(cached.Path, cached.ContentType) is { } hit) { _ledger?.Count(Interactions.HostOf(url), "media:hit"); - return PhysicalFile(cached.Path, cached.ContentType, enableRangeProcessing: true); + return hit; } if (Request.Headers.Range.Count == 0) { - var (path, contentType) = await _proxy.Fetch(signature, encoded, token); - if (path != default) - return PhysicalFile(path, contentType, enableRangeProcessing: true); + var (outcome, path, contentType) = await _proxy.Download(url, token); + if (outcome == ProxyOutcome.Cached && Serve(path, contentType) is { } fetched) + return fetched; + if (outcome == ProxyOutcome.Failed) + return NotFound(); } return await Stream(url, token); } + // a cached file, opened before it is answered: trimmed meanwhile, what is open still reads; cached by browsers only + // once there is something to cache + IActionResult Serve(string path, string contentType) + { + FileStream file; + try + { + file = new FileStream(path, FileMode.Open, FileAccess.Read, FileShare.ReadWrite | FileShare.Delete, 64 * 1024, useAsync: true); + } + catch (IOException) + { + return default; + } + Response.Headers["Cache-Control"] = "public, max-age=604800"; + return File(file, contentType, enableRangeProcessing: true); + } + async Task Stream(string url, CancellationToken token) { var range = System.Net.Http.Headers.RangeHeaderValue.TryParse(Request.Headers.Range.ToString(), out var asked) ? asked : default; using var upstream = await _proxy.Open(url, range, token); if (upstream == default) return NotFound(); + if (upstream.IsSuccessStatusCode) + Response.Headers["Cache-Control"] = "public, max-age=604800"; Response.StatusCode = (int)upstream.StatusCode; Response.ContentType = upstream.Content.Headers.ContentType?.ToString() ?? "application/octet-stream"; if (upstream.Content.Headers.ContentLength is { } length) diff --git a/PrivaPub/Controllers/ClientToServer/DomainBlockController.cs b/PrivaPub/Controllers/ClientToServer/DomainBlockController.cs index 0acbf42..ac4b5ac 100644 --- a/PrivaPub/Controllers/ClientToServer/DomainBlockController.cs +++ b/PrivaPub/Controllers/ClientToServer/DomainBlockController.cs @@ -19,11 +19,13 @@ namespace PrivaPub.Controllers.ClientToServer { readonly IDomainBlocks _domainBlocks; readonly IStringLocalizer _localizer; + readonly Domain.Media.IMediaProxy _proxy; - public DomainBlockController(IDomainBlocks domainBlocks, IStringLocalizer localizer) + public DomainBlockController(IDomainBlocks domainBlocks, IStringLocalizer localizer, Domain.Media.IMediaProxy proxy) { _domainBlocks = domainBlocks; _localizer = localizer; + _proxy = proxy; } [HttpGet, Route("/clientapi/admin/domainblocks/list")] @@ -48,6 +50,9 @@ namespace PrivaPub.Controllers.ClientToServer .Option(o => o.IsUpsert = true) .ExecuteAsync(token); await _domainBlocks.Reload(token); + // what the media proxy holds of it goes with the block + if (block.Severity == DomainBlockSeverity.Suspend || block.RejectMedia) + _proxy.Purge(domain); return Ok(ToView(block)); } diff --git a/PrivaPub/Domain/Media/MediaProxy.cs b/PrivaPub/Domain/Media/MediaProxy.cs index 920afc1..7d26f19 100644 --- a/PrivaPub/Domain/Media/MediaProxy.cs +++ b/PrivaPub/Domain/Media/MediaProxy.cs @@ -1,3 +1,5 @@ +using PrivaPub.Models.Federation; +using System.Collections.Concurrent; using Microsoft.Extensions.Options; using MongoDB.Entities; @@ -11,32 +13,76 @@ using System.Text; namespace PrivaPub.Domain.Media { + public enum ProxyOutcome + { + Cached, + TooBig, + Failed + } + public interface IMediaProxy { string Wrap(string remoteUrl); - Task<(string Path, string ContentType)> Fetch(string signature, string encodedUrl, CancellationToken token); string Verified(string signature, string encodedUrl); + /// Whether a remote file's server is suspended or has its media rejected: nothing of it is proxied. + bool Refuses(string url); (string Path, string ContentType) Cached(string url); + /// Downloads a remote file into the cache, once however many ask at the same time. + Task<(ProxyOutcome Outcome, string Path, string ContentType)> Download(string url, CancellationToken token); + Task<(string Path, string ContentType)> Fetch(string signature, string encodedUrl, CancellationToken token); Task Open(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token); + /// Deletes what the cache holds of a domain and its subdomains; how many files. + int Purge(string domain); + /// Trims the cache to its size, oldest first. + void Trim(); } + // Remote media, fetched for clients so that they never contact a remote server (CLAUDE.md, media invariants): + // - URLs are HMAC-signed, so only what PrivaPub showed is fetched; + // - a download is shared by everyone asking for the same URL at once, streamed into a temporary file beside its place + // and renamed there (a reader never sees half a file), and never held in memory; + // - what is too big to cache is remembered for an hour (it is streamed instead), what failed for five minutes; + // - the cache's size is counted as it grows, and trimmed as soon as it passes the cap; + // - nothing of a suspended server, or of one whose media are rejected, is proxied, and blocking one purges its files. public class MediaProxy : IMediaProxy { + static readonly TimeSpan TooBigFor = TimeSpan.FromHours(1); + static readonly TimeSpan FailedFor = TimeSpan.FromMinutes(5); + const int Downloads = 8;//remote files fetched at once, however many clients ask + readonly ILocalActorService _localActors; readonly IFederationHttp _http; readonly IMediaService _media; readonly IOptionsMonitor _options; + readonly Federation.Moderation.IDomainBlocks _domainBlocks; + readonly ConcurrentDictionary> _inFlight = new(); + readonly ConcurrentDictionary _refused = new(); + readonly SemaphoreSlim _downloads = new(Downloads); + readonly object _keyLock = new(); byte[] _key; + long _bytes = -1;//what the cache holds, counted once from the disk and then as it changes + int _trimming; - public MediaProxy(ILocalActorService localActors, IFederationHttp http, IMediaService media, IOptionsMonitor options) + public MediaProxy(ILocalActorService localActors, IFederationHttp http, IMediaService media, IOptionsMonitor options, + Federation.Moderation.IDomainBlocks domainBlocks = default) { _localActors = localActors; _http = http; _media = media; _options = options; + _domainBlocks = domainBlocks; } - byte[] Key => _key ??= LoadKey(); + byte[] Key + { + get + { + if (_key != default) + return _key; + lock (_keyLock) + return _key ??= LoadKey(); + } + } public string Wrap(string remoteUrl) { @@ -60,18 +106,180 @@ namespace PrivaPub.Domain.Media return CryptographicOperations.FixedTimeEquals(Encoding.ASCII.GetBytes(signature ?? string.Empty), Encoding.ASCII.GetBytes(Sign(url))) ? url : default; } + public bool Refuses(string url) + { + if (_domainBlocks == default || !Uri.TryCreate(url, UriKind.Absolute, out var target)) + return false; + var block = _domainBlocks.Find(target.Host); + return block is { Severity: DomainBlockSeverity.Suspend } or { RejectMedia: true }; + } + public (string Path, string ContentType) Cached(string url) { var (path, typePath) = CachePaths(url); - if (!File.Exists(path) || !File.Exists(typePath)) - return default; - File.SetLastWriteTimeUtc(path, DateTime.UtcNow); - return (path, File.ReadAllText(typePath)); + try + { + if (!File.Exists(path) || !File.Exists(typePath)) + return default; + File.SetLastWriteTimeUtc(path, DateTime.UtcNow); + return (path, File.ReadLines(typePath).FirstOrDefault()); + } + catch (IOException) + { + return default;//trimmed meanwhile + } } public Task Open(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token) => _http.OpenMedia(url, range, token); + public async Task<(string Path, string ContentType)> Fetch(string signature, string encodedUrl, CancellationToken token) + { + var url = Verified(signature, encodedUrl); + if (url == default || Refuses(url)) + return default; + var (outcome, path, contentType) = await Download(url, token); + return outcome == ProxyOutcome.Cached ? (path, contentType) : default; + } + + public async Task<(ProxyOutcome Outcome, string Path, string ContentType)> Download(string url, CancellationToken token) + { + if (Cached(url) is { Path: not null } cached) + return (ProxyOutcome.Cached, cached.Path, cached.ContentType); + if (_refused.TryGetValue(url, out var refused)) + { + if (refused.Until > DateTime.UtcNow) + return (refused.Outcome, default, default); + _refused.TryRemove(url, out _); + } + // shared by everyone asking for it now; the download itself isn't cancelled when one of them leaves + var download = _inFlight.GetOrAdd(url, key => DownloadOnce(key)); + try + { + return await download.WaitAsync(token); + } + finally + { + if (download.IsCompleted) + _inFlight.TryRemove(new KeyValuePair>(url, download)); + } + } + + async Task<(ProxyOutcome, string, string)> DownloadOnce(string url) + { + await Task.Yield(); + if (!await _downloads.WaitAsync(TimeSpan.FromSeconds(30))) + return (ProxyOutcome.Failed, default, default);//too busy: the client tries again + var (path, typePath) = CachePaths(url); + var part = $"{path}.{Guid.NewGuid():N}.part"; + try + { + Directory.CreateDirectory(System.IO.Path.GetDirectoryName(path)!); + string contentType, refusal; + await using (var file = new FileStream(part, FileMode.CreateNew, FileAccess.Write, FileShare.None, 64 * 1024, useAsync: true)) + (contentType, refusal) = await _http.DownloadMedia(url, _options.CurrentValue.MaxProxiedBytes, file, CancellationToken.None); + if (contentType == default) + { + var outcome = refusal == "too-large" ? ProxyOutcome.TooBig : ProxyOutcome.Failed; + _refused[url] = (outcome, DateTime.UtcNow + (outcome == ProxyOutcome.TooBig ? TooBigFor : FailedFor)); + return (outcome, default, default); + } + await File.WriteAllTextAsync(typePath, contentType + "\n" + new Uri(url).Host); + var size = new FileInfo(part).Length; + File.Move(part, path, overwrite: true); + Grew(size); + return (ProxyOutcome.Cached, path, contentType); + } + finally + { + _downloads.Release(); + if (File.Exists(part)) + File.Delete(part); + } + } + + void Grew(long bytes) + { + if (Interlocked.Read(ref _bytes) < 0) + Interlocked.CompareExchange(ref _bytes, Measure(), -1); + if (Interlocked.Add(ref _bytes, bytes) > _options.CurrentValue.ProxyCacheBytes) + _ = Task.Run(Trim); + } + + public void Trim() + { + if (Interlocked.Exchange(ref _trimming, 1) == 1) + return; + try + { + Interlocked.Exchange(ref _bytes, TrimDirectory(_media.ProxyRoot, _options.CurrentValue.ProxyCacheBytes)); + } + finally + { + Interlocked.Exchange(ref _trimming, 0); + } + } + + public int Purge(string domain) + { + var directory = new DirectoryInfo(_media.ProxyRoot); + if (string.IsNullOrEmpty(domain) || !directory.Exists) + return 0; + var purged = 0; + foreach (var type in directory.EnumerateFiles("*.type", SearchOption.AllDirectories)) + { + try + { + var host = File.ReadLines(type.FullName).Skip(1).FirstOrDefault(); + if (host == default || !(host.Equals(domain, StringComparison.OrdinalIgnoreCase) || host.EndsWith("." + domain, StringComparison.OrdinalIgnoreCase))) + continue; + var file = new FileInfo(type.FullName[..^".type".Length]); + if (file.Exists) + file.Delete(); + type.Delete(); + purged++; + } + catch (IOException) + { + } + } + Interlocked.Exchange(ref _bytes, -1); + return purged; + } + + long Measure() + { + var directory = new DirectoryInfo(_media.ProxyRoot); + return directory.Exists ? directory.EnumerateFiles("*", SearchOption.AllDirectories).Where(f => f.Extension is not (".type" or ".part")).Sum(f => f.Length) : 0; + } + + // deletes the oldest files until the cache is within its cap; what it holds afterwards + public static long TrimDirectory(string root, long cap) + { + var directory = new DirectoryInfo(root); + if (!directory.Exists) + return 0; + var files = directory.EnumerateFiles("*", SearchOption.AllDirectories).Where(f => f.Extension is not (".type" or ".part")).OrderBy(f => f.LastWriteTimeUtc).ToList(); + var total = files.Sum(f => f.Length); + foreach (var file in files) + { + if (total <= cap) + break; + try + { + total -= file.Length; + file.Delete(); + var type = new FileInfo(file.FullName + ".type"); + if (type.Exists) + type.Delete(); + } + catch (IOException) + { + } + } + return total; + } + (string Path, string TypePath) CachePaths(string url) { var name = Convert.ToHexStringLower(SHA256.HashData(Encoding.UTF8.GetBytes(url))); @@ -79,35 +287,16 @@ namespace PrivaPub.Domain.Media return (path, path + ".type"); } - public async Task<(string Path, string ContentType)> Fetch(string signature, string encodedUrl, CancellationToken token) - { - var url = Verified(signature, encodedUrl); - if (url == default) - return default; - if (Cached(url) is { Path: not null } cached) - return cached; - var (path, typePath) = CachePaths(url); - var directory = System.IO.Path.GetDirectoryName(path); - - var (bytes, contentType) = await _http.GetMedia(url, _options.CurrentValue.MaxProxiedBytes, token); - if (bytes == default) - return default; - Directory.CreateDirectory(directory); - await File.WriteAllBytesAsync(path, bytes, token); - await File.WriteAllTextAsync(typePath, contentType, token); - return (path, contentType); - } - string Sign(string url) => Base64Url(HMACSHA256.HashData(Key, Encoding.UTF8.GetBytes(url))[..16]); + // the oldest key, so that two first uses at once agree on one static byte[] LoadKey() { - var secret = DB.Default.Find().ExecuteFirstAsync().GetAwaiter().GetResult(); + var secret = DB.Default.Find().Sort(m => m.ID, Order.Ascending).ExecuteFirstAsync().GetAwaiter().GetResult(); if (secret == default) { - secret = new MediaSecret { Key = Convert.ToBase64String(RandomNumberGenerator.GetBytes(32)) }; - DB.Default.SaveAsync(secret).GetAwaiter().GetResult(); - secret = DB.Default.Find().ExecuteFirstAsync().GetAwaiter().GetResult(); + DB.Default.SaveAsync(new MediaSecret { Key = Convert.ToBase64String(RandomNumberGenerator.GetBytes(32)) }).GetAwaiter().GetResult(); + secret = DB.Default.Find().Sort(m => m.ID, Order.Ascending).ExecuteFirstAsync().GetAwaiter().GetResult(); } return Convert.FromBase64String(secret.Key); } @@ -192,23 +381,6 @@ namespace PrivaPub.Domain.Media TrimProxyCache(); } - void TrimProxyCache() - { - var directory = new DirectoryInfo(_media.ProxyRoot); - if (!directory.Exists) - return; - var files = directory.EnumerateFiles("*", SearchOption.AllDirectories).Where(f => f.Extension != ".type").OrderBy(f => f.LastWriteTimeUtc).ToList(); - var total = files.Sum(f => f.Length); - foreach (var file in files) - { - if (total <= _options.CurrentValue.ProxyCacheBytes) - break; - total -= file.Length; - file.Delete(); - var type = new FileInfo(file.FullName + ".type"); - if (type.Exists) - type.Delete(); - } - } + void TrimProxyCache() => MediaProxy.TrimDirectory(_media.ProxyRoot, _options.CurrentValue.ProxyCacheBytes); } } diff --git a/PrivaPub/Infrastructure/Http/FederationHttp.cs b/PrivaPub/Infrastructure/Http/FederationHttp.cs index ddf8e57..da6cb1b 100644 --- a/PrivaPub/Infrastructure/Http/FederationHttp.cs +++ b/PrivaPub/Infrastructure/Http/FederationHttp.cs @@ -29,6 +29,9 @@ namespace PrivaPub.Infrastructure.Http Task OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token); Task Send(HttpRequestMessage request, CancellationToken token); Task<(byte[] Bytes, string ContentType)> GetMedia(string url, long maxBytes, CancellationToken token); + /// Copies a media file into , up to maxBytes: its type, or why it was refused + /// ("too-large", "content-type", ...). + 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> GetStringArray(string url, int maxItems, CancellationToken token); } @@ -389,14 +392,29 @@ namespace PrivaPub.Infrastructure.Http 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> consume, 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) + var contentType = await ReadMedia(url, maxBytes, consume, exchange, token); + if (contentType == default && exchange.Outcome == Interactions.Ok) exchange.Refused("unusable"); - return media; + return (contentType, contentType == default ? exchange.Reason ?? "unusable" : default); } finally { @@ -404,7 +422,7 @@ namespace PrivaPub.Infrastructure.Http } } - async Task<(byte[] Bytes, string ContentType)> GetMedia(string url, long maxBytes, Exchange exchange, CancellationToken token) + async Task ReadMedia(string url, long maxBytes, Func> consume, Exchange exchange, CancellationToken token) { if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) { @@ -451,14 +469,14 @@ namespace PrivaPub.Infrastructure.Http exchange.Refused("too-large"); return default; } - var bytes = await ReadBounded(response.Content, (int)Math.Min(maxBytes, int.MaxValue), timeout.Token); - if (bytes == default) + var taken = await consume(response.Content, maxBytes, timeout.Token); + if (taken < 0) { exchange.Refused("too-large"); return default; } - exchange.Bytes = bytes.Length; - return (bytes, mediaType); + exchange.Bytes = taken; + return mediaType; } exchange.Refused("too-many-redirects"); return default; @@ -608,6 +626,23 @@ namespace PrivaPub.Infrastructure.Http 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 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 ReadBounded(HttpContent content, int limit, CancellationToken token) { await using var stream = await content.ReadAsStreamAsync(token); diff --git a/PrivaPub/Infrastructure/RateLimiting.cs b/PrivaPub/Infrastructure/RateLimiting.cs index ab731e3..99a3b58 100644 --- a/PrivaPub/Infrastructure/RateLimiting.cs +++ b/PrivaPub/Infrastructure/RateLimiting.cs @@ -17,6 +17,8 @@ namespace PrivaPub.Infrastructure public int InboxPerTenSeconds { get; set; } = 50;//and the rate it earns them back public int UploadsBurst { get; set; } = 30;//uploads (media, profile pictures) a session may make at once public int UploadsPerMinute { get; set; } = 10;//and the rate it earns them back + public int ProxyBurst { get; set; } = 1200;//remote media a client address may ask for at once (a timeline is many pictures) + public int ProxyPerMinute { get; set; } = 600;//and the rate it earns them back } public static class RateLimiting @@ -24,6 +26,7 @@ namespace PrivaPub.Infrastructure public const string Accounts = "accounts"; public const string Inbox = "inbox"; public const string Uploads = "uploads"; + public const string Proxy = "proxy"; static RateLimitOptions Limits(HttpContext context) => context.RequestServices.GetRequiredService>().Value; @@ -62,6 +65,15 @@ namespace PrivaPub.Infrastructure ReplenishmentPeriod = TimeSpan.FromSeconds(10), QueueLimit = 0 })); + options.AddPolicy(Proxy, context => RateLimitPartition.GetTokenBucketLimiter( + context.Connection.RemoteIpAddress?.ToString() ?? "unknown", + _ => new TokenBucketRateLimiterOptions + { + TokenLimit = Limits(context).ProxyBurst, + TokensPerPeriod = Limits(context).ProxyPerMinute, + ReplenishmentPeriod = TimeSpan.FromMinutes(1), + QueueLimit = 0 + })); // per session: the limiter runs before authentication, so the credential sent stands for whoever sends it options.AddPolicy(Uploads, context => RateLimitPartition.GetTokenBucketLimiter( Credential(context.Request) ?? "anonymous:" + context.Connection.RemoteIpAddress, diff --git a/PrivaPub/Services/GroupUsersService.cs b/PrivaPub/Services/GroupUsersService.cs index 1d5d49c..9028423 100644 --- a/PrivaPub/Services/GroupUsersService.cs +++ b/PrivaPub/Services/GroupUsersService.cs @@ -37,6 +37,7 @@ namespace PrivaPub.Services readonly IPasswordHasher _passwordHasher; readonly ILocalActorService _localActors; readonly IDeliveryService _delivery; + readonly Domain.Media.IMediaProxy _proxy; readonly IStringLocalizer _localizer; readonly ILogger _logger; @@ -45,8 +46,10 @@ namespace PrivaPub.Services ILocalActorService localActors, IDeliveryService delivery, IStringLocalizer localizer, - ILogger logger) + ILogger logger, + Domain.Media.IMediaProxy proxy) { + _proxy = proxy; _dbEntities = dbEntities; _passwordHasher = passwordHasher; _localActors = localActors; @@ -387,12 +390,13 @@ namespace PrivaPub.Services } } - static ViewGroupMember Remote(string actorUri, Models.User.ForeignAvatar foreign, GroupRole role, DateTime since, bool pending) => new() + // a remote member's picture through the media proxy, as everywhere else a client sees remote media + ViewGroupMember Remote(string actorUri, Models.User.ForeignAvatar foreign, GroupRole role, DateTime since, bool pending) => new() { ActorUri = actorUri, Handle = foreign == default ? actorUri : $"{foreign.UserName}@{foreign.Domain}", Name = foreign?.Name ?? foreign?.UserName, - PictureUrl = foreign?.PictureURL, + PictureUrl = _proxy.Wrap(foreign?.PictureURL), IsPending = pending, Role = role.ToString().ToLowerInvariant(), Since = since diff --git a/tools/pasture/appsettings.Pasture.json b/tools/pasture/appsettings.Pasture.json index b1aade4..eaf8857 100644 --- a/tools/pasture/appsettings.Pasture.json +++ b/tools/pasture/appsettings.Pasture.json @@ -26,7 +26,7 @@ "Relays": [ "https://relay.test/actor", "https://aoderelay.test/actor" ] }, "Media": { "Root": "/tmp/privapub-media" }, - "RateLimits": { "AccountsPerMinute": 1000, "UploadsBurst": 1000, "UploadsPerMinute": 1000 }, + "RateLimits": { "AccountsPerMinute": 1000, "UploadsBurst": 1000, "UploadsPerMinute": 1000, "ProxyBurst": 100000, "ProxyPerMinute": 100000 }, "Registrations": { "Mode": "Open" }, "Statistics": { "Geo": { "AutoUpdate": false } }, "Kestrel": { "Endpoints": { "Http": { "Url": "http://0.0.0.0:80", "Protocols": "Http1AndHttp2" } } },