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 Handle(Job job, CancellationToken token) { var claim = JsonSerializer.Deserialize(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(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().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> Listed(string url, CancellationToken token) { using var scope = HttpScope.For("followers-sync"); var items = new List(); 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 }; } }