Relays: PrivaPub reads from the relays its configuration names
Federation:Relays names relays by their actor (or inbox) address; the instance actor follows Public at each, as Mastodon subscribes, a minute after start and every six hours (asked again a day after no answer or a refusal, undone when a relay is no longer named). What an accepted relay passes on comes to the federated timeline and nobody's home: a public post it forwards (Activity-Relay), read again from its origin like any forwarded post, and a post it announces (aode-relay), kept as its author's and never as the relay's boost. Nothing of a persona's is sent to a relay; sending public posts there waits for the owner. The pasture gains both relays (peers/relay.sh, peers/aoderelay.sh) and scenarios/relay.sh, 13 checks. The village backlog is clean: 2574 checks pass, one known gap (Misskey's). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
This commit is contained in:
1 parent
1c7d1ece9c
commit
f6fc0c535c
22 files changed
+542
-24
No files matched your search
@@ -5,9 +5,10 @@ using PrivaPub.Models.Post;
|
||||
|
||||
namespace PrivaPub.Federation.Inbox
|
||||
{
|
||||
// Raw: the activity as it came, when its own author sent it (ReplyRelay passes it on as it is)
|
||||
// Raw: the activity as it came, when its own author sent it (ReplyRelay passes it on as it is); ForwardedBy: the server
|
||||
// that passed it on instead (a relay, a thread's server)
|
||||
public sealed record Arrival(string ActivityId, string ActivityType, string ActorURI, string Inbox, string KeyId, string Algorithm,
|
||||
IReadOnlyList<string> SignedHeaders, DateTime ReceivedAt, string Context = default, string Raw = default)
|
||||
IReadOnlyList<string> SignedHeaders, DateTime ReceivedAt, string Context = default, string Raw = default, string ForwardedBy = default)
|
||||
{
|
||||
static readonly AsyncLocal<Arrival> current = new();
|
||||
|
||||
|
||||
@@ -19,10 +19,12 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
readonly IQuoteService _quotes;
|
||||
readonly IInteractionApprovals _approvals;
|
||||
readonly IParticipations _participations;
|
||||
readonly Relays.IRelays _relays;
|
||||
|
||||
public AcceptHandler(DbEntities dbEntities, ILocalActorService localActors, IQuoteService quotes, IInteractionApprovals approvals = default,
|
||||
IParticipations participations = default)
|
||||
IParticipations participations = default, Relays.IRelays relays = default)
|
||||
{
|
||||
_relays = relays;
|
||||
_participations = participations;
|
||||
_dbEntities = dbEntities;
|
||||
_localActors = localActors;
|
||||
@@ -49,6 +51,11 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
Arrival.Accept("participation-answer");
|
||||
return;
|
||||
}
|
||||
if (_relays != default && await _relays.Answered(activity, actor, accepted: Type == "Accept", token))
|
||||
{
|
||||
Arrival.Accept("relay-answer");
|
||||
return;
|
||||
}
|
||||
var following = await FindFollowing(activity["object"], actor, _dbEntities, _localActors, token);
|
||||
if (following == default)
|
||||
{
|
||||
@@ -88,8 +95,8 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
public class RejectHandler : AcceptHandler
|
||||
{
|
||||
public RejectHandler(DbEntities dbEntities, ILocalActorService localActors, IQuoteService quotes, IInteractionApprovals approvals = default,
|
||||
IParticipations participations = default)
|
||||
: base(dbEntities, localActors, quotes, approvals, participations)
|
||||
IParticipations participations = default, Relays.IRelays relays = default)
|
||||
: base(dbEntities, localActors, quotes, approvals, participations, relays)
|
||||
{
|
||||
}
|
||||
|
||||
|
||||
@@ -38,9 +38,12 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
// handler is one of them; tests set it
|
||||
public IEnumerable<IActivityHandler> Relays { get; set; }
|
||||
|
||||
readonly PrivaPub.Federation.Relays.IRelays _relayService;
|
||||
|
||||
public AnnounceHandler(DbEntities dbEntities, ILocalActorService localActors, IRemotePosts remotePosts, IFanout fanout, IRemoteActorService remoteActors,
|
||||
IObjectRecords records, IQuoteService quotes, IServiceProvider services = default)
|
||||
IObjectRecords records, IQuoteService quotes, IServiceProvider services = default, PrivaPub.Federation.Relays.IRelays relays = default)
|
||||
{
|
||||
_relayService = relays;
|
||||
_services = services;
|
||||
_records = records;
|
||||
_quotes = quotes;
|
||||
@@ -74,6 +77,11 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
return;
|
||||
}
|
||||
|
||||
if (_relayService != default && await _relayService.Passes(actor.ActorURI, token))
|
||||
{
|
||||
await FromRelay(objectUri, token);
|
||||
return;
|
||||
}
|
||||
var isLocal = objectUri.StartsWith(_localActors.BaseAddress + "/", StringComparison.OrdinalIgnoreCase);
|
||||
var original = await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri && !p.DeletedAt.HasValue && p.ReblogOfPostId == null).ExecuteFirstAsync(token);
|
||||
var followed = await _dbEntities.Followings.Match(f => f.TargetActorURI == actor.ActorURI && f.State == FollowState.Accepted).ExecuteAnyAsync(token);
|
||||
@@ -139,6 +147,27 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
await _fanout.Distribute(reblog, token);
|
||||
}
|
||||
|
||||
// a post a relay we subscribe to announces: kept as its author's, for the federated timeline, never as the relay's
|
||||
// boost (Mastodon unwraps them the same way)
|
||||
async Task FromRelay(string objectUri, CancellationToken token)
|
||||
{
|
||||
var held = await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri).ExecuteFirstAsync(token);
|
||||
var post = held ?? await _remotePosts.StoreContext(objectUri, 0, token);
|
||||
if (post == default)
|
||||
{
|
||||
Arrival.Drop("fetch-failed");
|
||||
return;
|
||||
}
|
||||
Arrival.About(post.ObjectType ?? "Note", post.Visibility, post.CreationDate);
|
||||
if (held != default)
|
||||
{
|
||||
Arrival.Drop("duplicate");
|
||||
return;
|
||||
}
|
||||
Arrival.Accept("relayed");
|
||||
await _fanout.Distribute(post, token);
|
||||
}
|
||||
|
||||
async Task GroupActivity(JsonNode inner, ForeignAvatar group, CancellationToken token)
|
||||
{
|
||||
Arrival.About(Value(inner, "type"));
|
||||
|
||||
@@ -41,11 +41,13 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
readonly ILinkPreviews _previews;
|
||||
readonly IQuoteService _quotes;
|
||||
readonly IInteractionApprovals _approvals;
|
||||
readonly Relays.IRelays _relays;
|
||||
|
||||
public CreateHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery,
|
||||
IDomainBlocks domainBlocks, IFanout fanout, IRemotePosts remotePosts, IGroupDistributor groups, IObjectRecords records, IPollService polls, ILinkPreviews previews, IQuoteService quotes,
|
||||
IInteractionApprovals approvals = default)
|
||||
IInteractionApprovals approvals = default, Relays.IRelays relays = default)
|
||||
{
|
||||
_relays = relays;
|
||||
_approvals = approvals;
|
||||
_quotes = quotes;
|
||||
_previews = previews;
|
||||
@@ -183,8 +185,12 @@ namespace PrivaPub.Federation.Inbox.Handlers
|
||||
var inFollowedGroup = !followed && visibility is PostVisibility.Public or PostVisibility.Unlisted && note.Audience != default
|
||||
&& Origin.Same(note.Audience, author.ActorURI) && Origin.Same(note.Id, note.Audience)
|
||||
&& await _dbEntities.Followings.Match(f => f.TargetActorURI == note.Audience && f.State == FollowState.Accepted).ExecuteAnyAsync(token);
|
||||
// a public post a relay we subscribe to passes on, for the federated timeline (read again from its origin, as any
|
||||
// forwarded post)
|
||||
var relayed = !followed && visibility == PostVisibility.Public && _relays != default
|
||||
&& await _relays.Passes(Arrival.Current?.ForwardedBy, token);
|
||||
if (visibility == PostVisibility.Direct ? persons.Count == 0
|
||||
: group == default && persons.Count == 0 && !repliesToLocal && !followed && !repliesToFollowed && !inFollowedGroup)
|
||||
: group == default && persons.Count == 0 && !repliesToLocal && !followed && !repliesToFollowed && !inFollowedGroup && !relayed)
|
||||
{
|
||||
Arrival.Drop("not-addressed");
|
||||
return;
|
||||
|
||||
@@ -78,7 +78,7 @@ namespace PrivaPub.Federation.Inbox
|
||||
|
||||
var arrival = new Arrival(Id(activity), type, actor.ActorURI, payload.Inbox, payload.KeyId, payload.Algorithm,
|
||||
payload.SignedHeaders ?? Array.Empty<string>(), payload.ReceivedAt ?? job.CreatedAt, activity["@context"]?.ToJsonString(),
|
||||
payload.ForwardedBy == default ? payload.Activity : default);
|
||||
payload.ForwardedBy == default ? payload.Activity : default, payload.ForwardedBy);
|
||||
Arrival.Current = arrival;
|
||||
try
|
||||
{
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
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<bool> Answered(JsonNode activity, ForeignAvatar actor, bool accepted, CancellationToken token);
|
||||
Task<bool> Passes(string actorUri, 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. Nothing of a persona's is sent to a relay.
|
||||
public class Relays : IRelays
|
||||
{
|
||||
static readonly TimeSpan AskAgain = TimeSpan.FromDays(1);
|
||||
|
||||
readonly IOptions<FederationOptions> _options;
|
||||
readonly IRemoteActorService _remoteActors;
|
||||
readonly ILocalActorService _localActors;
|
||||
readonly IDeliveryService _delivery;
|
||||
volatile HashSet<string> _accepted;
|
||||
|
||||
public Relays(IOptions<FederationOptions> 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<RelaySubscription>().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<RelaySubscription>(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<bool> 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<RelaySubscription>().Match(s => s.FollowActivityURI == followId && s.ActorURI == actor.ActorURI)
|
||||
.ExecuteFirstAsync(token);
|
||||
if (subscription == default)
|
||||
return false;
|
||||
await DB.Default.Update<RelaySubscription>().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<bool> Passes(string actorUri, CancellationToken token)
|
||||
{
|
||||
if (string.IsNullOrEmpty(actorUri))
|
||||
return false;
|
||||
var accepted = _accepted ??= (await DB.Default.Find<RelaySubscription>().Match(s => s.State == RelayState.Accepted).ExecuteAsync(token))
|
||||
.Select(s => s.ActorURI).ToHashSet(StringComparer.Ordinal);
|
||||
return accepted.Contains(actorUri);
|
||||
}
|
||||
|
||||
// the relay's actor, from its address or, for a relay named by its inbox, from the /actor beside it
|
||||
async Task<ForeignAvatar> 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<RelaySubscriber> 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<IRelays>().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));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -122,6 +122,7 @@ namespace PrivaPub.Infrastructure.Data
|
||||
(() => DB.Default.Index<AccountDomainBlock>().Key(b => b.AvatarId, KeyType.Ascending).Key(b => b.Domain, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token), "domain block"),
|
||||
(() => DB.Default.Index<Bookmark>().Key(b => b.AvatarId, KeyType.Ascending).Key(b => b.PostId, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token), "bookmark"),
|
||||
(() => DB.Default.Index<Pin>().Key(p => p.AvatarId, KeyType.Ascending).Key(p => p.PostId, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token), "pin"),
|
||||
(() => DB.Default.Index<Models.Federation.RelaySubscription>().Key(r => r.Configured, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token), "relay"),
|
||||
(() => DB.Default.Index<RemoteFeatured>().Key(f => f.ActorURI, KeyType.Ascending).Key(f => f.ObjectURI, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token), "remote-featured")
|
||||
})
|
||||
await pair.Item1();
|
||||
|
||||
@@ -11,5 +11,8 @@ namespace PrivaPub.Infrastructure.Http
|
||||
// At two the inbox kept up with about 110 activities a second (tools/pasture/load.sh)
|
||||
public int InboxConcurrency { get; set; } = 8;
|
||||
public int DeliveryConcurrency { get; set; } = 8;
|
||||
// relays the instance actor subscribes to, by their actor's (or inbox's) address: their public posts come to the
|
||||
// federated timeline, and nothing is sent to them but the subscription
|
||||
public List<string> Relays { get; set; } = new();
|
||||
}
|
||||
}
|
||||
@@ -121,7 +121,9 @@ namespace PrivaPub.Middleware
|
||||
.AddSingleton<IHostCircuitBreaker, HostCircuitBreaker>()
|
||||
.AddSingleton<IJobHandler, DeliveryJobHandler>()
|
||||
.AddHostedService<JobWorker>()
|
||||
.AddHostedService<Domain.Social.FollowResender>();
|
||||
.AddHostedService<Domain.Social.FollowResender>()
|
||||
.AddSingleton<Federation.Relays.IRelays, Federation.Relays.Relays>()
|
||||
.AddHostedService<Federation.Relays.RelaySubscriber>();
|
||||
}
|
||||
public static IServiceCollection PrivaPubStatisticsConfiguration(this IServiceCollection service, IConfiguration configuration)
|
||||
{
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
using MongoDB.Entities;
|
||||
|
||||
namespace PrivaPub.Models.Federation
|
||||
{
|
||||
// a relay the instance actor subscribes to (Federation:Relays): it passes on the public posts its other subscribers
|
||||
// send it. PrivaPub only reads from it, and sends it nothing of its personas.
|
||||
public class RelaySubscription : Entity
|
||||
{
|
||||
public string Configured { get; set; }//as the configuration names it
|
||||
public string ActorURI { get; set; }
|
||||
public string InboxURL { get; set; }
|
||||
public string FollowActivityURI { get; set; }
|
||||
public RelayState State { get; set; } = RelayState.Pending;
|
||||
public DateTime RequestedAt { get; set; } = DateTime.UtcNow;
|
||||
public DateTime? AnsweredAt { get; set; }
|
||||
}
|
||||
|
||||
public enum RelayState
|
||||
{
|
||||
Pending,
|
||||
Accepted,
|
||||
Rejected
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user