What a restore lost of either side is not the other server's to mend: for RestoreRecord.FollowersGrace (14 days) deliveries carry no Collection-Synchronization header, and a follow only the remote remembers is adopted rather than undone. Without a restore nothing changes; the interop sweep (gts, mastodon, misskey, sharkey, akkoma, pixelfed, smithereen, followsync, pins, moves) passed, Smithereen's poll check on a second run. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
289 lines
12 KiB
C#
289 lines
12 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
|
|
{
|
|
// again: a delivery made once more on purpose, which an earlier one of the same activity does not stop; deliveries
|
|
// with the same again are made once
|
|
Task Enqueue(LocalActor signer, IEnumerable<string> inboxes, JsonObject activity, CancellationToken token, string again = default);
|
|
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, string again = default)
|
|
{
|
|
var body = activity.ToJsonString();
|
|
var activityId = activity["id"] is JsonValue id && id.TryGetValue<string>(out var text) ? text : default;
|
|
var targets = inboxes
|
|
.Where(i => !string.IsNullOrEmpty(i) && !i.StartsWith(signer.BaseAddress + "/", StringComparison.OrdinalIgnoreCase))
|
|
.Distinct(StringComparer.Ordinal)
|
|
.ToList();
|
|
// what goes to a relay carries the signer's proof (FEP-8b32): the relay passes it on as it came, signed with its own
|
|
// key, and Mastodon takes it on the proof's strength. Nothing else does: a server that knew the persona before it
|
|
// had its key (Mitra) refuses a proof by a key it does not hold, and does not read the actor again
|
|
var relayed = new HashSet<string>(StringComparer.Ordinal);
|
|
if (!string.IsNullOrEmpty(signer.SigningKey) && activity["actor"] is JsonValue actor && actor.TryGetValue<string>(out var actorUri) && actorUri == signer.Uri)
|
|
relayed = (await DB.Default.Find<RelaySubscription, string>().Match(s => targets.Contains(s.InboxURL)).Project(s => s.InboxURL).ExecuteAsync(token))
|
|
.ToHashSet(StringComparer.Ordinal);
|
|
var provenBody = default(string);
|
|
if (relayed.Count > 0)
|
|
{
|
|
var proven = activity.DeepClone().AsObject();
|
|
proven.Remove("proof");
|
|
proven["proof"] = IntegrityProofs.Create(proven, signer.AssertionKeyId, signer.SigningKey, DateTime.UtcNow);
|
|
provenBody = proven.ToJsonString();
|
|
}
|
|
var jobs = targets
|
|
.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 : again == default ? $"{activityId}|{target.inbox}" : $"{activityId}|{target.inbox}|{again}",
|
|
Payload = JsonSerializer.Serialize(new DeliveryPayload(signer.Id, signer.Kind, target.inbox, relayed.Contains(target.inbox) ? provenBody : 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;
|
|
|
|
readonly int _concurrency;
|
|
|
|
public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger<DeliveryJobHandler> logger,
|
|
IInteractionLedger ledger = default, Microsoft.Extensions.Options.IOptions<Infrastructure.Http.FederationOptions> federation = default)
|
|
{
|
|
_concurrency = Math.Max(1, federation?.Value.DeliveryConcurrency ?? 8);
|
|
_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 => _concurrency;
|
|
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",
|
|
{ Kind: LocalActorKind.Reporter } => "reporter",
|
|
_ => "application"
|
|
},
|
|
Signature = "cavage:rsa-sha256"
|
|
});
|
|
}
|
|
|
|
// FEP-8fcf: what a delivery to the persona's followers tells the receiving server of its followers there; nothing for
|
|
// a while after a restore, whose losses are not the other server's to mend
|
|
static async Task<(string, string)[]> Synchronization(LocalActor signer, string body, Uri inbox, CancellationToken token)
|
|
{
|
|
if (!FollowersSynchronization.Synchronizes(signer) || await Infrastructure.Backup.RestoreRecord.InFollowersGrace(token) || JsonNode.Parse(body) is not JsonObject activity || !FollowersSynchronization.ForFollowers(activity, signer))
|
|
return [];
|
|
var there = await FollowersSynchronization.On(signer, FollowersSynchronization.Origin(inbox.AbsoluteUri), token);
|
|
return [(FollowersSynchronization.Header, FollowersSynchronization.HeaderValue(signer, FollowersSynchronization.Digest(there)))];
|
|
}
|
|
|
|
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, await Synchronization(signer, payload.Body, inbox, token));
|
|
|
|
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}");
|
|
}
|
|
// a gateway's 502, 503 or 504 says the server behind it is down, which counts against the host; any other 5xx
|
|
// says the server is up and failed on this activity (Smithereen's 500 on a document its JSON-LD reader
|
|
// refused), retried on its own while the host's other deliveries go on
|
|
if (status is 502 or 503 or 504)
|
|
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;
|
|
}
|
|
}
|