From e99b76dbd736ee202822075990a2ad64a04a2ebc Mon Sep 17 00:00:00 2001 From: thepra Date: Thu, 1 Oct 2026 11:37:41 +0200 Subject: [PATCH] Home timelines, mention notifications, and threads that fetch their parents Fanout writes a TimelineEntry for every persona a post should reach: the author, local followers of a local author, local followers of a remote one, the members of a direct conversation or a circle. Mastodon's home rules apply when it is written: a reply shows only to followers of both sides (or to the one replied to), a reblog only where reblogs are wanted. A mention of a local persona becomes a Mention notification. Inbound posts are now also kept when a persona follows their author. RemotePosts holds what CreateHandler and backfill share: building a Post from a note, and fetching a public parent the first time a reply to it arrives, so its author is known (a reply to an unknown or non-public parent stays out of home timelines). Each fetched ancestor queues a FetchAncestors job for the next one, up to ten deep. /clientapi/timeline/home and /clientapi/notifications (with /notifications/read) page by max_id for the persona's own root only. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB --- PrivaPub.ClientModels/Post/ViewPost.cs | 1 + PrivaPub.Tests/Domain/TimelineTests.cs | 142 ++++++++++++++ .../Federation/InboxScenarioTests.cs | 3 +- PrivaPub.Tests/Support/Harness.cs | 11 +- .../ClientToServer/TimelineController.cs | 38 ++++ PrivaPub/Domain/Timelines/Fanout.cs | 93 ++++++++++ PrivaPub/Domain/Timelines/TimelineService.cs | 119 ++++++++++++ .../Inbox/Handlers/CreateHandler.cs | 66 ++----- PrivaPub/Federation/Inbox/RemotePosts.cs | 173 ++++++++++++++++++ .../Middleware/SocialPubConfigurations.cs | 5 + PrivaPub/Models/Jobs/Job.cs | 3 +- PrivaPub/Services/PostsService.cs | 7 + 12 files changed, 609 insertions(+), 52 deletions(-) create mode 100644 PrivaPub.Tests/Domain/TimelineTests.cs create mode 100644 PrivaPub/Controllers/ClientToServer/TimelineController.cs create mode 100644 PrivaPub/Domain/Timelines/Fanout.cs create mode 100644 PrivaPub/Domain/Timelines/TimelineService.cs create mode 100644 PrivaPub/Federation/Inbox/RemotePosts.cs diff --git a/PrivaPub.ClientModels/Post/ViewPost.cs b/PrivaPub.ClientModels/Post/ViewPost.cs index 4174503..9f7a833 100644 --- a/PrivaPub.ClientModels/Post/ViewPost.cs +++ b/PrivaPub.ClientModels/Post/ViewPost.cs @@ -16,6 +16,7 @@ namespace PrivaPub.ClientModels.Post public bool HasContentWarning { get; set; } public bool IsFederatedCopy { get; set; } public DateTime CreationDate { get; set; } + public ViewPost ReblogOf { get; set; } } public class ViewDmGroup diff --git a/PrivaPub.Tests/Domain/TimelineTests.cs b/PrivaPub.Tests/Domain/TimelineTests.cs new file mode 100644 index 0000000..6420503 --- /dev/null +++ b/PrivaPub.Tests/Domain/TimelineTests.cs @@ -0,0 +1,142 @@ +using MongoDB.Entities; + +using PrivaPub.ClientModels.Post; +using PrivaPub.ClientModels.Social; +using PrivaPub.Federation.Objects; +using PrivaPub.Models.Social; +using PrivaPub.Tests.Support; + +using System.Text.Json.Nodes; + +namespace PrivaPub.Tests.Domain +{ + [Trait("Category", "Integration")] + public sealed class TimelineTests : IAsyncLifetime + { + Harness _harness; + + public async ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + _harness = await Harness.Start(); + } + + public async ValueTask DisposeAsync() + { + if (_harness != default) + await _harness.DisposeAsync(); + } + + async Task> Home(string root, string avatarId) => + (List)(await _harness.Timelines.Home(root, avatarId, default, 40, TestContext.Current.CancellationToken)).Data; + + async Task Accepted(string root, string avatarId, RemoteActor target) + { + await _harness.Follows.Follow(root, new FollowForm { AvatarId = avatarId, Target = target.Id }, TestContext.Current.CancellationToken); + await DB.Default.Update().Match(f => f.AvatarId == avatarId && f.TargetActorURI == target.Id) + .Modify(f => f.State, FollowState.Accepted).ExecuteAsync(TestContext.Current.CancellationToken); + } + + static JsonObject Create(RemoteActor author, string inReplyTo = default, string text = "hello", params (string Uri, string Name)[] mentions) + { + var origin = new Uri(author.Id).GetLeftPart(UriPartial.Authority); + var note = new JsonObject + { + ["id"] = $"{origin}/notes/{Guid.NewGuid():N}", + ["type"] = "Note", + ["attributedTo"] = author.Id, + ["content"] = $"

{text}

", + ["to"] = new JsonArray(Addressing.Public), + ["cc"] = new JsonArray(mentions.Select(m => (JsonNode)m.Uri).Prepend(author.Id + "/followers").ToArray()), + ["tag"] = new JsonArray(mentions.Select(m => (JsonNode)new JsonObject { ["type"] = "Mention", ["href"] = m.Uri, ["name"] = m.Name }).ToArray()) + }; + if (inReplyTo != default) + note["inReplyTo"] = inReplyTo; + return new JsonObject { ["id"] = $"{origin}/activities/{Guid.NewGuid():N}", ["type"] = "Create", ["actor"] = author.Id, ["object"] = note }; + } + + [Fact] + public async Task A_local_post_reaches_the_author_and_local_followers_only() + { + var token = TestContext.Current.CancellationToken; + var (aliceRoot, alice) = await _harness.Persona("alice"); + var (bobRoot, bob) = await _harness.Persona("bob"); + var (carolRoot, carol) = await _harness.Persona("carol"); + await _harness.Follows.Follow(bobRoot, new FollowForm { AvatarId = bob.Id, Target = alice.UserName }, token); + + await _harness.Posts.InsertPost(aliceRoot, new InsertPostForm { AvatarId = alice.Id, Text = "for my followers" }, token); + + Assert.Single(await Home(aliceRoot, alice.Id)); + Assert.Equal("

for my followers

", Assert.Single(await Home(bobRoot, bob.Id)).ContentHtml); + Assert.Empty(await Home(carolRoot, carol.Id)); + Assert.Equal(404, (await _harness.Timelines.Home(carolRoot, bob.Id, default, 20, token)).StatusCode); + } + + [Fact] + public async Task A_followed_remote_account_is_kept_and_its_replies_follow_mastodons_rule() + { + var (bobRoot, bob) = await _harness.Persona("bob"); + var carol = new RemoteActor(_harness.Peer, "carol"); + var dave = new RemoteActor(_harness.Peer, "dave"); + await Accepted(bobRoot, bob.Id, carol); + + var top = Create(carol, text: "top"); + await _harness.Deliver(carol, "/human-centipede", top); + await _harness.Deliver(carol, "/human-centipede", Create(carol, top["object"]!["id"]!.GetValue(), "self reply")); + await _harness.Deliver(dave, "/human-centipede", Create(dave, text: "dave alone")); + var daveRoot = Create(dave, text: "unfollowed root"); + await _harness.Deliver(carol, "/human-centipede", Create(carol, $"{new Uri(dave.Id).GetLeftPart(UriPartial.Authority)}/notes/unknown", "reply to a stranger")); + + var home = (await Home(bobRoot, bob.Id)).Select(p => p.ContentHtml).ToList(); + Assert.Contains("

top

", home); + Assert.Contains("

self reply

", home); + Assert.DoesNotContain("

dave alone

", home); + Assert.DoesNotContain("

reply to a stranger

", home); + } + + [Fact] + public async Task A_reply_brings_its_public_parent_and_the_grandparent_is_queued() + { + var token = TestContext.Current.CancellationToken; + var (bobRoot, bob) = await _harness.Persona("bob"); + var carol = new RemoteActor(_harness.Peer, "carol"); + var dave = new RemoteActor(_harness.Peer, "dave"); + await Accepted(bobRoot, bob.Id, carol); + var origin = new Uri(dave.Id).GetLeftPart(UriPartial.Authority); + var grandparentId = $"{origin}/notes/{Guid.NewGuid():N}"; + var parentId = $"{origin}/notes/{Guid.NewGuid():N}"; + _harness.Peer.Serve(new Uri(parentId).AbsolutePath, new JsonObject + { + ["id"] = parentId, ["type"] = "Note", ["attributedTo"] = dave.Id, ["content"] = "

parent

", + ["to"] = new JsonArray(Addressing.Public), ["inReplyTo"] = grandparentId + }.ToJsonString()); + + await _harness.Deliver(carol, "/human-centipede", Create(carol, parentId, "reply to dave")); + + var parent = await DB.Default.Find().Match(p => p.ObjectURI == parentId).ExecuteSingleAsync(token); + var reply = await DB.Default.Find().Match(p => p.InReplyToURI == parentId).ExecuteSingleAsync(token); + Assert.Equal(parent.ID, reply.AnsweringToPostId); + Assert.Equal(parent.AuthorAccountId, reply.InReplyToAccountId); + Assert.True(await DB.Default.Find().Match(j => j.DedupeKey == "ancestors|" + parent.ID).ExecuteAnyAsync(token)); + Assert.DoesNotContain("

reply to dave

", (await Home(bobRoot, bob.Id)).Select(p => p.ContentHtml)); + Assert.DoesNotContain("

parent

", (await Home(bobRoot, bob.Id)).Select(p => p.ContentHtml)); + } + + [Fact] + public async Task A_mention_notifies_the_persona() + { + var token = TestContext.Current.CancellationToken; + var (bobRoot, bob) = await _harness.Persona("bob"); + var carol = new RemoteActor(_harness.Peer, "carol"); + + await _harness.Deliver(carol, "/human-centipede", Create(carol, mentions: (bob.Uri, "@" + bob.Handle))); + + var notifications = (List)(await _harness.Timelines.Notifications(bobRoot, bob.Id, default, 20, token)).Data; + var mention = Assert.Single(notifications); + Assert.Equal("Mention", mention.Type); + Assert.Equal(carol.Id, mention.FromActorURI); + await _harness.Timelines.MarkNotificationsRead(bobRoot, bob.Id, token); + Assert.True(((List)(await _harness.Timelines.Notifications(bobRoot, bob.Id, default, 20, token)).Data).Single().IsRead); + } + } +} diff --git a/PrivaPub.Tests/Federation/InboxScenarioTests.cs b/PrivaPub.Tests/Federation/InboxScenarioTests.cs index 82772c4..d27236d 100644 --- a/PrivaPub.Tests/Federation/InboxScenarioTests.cs +++ b/PrivaPub.Tests/Federation/InboxScenarioTests.cs @@ -3,6 +3,7 @@ using Microsoft.Extensions.Logging.Abstractions; using MongoDB.Entities; +using PrivaPub.Domain.Timelines; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Inbox.Handlers; using PrivaPub.Federation.Inbox; @@ -53,7 +54,7 @@ namespace PrivaPub.Tests.Federation { new FollowHandler(db, _local, remote, delivery), new UndoHandler(db, _local, remote, delivery), - new CreateHandler(db, _local, remote, delivery, _blocks), + new CreateHandler(db, _local, remote, delivery, _blocks, new Fanout(db), new RemotePosts(db, _local, remote, _blocks, queue)), new DeleteHandler(db, _local, remote, delivery), new UpdateHandler(db, _local, remote) }, NullLogger.Instance); diff --git a/PrivaPub.Tests/Support/Harness.cs b/PrivaPub.Tests/Support/Harness.cs index a41cc4b..656ec9f 100644 --- a/PrivaPub.Tests/Support/Harness.cs +++ b/PrivaPub.Tests/Support/Harness.cs @@ -6,6 +6,7 @@ using MongoDB.Entities; using PrivaPub.Domain.Content; using PrivaPub.Domain.Social; +using PrivaPub.Domain.Timelines; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Inbox; using PrivaPub.Federation.Inbox.Handlers; @@ -36,6 +37,8 @@ namespace PrivaPub.Tests.Support Local = new LocalActorService(Db, new StaticOptions(new AppConfiguration { BackendBaseAddress = Base })); Remote = new RemoteActorService(Peer.Http(cache), Local, cache, Db); Delivery = new DeliveryService(Db, Queue); + Fanout = new Fanout(Db); + RemotePosts = new RemotePosts(Db, Local, Remote, new NoBlocks(), Queue); Receiver = new InboxReceiver(Local, Remote, Queue, new NoBlocks(), NullLogger.Instance); Processor = new InboxProcessor(Remote, new IActivityHandler[] { @@ -43,14 +46,15 @@ namespace PrivaPub.Tests.Support new AcceptHandler(Db, Local), new RejectHandler(Db, Local), new UndoHandler(Db, Local, Remote, Delivery), - new CreateHandler(Db, Local, Remote, Delivery, new NoBlocks()), + new CreateHandler(Db, Local, Remote, Delivery, new NoBlocks(), Fanout, RemotePosts), new DeleteHandler(Db, Local, Remote, Delivery), new UpdateHandler(Db, Local, Remote) }, NullLogger.Instance); Follows = new FollowService(Db, Local, Remote, Delivery, new KeyLocalizer(), NullLogger.Instance); Content = new ContentRenderer(Local, Remote); Outbox = new OutboxPublisher(Db, Local, Delivery); - Posts = new PostsService(Db, Local, Remote, Delivery, Content, Outbox, new KeyLocalizer(), NullLogger.Instance); + Posts = new PostsService(Db, Local, Remote, Delivery, Content, Outbox, Fanout, new KeyLocalizer(), NullLogger.Instance); + Timelines = new TimelineService(Db, new KeyLocalizer()); } public Peer Peer { get; } @@ -65,6 +69,9 @@ namespace PrivaPub.Tests.Support public ContentRenderer Content { get; } public OutboxPublisher Outbox { get; } public PostsService Posts { get; } + public Fanout Fanout { get; } + public RemotePosts RemotePosts { get; } + public TimelineService Timelines { get; } public async Task FollowedBy(LocalActor local, RemoteActor follower) => await DB.Default.SaveAsync(new Follower diff --git a/PrivaPub/Controllers/ClientToServer/TimelineController.cs b/PrivaPub/Controllers/ClientToServer/TimelineController.cs new file mode 100644 index 0000000..fad1418 --- /dev/null +++ b/PrivaPub/Controllers/ClientToServer/TimelineController.cs @@ -0,0 +1,38 @@ +using Microsoft.AspNetCore.Authorization; +using Microsoft.AspNetCore.Mvc; + +using PrivaPub.ClientModels; +using PrivaPub.Domain.Timelines; +using PrivaPub.Extensions; + +namespace PrivaPub.Controllers.ClientToServer +{ + [ApiController, + Route("clientapi/timeline"), + Authorize(Policy = Policies.IsUser)] + public class TimelineController : ControllerBase + { + readonly ITimelineService _timelines; + + public TimelineController(ITimelineService timelines) + { + _timelines = timelines; + } + + [HttpGet, Route("/clientapi/timeline/home")] + public async Task Home([FromQuery] string avatarId, [FromQuery(Name = "max_id")] string maxId, [FromQuery] int limit = 20, + CancellationToken token = default) => + Answer(await _timelines.Home(User.GetUserId(), avatarId, maxId, limit, token)); + + [HttpGet, Route("/clientapi/notifications")] + public async Task Notifications([FromQuery] string avatarId, [FromQuery(Name = "max_id")] string maxId, [FromQuery] int limit = 20, + CancellationToken token = default) => + Answer(await _timelines.Notifications(User.GetUserId(), avatarId, maxId, limit, token)); + + [HttpPost, Route("/clientapi/notifications/read")] + public async Task Read([FromQuery] string avatarId, CancellationToken token) => + Answer(await _timelines.MarkNotificationsRead(User.GetUserId(), avatarId, token)); + + IActionResult Answer(WebResult result) => result.IsValid ? Ok(result.Data) : StatusCode(result.StatusCode, result); + } +} diff --git a/PrivaPub/Domain/Timelines/Fanout.cs b/PrivaPub/Domain/Timelines/Fanout.cs new file mode 100644 index 0000000..db0676b --- /dev/null +++ b/PrivaPub/Domain/Timelines/Fanout.cs @@ -0,0 +1,93 @@ +using MongoDB.Driver; +using MongoDB.Entities; + +using PrivaPub.Domain.Social; +using PrivaPub.Models.Group; +using PrivaPub.Models.Post; +using PrivaPub.Models.Social; +using PrivaPub.StaticServices; + +using GroupEntity = PrivaPub.Models.Group.Group; + +namespace PrivaPub.Domain.Timelines +{ + public interface IFanout + { + Task Distribute(Post post, CancellationToken token); + } + + public class Fanout : IFanout + { + readonly DbEntities _dbEntities; + + public Fanout(DbEntities dbEntities) + { + _dbEntities = dbEntities; + } + + public async Task Distribute(Post post, CancellationToken token) + { + var recipients = new HashSet(StringComparer.Ordinal); + if (!post.IsFederatedCopy && !string.IsNullOrEmpty(post.GroupUserId)) + recipients.Add(post.GroupUserId); + + switch (post.Visibility) + { + case PostVisibility.Public or PostVisibility.Unlisted or PostVisibility.FollowersOnly: + var followers = post.IsFederatedCopy + ? await _dbEntities.Followings.Match(f => f.TargetActorURI == post.ActorURI && f.State == FollowState.Accepted).ExecuteAsync(token) + : await _dbEntities.Followings.Match(f => f.TargetAccountId == post.GroupUserId && f.TargetIsLocal && f.State == FollowState.Accepted).ExecuteAsync(token); + foreach (var following in followers) + if (await Shows(following, post, token)) + recipients.Add(following.AvatarId); + break; + case PostVisibility.Direct when !string.IsNullOrEmpty(post.ConversationId): + var conversation = await _dbEntities.DmGroups.MatchID(post.ConversationId).ExecuteFirstAsync(token); + foreach (var member in conversation?.Members.Where(m => !m.IsForeign) ?? Enumerable.Empty()) + recipients.Add(member.AvatarId); + break; + case PostVisibility.Circle when !string.IsNullOrEmpty(post.GroupId): + var circle = await DB.Default.Find().MatchID(post.GroupId).ExecuteFirstAsync(token); + foreach (var member in circle?.Members.Where(m => !m.IsForeign) ?? Enumerable.Empty()) + recipients.Add(member.AvatarId); + break; + } + + foreach (var avatarId in recipients) + { + try + { + await DB.Default.SaveAsync(new TimelineEntry + { + AvatarId = avatarId, + PostId = post.ID, + AuthorAccountId = post.AuthorAccountId ?? post.GroupUserId, + ReblogOfPostId = post.ReblogOfPostId + }, token); + } + catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey) + { + } + } + + if (string.IsNullOrEmpty(post.ReblogOfPostId)) + foreach (var mention in post.Mentions.Where(m => m.IsLocal && !string.IsNullOrEmpty(m.AccountId))) + await Notifications.Add(mention.AccountId, NotificationType.Mention, post.AuthorAccountId ?? post.GroupUserId, post.ActorURI, post.ID, token); + } + + async Task Shows(Following following, Post post, CancellationToken token) + { + if (!string.IsNullOrEmpty(post.ReblogOfPostId)) + return following.ShowReblogs; + var repliedTo = post.InReplyToAccountId; + var author = post.AuthorAccountId ?? post.GroupUserId; + if (string.IsNullOrEmpty(post.InReplyToURI) || repliedTo == author || repliedTo == following.AvatarId) + return true; + if (string.IsNullOrEmpty(repliedTo)) + return false; + return await _dbEntities.Followings + .Match(f => f.AvatarId == following.AvatarId && f.TargetAccountId == repliedTo && f.State == FollowState.Accepted) + .ExecuteAnyAsync(token); + } + } +} diff --git a/PrivaPub/Domain/Timelines/TimelineService.cs b/PrivaPub/Domain/Timelines/TimelineService.cs new file mode 100644 index 0000000..642f95c --- /dev/null +++ b/PrivaPub/Domain/Timelines/TimelineService.cs @@ -0,0 +1,119 @@ +using Microsoft.Extensions.Localization; + +using MongoDB.Entities; + +using PrivaPub.ClientModels; +using PrivaPub.ClientModels.Post; +using PrivaPub.ClientModels.Social; +using PrivaPub.Models.Post; +using PrivaPub.Models.Social; +using PrivaPub.Resources; +using PrivaPub.StaticServices; + +namespace PrivaPub.Domain.Timelines +{ + public interface ITimelineService + { + Task Home(string rootUserId, string avatarId, string maxId, int limit, CancellationToken token); + Task Notifications(string rootUserId, string avatarId, string maxId, int limit, CancellationToken token); + Task MarkNotificationsRead(string rootUserId, string avatarId, CancellationToken token); + } + + public class TimelineService : ITimelineService + { + const int MaxLimit = 40; + + readonly DbEntities _dbEntities; + readonly IStringLocalizer _localizer; + + public TimelineService(DbEntities dbEntities, IStringLocalizer localizer) + { + _dbEntities = dbEntities; + _localizer = localizer; + } + + public async Task Home(string rootUserId, string avatarId, string maxId, int limit, CancellationToken token) + { + var result = new WebResult(); + if (!await Owns(rootUserId, avatarId, token)) + return result.Invalidate(_localizer["Avatar not found."], StatusCodes.Status404NotFound); + + var query = _dbEntities.TimelineEntries.Match(e => e.AvatarId == avatarId); + if (!string.IsNullOrEmpty(maxId)) + query.Match(f => f.Lt(e => e.PostId, maxId)); + var entries = await query.Sort(e => e.PostId, Order.Descending).Limit(Math.Clamp(limit, 1, MaxLimit)).ExecuteAsync(token); + + var ids = entries.Select(e => e.PostId).Concat(entries.Where(e => e.ReblogOfPostId != default).Select(e => e.ReblogOfPostId)).Distinct().ToList(); + var posts = ids.Count == 0 + ? new Dictionary() + : (await _dbEntities.Posts.Match(p => ids.Contains(p.ID) && !p.DeletedAt.HasValue).ExecuteAsync(token)).ToDictionary(p => p.ID); + + result.Data = entries + .Where(e => posts.ContainsKey(e.PostId) && (e.ReblogOfPostId == default || posts.ContainsKey(e.ReblogOfPostId))) + .Select(e => + { + var view = ToView(posts[e.PostId]); + if (e.ReblogOfPostId != default) + view.ReblogOf = ToView(posts[e.ReblogOfPostId]); + return view; + }) + .ToList(); + return result; + } + + public async Task Notifications(string rootUserId, string avatarId, string maxId, int limit, CancellationToken token) + { + var result = new WebResult(); + if (!await Owns(rootUserId, avatarId, token)) + return result.Invalidate(_localizer["Avatar not found."], StatusCodes.Status404NotFound); + + var query = _dbEntities.Notifications.Match(n => n.AvatarId == avatarId); + if (!string.IsNullOrEmpty(maxId)) + query.Match(f => f.Lt(n => n.ID, maxId)); + result.Data = (await query.Sort(n => n.ID, Order.Descending).Limit(Math.Clamp(limit, 1, MaxLimit)).ExecuteAsync(token)) + .Select(n => new ViewNotification + { + Id = n.ID, + Type = n.Type.ToString(), + FromAccountId = n.FromAccountId, + FromActorURI = n.FromActorURI, + PostId = n.PostId, + IsRead = n.IsRead, + CreatedAt = n.CreatedAt + }) + .ToList(); + return result; + } + + public async Task MarkNotificationsRead(string rootUserId, string avatarId, CancellationToken token) + { + var result = new WebResult(); + if (!await Owns(rootUserId, avatarId, token)) + return result.Invalidate(_localizer["Avatar not found."], StatusCodes.Status404NotFound); + await DB.Default.Update().Match(n => n.AvatarId == avatarId && !n.IsRead).Modify(n => n.IsRead, true).ExecuteAsync(token); + return result; + } + + async Task Owns(string rootUserId, string avatarId, CancellationToken token) => + !string.IsNullOrEmpty(rootUserId) && !string.IsNullOrEmpty(avatarId) + && await _dbEntities.RootToAvatars.Match(ra => ra.RootId == rootUserId && ra.AvatarId == avatarId).ExecuteAnyAsync(token); + + static ViewPost ToView(Post post) => new() + { + Id = post.ID, + ObjectURI = post.ObjectURI, + AuthorAvatarId = post.IsFederatedCopy ? default : post.GroupUserId, + AuthorActorURI = post.ActorURI, + GroupId = post.GroupId, + DmGroupId = post.ConversationId, + Visibility = post.Visibility.ToString(), + AnsweringToPostId = post.AnsweringToPostId, + Title = post.Title, + Text = post.Text, + ContentHtml = post.ContentHtml, + HasContentWarning = post.HasContentWarning, + IsFederatedCopy = post.IsFederatedCopy, + CreationDate = post.CreationDate + }; + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs b/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs index b6b6501..729c463 100644 --- a/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs +++ b/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs @@ -1,6 +1,7 @@ using MongoDB.Driver; using MongoDB.Entities; +using PrivaPub.Domain.Timelines; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Objects; @@ -10,6 +11,7 @@ using PrivaPub.Federation.Rendering; using PrivaPub.Models.Federation; using PrivaPub.Models.Group; using PrivaPub.Models.Post; +using PrivaPub.Models.Social; using PrivaPub.Models.User; using PrivaPub.StaticServices; @@ -29,10 +31,14 @@ namespace PrivaPub.Federation.Inbox.Handlers readonly IRemoteActorService _remoteActors; readonly IDeliveryService _delivery; readonly IDomainBlocks _domainBlocks; + readonly IFanout _fanout; + readonly IRemotePosts _remotePosts; public CreateHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery, - IDomainBlocks domainBlocks) + IDomainBlocks domainBlocks, IFanout fanout, IRemotePosts remotePosts) { + _fanout = fanout; + _remotePosts = remotePosts; _dbEntities = dbEntities; _localActors = localActors; _remoteActors = remoteActors; @@ -86,21 +92,13 @@ namespace PrivaPub.Federation.Inbox.Handlers group = default; var repliesToLocal = parent is { IsFederatedCopy: false } && visibility != PostVisibility.Direct; - if (visibility == PostVisibility.Direct ? persons.Count == 0 : group == default && persons.Count == 0 && !repliesToLocal) + var followed = visibility != PostVisibility.Direct && await _dbEntities.Followings + .Match(f => f.TargetActorURI == author.ActorURI && f.State == FollowState.Accepted) + .ExecuteAnyAsync(token); + if (visibility == PostVisibility.Direct ? persons.Count == 0 : group == default && persons.Count == 0 && !repliesToLocal && !followed) return; - - var mentions = new List(); - foreach (var mention in note.Mentions) - { - var local = await _localActors.FindByUri(mention.ActorURI, token); - mentions.Add(new PostMention - { - ActorURI = mention.ActorURI, - Handle = mention.Handle, - IsLocal = local != default, - AccountId = local?.Id - }); - } + if (parent == default && visibility != PostVisibility.Direct) + parent = await _remotePosts.Parent(note.InReplyTo, 1, token); DmGroup conversation = default; if (visibility == PostVisibility.Direct) @@ -112,39 +110,10 @@ namespace PrivaPub.Federation.Inbox.Handlers conversation = await FindOrCreateConversation(participants, Origin.Same(note.Context, author.ActorURI) ? note.Context : default, token); } - var post = new PostEntity - { - ID = PrivacyIds.At(note.Published), - ObjectURI = note.Id, - ActivityURI = Id(activity), - ActorURI = author.ActorURI, - AuthorAccountId = author.ID, - Visibility = visibility, - To = to, - Cc = cc, - Url = note.Url, - ContextURI = note.Context, - QuoteURI = note.QuoteUri, - GroupId = group?.Id, - ConversationId = conversation?.ID, - Title = note.Title, - SpoilerText = note.SpoilerText, - HasContentWarning = note.Sensitive, - Text = note.ContentHtml, - ContentHtml = note.ContentHtml, - ContentFormat = ContentFormat.Html, - Language = note.Language, - Mentions = mentions, - Tags = note.Tags.ToList(), - Media = _domainBlocks.Find(new Uri(author.ActorURI).Host)?.RejectMedia == true ? new() : note.Attachments.ToList(), - InReplyToURI = note.InReplyTo, - AnsweringToPostId = parent?.ID, - InReplyToAccountId = parent?.AuthorAccountId ?? parent?.GroupUserId, - IsFederatedCopy = true, - CreationDate = note.Published, - UpdateDate = note.Updated ?? note.Published, - EditedAt = note.Updated - }; + var post = await _remotePosts.Build(note, author, visibility, to, cc, parent, token); + post.ActivityURI = Id(activity); + post.GroupId = group?.Id; + post.ConversationId = conversation?.ID; try { await DB.Default.SaveAsync(post, token); @@ -156,6 +125,7 @@ namespace PrivaPub.Federation.Inbox.Handlers if (parent != default) await DB.Default.Update().MatchID(parent.ID).Modify(b => b.Inc(p => p.RepliesCount, 1)).ExecuteAsync(token); + await _fanout.Distribute(post, token); if (conversation != default) await DB.Default.Update().MatchID(conversation.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token); if (group != default) diff --git a/PrivaPub/Federation/Inbox/RemotePosts.cs b/PrivaPub/Federation/Inbox/RemotePosts.cs new file mode 100644 index 0000000..accdb12 --- /dev/null +++ b/PrivaPub/Federation/Inbox/RemotePosts.cs @@ -0,0 +1,173 @@ +using MongoDB.Driver; +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Moderation; +using PrivaPub.Federation.Objects; +using PrivaPub.Infrastructure.Ids; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Post; +using PrivaPub.Models.User; +using PrivaPub.StaticServices; + +using System.Text.Json; +using System.Text.Json.Nodes; + +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Federation.Inbox +{ + public sealed record AncestorsPayload(string PostId, int Depth); + + public interface IRemotePosts + { + Task Build(NoteDocument note, ForeignAvatar author, PostVisibility visibility, IReadOnlyList to, IReadOnlyList cc, + PostEntity parent, CancellationToken token); + Task Parent(string inReplyTo, int depth, CancellationToken token); + Task StoreContext(string objectUri, int depth, CancellationToken token); + } + + public class RemotePosts : IRemotePosts + { + public const int MaxDepth = 10; + + readonly DbEntities _dbEntities; + readonly ILocalActorService _localActors; + readonly IRemoteActorService _remoteActors; + readonly IDomainBlocks _domainBlocks; + readonly IJobQueue _queue; + + public RemotePosts(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDomainBlocks domainBlocks, + IJobQueue queue) + { + _dbEntities = dbEntities; + _localActors = localActors; + _remoteActors = remoteActors; + _domainBlocks = domainBlocks; + _queue = queue; + } + + public async Task Build(NoteDocument note, ForeignAvatar author, PostVisibility visibility, IReadOnlyList to, + IReadOnlyList cc, PostEntity parent, CancellationToken token) + { + var mentions = new List(); + foreach (var mention in note.Mentions) + { + var local = await _localActors.FindByUri(mention.ActorURI, token); + mentions.Add(new PostMention { ActorURI = mention.ActorURI, Handle = mention.Handle, IsLocal = local != default, AccountId = local?.Id }); + } + + return new PostEntity + { + ID = PrivacyIds.At(note.Published), + ObjectURI = note.Id, + ActorURI = author.ActorURI, + AuthorAccountId = author.ID, + Visibility = visibility, + To = to.ToList(), + Cc = cc.ToList(), + Url = note.Url, + ContextURI = note.Context, + QuoteURI = note.QuoteUri, + Title = note.Title, + SpoilerText = note.SpoilerText, + HasContentWarning = note.Sensitive, + Text = note.ContentHtml, + ContentHtml = note.ContentHtml, + ContentFormat = ContentFormat.Html, + Language = note.Language, + Mentions = mentions, + Tags = note.Tags.ToList(), + Media = _domainBlocks.Find(new Uri(author.ActorURI).Host)?.RejectMedia == true ? new() : note.Attachments.ToList(), + InReplyToURI = note.InReplyTo, + AnsweringToPostId = parent?.ID, + InReplyToAccountId = parent?.AuthorAccountId ?? parent?.GroupUserId, + IsFederatedCopy = true, + CreationDate = note.Published, + UpdateDate = note.Updated ?? note.Published, + EditedAt = note.Updated + }; + } + + public async Task Parent(string inReplyTo, int depth, CancellationToken token) + { + if (string.IsNullOrEmpty(inReplyTo)) + return default; + var known = await _dbEntities.Posts.Match(p => p.ObjectURI == inReplyTo && !p.DeletedAt.HasValue).ExecuteFirstAsync(token); + if (known != default || inReplyTo.StartsWith(_localActors.BaseAddress + "/", StringComparison.OrdinalIgnoreCase)) + return known; + return await StoreContext(inReplyTo, depth, token); + } + + public async Task StoreContext(string objectUri, int depth, CancellationToken token) + { + var existing = await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri).ExecuteFirstAsync(token); + if (existing != default || depth > MaxDepth) + return existing; + + using var fetched = await _remoteActors.FetchObject(objectUri, token); + var note = fetched == default ? default : NoteParser.Parse(JsonNode.Parse(fetched.Root.GetRawText())); + if (note == default || !Origin.Same(note.Id, note.AttributedTo)) + return default; + var visibility = Addressing.Classify(note.To, note.Cc, default); + if (visibility is not (PostVisibility.Public or PostVisibility.Unlisted)) + return default; + var author = await _remoteActors.GetActor(note.AttributedTo, refresh: false, token); + if (author == default) + return default; + + var grandparent = string.IsNullOrEmpty(note.InReplyTo) + ? default + : await _dbEntities.Posts.Match(p => p.ObjectURI == note.InReplyTo && !p.DeletedAt.HasValue).ExecuteFirstAsync(token); + var post = await Build(note, author, visibility, note.To, note.Cc, grandparent, token); + try + { + await DB.Default.SaveAsync(post, token); + } + catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey) + { + return await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri).ExecuteFirstAsync(token); + } + + if (grandparent == default && !string.IsNullOrEmpty(note.InReplyTo) && depth < MaxDepth) + await _queue.Enqueue(JobKind.FetchAncestors, JsonSerializer.Serialize(new AncestorsPayload(post.ID, depth + 1)), new Uri(note.InReplyTo).Host, + "ancestors|" + post.ID, token); + return post; + } + } + + public class AncestorsJobHandler : IJobHandler + { + readonly DbEntities _dbEntities; + readonly IRemotePosts _remotePosts; + + public AncestorsJobHandler(DbEntities dbEntities, IRemotePosts remotePosts) + { + _dbEntities = dbEntities; + _remotePosts = remotePosts; + } + + public JobKind Kind => JobKind.FetchAncestors; + public int Concurrency => 1; + public int MaxAttempts => 3; + public int PerHostLimit => 1; + + public async Task Handle(Job job, CancellationToken token) + { + var payload = JsonSerializer.Deserialize(job.Payload); + var child = await _dbEntities.Posts.MatchID(payload.PostId).ExecuteFirstAsync(token); + if (child == default || string.IsNullOrEmpty(child.InReplyToURI) || !string.IsNullOrEmpty(child.AnsweringToPostId)) + return JobOutcome.Done; + + var parent = await _remotePosts.Parent(child.InReplyToURI, payload.Depth, token); + if (parent == default) + return JobOutcome.Done; + await DB.Default.Update().MatchID(child.ID) + .Modify(p => p.AnsweringToPostId, parent.ID) + .Modify(p => p.InReplyToAccountId, parent.AuthorAccountId ?? parent.GroupUserId) + .ExecuteAsync(token); + return JobOutcome.Done; + } + } +} diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index 9909abf..0245530 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -18,6 +18,7 @@ using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Inbox.Handlers; using PrivaPub.Domain.Content; using PrivaPub.Domain.Social; +using PrivaPub.Domain.Timelines; using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Jobs; using Microsoft.Extensions.Options; @@ -58,6 +59,7 @@ namespace PrivaPub.Middleware .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton() .AddSingleton() @@ -67,6 +69,8 @@ namespace PrivaPub.Middleware .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton() .AddSingleton() @@ -123,6 +127,7 @@ namespace PrivaPub.Middleware .AddTransient() .AddTransient() .AddTransient() + .AddTransient() .AddSingleton() .AddHttpContextAccessor() .AddMemoryCache() diff --git a/PrivaPub/Models/Jobs/Job.cs b/PrivaPub/Models/Jobs/Job.cs index a31e6f8..5457803 100644 --- a/PrivaPub/Models/Jobs/Job.cs +++ b/PrivaPub/Models/Jobs/Job.cs @@ -21,7 +21,8 @@ namespace PrivaPub.Models.Jobs public enum JobKind { Deliver, - ProcessInbox + ProcessInbox, + FetchAncestors } public enum JobState diff --git a/PrivaPub/Services/PostsService.cs b/PrivaPub/Services/PostsService.cs index a607af8..844224a 100644 --- a/PrivaPub/Services/PostsService.cs +++ b/PrivaPub/Services/PostsService.cs @@ -15,6 +15,7 @@ using System.Text.Json.Nodes; using PostEntity = PrivaPub.Models.Post.Post; using PrivaPub.Domain.Content; +using PrivaPub.Domain.Timelines; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Rendering; using PrivaPub.Federation.Outbox; @@ -41,6 +42,7 @@ namespace PrivaPub.Services readonly IDeliveryService _delivery; readonly IContentRenderer _content; readonly IOutboxPublisher _outbox; + readonly IFanout _fanout; readonly IStringLocalizer _localizer; readonly ILogger _logger; @@ -50,6 +52,7 @@ namespace PrivaPub.Services IDeliveryService delivery, IContentRenderer content, IOutboxPublisher outbox, + IFanout fanout, IStringLocalizer localizer, ILogger logger) { @@ -59,6 +62,7 @@ namespace PrivaPub.Services _delivery = delivery; _content = content; _outbox = outbox; + _fanout = fanout; _localizer = localizer; _logger = logger; } @@ -113,6 +117,7 @@ namespace PrivaPub.Services if (isLocalOnly) { await DB.Default.SaveAsync(post, token); + await _fanout.Distribute(post, token); result.Data = ToView(post); return result; } @@ -126,6 +131,7 @@ namespace PrivaPub.Services if (parent != default) await DB.Default.Update().MatchID(parent.ID).Modify(b => b.Inc(p => p.RepliesCount, 1)).ExecuteAsync(token); + await _fanout.Distribute(post, token); await _outbox.Publish(author, post, create, token); if (group is { IsFederated: true } && post.Visibility is PostVisibility.Public or PostVisibility.Unlisted) await _delivery.EnqueueToFollowers(group, ActivityPubRenderer.Announce(group, post.ObjectURI, $"announce-{post.ID}"), token); @@ -367,6 +373,7 @@ namespace PrivaPub.Services dm.Mentions = remote.Select(r => new PostMention { ActorURI = r.Uri, Handle = "@" + r.Handle, IsLocal = r.Uri.StartsWith(author.BaseAddress + "/") }).ToList(); await DB.Default.SaveAsync(dm, token); await DB.Default.Update().MatchID(dmGroup.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token); + await _fanout.Distribute(dm, token); if (inboxes.Count > 0) await _delivery.Enqueue(author, inboxes, directCreate, token);