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 { [Xunit.Collection(nameof(Exclusive))] [Trait("Category", "Integration")] public sealed class DeliveryWorkerTests : IAsyncLifetime { static readonly string[] PeerHosts = { "localhost", "127.0.0.1" }; Peer _peer; public async ValueTask InitializeAsync() { Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); _peer = await Peer.Start(); await DB.Default.DeleteAsync(i => PeerHosts.Contains(i.Host)); } public async ValueTask DisposeAsync() { if (_peer == default) return; await _peer.DisposeAsync(); await DB.Default.DeleteAsync(i => PeerHosts.Contains(i.Host)); } [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 run = $"https://privapub.test/a/{Guid.NewGuid():N}/"; var queue = new JobQueue(j => j.DedupeKey.StartsWith(run)); 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"] = run + "dead", ["type"] = "Create" }, token); await delivery.Enqueue(sender, new[] { $"{_peer.A}/live/inbox" }, new JsonObject { ["id"] = run + "live", ["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); try { await Until(() => Task.FromResult(_peer.Requests.Any(r => r.Path == "/live/inbox")), token); var deadSoFar = _peer.Requests.Count(r => r.Path == "/dead/inbox"); var live = Assert.Single(_peer.Requests, r => r.Path == "/live/inbox"); Assert.Contains($"keyId=\"{sender.KeyId}\"", live.Signature); Assert.True(deadSoFar < 30, $"{deadSoFar} dead deliveries were tried before the live one"); await Until(() => DB.Default.Find().Match(j => j.DedupeKey.StartsWith(run + "dead") && j.State == JobState.Pending && j.Attempts > 0).ExecuteAnyAsync(token), token); } finally { await worker.StopAsync(CancellationToken.None); } var retried = await DB.Default.Find().Match(j => j.DedupeKey.StartsWith(run + "dead") && j.Attempts > 0 && j.State == JobState.Pending).ExecuteFirstAsync(token); Assert.Equal("503", retried.LastError?.Split(' ')[0]); Assert.True(retried.RunAt > DateTime.UtcNow); Assert.True((await DB.Default.Find().Match(i => i.Host == "localhost").ExecuteFirstAsync(token))?.ConsecutiveFailures > 0); } static async Task Until(Func> condition, CancellationToken token) { var deadline = DateTime.UtcNow.AddSeconds(30); while (!await condition()) { Assert.True(DateTime.UtcNow < deadline, "timed out"); await Task.Delay(50, token); } } } }