Files
SocialPub/PrivaPub/Federation/Outbox/DeliveryService.cs
T
thepraandClaude Opus 5.5 b9286c2c6e Deliveries refused with 409 or 422 get three tries
Mastodon answers 422 when two first contacts from one actor race to
create its account (ActiveRecord::RecordInvalid on the unique uri), and
409 while another worker holds its lock. A persona that followed two
Mastodon accounts at once had one Follow refused that way; the job died
on its first attempt and the persona waited on "requested" forever.
Both answers are now retried twice on the usual backoff before they
count as refusals.

Found by the town (a village of 23 accounts on seven servers).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
2026-10-04 23:42:23 +02:00

252 lines
9.1 KiB
C#

using MongoDB.Entities;
using PrivaPub.Models.Federation;
using PrivaPub.StaticServices;
using System.Net;
using System.Net.Http.Headers;
using System.Text;
using System.Text.Json.Nodes;
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
{
public interface IDeliveryService
{
Task Enqueue(LocalActor signer, IEnumerable<string> inboxes, JsonObject activity, CancellationToken token);
Task EnqueueToFollowers(LocalActor signer, JsonObject activity, CancellationToken token, IEnumerable<string> extraInboxes = default);
Task<IReadOnlyList<string>> 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, IJobQueue queue)
{
_dbEntities = dbEntities;
_queue = queue;
}
public async Task Enqueue(LocalActor signer, IEnumerable<string> inboxes, JsonObject activity, CancellationToken token)
{
var body = activity.ToJsonString();
var activityId = activity["id"] is JsonValue id && id.TryGetValue<string>(out var text) ? text : default;
var jobs = inboxes
.Where(i => !string.IsNullOrEmpty(i) && !i.StartsWith(signer.BaseAddress + "/", StringComparison.OrdinalIgnoreCase))
.Distinct(StringComparer.Ordinal)
.Select(inbox => Uri.TryCreate(inbox, UriKind.Absolute, out var uri) ? (inbox, uri) : default)
.Where(target => target.uri != default)
.Select(target => new Job
{
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 (jobs.Count > 0)
await _queue.EnqueueMany(jobs, token);
}
public async Task EnqueueToFollowers(LocalActor signer, JsonObject activity, CancellationToken token, IEnumerable<string> extraInboxes = default)
{
var inboxes = (await FollowerInboxes(signer, token)).Concat(extraInboxes ?? Enumerable.Empty<string>());
await Enqueue(signer, inboxes, activity, token);
}
public async Task<IReadOnlyList<string>> FollowerInboxes(LocalActor actor, CancellationToken token)
{
var followers = await _dbEntities.Followers
.Match(f => f.LocalActorId == actor.Id && f.LocalActorKind == actor.Kind && f.IsAccepted)
.ExecuteAsync(token);
return followers.Select(f => string.IsNullOrEmpty(f.SharedInboxURL) ? f.InboxURL : f.SharedInboxURL)
.Distinct(StringComparer.Ordinal)
.ToList();
}
}
public class DeliveryJobHandler : IJobHandler
{
readonly ILocalActorService _actors;
readonly IFederationHttp _http;
readonly IHostCircuitBreaker _breaker;
readonly ILogger<DeliveryJobHandler> _logger;
readonly IInteractionLedger _ledger;
public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger<DeliveryJobHandler> 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;
public int Concurrency => 8;
public int MaxAttempts => 16;
public int PerHostLimit => 2;
public async Task<JobOutcome> 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<JobOutcome> Deliver(Job job, Attempt attempt, CancellationToken token)
{
var payload = JsonSerializer.Deserialize<DeliveryPayload>(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);
return JobOutcome.Done;
}
if (response.StatusCode == HttpStatusCode.TooManyRequests
|| response.StatusCode == HttpStatusCode.ServiceUnavailable && response.Headers.RetryAfter != default)
{
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)), $"{status}")
: JobOutcome.Retry($"{status}");
}
// Mastodon answers 422 when two first contacts from one actor race to create its account (RecordInvalid on the
// unique uri) and 409 while another worker holds its lock: a Follow refused that way was lost for good, and the
// persona waited on "requested" forever (found by the town). Both get a few tries before they count as refusals.
if (status is 409 or 422 && job.Attempts < TransientRefusalAttempts)
return JobOutcome.Retry($"{status} {response.ReasonPhrase}");
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 (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);
}
}
const int TransientRefusalAttempts = 3;
static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b;
}
}