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.Models.Jobs; 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; public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger logger) { _actors = actors; _http = http; _breaker = breaker; _logger = logger; } 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 payload = JsonSerializer.Deserialize(job.Payload); if (payload == default || !Uri.TryCreate(payload.Inbox, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox)) return JobOutcome.Dead("not a deliverable inbox"); var unavailableUntil = await _breaker.UnavailableUntil(job.Host, token); if (unavailableUntil.HasValue) return JobOutcome.Defer(unavailableUntil.Value, "the host is unavailable"); var signer = await _actors.FindById(payload.SignerKind, payload.SignerId, token); if (signer == default) 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); try { using var response = await _http.Send(request, token); var status = (int)response.StatusCode; if (response.IsSuccessStatusCode) { await _breaker.Succeeded(job.Host, token); return JobOutcome.Done; } if (response.StatusCode == HttpStatusCode.TooManyRequests) { 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)), "429") : JobOutcome.Retry("429"); } 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) { return JobOutcome.Dead(ex.Message); } catch (Exception ex) when (ex is HttpRequestException or TaskCanceledException && !token.IsCancellationRequested) { 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; } }