using System.Text.Json; using System.Text.Json.Nodes; using MongoDB.Entities; using PrivaPub.Federation.Actors; using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Jobs; using PrivaPub.Models.Federation; using PrivaPub.Models.Jobs; using PrivaPub.Models.Post; using PrivaPub.StaticServices; using static PrivaPub.Federation.Objects.ActivityJson; using PostEntity = PrivaPub.Models.Post.Post; namespace PrivaPub.Federation.Inbox { public sealed record RepliesPayload(string PostId, int Depth); // What was said under a remote post on other servers, read when a persona opens its thread, so the thread shows more // than what happened to reach us. The thread's own collection comes first (FEP-7888's `context`, which Mastodon 4.5+ // serves with every reply at any depth); otherwise the post's `replies` (PeerTube's `comments`), and the replies' // own, two levels down. Read with the instance actor's signature, never a persona's, a few pages and at most // MaxItems posts a job; each one is fetched from its own server and stored like any post fetched for a thread // (RemotePosts.StoreContext: public and unlisted only, its author checked against its origin). public class RepliesJobHandler : IJobHandler { public const int MaxPages = 5; public const int MaxItems = 100; public const int MaxDepth = 2; public const int MaxBranches = 20; static readonly HashSet CollectionTypes = new(StringComparer.Ordinal) { "Collection", "OrderedCollection", "CollectionPage", "OrderedCollectionPage" }; readonly DbEntities _dbEntities; readonly IRemoteActorService _remoteActors; readonly IRemotePosts _remotePosts; readonly IJobQueue _queue; public RepliesJobHandler(DbEntities dbEntities, IRemoteActorService remoteActors, IRemotePosts remotePosts, IJobQueue queue) { _dbEntities = dbEntities; _remoteActors = remoteActors; _remotePosts = remotePosts; _queue = queue; } public JobKind Kind => JobKind.FetchReplies; public int Concurrency => 2; public int MaxAttempts => 2; public int PerHostLimit => 1; // a thread is read again at most once an hour public static Job For(PostEntity post, int depth) => new() { Kind = JobKind.FetchReplies, Payload = JsonSerializer.Serialize(new RepliesPayload(post.ID, depth)), Host = new Uri(post.ObjectURI).Host, DedupeKey = $"replies|{post.ID}|{DateTime.UtcNow:yyyyMMddHH}" }; public async Task Handle(Job job, CancellationToken token) { var payload = JsonSerializer.Deserialize(job.Payload); var post = await _dbEntities.Posts.MatchID(payload.PostId).ExecuteFirstAsync(token); if (post is not { IsFederatedCopy: true, Visibility: PostVisibility.Public or PostVisibility.Unlisted } || post.DeletedAt.HasValue || string.IsNullOrEmpty(post.ObjectURI)) return JobOutcome.Done; if (await Document(post, token) is not { } document) return JobOutcome.Done; if (await Collection(document["context"], token) is { } thread) { await StoreAll(thread, token); return JobOutcome.Done; } if (await Collection(document["replies"] ?? document["comments"], token) is not { } replies) return JobOutcome.Done; var answers = await StoreAll(replies, token); if (payload.Depth < MaxDepth) await _queue.EnqueueMany(answers.Where(a => a.IsFederatedCopy && a.AnsweringToPostId == post.ID).Take(MaxBranches) .Select(a => For(a, payload.Depth + 1)), token); return JobOutcome.Done; } // the post as it reached us, or as its server serves it now when we kept no copy async Task Document(PostEntity post, CancellationToken token) { var record = await DB.Default.Find().Match(r => r.ObjectURI == post.ObjectURI && r.Raw != null && !r.RawTruncated) .Project(r => new ObjectRecord { Raw = r.Raw }).ExecuteFirstAsync(token); if (record != default) return JsonNode.Parse(record.Raw) as JsonObject; using var scope = HttpScope.For("replies"); using var fetched = await _remoteActors.FetchObject(post.ObjectURI, token); return fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()) as JsonObject; } // a collection as its server serves it now: what a post carries is read again by its id, since an embedded first // page is as old as the post; a value that is no web address (Pleroma's `context` is a tag: URI), or no collection, // gives nothing async Task Collection(JsonNode node, CancellationToken token) { var uri = Id(node); if (uri != default && Uri.TryCreate(uri, UriKind.Absolute, out var parsed) && (parsed.Scheme == Uri.UriSchemeHttps || parsed.Scheme == Uri.UriSchemeHttp)) { using var scope = HttpScope.For("replies"); using var fetched = await _remoteActors.FetchObject(uri, token); node = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()); } return node is JsonObject collection && Value(collection, "type") is { } type && CollectionTypes.Contains(type) ? collection : default; } // every post the collection lists, stored when we lack it: its pages in order, at most MaxPages and MaxItems async Task> StoreAll(JsonObject collection, CancellationToken token) { var posts = new List(); var seen = 0; var page = collection["orderedItems"] is JsonArray || collection["items"] is JsonArray ? collection : await Page(collection["first"], token); for (var pages = 0; page != default && pages < MaxPages && seen < MaxItems; pages++) { foreach (var item in Items(page)) { if (++seen > MaxItems) break; // a context may list the activities that made the thread (FEP-f228) rather than its posts var id = Value(item, "type") is "Create" or "Update" ? Id(item["object"]) : Id(item); if (id == default) continue; var post = await _dbEntities.Posts.Match(p => p.ObjectURI == id && !p.DeletedAt.HasValue).ExecuteFirstAsync(token) ?? await _remotePosts.StoreContext(id, 0, token); if (post != default) posts.Add(post); } page = page["next"] is { } next ? await Page(next, token) : default; } return posts; } async Task Page(JsonNode node, CancellationToken token) { if (node is JsonValue && Id(node) is { } uri) { using var scope = HttpScope.For("replies"); using var fetched = await _remoteActors.FetchObject(uri, token); node = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()); } return node as JsonObject; } static IEnumerable Items(JsonNode page) => (page["orderedItems"] as JsonArray ?? page["items"] as JsonArray ?? new JsonArray()).Where(i => i != default); } }