Files
SocialPub/PrivaPub/Federation/Outbox/DeliveryService.cs
T

244 lines
8.5 KiB
C#
Raw Normal View History

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;
2026-10-03 11:06:24 +02:00
using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Federation.Objects;
using PrivaPub.Models.Statistics;
using PrivaPub.Models.Jobs;
2026-10-03 11:06:24 +02:00
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;
2026-10-03 11:06:24 +02:00
readonly IInteractionLedger _ledger;
2026-10-03 11:06:24 +02:00
public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger<DeliveryJobHandler> logger,
IInteractionLedger ledger = default)
{
_actors = actors;
_http = http;
_breaker = breaker;
_logger = logger;
2026-10-03 11:06:24 +02:00
_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)
2026-10-03 11:06:24 +02:00
{
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);
2026-10-03 11:06:24 +02:00
attempt.Payload = payload;
if (payload == default || !Uri.TryCreate(payload.Inbox, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox))
2026-10-03 11:06:24 +02:00
{
attempt.Reason = "not-deliverable";
return JobOutcome.Dead("not a deliverable inbox");
2026-10-03 11:06:24 +02:00
}
var unavailableUntil = await _breaker.UnavailableUntil(job.Host, token);
if (unavailableUntil.HasValue)
2026-10-03 11:06:24 +02:00
{
attempt.Reason = "host-unavailable";
return JobOutcome.Defer(unavailableUntil.Value, "the host is unavailable");
2026-10-03 11:06:24 +02:00
}
var signer = await _actors.FindById(payload.SignerKind, payload.SignerId, token);
2026-10-03 11:06:24 +02:00
attempt.Signer = signer;
if (signer == default)
2026-10-03 11:06:24 +02:00
{
attempt.Reason = "signer-gone";
return JobOutcome.Dead("the signing actor no longer exists");
2026-10-03 11:06:24 +02:00
}
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);
2026-10-03 11:06:24 +02:00
var started = Stopwatch.GetTimestamp();
try
{
using var response = await _http.Send(request, token);
2026-10-03 11:06:24 +02:00
attempt.LatencyMs = (int)Stopwatch.GetElapsedTime(started).TotalMilliseconds;
var status = (int)response.StatusCode;
2026-10-03 11:06:24 +02:00
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}");
}
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)
{
2026-10-03 11:06:24 +02:00
attempt.Reason = "private-address";
return JobOutcome.Dead(ex.Message);
}
catch (Exception ex) when (ex is HttpRequestException or TaskCanceledException && !token.IsCancellationRequested)
{
2026-10-03 11:06:24 +02:00
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);
}
}
static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b;
}
}