A sender's digest of its followers here is compared and mended
FEP-8fcf, received. When a delivery's Collection-Synchronization header digests the sender's followers on PrivaPub
otherwise than the personas following it, a job reads the list the header names (on the sender's origin, signed by the
instance actor): a follow the list leaves out ends, only when the list is the one the digest describes; a request it
lists is taken as accepted; a persona it lists that follows nothing there sends Undo{Follow}, as Mastodon does. Each
claiming delivery is compared once.
Checked live (scenarios/followsync.sh, now 15 checks): PrivaPub ends a follow Mastodon lost and undoes one only
Mastodon remembered.
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
bd18de35a6
commit
e0682fa868
11 files changed
+291
-12
No files matched your search
@@ -55,6 +55,23 @@ namespace PrivaPub.Federation.Actors
|
||||
public static string HeaderValue(LocalActor actor, string digest) =>
|
||||
$"collectionId=\"{actor.Followers}\", url=\"{RollCall(actor)}\", digest=\"{digest}\"";
|
||||
|
||||
// what a Collection-Synchronization header claims, written as a Signature header's parameters; default when it
|
||||
// names no collection, list or digest
|
||||
public static SynchronizationClaim Claim(string actorUri, string header)
|
||||
{
|
||||
var values = new Dictionary<string, string>(StringComparer.Ordinal);
|
||||
foreach (var part in (header ?? "").Split(',', StringSplitOptions.TrimEntries | StringSplitOptions.RemoveEmptyEntries))
|
||||
{
|
||||
var equals = part.IndexOf('=');
|
||||
if (equals > 0)
|
||||
values[part[..equals].Trim()] = part[(equals + 1)..].Trim().Trim('"');
|
||||
}
|
||||
return values.TryGetValue("collectionId", out var collection) && values.TryGetValue("url", out var url) && values.TryGetValue("digest", out var digest)
|
||||
&& digest.Length == 64 && digest.All(Uri.IsHexDigit)
|
||||
? new SynchronizationClaim(actorUri, collection, url, digest.ToLowerInvariant())
|
||||
: default;
|
||||
}
|
||||
|
||||
// whether the activity, or the object it carries, is addressed to the actor's followers
|
||||
public static bool ForFollowers(JsonObject activity, LocalActor actor) =>
|
||||
Addressed(activity, actor.Followers) || activity["object"] is JsonObject inner && Addressed(inner, actor.Followers);
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
using MongoDB.Entities;
|
||||
|
||||
using PrivaPub.Domain.Social;
|
||||
using PrivaPub.Federation.Objects;
|
||||
using PrivaPub.Federation.Outbox;
|
||||
using PrivaPub.Federation.Rendering;
|
||||
using PrivaPub.Infrastructure.Http;
|
||||
using PrivaPub.Infrastructure.Jobs;
|
||||
using PrivaPub.Models.Federation;
|
||||
using PrivaPub.Models.Jobs;
|
||||
using PrivaPub.Models.Social;
|
||||
using PrivaPub.StaticServices;
|
||||
|
||||
using System.Text.Json;
|
||||
using System.Text.Json.Nodes;
|
||||
|
||||
namespace PrivaPub.Federation.Actors
|
||||
{
|
||||
public sealed record SynchronizationClaim(string Actor, string CollectionId, string Url, string Digest);
|
||||
|
||||
// FEP-8fcf, received: an account elsewhere tells, with its delivery, a digest of its followers here. When the personas
|
||||
// following it here digest otherwise, its partial list (read signed by the instance actor) says who they are: a follow
|
||||
// only PrivaPub remembers ends, a request it answered without our knowing is accepted, and a follow only it remembers
|
||||
// is undone by the persona, as Mastodon does. A follow ends only when the list is the one the digest describes.
|
||||
public class FollowingSynchronization : IJobHandler
|
||||
{
|
||||
const int MaxPages = 10;
|
||||
|
||||
readonly DbEntities _dbEntities;
|
||||
readonly ILocalActorService _localActors;
|
||||
readonly IRemoteActorService _remoteActors;
|
||||
readonly IFollowService _follows;
|
||||
readonly IDeliveryService _delivery;
|
||||
|
||||
public FollowingSynchronization(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IFollowService follows,
|
||||
IDeliveryService delivery)
|
||||
{
|
||||
_dbEntities = dbEntities;
|
||||
_localActors = localActors;
|
||||
_remoteActors = remoteActors;
|
||||
_follows = follows;
|
||||
_delivery = delivery;
|
||||
}
|
||||
|
||||
public JobKind Kind => JobKind.SynchronizeFollowing;
|
||||
public int Concurrency => 1;
|
||||
public int MaxAttempts => 3;
|
||||
public int PerHostLimit => 1;
|
||||
|
||||
// once for each delivery that claims it, as Mastodon compares each: cheap while the digests agree
|
||||
public static string DedupeKey(SynchronizationClaim claim, string activityId) =>
|
||||
activityId == default ? default : $"followsync|{activityId}|{claim.Digest}";
|
||||
|
||||
public async Task<JobOutcome> Handle(Job job, CancellationToken token)
|
||||
{
|
||||
var claim = JsonSerializer.Deserialize<SynchronizationClaim>(job.Payload);
|
||||
var actor = claim == default ? default : await _dbEntities.ForeignAvatars.Match(a => a.ActorURI == claim.Actor && !a.DeletionAt.HasValue).ExecuteFirstAsync(token);
|
||||
if (actor == default || string.IsNullOrEmpty(actor.FollowersURL) || claim.CollectionId != actor.FollowersURL || !Origin.Same(claim.Url, actor.ActorURI))
|
||||
return JobOutcome.Done;
|
||||
|
||||
var followings = await _dbEntities.Followings.Match(f => f.TargetActorURI == actor.ActorURI && !f.TargetIsLocal).ExecuteAsync(token);
|
||||
var personas = new Dictionary<string, (LocalActor Persona, Following Following)>(StringComparer.Ordinal);
|
||||
foreach (var following in followings)
|
||||
if (await _localActors.FindById(LocalActorKind.Person, following.AvatarId, token) is { } persona)
|
||||
personas[persona.Uri] = (persona, following);
|
||||
var accepted = personas.Values.Where(p => p.Following.State == FollowState.Accepted).Select(p => p.Persona.Uri).ToList();
|
||||
if (FollowersSynchronization.Digest(accepted) == claim.Digest)
|
||||
return JobOutcome.Done;
|
||||
|
||||
var listed = await Listed(claim.Url, token);
|
||||
if (listed == default)
|
||||
return JobOutcome.Retry("the partial followers collection could not be read");
|
||||
|
||||
var complete = FollowersSynchronization.Digest(listed) == claim.Digest;
|
||||
var named = listed.ToHashSet(StringComparer.Ordinal);
|
||||
if (complete)
|
||||
foreach (var uri in accepted.Where(uri => !named.Contains(uri)))
|
||||
await _follows.UnfollowAs(personas[uri].Persona, actor.ActorURI, token);
|
||||
|
||||
foreach (var uri in named)
|
||||
{
|
||||
if (personas.TryGetValue(uri, out var known))
|
||||
{
|
||||
if (known.Following.State == FollowState.Requested)
|
||||
await DB.Default.Update<Following>().MatchID(known.Following.ID).Modify(f => f.State, FollowState.Accepted).ExecuteAsync(token);
|
||||
continue;
|
||||
}
|
||||
if (await _localActors.FindByUri(uri, token) is not { Kind: LocalActorKind.Person } stranger)
|
||||
continue;
|
||||
// (it has no id for a follow PrivaPub never made; it is undone by actor and object)
|
||||
var undo = new JsonObject
|
||||
{
|
||||
["@context"] = ActivityPubRenderer.ActivityStreams,
|
||||
["id"] = stranger.ActivityUri($"undo-follow-sync-{Guid.NewGuid():N}"),
|
||||
["type"] = "Undo",
|
||||
["actor"] = stranger.Uri,
|
||||
["object"] = new JsonObject { ["type"] = "Follow", ["actor"] = stranger.Uri, ["object"] = actor.ActorURI }
|
||||
};
|
||||
await _delivery.Enqueue(stranger, new[] { string.IsNullOrEmpty(actor.InboxURL) ? actor.SharedInboxURL : actor.InboxURL }, undo, token);
|
||||
}
|
||||
return JobOutcome.Done;
|
||||
}
|
||||
|
||||
// the ids the collection lists, its pages followed on its own server; default when it cannot be read
|
||||
async Task<List<string>> Listed(string url, CancellationToken token)
|
||||
{
|
||||
using var scope = HttpScope.For("followers-sync");
|
||||
var items = new List<string>();
|
||||
var next = url;
|
||||
for (var page = 0; page < MaxPages && next != default; page++)
|
||||
{
|
||||
using var fetched = await _remoteActors.FetchObject(next, token);
|
||||
if (fetched == default || fetched.Root.ValueKind != JsonValueKind.Object)
|
||||
return page == 0 ? default : items;
|
||||
var root = fetched.Root;
|
||||
if (page == 0 && !Has(root, "orderedItems") && !Has(root, "items") && Link(root, "first") is { } first)
|
||||
{
|
||||
next = Origin.Same(first, url) ? first : default;
|
||||
continue;
|
||||
}
|
||||
foreach (var name in new[] { "orderedItems", "items" })
|
||||
if (root.TryGetProperty(name, out var list) && list.ValueKind == JsonValueKind.Array)
|
||||
items.AddRange(list.EnumerateArray().Select(Id).Where(id => id != default));
|
||||
next = Link(root, "next") is { } after && Origin.Same(after, url) ? after : default;
|
||||
}
|
||||
return items;
|
||||
}
|
||||
|
||||
static bool Has(JsonElement node, string name) => node.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.Array;
|
||||
|
||||
static string Link(JsonElement node, string name) => node.TryGetProperty(name, out var value) ? Id(value) : default;
|
||||
|
||||
static string Id(JsonElement node) => node.ValueKind switch
|
||||
{
|
||||
JsonValueKind.String => node.GetString(),
|
||||
JsonValueKind.Object when node.TryGetProperty("id", out var id) && id.ValueKind == JsonValueKind.String => id.GetString(),
|
||||
_ => default
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -204,6 +204,9 @@ namespace PrivaPub.Federation.Inbox
|
||||
// once per follow, so a retry still counts as a copy
|
||||
if (!queued && dedupe != default && type == "Undo" && await FollowAgain(activity, actorUri, token) is { } follow)
|
||||
queued = await _queue.Enqueue(JobKind.ProcessInbox, queuedPayload, host, $"{dedupe}|{Digest(activity)}|{follow}", token);
|
||||
// FEP-8fcf: what the sender says of its followers here (FollowingSynchronization)
|
||||
if (queued && forwardedBy == default && FollowersSynchronization.Claim(actorUri, request.Headers[FollowersSynchronization.Header].ToString()) is { } claim)
|
||||
await _queue.Enqueue(JobKind.SynchronizeFollowing, JsonSerializer.Serialize(claim), host, FollowingSynchronization.DedupeKey(claim, activityId), token);
|
||||
_logger.LogInformation("Inbox {Recipient}: {Type} from {Actor} queued", recipient?.Handle ?? "shared", type, actorUri);
|
||||
return new(StatusCodes.Status202Accepted, Reason: !queued ? "duplicate" : forwardedBy != default ? "forwarded" : "queued");
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user