From e247001bdba382d64437100c16e4ecdffb260346 Mon Sep 17 00:00:00 2001 From: thepra Date: Sat, 3 Oct 2026 11:06:24 +0200 Subject: [PATCH] M4: every delivery attempt is recorded DeliveryJobHandler records each attempt as an 'out' event: the activity, object type and audience read from the body (ActivityShape, shared with the inbox side), the receiving server, the status, the time the POST took, the wait since it was queued, the bytes, the attempt number and the signer's kind (none for circles or private traffic). The outcome is ok, deferred (429 or 503 with Retry-After, or an open breaker: host-unavailable), retry (5xx, 408, timeout, network), or dead: another 4xx, private-address, not-deliverable, signer-gone, or a retry at the last attempt. Until now none of this outlived the job's seven-day TTL, and the last error was overwritten on success. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2 --- .../Statistics/DeliveryAttemptTests.cs | 111 ++++++++++++++++++ PrivaPub/Federation/Inbox/InboxProcessor.cs | 21 +--- PrivaPub/Federation/Objects/ActivityShape.cs | 48 ++++++++ PrivaPub/Federation/Outbox/DeliveryService.cs | 90 +++++++++++++- 4 files changed, 249 insertions(+), 21 deletions(-) create mode 100644 PrivaPub.Tests/Statistics/DeliveryAttemptTests.cs create mode 100644 PrivaPub/Federation/Objects/ActivityShape.cs diff --git a/PrivaPub.Tests/Statistics/DeliveryAttemptTests.cs b/PrivaPub.Tests/Statistics/DeliveryAttemptTests.cs new file mode 100644 index 0000000..e48458e --- /dev/null +++ b/PrivaPub.Tests/Statistics/DeliveryAttemptTests.cs @@ -0,0 +1,111 @@ +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); + } + } +} diff --git a/PrivaPub/Federation/Inbox/InboxProcessor.cs b/PrivaPub/Federation/Inbox/InboxProcessor.cs index 280ff7d..0951558 100644 --- a/PrivaPub/Federation/Inbox/InboxProcessor.cs +++ b/PrivaPub/Federation/Inbox/InboxProcessor.cs @@ -97,7 +97,7 @@ namespace PrivaPub.Federation.Inbox LatencyMs = (int)Stopwatch.GetElapsedTime(started).TotalMilliseconds, WaitMs = (int)Math.Min(int.MaxValue, Math.Max(0, (DateTime.UtcNow - (payload.ReceivedAt ?? job.CreatedAt)).TotalMilliseconds)), Attempt = job.Attempts, - Audience = verdict.Audience ?? Audience(activity), + Audience = verdict.Audience ?? ActivityShape.Of(activity).Audience, LocalKind = verdict.LocalKind, AgeSeconds = verdict.AgeSeconds, Inbox = payload.Inbox, @@ -105,24 +105,5 @@ namespace PrivaPub.Federation.Inbox Features = type is "Create" or "Update" && embedded != default ? ObjectFeatures.Detect(embedded) : default }, payload.ActorURI); } - - static string Audience(JsonNode activity) - { - var source = activity?["to"] != default || activity?["cc"] != default ? activity : activity?["object"] as JsonObject; - var to = Addresses(source?["to"]); - var cc = Addresses(source?["cc"]); - if (to.Any(Addressing.IsPublic)) - return Interactions.Public; - if (cc.Any(Addressing.IsPublic)) - return Interactions.Unlisted; - return to.Count + cc.Count == 0 ? Interactions.None : Interactions.Private; - } - - static List Addresses(JsonNode node) => node switch - { - JsonArray array => array.Select(Id).Where(id => id != default).ToList(), - JsonNode single when Id(single) is { } id => new List { id }, - _ => new List() - }; } } diff --git a/PrivaPub/Federation/Objects/ActivityShape.cs b/PrivaPub/Federation/Objects/ActivityShape.cs new file mode 100644 index 0000000..26719c6 --- /dev/null +++ b/PrivaPub/Federation/Objects/ActivityShape.cs @@ -0,0 +1,48 @@ +using PrivaPub.Infrastructure.Statistics; + +using System.Text.Json.Nodes; + +using static PrivaPub.Federation.Objects.ActivityJson; + +namespace PrivaPub.Federation.Objects +{ + public sealed record ActivityShape(string Type, string ObjectType, string Audience) + { + public static ActivityShape Of(JsonNode activity) + { + var embedded = activity?["object"] as JsonObject; + return new ActivityShape(Value(activity, "type"), Value(embedded, "type"), AudienceOf(activity)); + } + + public static ActivityShape Of(string json) + { + try + { + return Of(JsonNode.Parse(json)); + } + catch (System.Text.Json.JsonException) + { + return new ActivityShape(default, default, default); + } + } + + static string AudienceOf(JsonNode activity) + { + var source = activity?["to"] != default || activity?["cc"] != default ? activity : activity?["object"] as JsonObject; + var to = Addresses(source?["to"]); + var cc = Addresses(source?["cc"]); + if (to.Any(Addressing.IsPublic)) + return Interactions.Public; + if (cc.Any(Addressing.IsPublic)) + return Interactions.Unlisted; + return to.Count + cc.Count == 0 ? Interactions.None : Interactions.Private; + } + + static List Addresses(JsonNode node) => node switch + { + JsonArray array => array.Select(Id).Where(id => id != default).ToList(), + JsonNode single when Id(single) is { } id => new List { id }, + _ => new List() + }; + } +} diff --git a/PrivaPub/Federation/Outbox/DeliveryService.cs b/PrivaPub/Federation/Outbox/DeliveryService.cs index 3f92f3d..1e57f52 100644 --- a/PrivaPub/Federation/Outbox/DeliveryService.cs +++ b/PrivaPub/Federation/Outbox/DeliveryService.cs @@ -11,8 +11,12 @@ using PrivaPub.Federation.Actors; using PrivaPub.Federation.Signing; using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Infrastructure.Statistics; +using PrivaPub.Federation.Objects; +using PrivaPub.Models.Statistics; using PrivaPub.Models.Jobs; +using System.Diagnostics; using System.Text.Json; namespace PrivaPub.Federation.Outbox @@ -81,13 +85,25 @@ namespace PrivaPub.Federation.Outbox readonly IFederationHttp _http; readonly IHostCircuitBreaker _breaker; readonly ILogger _logger; + readonly IInteractionLedger _ledger; - public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger logger) + public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger logger, + IInteractionLedger ledger = default) { _actors = actors; _http = http; _breaker = breaker; _logger = logger; + _ledger = ledger; + } + + sealed class Attempt + { + public DeliveryPayload Payload; + public LocalActor Signer; + public int? Status; + public string Reason; + public int? LatencyMs; } public JobKind Kind => JobKind.Deliver; @@ -96,28 +112,97 @@ namespace PrivaPub.Federation.Outbox public int PerHostLimit => 2; public async Task Handle(Job job, CancellationToken token) + { + var attempt = new Attempt(); + JobOutcome outcome = default; + try + { + outcome = await Deliver(job, attempt, token); + return outcome; + } + finally + { + Record(job, attempt, outcome); + } + } + + void Record(Job job, Attempt attempt, JobOutcome outcome) + { + if (_ledger == default) + return; + var shape = attempt.Payload == default ? default : ActivityShape.Of(attempt.Payload.Body); + _ledger.Record(new InteractionEvent + { + Channel = Interactions.Out, + Host = job.Host, + Activity = shape?.Type, + Object = shape?.ObjectType, + Outcome = outcome?.Result switch + { + JobResult.Done => Interactions.Ok, + JobResult.Defer => Interactions.Deferred, + JobResult.Retry when job.Attempts >= MaxAttempts => Interactions.Dead, + JobResult.Retry => Interactions.Retry, + JobResult.Dead => Interactions.Dead, + _ => Interactions.Failed + }, + Reason = attempt.Reason, + Status = attempt.Status, + LatencyMs = attempt.LatencyMs, + WaitMs = (int)Math.Min(int.MaxValue, Math.Max(0, (DateTime.UtcNow - job.CreatedAt).TotalMilliseconds)), + Bytes = attempt.Payload?.Body == default ? default : Encoding.UTF8.GetByteCount(attempt.Payload.Body), + Attempt = job.Attempts, + Audience = shape?.Audience, + LocalKind = attempt.Signer switch + { + null => default, + { IsCircle: true } => default, + { Kind: LocalActorKind.Person } => "person", + { Kind: LocalActorKind.Group } => "group", + _ => "application" + }, + Signature = "cavage:rsa-sha256" + }); + } + + async Task Deliver(Job job, Attempt attempt, CancellationToken token) { var payload = JsonSerializer.Deserialize(job.Payload); + attempt.Payload = payload; if (payload == default || !Uri.TryCreate(payload.Inbox, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox)) + { + attempt.Reason = "not-deliverable"; return JobOutcome.Dead("not a deliverable inbox"); + } var unavailableUntil = await _breaker.UnavailableUntil(job.Host, token); if (unavailableUntil.HasValue) + { + attempt.Reason = "host-unavailable"; return JobOutcome.Defer(unavailableUntil.Value, "the host is unavailable"); + } var signer = await _actors.FindById(payload.SignerKind, payload.SignerId, token); + attempt.Signer = signer; if (signer == default) + { + attempt.Reason = "signer-gone"; 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); + var started = Stopwatch.GetTimestamp(); try { using var response = await _http.Send(request, token); + attempt.LatencyMs = (int)Stopwatch.GetElapsedTime(started).TotalMilliseconds; var status = (int)response.StatusCode; + attempt.Status = status; + attempt.Reason = response.IsSuccessStatusCode ? default : status.ToString(System.Globalization.CultureInfo.InvariantCulture); if (response.IsSuccessStatusCode) { await _breaker.Succeeded(job.Host, token); @@ -142,10 +227,13 @@ namespace PrivaPub.Federation.Outbox } catch (BlockedDestinationException ex) { + attempt.Reason = "private-address"; return JobOutcome.Dead(ex.Message); } catch (Exception ex) when (ex is HttpRequestException or TaskCanceledException && !token.IsCancellationRequested) { + attempt.LatencyMs = (int)Stopwatch.GetElapsedTime(started).TotalMilliseconds; + attempt.Reason = ex is TaskCanceledException ? "timeout" : "network"; await _breaker.Failed(job.Host, ex.GetType().Name, token); return JobOutcome.Retry(ex.Message); }