diff --git a/PrivaPub.Tests/Federation/InboxScenarioTests.cs b/PrivaPub.Tests/Federation/InboxScenarioTests.cs index 1dc5d19..40779ed 100644 --- a/PrivaPub.Tests/Federation/InboxScenarioTests.cs +++ b/PrivaPub.Tests/Federation/InboxScenarioTests.cs @@ -6,6 +6,7 @@ using MongoDB.Entities; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Inbox; using PrivaPub.Federation.Outbox; +using PrivaPub.Infrastructure.Jobs; using PrivaPub.Models; using PrivaPub.Models.Group; using PrivaPub.Models.Post; @@ -36,7 +37,7 @@ namespace PrivaPub.Tests.Federation var cache = new MemoryCache(new MemoryCacheOptions()); _local = new LocalActorService(new DbEntities(), new StaticOptions(new AppConfiguration { BackendBaseAddress = Base })); var remote = new RemoteActorService(Peer.Http(cache), _local, cache, new DbEntities()); - _inbox = new InboxService(new DbEntities(), _local, remote, new DeliveryService(new DbEntities()), NullLogger.Instance); + _inbox = new InboxService(new DbEntities(), _local, remote, new DeliveryService(new DbEntities(), new JobQueue()), NullLogger.Instance); } public async ValueTask DisposeAsync() diff --git a/PrivaPub.Tests/Infrastructure/IndexTests.cs b/PrivaPub.Tests/Infrastructure/IndexTests.cs index 072216e..4ed1741 100644 --- a/PrivaPub.Tests/Infrastructure/IndexTests.cs +++ b/PrivaPub.Tests/Infrastructure/IndexTests.cs @@ -20,6 +20,8 @@ namespace PrivaPub.Tests.Infrastructure { Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); var token = TestContext.Current.CancellationToken; + await DB.Default.Index().DropAllAsync(token); + await DB.Default.Index().DropAllAsync(token); var actorUri = $"https://r.example/users/{Guid.NewGuid():N}"; var localId = Guid.NewGuid().ToString("N")[..24]; await DB.Default.SaveAsync(new[] @@ -48,7 +50,6 @@ namespace PrivaPub.Tests.Infrastructure { Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); var token = TestContext.Current.CancellationToken; - await Indexes.Create(token); var local = new LocalActorService(new DbEntities(), new StaticOptions(new AppConfiguration { BackendBaseAddress = "https://privapub.test" })); var name = $"n{Guid.NewGuid():N}"[..20]; diff --git a/PrivaPub.Tests/Infrastructure/JobQueueTests.cs b/PrivaPub.Tests/Infrastructure/JobQueueTests.cs new file mode 100644 index 0000000..87f40db --- /dev/null +++ b/PrivaPub.Tests/Infrastructure/JobQueueTests.cs @@ -0,0 +1,145 @@ +using Microsoft.Extensions.Caching.Memory; +using Microsoft.Extensions.Logging.Abstractions; + +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Outbox; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.User; +using PrivaPub.StaticServices; +using PrivaPub.Tests.Support; + +using System.Text.Json.Nodes; + +namespace PrivaPub.Tests.Infrastructure +{ + public class BackoffTests + { + [Fact] + public void Grows_like_mastodons() + { + Assert.InRange(Backoff.After(1).TotalSeconds, 16, 16 + 9 * 2); + Assert.InRange(Backoff.After(10).TotalSeconds, 10015, 10015 + 9 * 11); + } + + [Fact] + public void Quarantines_a_host_only_past_the_threshold() + { + Assert.Equal(TimeSpan.Zero, Backoff.HostQuarantine(9, 10)); + Assert.Equal(TimeSpan.FromHours(1), Backoff.HostQuarantine(10, 10)); + Assert.Equal(TimeSpan.FromHours(4), Backoff.HostQuarantine(12, 10)); + Assert.Equal(TimeSpan.FromDays(7), Backoff.HostQuarantine(40, 10)); + } + } + + [Trait("Category", "Integration")] + public sealed class JobQueueTests : IAsyncLifetime + { + Peer _peer; + + public async ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + _peer = await Peer.Start(); + } + + public async ValueTask DisposeAsync() + { + if (_peer != default) + await _peer.DisposeAsync(); + } + + static Job NewJob(string host = default, string dedupe = default) => new() + { + Kind = (JobKind)99, + Host = host, + DedupeKey = dedupe, + Payload = "{}" + }; + + [Fact] + public async Task A_dedupe_key_is_queued_once() + { + var queue = new JobQueue(); + var key = Guid.NewGuid().ToString("N"); + + var inserted = await queue.EnqueueMany(new[] { NewJob(dedupe: key), NewJob(dedupe: key) }, TestContext.Current.CancellationToken); + + Assert.Equal(1, inserted); + } + + [Fact] + public async Task A_busy_host_is_skipped_and_a_failed_job_backs_off() + { + var token = TestContext.Current.CancellationToken; + var queue = new JobQueue(); + var kind = (JobKind)(100 + Random.Shared.Next(1000)); + var busy = $"busy{Guid.NewGuid():N}.example"; + var free = $"free{Guid.NewGuid():N}.example"; + await queue.EnqueueMany(new[] { new Job { Kind = kind, Host = busy }, new Job { Kind = kind, Host = free } }, token); + + var leased = await queue.Lease(kind, new[] { busy }, token); + Assert.Equal(free, leased.Host); + Assert.Equal(1, leased.Attempts); + + await queue.Finish(leased, JobOutcome.Retry("boom"), maxAttempts: 3, token); + var after = await DB.Default.Find().OneAsync(leased.ID, token); + Assert.Equal(JobState.Pending, after.State); + Assert.True(after.RunAt > DateTime.UtcNow.AddSeconds(10)); + + leased.Attempts = 3; + await queue.Finish(leased, JobOutcome.Retry("boom"), maxAttempts: 3, token); + Assert.Equal(JobState.Dead, (await DB.Default.Find().OneAsync(leased.ID, token)).State); + } + + [Fact] + public async Task An_expired_lease_returns_to_the_queue() + { + var token = TestContext.Current.CancellationToken; + var job = new Job { Kind = (JobKind)98, State = JobState.Running, LeasedUntil = DateTime.UtcNow.AddMinutes(-1) }; + await DB.Default.SaveAsync(job, token); + + await new JobQueue().Reap(token); + + Assert.Equal(JobState.Pending, (await DB.Default.Find().OneAsync(job.ID, token)).State); + } + + [Fact] + public async Task A_dead_host_does_not_hold_up_deliveries_to_live_ones() + { + var token = TestContext.Current.CancellationToken; + var (privateKey, publicKey) = Keys.NewKeyPair(); + var avatar = new Avatar { UserName = $"sender{Guid.NewGuid():N}"[..20], PrivateKey = privateKey, PublicKey = publicKey }; + await DB.Default.SaveAsync(avatar, token); + var local = new LocalActorService(new DbEntities(), new StaticOptions(new AppConfiguration { BackendBaseAddress = "https://privapub.test" })); + var sender = local.FromAvatar(avatar); + var queue = new JobQueue(); + var delivery = new DeliveryService(new DbEntities(), queue); + _peer.Answer("/dead/inbox", 503, TimeSpan.FromMilliseconds(300)); + _peer.Answer("/live/inbox", 202); + + var deadInboxes = Enumerable.Range(0, 60).Select(i => $"{_peer.B}/dead/inbox?{i}"); + await delivery.Enqueue(sender, deadInboxes, new JsonObject { ["id"] = $"https://privapub.test/a/{Guid.NewGuid():N}", ["type"] = "Create" }, token); + await delivery.Enqueue(sender, new[] { $"{_peer.A}/live/inbox" }, new JsonObject { ["id"] = $"https://privapub.test/a/{Guid.NewGuid():N}", ["type"] = "Create" }, token); + + var handler = new DeliveryJobHandler(local, Peer.Http(), new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())), NullLogger.Instance); + using var worker = new JobWorker(queue, new IJobHandler[] { handler }, NullLogger.Instance); + await worker.StartAsync(token); + var deadline = DateTime.UtcNow.AddSeconds(5); + while (DateTime.UtcNow < deadline && !_peer.Requests.Any(r => r.Path == "/live/inbox")) + await Task.Delay(50, token); + var live = Assert.Single(_peer.Requests, r => r.Path == "/live/inbox"); + Assert.True(_peer.Requests.Count(r => r.Path == "/dead/inbox") < 10); + while (DateTime.UtcNow < deadline && !await DB.Default.Find().Match(i => i.Host == "localhost").ExecuteAnyAsync(token)) + await Task.Delay(50, token); + await worker.StopAsync(token); + + Assert.Contains($"keyId=\"{sender.KeyId}\"", live.Signature); + var instance = await DB.Default.Find().Match(i => i.Host == "localhost").ExecuteFirstAsync(token); + Assert.True(instance.ConsecutiveFailures > 0); + } + } +} diff --git a/PrivaPub.Tests/Support/MongoFixture.cs b/PrivaPub.Tests/Support/MongoFixture.cs index 2eeb807..72f065c 100644 --- a/PrivaPub.Tests/Support/MongoFixture.cs +++ b/PrivaPub.Tests/Support/MongoFixture.cs @@ -32,6 +32,7 @@ namespace PrivaPub.Tests.Support var connection = Environment.GetEnvironmentVariable("PRIVAPUB_TEST_MONGO") ?? "mongodb://127.0.0.1:27017"; await DB.InitAsync(Database, MongoClientSettings.FromConnectionString(connection)); EntityMaps.Warm(); + await Indexes.Create(); } public async ValueTask DisposeAsync() diff --git a/PrivaPub.Tests/Support/Peer.cs b/PrivaPub.Tests/Support/Peer.cs index 98ce2c6..d3ac7bb 100644 --- a/PrivaPub.Tests/Support/Peer.cs +++ b/PrivaPub.Tests/Support/Peer.cs @@ -16,6 +16,7 @@ namespace PrivaPub.Tests.Support { readonly WebApplication _app; readonly ConcurrentDictionary _documents = new(); + readonly ConcurrentDictionary _answers = new(); public int Port { get; } public string A => $"http://127.0.0.1:{Port}"; @@ -38,6 +39,13 @@ namespace PrivaPub.Tests.Support { peer.Requests.Enqueue(new(context.Request.Method, context.Request.Path, context.Request.Headers["Signature"].ToString())); var key = context.Request.Path.Value; + if (peer._answers.TryGetValue(key, out var answer)) + { + if (answer.Delay > TimeSpan.Zero) + await Task.Delay(answer.Delay); + context.Response.StatusCode = answer.Status; + return; + } if (!peer._documents.TryGetValue(key, out var document)) { context.Response.StatusCode = StatusCodes.Status404NotFound; @@ -53,6 +61,8 @@ namespace PrivaPub.Tests.Support public void Serve(string path, string json) => _documents[path] = json; + public void Answer(string path, int status, TimeSpan delay = default) => _answers[path] = (status, delay); + public static FederationHttp Http(IMemoryCache cache = default) { var options = new FederationOptions { AllowPrivateNetworks = true, AllowPlainHttp = true }; diff --git a/PrivaPub/Federation/Outbox/DeliveryService.cs b/PrivaPub/Federation/Outbox/DeliveryService.cs index 9d13f73..c8a3e1d 100644 --- a/PrivaPub/Federation/Outbox/DeliveryService.cs +++ b/PrivaPub/Federation/Outbox/DeliveryService.cs @@ -10,6 +10,10 @@ using System.Text.Json.Nodes; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Signing; using PrivaPub.Infrastructure.Http; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Jobs; + +using System.Text.Json; namespace PrivaPub.Federation.Outbox { @@ -20,31 +24,38 @@ namespace PrivaPub.Federation.Outbox Task> FollowerInboxes(LocalActor actor, CancellationToken token); } + public sealed record DeliveryPayload(string SignerId, LocalActorKind SignerKind, string Inbox, string Body); + public class DeliveryService : IDeliveryService { readonly DbEntities _dbEntities; + readonly IJobQueue _queue; - public DeliveryService(DbEntities dbEntities) + public DeliveryService(DbEntities dbEntities, IJobQueue queue) { _dbEntities = dbEntities; + _queue = queue; } public async Task Enqueue(LocalActor signer, IEnumerable inboxes, JsonObject activity, CancellationToken token) { var body = activity.ToJsonString(); - var deliveries = inboxes + var activityId = activity["id"] is JsonValue id && id.TryGetValue(out var text) ? text : default; + var jobs = inboxes .Where(i => !string.IsNullOrEmpty(i) && !i.StartsWith(signer.BaseAddress + "/", StringComparison.OrdinalIgnoreCase)) .Distinct(StringComparer.Ordinal) - .Select(inbox => new Delivery + .Select(inbox => Uri.TryCreate(inbox, UriKind.Absolute, out var uri) ? (inbox, uri) : default) + .Where(target => target.uri != default) + .Select(target => new Job { - SignerId = signer.Id, - SignerKind = signer.Kind, - InboxURL = inbox, - Body = body + Kind = JobKind.Deliver, + Host = target.uri.Host.ToLowerInvariant(), + DedupeKey = activityId == default ? default : $"{activityId}|{target.inbox}", + Payload = JsonSerializer.Serialize(new DeliveryPayload(signer.Id, signer.Kind, target.inbox, body)) }) .ToList(); - if (deliveries.Count > 0) - await DB.Default.SaveAsync(deliveries, token); + if (jobs.Count > 0) + await _queue.EnqueueMany(jobs, token); } public async Task EnqueueToFollowers(LocalActor signer, JsonObject activity, CancellationToken token, IEnumerable extraInboxes = default) @@ -64,118 +75,81 @@ namespace PrivaPub.Federation.Outbox } } - public class DeliveryWorker : BackgroundService + public class DeliveryJobHandler : IJobHandler { - const int MaxAttempts = 8; - static readonly TimeSpan Poll = TimeSpan.FromSeconds(3); - - readonly IServiceProvider _services; + readonly ILocalActorService _actors; readonly IFederationHttp _http; - readonly ILogger _logger; + readonly IHostCircuitBreaker _breaker; + readonly ILogger _logger; - public DeliveryWorker(IServiceProvider services, IFederationHttp http, ILogger logger) + public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger logger) { - _services = services; + _actors = actors; _http = http; + _breaker = breaker; _logger = logger; } - protected override async Task ExecuteAsync(CancellationToken stoppingToken) - { - while (!stoppingToken.IsCancellationRequested) - { - try - { - await DeliverDue(stoppingToken); - } - catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) - { - return; - } - catch (Exception ex) - { - _logger.LogError(ex, "{Worker} pass failed", nameof(DeliveryWorker)); - } - await Task.Delay(Poll, stoppingToken); - } - } + public JobKind Kind => JobKind.Deliver; + public int Concurrency => 8; + public int MaxAttempts => 16; + public int PerHostLimit => 2; - async Task DeliverDue(CancellationToken token) + public async Task Handle(Job job, CancellationToken token) { - using var scope = _services.CreateScope(); - var dbEntities = scope.ServiceProvider.GetRequiredService(); - var actors = scope.ServiceProvider.GetRequiredService(); - var now = DateTime.UtcNow; - var due = await dbEntities.Deliveries - .Match(d => !d.DeliveredAt.HasValue && !d.AbandonedAt.HasValue && d.NextAttemptAt <= now) - .Sort(d => d.NextAttemptAt, Order.Ascending) - .Limit(20) - .ExecuteAsync(token); + var payload = JsonSerializer.Deserialize(job.Payload); + if (payload == default || !Uri.TryCreate(payload.Inbox, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox)) + return JobOutcome.Dead("not a deliverable inbox"); - foreach (var delivery in due) - { - var signer = await actors.FindById(delivery.SignerKind, delivery.SignerId, token); - if (signer == default) - { - delivery.AbandonedAt = DateTime.UtcNow; - delivery.LastError = "the signing actor no longer exists"; - await DB.Default.SaveAsync(delivery, token); - continue; - } - await Deliver(delivery, signer, token); - await DB.Default.SaveAsync(delivery, token); - } - } + var unavailableUntil = await _breaker.UnavailableUntil(job.Host, token); + if (unavailableUntil.HasValue) + return JobOutcome.Defer(unavailableUntil.Value, "the host is unavailable"); + + var signer = await _actors.FindById(payload.SignerKind, payload.SignerId, token); + if (signer == default) + return JobOutcome.Dead("the signing actor no longer exists"); + + var body = Encoding.UTF8.GetBytes(payload.Body); + using var request = new HttpRequestMessage(HttpMethod.Post, inbox) { Content = new ByteArrayContent(body) }; + request.Content.Headers.ContentType = MediaTypeHeaderValue.Parse(RemoteActorService.ActivityJson); + HttpSignatures.Sign(request, signer, body); - async Task Deliver(Delivery delivery, LocalActor signer, CancellationToken token) - { - delivery.Attempts++; try { - if (!Uri.TryCreate(delivery.InboxURL, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox)) - { - delivery.AbandonedAt = DateTime.UtcNow; - delivery.LastError = "not a deliverable inbox"; - return; - } - - var body = Encoding.UTF8.GetBytes(delivery.Body); - using var request = new HttpRequestMessage(HttpMethod.Post, inbox) - { - Content = new ByteArrayContent(body) - }; - request.Content.Headers.ContentType = MediaTypeHeaderValue.Parse(RemoteActorService.ActivityJson); - HttpSignatures.Sign(request, signer, body); - using var response = await _http.Send(request, token); + var status = (int)response.StatusCode; if (response.IsSuccessStatusCode) { - delivery.DeliveredAt = DateTime.UtcNow; - delivery.LastError = default; - return; + await _breaker.Succeeded(job.Host, token); + return JobOutcome.Done; } - - delivery.LastError = $"{(int)response.StatusCode} {response.ReasonPhrase}"; - if (response.StatusCode is HttpStatusCode.Gone or HttpStatusCode.NotFound or HttpStatusCode.BadRequest or HttpStatusCode.Forbidden) + if (response.StatusCode == HttpStatusCode.TooManyRequests) { - delivery.AbandonedAt = DateTime.UtcNow; - _logger.LogWarning("Delivery {Id} to {Inbox} abandoned: {Status}", delivery.ID, delivery.InboxURL, delivery.LastError); - return; + var retryAfter = response.Headers.RetryAfter?.Delta ?? (response.Headers.RetryAfter?.Date - DateTimeOffset.UtcNow); + return retryAfter > TimeSpan.Zero + ? JobOutcome.Defer(DateTime.UtcNow + Min(retryAfter.Value, TimeSpan.FromHours(6)), "429") + : JobOutcome.Retry("429"); } + if (status is >= 400 and < 500 && status != 408) + { + await _breaker.Succeeded(job.Host, token); + _logger.LogInformation("Delivery to {Inbox} refused with {Status}", payload.Inbox, status); + return JobOutcome.Dead($"{status} {response.ReasonPhrase}"); + } + await _breaker.Failed(job.Host, $"{status}", token); + return JobOutcome.Retry($"{status} {response.ReasonPhrase}"); } - catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or TaskCanceledException && !token.IsCancellationRequested) + catch (BlockedDestinationException ex) { - delivery.LastError = ex.Message; + return JobOutcome.Dead(ex.Message); } - - if (delivery.Attempts >= MaxAttempts) + catch (Exception ex) when (ex is HttpRequestException or TaskCanceledException && !token.IsCancellationRequested) { - delivery.AbandonedAt = DateTime.UtcNow; - _logger.LogWarning("Delivery {Id} to {Inbox} abandoned after {Attempts} attempts: {Error}", - delivery.ID, delivery.InboxURL, delivery.Attempts, delivery.LastError); - return; + await _breaker.Failed(job.Host, ex.GetType().Name, token); + return JobOutcome.Retry(ex.Message); } - delivery.NextAttemptAt = DateTime.UtcNow.AddMinutes(Math.Pow(2, delivery.Attempts)); } + + static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b; } } diff --git a/PrivaPub/Infrastructure/Data/Indexes.cs b/PrivaPub/Infrastructure/Data/Indexes.cs index 428d207..d690651 100644 --- a/PrivaPub/Infrastructure/Data/Indexes.cs +++ b/PrivaPub/Infrastructure/Data/Indexes.cs @@ -4,6 +4,7 @@ using MongoDB.Entities; using PrivaPub.Models.Federation; using PrivaPub.Models.Group; +using PrivaPub.Models.Jobs; using PrivaPub.Models.Post; using PrivaPub.Models.User; @@ -48,6 +49,15 @@ namespace PrivaPub.Infrastructure.Data await Plain(token, g => g.ParticipantsKey); await Plain(token, d => d.DeliveredAt, d => d.AbandonedAt, d => d.NextAttemptAt); + + await Plain(token, j => j.Kind, j => j.State, j => j.RunAt); + await Plain(token, j => j.State, j => j.LeasedUntil); + await Unique(j => j.DedupeKey, Builders.Filter.Type(j => j.DedupeKey, BsonType.String), token); + await DB.Default.Index() + .Key(j => j.FinishedAt, KeyType.Ascending) + .Option(o => o.ExpireAfter = TimeSpan.FromDays(7)) + .CreateAsync(token); + await Unique(i => i.Host, Builders.Filter.Type(i => i.Host, BsonType.String), token); } static async Task Unique(System.Linq.Expressions.Expression> key, FilterDefinition partial, diff --git a/PrivaPub/Infrastructure/Data/Migrations/_004_pending_deliveries_become_jobs.cs b/PrivaPub/Infrastructure/Data/Migrations/_004_pending_deliveries_become_jobs.cs new file mode 100644 index 0000000..7a35839 --- /dev/null +++ b/PrivaPub/Infrastructure/Data/Migrations/_004_pending_deliveries_become_jobs.cs @@ -0,0 +1,43 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Outbox; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Jobs; + +using System.Text.Json; +using System.Text.Json.Nodes; + +namespace PrivaPub.Infrastructure.Data.Migrations +{ + public class _004_pending_deliveries_become_jobs : IMigration + { + public async Task UpgradeAsync() + { + var pending = await DB.Default.Find() + .Match(d => !d.DeliveredAt.HasValue && !d.AbandonedAt.HasValue) + .ExecuteAsync(); + foreach (var delivery in pending) + { + if (Uri.TryCreate(delivery.InboxURL, UriKind.Absolute, out var inbox)) + { + var activityId = JsonNode.Parse(delivery.Body)?["id"]?.GetValue(); + var job = new Job + { + Kind = JobKind.Deliver, + Host = inbox.Host.ToLowerInvariant(), + DedupeKey = activityId == default ? default : $"{activityId}|{delivery.InboxURL}", + Attempts = delivery.Attempts, + RunAt = delivery.NextAttemptAt, + Payload = JsonSerializer.Serialize(new DeliveryPayload(delivery.SignerId, delivery.SignerKind, delivery.InboxURL, delivery.Body)) + }; + if (!await DB.Default.Find().Match(j => j.DedupeKey == job.DedupeKey && job.DedupeKey != null).ExecuteAnyAsync()) + await DB.Default.SaveAsync(job); + } + await DB.Default.Update().MatchID(delivery.ID) + .Modify(d => d.AbandonedAt, DateTime.UtcNow) + .Modify(d => d.LastError, "moved to the job queue") + .ExecuteAsync(); + } + } + } +} diff --git a/PrivaPub/Infrastructure/Jobs/Backoff.cs b/PrivaPub/Infrastructure/Jobs/Backoff.cs new file mode 100644 index 0000000..aecbe8f --- /dev/null +++ b/PrivaPub/Infrastructure/Jobs/Backoff.cs @@ -0,0 +1,13 @@ +namespace PrivaPub.Infrastructure.Jobs +{ + public static class Backoff + { + public static TimeSpan After(int attempt) => + TimeSpan.FromSeconds(Math.Pow(attempt, 4) + 15 + Random.Shared.Next(10) * (attempt + 1)); + + public static TimeSpan HostQuarantine(int consecutiveFailures, int threshold) => + consecutiveFailures < threshold + ? TimeSpan.Zero + : TimeSpan.FromHours(Math.Min(Math.Pow(2, consecutiveFailures - threshold), 24 * 7)); + } +} diff --git a/PrivaPub/Infrastructure/Jobs/HostCircuitBreaker.cs b/PrivaPub/Infrastructure/Jobs/HostCircuitBreaker.cs new file mode 100644 index 0000000..84fe1f6 --- /dev/null +++ b/PrivaPub/Infrastructure/Jobs/HostCircuitBreaker.cs @@ -0,0 +1,76 @@ +using Microsoft.Extensions.Caching.Memory; + +using MongoDB.Entities; + +using PrivaPub.Models.Jobs; + +namespace PrivaPub.Infrastructure.Jobs +{ + public interface IHostCircuitBreaker + { + Task UnavailableUntil(string host, CancellationToken token); + Task Succeeded(string host, CancellationToken token); + Task Failed(string host, string error, CancellationToken token); + } + + public class HostCircuitBreaker : IHostCircuitBreaker + { + public const int Threshold = 10; + static readonly TimeSpan CacheLifetime = TimeSpan.FromSeconds(30); + + readonly IMemoryCache _cache; + + public HostCircuitBreaker(IMemoryCache cache) + { + _cache = cache; + } + + public async Task UnavailableUntil(string host, CancellationToken token) + { + var instance = await Instance(host, token); + return instance?.UnavailableUntil > DateTime.UtcNow ? instance.UnavailableUntil : default; + } + + public async Task Succeeded(string host, CancellationToken token) + { + var instance = await Instance(host, token); + if (instance is not { ConsecutiveFailures: > 0 } && instance?.LastSuccessAt > DateTime.UtcNow.AddHours(-1)) + return; + await DB.Default.Update() + .Match(i => i.Host == host) + .Modify(i => i.ConsecutiveFailures, 0) + .Modify(i => i.UnavailableUntil, null) + .Modify(i => i.LastSuccessAt, DateTime.UtcNow) + .Option(o => o.IsUpsert = true) + .ExecuteAsync(token); + _cache.Remove(Key(host)); + } + + public async Task Failed(string host, string error, CancellationToken token) + { + var now = DateTime.UtcNow; + var instance = await DB.Default.UpdateAndGet() + .Match(i => i.Host == host) + .Modify(b => b.Inc(i => i.ConsecutiveFailures, 1)) + .Modify(i => i.LastFailureAt, now) + .Modify(i => i.LastError, error) + .Option(o => o.IsUpsert = true) + .ExecuteAsync(token); + var quarantine = Backoff.HostQuarantine(instance.ConsecutiveFailures, Threshold); + if (quarantine > TimeSpan.Zero) + await DB.Default.Update().MatchID(instance.ID) + .Modify(i => i.UnavailableUntil, now + quarantine) + .ExecuteAsync(token); + _cache.Remove(Key(host)); + } + + async Task Instance(string host, CancellationToken token) => + await _cache.GetOrCreateAsync(Key(host), async entry => + { + entry.AbsoluteExpirationRelativeToNow = CacheLifetime; + return await DB.Default.Find().Match(i => i.Host == host).ExecuteFirstAsync(token); + }); + + static string Key(string host) => "remote-instance:" + host; + } +} diff --git a/PrivaPub/Infrastructure/Jobs/JobQueue.cs b/PrivaPub/Infrastructure/Jobs/JobQueue.cs new file mode 100644 index 0000000..9d2f655 --- /dev/null +++ b/PrivaPub/Infrastructure/Jobs/JobQueue.cs @@ -0,0 +1,142 @@ +using MongoDB.Driver; +using MongoDB.Entities; + +using PrivaPub.Models.Jobs; + +using System.Collections.Concurrent; + +namespace PrivaPub.Infrastructure.Jobs +{ + public sealed record JobOutcome(JobResult Result, string Error = default, DateTime? RetryAt = default) + { + public static readonly JobOutcome Done = new(JobResult.Done); + public static JobOutcome Retry(string error) => new(JobResult.Retry, error); + public static JobOutcome Dead(string error) => new(JobResult.Dead, error); + public static JobOutcome Defer(DateTime until, string error) => new(JobResult.Defer, error, until); + } + + public enum JobResult + { + Done, + Retry, + Dead, + Defer + } + + public interface IJobQueue + { + Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token); + Task EnqueueMany(IEnumerable jobs, CancellationToken token); + Task Lease(JobKind kind, IReadOnlyCollection busyHosts, CancellationToken token); + Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token); + Task Reap(CancellationToken token); + Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token); + } + + public class JobQueue : IJobQueue + { + public static readonly TimeSpan LeaseTime = TimeSpan.FromMinutes(2); + + readonly string _owner = $"{Environment.MachineName}:{Environment.ProcessId}"; + readonly ConcurrentDictionary _signals = new(); + + public async Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token) => + await EnqueueMany(new[] { new Job { Kind = kind, Payload = payload, Host = host, DedupeKey = dedupeKey } }, token) == 1; + + public async Task EnqueueMany(IEnumerable jobs, CancellationToken token) + { + var inserted = 0; + foreach (var job in jobs) + { + try + { + await DB.Default.SaveAsync(job, token); + inserted++; + Signal(job.Kind); + } + catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey) + { + } + } + return inserted; + } + + public async Task Lease(JobKind kind, IReadOnlyCollection busyHosts, CancellationToken token) + { + var now = DateTime.UtcNow; + var busy = busyHosts?.ToList() ?? new List(); + return await DB.Default.UpdateAndGet() + .Match(j => j.Kind == kind && j.State == JobState.Pending && j.RunAt <= now && !busy.Contains(j.Host)) + .Modify(j => j.State, JobState.Running) + .Modify(j => j.LeasedUntil, now + LeaseTime) + .Modify(j => j.LeaseOwner, _owner) + .Modify(b => b.Inc(j => j.Attempts, 1)) + .Option(o => o.Sort = Builders.Sort.Ascending(j => j.RunAt)) + .ExecuteAsync(token); + } + + public async Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token) + { + var now = DateTime.UtcNow; + var update = DB.Default.Update().MatchID(job.ID) + .Modify(j => j.LeasedUntil, null) + .Modify(j => j.LeaseOwner, null) + .Modify(j => j.LastError, outcome.Error); + switch (outcome.Result) + { + case JobResult.Done: + update.Modify(j => j.State, JobState.Done).Modify(j => j.FinishedAt, now); + break; + case JobResult.Retry when job.Attempts < maxAttempts: + update.Modify(j => j.State, JobState.Pending).Modify(j => j.RunAt, now + Backoff.After(job.Attempts)); + break; + case JobResult.Defer: + update.Modify(j => j.State, JobState.Pending) + .Modify(j => j.RunAt, outcome.RetryAt ?? now + Backoff.After(job.Attempts)) + .Modify(b => b.Inc(j => j.Attempts, -1)); + break; + default: + update.Modify(j => j.State, JobState.Dead).Modify(j => j.FinishedAt, now); + break; + } + await update.ExecuteAsync(token); + } + + public async Task Reap(CancellationToken token) + { + var now = DateTime.UtcNow; + var result = await DB.Default.Update() + .Match(j => j.State == JobState.Running && j.LeasedUntil < now) + .Modify(j => j.State, JobState.Pending) + .Modify(j => j.RunAt, now) + .Modify(j => j.LeasedUntil, null) + .Modify(j => j.LeaseOwner, null) + .ExecuteAsync(token); + return result.ModifiedCount; + } + + public async Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token) + { + try + { + await SignalFor(kind).WaitAsync(poll, token); + } + catch (OperationCanceledException) when (token.IsCancellationRequested) + { + } + } + + void Signal(JobKind kind) + { + try + { + SignalFor(kind).Release(); + } + catch (SemaphoreFullException) + { + } + } + + SemaphoreSlim SignalFor(JobKind kind) => _signals.GetOrAdd(kind, _ => new SemaphoreSlim(0, 64)); + } +} diff --git a/PrivaPub/Infrastructure/Jobs/JobWorker.cs b/PrivaPub/Infrastructure/Jobs/JobWorker.cs new file mode 100644 index 0000000..da6dd15 --- /dev/null +++ b/PrivaPub/Infrastructure/Jobs/JobWorker.cs @@ -0,0 +1,137 @@ +using PrivaPub.Models.Jobs; + +using System.Collections.Concurrent; + +namespace PrivaPub.Infrastructure.Jobs +{ + public interface IJobHandler + { + JobKind Kind { get; } + int Concurrency { get; } + int MaxAttempts { get; } + int PerHostLimit { get; } + Task Handle(Job job, CancellationToken token); + } + + public class JobWorker : BackgroundService + { + static readonly TimeSpan Poll = TimeSpan.FromSeconds(5); + static readonly TimeSpan ReapInterval = TimeSpan.FromSeconds(30); + + readonly IJobQueue _queue; + readonly IEnumerable _handlers; + readonly ILogger _logger; + + public JobWorker(IJobQueue queue, IEnumerable handlers, ILogger logger) + { + _queue = queue; + _handlers = handlers; + _logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + var loops = _handlers.SelectMany(handler => + { + var inFlight = new ConcurrentDictionary(StringComparer.OrdinalIgnoreCase); + return Enumerable.Range(0, handler.Concurrency).Select(_ => Run(handler, inFlight, stoppingToken)); + }).Append(Reap(stoppingToken)); + await Task.WhenAll(loops); + } + + async Task Run(IJobHandler handler, ConcurrentDictionary inFlight, CancellationToken stoppingToken) + { + while (!stoppingToken.IsCancellationRequested) + { + Job job; + try + { + var busy = inFlight.Where(h => h.Value >= handler.PerHostLimit).Select(h => h.Key).ToList(); + job = await _queue.Lease(handler.Kind, busy, stoppingToken); + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + return; + } + catch (Exception ex) + { + _logger.LogError(ex, "{Worker} could not lease a {Kind} job", nameof(JobWorker), handler.Kind); + await Delay(Poll, stoppingToken); + continue; + } + + if (job == default) + { + await _queue.WaitForWork(handler.Kind, Poll, stoppingToken); + continue; + } + + var host = job.Host ?? string.Empty; + inFlight.AddOrUpdate(host, 1, (_, count) => count + 1); + try + { + JobOutcome outcome; + try + { + outcome = await handler.Handle(job, stoppingToken); + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + return; + } + catch (Exception ex) + { + _logger.LogError(ex, "{Kind} job {Id} threw", handler.Kind, job.ID); + outcome = JobOutcome.Retry(ex.GetType().Name); + } + + if (outcome.Result != JobResult.Done && job.Attempts >= handler.MaxAttempts && outcome.Result == JobResult.Retry) + _logger.LogWarning("{Kind} job {Id} for {Host} is dead after {Attempts} attempts: {Error}", + handler.Kind, job.ID, job.Host, job.Attempts, outcome.Error); + await _queue.Finish(job, outcome, handler.MaxAttempts, CancellationToken.None); + } + catch (Exception ex) + { + _logger.LogError(ex, "{Kind} job {Id} could not be finished; its lease will expire", handler.Kind, job.ID); + } + finally + { + inFlight.AddOrUpdate(host, 0, (_, count) => count - 1); + } + } + } + + async Task Reap(CancellationToken stoppingToken) + { + while (!stoppingToken.IsCancellationRequested) + { + try + { + var reaped = await _queue.Reap(stoppingToken); + if (reaped > 0) + _logger.LogWarning("{Worker} returned {Count} expired leases to the queue", nameof(JobWorker), reaped); + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) + { + return; + } + catch (Exception ex) + { + _logger.LogError(ex, "{Worker} reaper failed", nameof(JobWorker)); + } + await Delay(ReapInterval, stoppingToken); + } + } + + static async Task Delay(TimeSpan delay, CancellationToken token) + { + try + { + await Task.Delay(delay, token); + } + catch (OperationCanceledException) + { + } + } + } +} diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index acf60d0..a2b87d2 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -15,6 +15,7 @@ using PrivaPub.Federation.Actors; using PrivaPub.Federation.Outbox; using PrivaPub.Federation.Inbox; using PrivaPub.Infrastructure.Http; +using PrivaPub.Infrastructure.Jobs; using Microsoft.Extensions.Options; namespace PrivaPub.Middleware @@ -51,7 +52,10 @@ namespace PrivaPub.Middleware .AddSingleton() .AddSingleton() .AddTransient() - .AddHostedService(); + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddHostedService(); } public static IServiceCollection PrivaPubAuthServicesConfiguration(this IServiceCollection service, IConfiguration configuration) { diff --git a/PrivaPub/Models/Jobs/Job.cs b/PrivaPub/Models/Jobs/Job.cs new file mode 100644 index 0000000..a31e6f8 --- /dev/null +++ b/PrivaPub/Models/Jobs/Job.cs @@ -0,0 +1,34 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Jobs +{ + public class Job : Entity + { + public JobKind Kind { get; set; } + public JobState State { get; set; } + public string Payload { get; set; } + public string Host { get; set; } + public string DedupeKey { get; set; } + public int Attempts { get; set; } + public DateTime RunAt { get; set; } = DateTime.UtcNow; + public DateTime? LeasedUntil { get; set; } + public string LeaseOwner { get; set; } + public string LastError { get; set; } + public DateTime CreatedAt { get; set; } = DateTime.UtcNow; + public DateTime? FinishedAt { get; set; } + } + + public enum JobKind + { + Deliver, + ProcessInbox + } + + public enum JobState + { + Pending, + Running, + Done, + Dead + } +} diff --git a/PrivaPub/Models/Jobs/RemoteInstance.cs b/PrivaPub/Models/Jobs/RemoteInstance.cs new file mode 100644 index 0000000..3241987 --- /dev/null +++ b/PrivaPub/Models/Jobs/RemoteInstance.cs @@ -0,0 +1,14 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Jobs +{ + public class RemoteInstance : Entity + { + public string Host { get; set; } + public int ConsecutiveFailures { get; set; } + public DateTime? UnavailableUntil { get; set; } + public DateTime? LastSuccessAt { get; set; } + public DateTime? LastFailureAt { get; set; } + public string LastError { get; set; } + } +}