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 inboxes, JsonObject activity, CancellationToken token); Task EnqueueToFollowers(LocalActor signer, JsonObject activity, CancellationToken token, IEnumerable extraInboxes = default); Task> 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 inboxes, JsonObject activity, CancellationToken token) { var body = activity.ToJsonString(); var activityId = activity["id"] is JsonValue id && id.TryGetValue(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 extraInboxes = default) { var inboxes = (await FollowerInboxes(signer, token)).Concat(extraInboxes ?? Enumerable.Empty()); await Enqueue(signer, inboxes, activity, token); } public async Task> 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 _logger; readonly IInteractionLedger _ledger; 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; public int Concurrency => 8; public int MaxAttempts => 16; 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); 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) { 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); } } static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b; } }