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; namespace PrivaPub.Services.Federation { 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 class DeliveryService : IDeliveryService { readonly DbEntities _dbEntities; public DeliveryService(DbEntities dbEntities) { _dbEntities = dbEntities; } public async Task Enqueue(LocalActor signer, IEnumerable inboxes, JsonObject activity, CancellationToken token) { var body = activity.ToJsonString(); var deliveries = inboxes .Where(i => !string.IsNullOrEmpty(i) && !i.StartsWith(signer.BaseAddress + "/", StringComparison.OrdinalIgnoreCase)) .Distinct(StringComparer.Ordinal) .Select(inbox => new Delivery { SignerId = signer.Id, SignerKind = signer.Kind, InboxURL = inbox, Body = body }) .ToList(); if (deliveries.Count > 0) await DB.Default.SaveAsync(deliveries, 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 DeliveryWorker : BackgroundService { const int MaxAttempts = 8; static readonly TimeSpan Poll = TimeSpan.FromSeconds(3); readonly IServiceProvider _services; readonly IHttpClientFactory _httpClientFactory; readonly ILogger _logger; public DeliveryWorker(IServiceProvider services, IHttpClientFactory httpClientFactory, ILogger logger) { _services = services; _httpClientFactory = httpClientFactory; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested) { try { await DeliverDue(stoppingToken); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } catch (Exception ex) { _logger.LogError(ex, "{Worker} pass failed", nameof(DeliveryWorker)); } await Task.Delay(Poll, stoppingToken); } } async Task DeliverDue(CancellationToken token) { using var scope = _services.CreateScope(); var dbEntities = scope.ServiceProvider.GetRequiredService(); var actors = scope.ServiceProvider.GetRequiredService(); var now = DateTime.UtcNow; var due = await dbEntities.Deliveries .Match(d => !d.DeliveredAt.HasValue && !d.AbandonedAt.HasValue && d.NextAttemptAt <= now) .Sort(d => d.NextAttemptAt, Order.Ascending) .Limit(20) .ExecuteAsync(token); foreach (var delivery in due) { var signer = await actors.FindById(delivery.SignerKind, delivery.SignerId, token); if (signer == default) { delivery.AbandonedAt = DateTime.UtcNow; delivery.LastError = "the signing actor no longer exists"; await DB.Default.SaveAsync(delivery, token); continue; } await Deliver(delivery, signer, token); await DB.Default.SaveAsync(delivery, token); } } async Task Deliver(Delivery delivery, LocalActor signer, CancellationToken token) { delivery.Attempts++; try { if (!Uri.TryCreate(delivery.InboxURL, UriKind.Absolute, out var inbox) || !RemoteActorService.IsFetchable(inbox)) { delivery.AbandonedAt = DateTime.UtcNow; delivery.LastError = "not a deliverable inbox"; return; } var body = Encoding.UTF8.GetBytes(delivery.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); using var response = await _httpClientFactory.CreateClient(RemoteActorService.HttpClientName).SendAsync(request, token); if (response.IsSuccessStatusCode) { delivery.DeliveredAt = DateTime.UtcNow; delivery.LastError = default; return; } delivery.LastError = $"{(int)response.StatusCode} {response.ReasonPhrase}"; if (response.StatusCode is HttpStatusCode.Gone or HttpStatusCode.NotFound or HttpStatusCode.BadRequest or HttpStatusCode.Forbidden) { delivery.AbandonedAt = DateTime.UtcNow; _logger.LogWarning("Delivery {Id} to {Inbox} abandoned: {Status}", delivery.ID, delivery.InboxURL, delivery.LastError); return; } } catch (Exception ex) when (ex is HttpRequestException or TaskCanceledException && !token.IsCancellationRequested) { delivery.LastError = ex.Message; } if (delivery.Attempts >= MaxAttempts) { delivery.AbandonedAt = DateTime.UtcNow; _logger.LogWarning("Delivery {Id} to {Inbox} abandoned after {Attempts} attempts: {Error}", delivery.ID, delivery.InboxURL, delivery.Attempts, delivery.LastError); return; } delivery.NextAttemptAt = DateTime.UtcNow.AddMinutes(Math.Pow(2, delivery.Attempts)); } } }