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 collections come first: its conversation container (FEP-171b's // `contextHistory`, every activity its owner added to the thread), then its posts (FEP-7888's `context`, which // Mastodon 4.5+, Mitra and NodeBB serve with every reply at any depth), sent the ETag of the last whole read so an // unchanged thread costs one 304; otherwise the post's `replies` (PeerTube's `comments`, Forte's whole thread), 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; foreach (var name in new[] { "contextHistory", "context" }) { var uri = Id(document[name]); // Forte's `context` is the thread's first post if (uri != default && await _dbEntities.Posts.Match(p => p.ObjectURI == uri).ExecuteAnyAsync(token)) continue; var known = uri == default ? default : await DB.Default.Find().Match(c => c.URI == uri).ExecuteFirstAsync(token); var (thread, etag, unchanged) = await Collection(document[name], known?.ETag, token); if (unchanged) { await DB.Default.Update().MatchID(known.ID).Modify(c => c.ReadAt, DateTime.UtcNow).ExecuteAsync(token); return JobOutcome.Done; } if (thread == default) continue; var (_, read) = await StoreAll(thread, name == "contextHistory", token); var whole = read && Counted(thread); if (uri != default && (known != default || whole && etag != default)) await DB.Default.Update() .Match(c => c.URI == uri) .Modify(c => c.ETag, whole ? etag : default) .Modify(c => c.ReadAt, DateTime.UtcNow) .Option(o => o.IsUpsert = true) .ExecuteAsync(token); return JobOutcome.Done; } if ((await Collection(document["replies"] ?? document["comments"], default, token)).Collection is not { } replies) return JobOutcome.Done; var (answers, _) = await StoreAll(replies, false, 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, and its ETag: 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, and one whose server still has the version `etag` names is unchanged async Task<(JsonObject Collection, string ETag, bool Unchanged)> Collection(JsonNode node, string etag, CancellationToken token) { var uri = Id(node); string served = default; 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, etag, token); if (fetched?.NotModified == true) return (default, etag, true); node = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()); served = fetched?.ETag; } return (node is JsonObject collection && Value(collection, "type") is { } type && CollectionTypes.Contains(type) ? collection : default, served, false); } // every post the collection lists, stored when we lack it: its pages in order, at most MaxPages and MaxItems; whether // it was read to its end. A conversation container's items are its owner's Adds of the thread's activities. async Task<(List Posts, bool Whole)> StoreAll(JsonObject collection, bool activities, 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++) { if (pages == MaxPages) return (posts, false); foreach (var item in Items(page)) { if (++seen > MaxItems) return (posts, false); var id = activities ? await Added(item, token) : PostId(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, true); } // whether the document itself changes with every post added to the thread, so its ETag is the thread's: it counts // them (NodeBB), or holds every one with no further page (a short Mastodon thread); a long Mastodon thread's first // page stays as it was while its last one grows static bool Counted(JsonObject collection) { if (collection["totalItems"] is JsonValue) return true; var page = collection["orderedItems"] is JsonArray || collection["items"] is JsonArray ? collection : collection["first"] as JsonObject; return page != default && (page["orderedItems"] is JsonArray || page["items"] is JsonArray) && page["next"] == default; } // a context may list the activities that made the thread (FEP-f228) rather than its posts static string PostId(JsonNode item) => Value(item, "type") is "Create" or "Update" ? Id(item["object"]) : Id(item); // the post an activity in a conversation container brought into the thread, added by its owner or (Forte's) not: a // Create's or an Update's object; a like, a removal or anything else adds no post. Only its id is taken, and the post // is read from its own server. async Task Added(JsonNode item, CancellationToken token) { var activity = await Page(item, token); if (activity != default && Value(activity, "type") == "Add") activity = await Page(activity["object"], token); return activity != default && Value(activity, "type") is "Create" or "Update" ? Id(activity["object"]) : default; } 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); } }