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.Federation; using PrivaPub.Models.Jobs; using PrivaPub.Models.Statistics; using PrivaPub.Tests.Support; using System.Text.Json.Nodes; namespace PrivaPub.Tests.Statistics { [Xunit.Collection(nameof(Exclusive))] [Trait("Category", "Integration")] public sealed class DeliveryAttemptTests : IAsyncLifetime { static readonly string[] PeerHosts = { "localhost", "127.0.0.1" }; Harness _harness; public async ValueTask InitializeAsync() { Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); _harness = await Harness.Start(); await DB.Default.DeleteAsync(i => PeerHosts.Contains(i.Host)); } public async ValueTask DisposeAsync() { if (_harness == default) return; await _harness.DisposeAsync(); await DB.Default.DeleteAsync(i => PeerHosts.Contains(i.Host)); } async Task<(JobOutcome Outcome, InteractionEvent Event)> Attempt(string path, JsonObject activity, LocalActor signer = default) { var ledger = new MemoryLedger(); var handler = new DeliveryJobHandler(_harness.Local, Peer.Http(), new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())), NullLogger.Instance, ledger); signer ??= (await _harness.Persona("sender")).Avatar; var id = activity["id"]!.GetValue(); await _harness.Delivery.Enqueue(signer, new[] { _harness.Peer.A + path }, activity, TestContext.Current.CancellationToken); var job = await DB.Default.Find().Match(j => j.DedupeKey == $"{id}|{_harness.Peer.A + path}").ExecuteFirstAsync(TestContext.Current.CancellationToken); job.Attempts = 1; var outcome = await handler.Handle(job, TestContext.Current.CancellationToken); return (outcome, Assert.Single(ledger.Of("out"))); } static JsonObject Create(string audience) => new() { ["id"] = $"https://privapub.test/a/{Guid.NewGuid():N}", ["type"] = "Create", ["to"] = new JsonArray(audience), ["object"] = new JsonObject { ["type"] = "Note", ["content"] = "hi" } }; [Fact] public async Task A_delivered_public_post_is_ok_with_its_shape_size_and_time() { _harness.Peer.Answer("/ok/inbox", 202); var (outcome, e) = await Attempt("/ok/inbox", Create("https://www.w3.org/ns/activitystreams#Public")); Assert.Equal(JobResult.Done, outcome.Result); Assert.Equal(("ok", "Create", "Note", "public", "person"), (e.Outcome, e.Activity, e.Object, e.Audience, e.LocalKind)); Assert.Equal(202, e.Status); Assert.Equal("127.0.0.1", e.Host); Assert.Null(e.Reason); Assert.True(e.Bytes > 0); Assert.NotNull(e.LatencyMs); Assert.Equal(1, e.Attempt); } [Fact] public async Task Refusals_and_failures_carry_their_status_and_private_traffic_has_no_kind() { _harness.Peer.Answer("/gone/inbox", 410); _harness.Peer.Answer("/busy/inbox", 429); _harness.Peer.Answer("/broken/inbox", 500); var (_, gone) = await Attempt("/gone/inbox", Create("https://elsewhere.example/users/x")); var (_, busy) = await Attempt("/busy/inbox", Create("https://elsewhere.example/users/x")); var (_, broken) = await Attempt("/broken/inbox", Create("https://elsewhere.example/users/x")); Assert.Equal(("dead", "410", 410), (gone.Outcome, gone.Reason, gone.Status!.Value)); Assert.Equal(("retry", "429"), (busy.Outcome, busy.Reason)); Assert.Equal(("retry", "500"), (broken.Outcome, broken.Reason)); Assert.All(new[] { gone, busy, broken }, e => Assert.Equal(("private", (string)null), (e.Audience, e.LocalKind))); } [Fact] public async Task A_signer_that_no_longer_exists_is_dead() { var ghost = new LocalActor { Id = "000000000000000000000000", Kind = LocalActorKind.Person, UserName = "ghost", BaseAddress = Harness.Base }; var (outcome, e) = await Attempt("/never/inbox", Create("https://www.w3.org/ns/activitystreams#Public"), ghost); Assert.Equal(JobResult.Dead, outcome.Result); Assert.Equal(("dead", "signer-gone"), (e.Outcome, e.Reason)); Assert.Null(e.Status); } } }