using Microsoft.Extensions.Options; using MongoDB.Entities; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Objects; using PrivaPub.Federation.Outbox; using PrivaPub.Federation.Rendering; using PrivaPub.Infrastructure.Http; using PrivaPub.Models.Federation; using PrivaPub.Models.User; using System.Text.Json.Nodes; using static PrivaPub.Federation.Objects.ActivityJson; namespace PrivaPub.Federation.Relays { public interface IRelays { Task Reconcile(CancellationToken token); Task Answered(JsonNode activity, ForeignAvatar actor, bool accepted, CancellationToken token); Task Passes(string actorUri, CancellationToken token); Task> Inboxes(CancellationToken token); } // Relays (Activity-Relay, aode-relay, pub-relay): the instance actor follows Public at each relay the configuration // names, as Mastodon subscribes, and takes what they pass on: public posts, forwarded as their authors sent them (read // again from their origin, as any forwarded post) or announced by the relay. A persona's public posts go to the relays // that accepted us (owner decision 2026-10-06), as Mastodon sends them; nothing else of a persona's does. public class Relays : IRelays { static readonly TimeSpan AskAgain = TimeSpan.FromDays(1); readonly IOptions _options; readonly IRemoteActorService _remoteActors; readonly ILocalActorService _localActors; readonly IDeliveryService _delivery; volatile HashSet _accepted; public Relays(IOptions options, IRemoteActorService remoteActors, ILocalActorService localActors, IDeliveryService delivery) { _options = options; _remoteActors = remoteActors; _localActors = localActors; _delivery = delivery; } // subscribes to the relays the configuration names (again a day after an unanswered or refused request), and // unsubscribes from those it no longer names public async Task Reconcile(CancellationToken token) { var configured = (_options.Value.Relays ?? new()).Where(r => !string.IsNullOrWhiteSpace(r)).Select(r => r.Trim()).Distinct().ToList(); var known = await DB.Default.Find().ExecuteAsync(token); if (configured.Count == 0 && known.Count == 0) return; var instance = await _localActors.GetInstanceActor(token); foreach (var gone in known.Where(k => !configured.Contains(k.Configured))) { if (gone.State != RelayState.Rejected && !string.IsNullOrEmpty(gone.InboxURL)) await _delivery.Enqueue(instance, new[] { gone.InboxURL }, Undo(instance, gone.FollowActivityURI), token); await DB.Default.DeleteAsync(gone.ID); } foreach (var address in configured) { var subscription = known.FirstOrDefault(k => k.Configured == address); if (subscription is { State: RelayState.Accepted } || subscription != default && DateTime.UtcNow - subscription.RequestedAt < AskAgain) continue; var relay = await Resolve(address, token); if (relay == default || string.IsNullOrEmpty(relay.InboxURL)) continue; subscription ??= new RelaySubscription { Configured = address }; subscription.ActorURI = relay.ActorURI; subscription.InboxURL = relay.InboxURL; subscription.FollowActivityURI = instance.ActivityUri($"relay-{Guid.NewGuid():N}"); subscription.State = RelayState.Pending; subscription.RequestedAt = DateTime.UtcNow; subscription.AnsweredAt = default; await DB.Default.SaveAsync(subscription, token); await _delivery.Enqueue(instance, new[] { relay.InboxURL }, Follow(instance, subscription.FollowActivityURI), token); } _accepted = default; } // the relay's Accept or Reject of the subscription, by the Follow it answers public async Task Answered(JsonNode activity, ForeignAvatar actor, bool accepted, CancellationToken token) { var followId = Id(activity["object"]); if (followId == default) return false; var subscription = await DB.Default.Find().Match(s => s.FollowActivityURI == followId && s.ActorURI == actor.ActorURI) .ExecuteFirstAsync(token); if (subscription == default) return false; await DB.Default.Update().MatchID(subscription.ID) .Modify(s => s.State, accepted ? RelayState.Accepted : RelayState.Rejected) .Modify(s => s.AnsweredAt, DateTime.UtcNow) .ExecuteAsync(token); _accepted = default; return true; } // whether the actor is a relay that accepted our subscription public async Task Passes(string actorUri, CancellationToken token) { if (string.IsNullOrEmpty(actorUri)) return false; var accepted = _accepted ??= (await DB.Default.Find().Match(s => s.State == RelayState.Accepted).ExecuteAsync(token)) .Select(s => s.ActorURI).ToHashSet(StringComparer.Ordinal); return accepted.Contains(actorUri); } // the inboxes of the relays that accepted our subscription, where a persona's public posts also go public async Task> Inboxes(CancellationToken token) => (await DB.Default.Find().Match(s => s.State == RelayState.Accepted).ExecuteAsync(token)) .Select(s => s.InboxURL).Where(i => !string.IsNullOrEmpty(i)).Distinct(StringComparer.Ordinal).ToList(); // the relay's actor, from its address or, for a relay named by its inbox, from the /actor beside it async Task Resolve(string address, CancellationToken token) { var actor = await _remoteActors.GetActor(address, refresh: false, token); if (actor != default || !address.EndsWith("/inbox", StringComparison.Ordinal)) return actor; return await _remoteActors.GetActor(address[..^"/inbox".Length] + "/actor", refresh: false, token); } static JsonObject Follow(LocalActor instance, string id) => new() { ["@context"] = ActivityPubRenderer.Context(), ["id"] = id, ["type"] = "Follow", ["actor"] = instance.Uri, ["object"] = Addressing.Public }; static JsonObject Undo(LocalActor instance, string followId) => new() { ["@context"] = ActivityPubRenderer.Context(), ["id"] = instance.ActivityUri($"relay-undo-{Guid.NewGuid():N}"), ["type"] = "Undo", ["actor"] = instance.Uri, ["object"] = new JsonObject { ["id"] = followId, ["type"] = "Follow", ["actor"] = instance.Uri, ["object"] = Addressing.Public } }; } // subscribes to the configured relays a minute after start, then every six hours public sealed class RelaySubscriber(IServiceScopeFactory scopes, ILogger logger) : BackgroundService { static readonly TimeSpan Every = TimeSpan.FromHours(6); protected override async Task ExecuteAsync(CancellationToken token) { await Task.Delay(TimeSpan.FromMinutes(1), token); using var timer = new PeriodicTimer(Every); do { try { using var scope = scopes.CreateScope(); await scope.ServiceProvider.GetRequiredService().Reconcile(token); } catch (Exception ex) when (ex is not OperationCanceledException) { logger.LogWarning(ex, "Relay subscriptions could not be brought up to date"); } } while (await timer.WaitForNextTickAsync(token)); } } }