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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB
This commit is contained in:
thepraandClaude Opus 5.5 committed 2026-10-01 11:37:41 +02:00
1 parent a76cc34557
commit e99b76dbd7
12 files changed
+609 -52

No files matched your search

+1
View File
@@ -16,6 +16,7 @@ namespace PrivaPub.ClientModels.Post
public bool HasContentWarning { get; set; } public bool HasContentWarning { get; set; }
public bool IsFederatedCopy { get; set; } public bool IsFederatedCopy { get; set; }
public DateTime CreationDate { get; set; } public DateTime CreationDate { get; set; }
public ViewPost ReblogOf { get; set; }
} }
public class ViewDmGroup public class ViewDmGroup
+142
View File
@@ -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<List<ViewPost>> Home(string root, string avatarId) =>
(List<ViewPost>)(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<Following>().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"] = $"<p>{text}</p>",
["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("<p>for my followers</p>", 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<string>(), "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("<p>top</p>", home);
Assert.Contains("<p>self reply</p>", home);
Assert.DoesNotContain("<p>dave alone</p>", home);
Assert.DoesNotContain("<p>reply to a stranger</p>", 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"] = "<p>parent</p>",
["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<PrivaPub.Models.Post.Post>().Match(p => p.ObjectURI == parentId).ExecuteSingleAsync(token);
var reply = await DB.Default.Find<PrivaPub.Models.Post.Post>().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<PrivaPub.Models.Jobs.Job>().Match(j => j.DedupeKey == "ancestors|" + parent.ID).ExecuteAnyAsync(token));
Assert.DoesNotContain("<p>reply to dave</p>", (await Home(bobRoot, bob.Id)).Select(p => p.ContentHtml));
Assert.DoesNotContain("<p>parent</p>", (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<ViewNotification>)(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<ViewNotification>)(await _harness.Timelines.Notifications(bobRoot, bob.Id, default, 20, token)).Data).Single().IsRead);
}
}
}
@@ -3,6 +3,7 @@ using Microsoft.Extensions.Logging.Abstractions;
using MongoDB.Entities; using MongoDB.Entities;
using PrivaPub.Domain.Timelines;
using PrivaPub.Federation.Actors; using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Inbox.Handlers; using PrivaPub.Federation.Inbox.Handlers;
using PrivaPub.Federation.Inbox; using PrivaPub.Federation.Inbox;
@@ -53,7 +54,7 @@ namespace PrivaPub.Tests.Federation
{ {
new FollowHandler(db, _local, remote, delivery), new FollowHandler(db, _local, remote, delivery),
new UndoHandler(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 DeleteHandler(db, _local, remote, delivery),
new UpdateHandler(db, _local, remote) new UpdateHandler(db, _local, remote)
}, NullLogger<InboxProcessor>.Instance); }, NullLogger<InboxProcessor>.Instance);
+9 -2
View File
@@ -6,6 +6,7 @@ using MongoDB.Entities;
using PrivaPub.Domain.Content; using PrivaPub.Domain.Content;
using PrivaPub.Domain.Social; using PrivaPub.Domain.Social;
using PrivaPub.Domain.Timelines;
using PrivaPub.Federation.Actors; using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Inbox; using PrivaPub.Federation.Inbox;
using PrivaPub.Federation.Inbox.Handlers; using PrivaPub.Federation.Inbox.Handlers;
@@ -36,6 +37,8 @@ namespace PrivaPub.Tests.Support
Local = new LocalActorService(Db, new StaticOptions<AppConfiguration>(new AppConfiguration { BackendBaseAddress = Base })); Local = new LocalActorService(Db, new StaticOptions<AppConfiguration>(new AppConfiguration { BackendBaseAddress = Base }));
Remote = new RemoteActorService(Peer.Http(cache), Local, cache, Db); Remote = new RemoteActorService(Peer.Http(cache), Local, cache, Db);
Delivery = new DeliveryService(Db, Queue); 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<InboxReceiver>.Instance); Receiver = new InboxReceiver(Local, Remote, Queue, new NoBlocks(), NullLogger<InboxReceiver>.Instance);
Processor = new InboxProcessor(Remote, new IActivityHandler[] Processor = new InboxProcessor(Remote, new IActivityHandler[]
{ {
@@ -43,14 +46,15 @@ namespace PrivaPub.Tests.Support
new AcceptHandler(Db, Local), new AcceptHandler(Db, Local),
new RejectHandler(Db, Local), new RejectHandler(Db, Local),
new UndoHandler(Db, Local, Remote, Delivery), 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 DeleteHandler(Db, Local, Remote, Delivery),
new UpdateHandler(Db, Local, Remote) new UpdateHandler(Db, Local, Remote)
}, NullLogger<InboxProcessor>.Instance); }, NullLogger<InboxProcessor>.Instance);
Follows = new FollowService(Db, Local, Remote, Delivery, new KeyLocalizer<GenericRes>(), NullLogger<FollowService>.Instance); Follows = new FollowService(Db, Local, Remote, Delivery, new KeyLocalizer<GenericRes>(), NullLogger<FollowService>.Instance);
Content = new ContentRenderer(Local, Remote); Content = new ContentRenderer(Local, Remote);
Outbox = new OutboxPublisher(Db, Local, Delivery); Outbox = new OutboxPublisher(Db, Local, Delivery);
Posts = new PostsService(Db, Local, Remote, Delivery, Content, Outbox, new KeyLocalizer<GenericRes>(), NullLogger<PostsService>.Instance); Posts = new PostsService(Db, Local, Remote, Delivery, Content, Outbox, Fanout, new KeyLocalizer<GenericRes>(), NullLogger<PostsService>.Instance);
Timelines = new TimelineService(Db, new KeyLocalizer<GenericRes>());
} }
public Peer Peer { get; } public Peer Peer { get; }
@@ -65,6 +69,9 @@ namespace PrivaPub.Tests.Support
public ContentRenderer Content { get; } public ContentRenderer Content { get; }
public OutboxPublisher Outbox { get; } public OutboxPublisher Outbox { get; }
public PostsService Posts { 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) => public async Task FollowedBy(LocalActor local, RemoteActor follower) =>
await DB.Default.SaveAsync(new Follower await DB.Default.SaveAsync(new Follower
@@ -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<IActionResult> 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<IActionResult> 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<IActionResult> 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);
}
}
+93
View File
@@ -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<string>(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<GroupMember>())
recipients.Add(member.AvatarId);
break;
case PostVisibility.Circle when !string.IsNullOrEmpty(post.GroupId):
var circle = await DB.Default.Find<GroupEntity>().MatchID(post.GroupId).ExecuteFirstAsync(token);
foreach (var member in circle?.Members.Where(m => !m.IsForeign) ?? Enumerable.Empty<GroupMember>())
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<bool> 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);
}
}
}
@@ -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<WebResult> Home(string rootUserId, string avatarId, string maxId, int limit, CancellationToken token);
Task<WebResult> Notifications(string rootUserId, string avatarId, string maxId, int limit, CancellationToken token);
Task<WebResult> MarkNotificationsRead(string rootUserId, string avatarId, CancellationToken token);
}
public class TimelineService : ITimelineService
{
const int MaxLimit = 40;
readonly DbEntities _dbEntities;
readonly IStringLocalizer<GenericRes> _localizer;
public TimelineService(DbEntities dbEntities, IStringLocalizer<GenericRes> localizer)
{
_dbEntities = dbEntities;
_localizer = localizer;
}
public async Task<WebResult> 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<string, Post>()
: (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<WebResult> 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<WebResult> 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<Notification>().Match(n => n.AvatarId == avatarId && !n.IsRead).Modify(n => n.IsRead, true).ExecuteAsync(token);
return result;
}
async Task<bool> 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
};
}
}
@@ -1,6 +1,7 @@
using MongoDB.Driver; using MongoDB.Driver;
using MongoDB.Entities; using MongoDB.Entities;
using PrivaPub.Domain.Timelines;
using PrivaPub.Federation.Actors; using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Moderation;
using PrivaPub.Federation.Objects; using PrivaPub.Federation.Objects;
@@ -10,6 +11,7 @@ using PrivaPub.Federation.Rendering;
using PrivaPub.Models.Federation; using PrivaPub.Models.Federation;
using PrivaPub.Models.Group; using PrivaPub.Models.Group;
using PrivaPub.Models.Post; using PrivaPub.Models.Post;
using PrivaPub.Models.Social;
using PrivaPub.Models.User; using PrivaPub.Models.User;
using PrivaPub.StaticServices; using PrivaPub.StaticServices;
@@ -29,10 +31,14 @@ namespace PrivaPub.Federation.Inbox.Handlers
readonly IRemoteActorService _remoteActors; readonly IRemoteActorService _remoteActors;
readonly IDeliveryService _delivery; readonly IDeliveryService _delivery;
readonly IDomainBlocks _domainBlocks; readonly IDomainBlocks _domainBlocks;
readonly IFanout _fanout;
readonly IRemotePosts _remotePosts;
public CreateHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery, public CreateHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery,
IDomainBlocks domainBlocks) IDomainBlocks domainBlocks, IFanout fanout, IRemotePosts remotePosts)
{ {
_fanout = fanout;
_remotePosts = remotePosts;
_dbEntities = dbEntities; _dbEntities = dbEntities;
_localActors = localActors; _localActors = localActors;
_remoteActors = remoteActors; _remoteActors = remoteActors;
@@ -86,21 +92,13 @@ namespace PrivaPub.Federation.Inbox.Handlers
group = default; group = default;
var repliesToLocal = parent is { IsFederatedCopy: false } && visibility != PostVisibility.Direct; 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; return;
if (parent == default && visibility != PostVisibility.Direct)
var mentions = new List<PostMention>(); parent = await _remotePosts.Parent(note.InReplyTo, 1, token);
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
});
}
DmGroup conversation = default; DmGroup conversation = default;
if (visibility == PostVisibility.Direct) 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); conversation = await FindOrCreateConversation(participants, Origin.Same(note.Context, author.ActorURI) ? note.Context : default, token);
} }
var post = new PostEntity var post = await _remotePosts.Build(note, author, visibility, to, cc, parent, token);
{ post.ActivityURI = Id(activity);
ID = PrivacyIds.At(note.Published), post.GroupId = group?.Id;
ObjectURI = note.Id, post.ConversationId = conversation?.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
};
try try
{ {
await DB.Default.SaveAsync(post, token); await DB.Default.SaveAsync(post, token);
@@ -156,6 +125,7 @@ namespace PrivaPub.Federation.Inbox.Handlers
if (parent != default) if (parent != default)
await DB.Default.Update<PostEntity>().MatchID(parent.ID).Modify(b => b.Inc(p => p.RepliesCount, 1)).ExecuteAsync(token); await DB.Default.Update<PostEntity>().MatchID(parent.ID).Modify(b => b.Inc(p => p.RepliesCount, 1)).ExecuteAsync(token);
await _fanout.Distribute(post, token);
if (conversation != default) if (conversation != default)
await DB.Default.Update<DmGroup>().MatchID(conversation.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token); await DB.Default.Update<DmGroup>().MatchID(conversation.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token);
if (group != default) if (group != default)
+173
View File
@@ -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<PostEntity> Build(NoteDocument note, ForeignAvatar author, PostVisibility visibility, IReadOnlyList<string> to, IReadOnlyList<string> cc,
PostEntity parent, CancellationToken token);
Task<PostEntity> Parent(string inReplyTo, int depth, CancellationToken token);
Task<PostEntity> 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<PostEntity> Build(NoteDocument note, ForeignAvatar author, PostVisibility visibility, IReadOnlyList<string> to,
IReadOnlyList<string> cc, PostEntity parent, CancellationToken token)
{
var mentions = new List<PostMention>();
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<PostEntity> 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<PostEntity> 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<JobOutcome> Handle(Job job, CancellationToken token)
{
var payload = JsonSerializer.Deserialize<AncestorsPayload>(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<PostEntity>().MatchID(child.ID)
.Modify(p => p.AnsweringToPostId, parent.ID)
.Modify(p => p.InReplyToAccountId, parent.AuthorAccountId ?? parent.GroupUserId)
.ExecuteAsync(token);
return JobOutcome.Done;
}
}
}
@@ -18,6 +18,7 @@ using PrivaPub.Federation.Moderation;
using PrivaPub.Federation.Inbox.Handlers; using PrivaPub.Federation.Inbox.Handlers;
using PrivaPub.Domain.Content; using PrivaPub.Domain.Content;
using PrivaPub.Domain.Social; using PrivaPub.Domain.Social;
using PrivaPub.Domain.Timelines;
using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Http;
using PrivaPub.Infrastructure.Jobs; using PrivaPub.Infrastructure.Jobs;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
@@ -58,6 +59,7 @@ namespace PrivaPub.Middleware
.AddSingleton<IRemoteActorService, RemoteActorService>() .AddSingleton<IRemoteActorService, RemoteActorService>()
.AddSingleton<IDeliveryService, DeliveryService>() .AddSingleton<IDeliveryService, DeliveryService>()
.AddSingleton<IOutboxPublisher, OutboxPublisher>() .AddSingleton<IOutboxPublisher, OutboxPublisher>()
.AddSingleton<IFanout, Fanout>()
.AddSingleton<IInboxReceiver, InboxReceiver>() .AddSingleton<IInboxReceiver, InboxReceiver>()
.AddSingleton<IActivityHandler, FollowHandler>() .AddSingleton<IActivityHandler, FollowHandler>()
.AddSingleton<IActivityHandler, AcceptHandler>() .AddSingleton<IActivityHandler, AcceptHandler>()
@@ -67,6 +69,8 @@ namespace PrivaPub.Middleware
.AddSingleton<IActivityHandler, DeleteHandler>() .AddSingleton<IActivityHandler, DeleteHandler>()
.AddSingleton<IActivityHandler, UpdateHandler>() .AddSingleton<IActivityHandler, UpdateHandler>()
.AddSingleton<IJobHandler, InboxProcessor>() .AddSingleton<IJobHandler, InboxProcessor>()
.AddSingleton<IRemotePosts, RemotePosts>()
.AddSingleton<IJobHandler, AncestorsJobHandler>()
.AddSingleton<IJobQueue, JobQueue>() .AddSingleton<IJobQueue, JobQueue>()
.AddSingleton<IHostCircuitBreaker, HostCircuitBreaker>() .AddSingleton<IHostCircuitBreaker, HostCircuitBreaker>()
.AddSingleton<IJobHandler, DeliveryJobHandler>() .AddSingleton<IJobHandler, DeliveryJobHandler>()
@@ -123,6 +127,7 @@ namespace PrivaPub.Middleware
.AddTransient<IGroupUsersService, GroupUsersService>() .AddTransient<IGroupUsersService, GroupUsersService>()
.AddTransient<IPostsService, PostsService>() .AddTransient<IPostsService, PostsService>()
.AddTransient<IFollowService, FollowService>() .AddTransient<IFollowService, FollowService>()
.AddTransient<ITimelineService, TimelineService>()
.AddSingleton<AppConfigurationService>() .AddSingleton<AppConfigurationService>()
.AddHttpContextAccessor() .AddHttpContextAccessor()
.AddMemoryCache() .AddMemoryCache()
+2 -1
View File
@@ -21,7 +21,8 @@ namespace PrivaPub.Models.Jobs
public enum JobKind public enum JobKind
{ {
Deliver, Deliver,
ProcessInbox ProcessInbox,
FetchAncestors
} }
public enum JobState public enum JobState
+7
View File
@@ -15,6 +15,7 @@ using System.Text.Json.Nodes;
using PostEntity = PrivaPub.Models.Post.Post; using PostEntity = PrivaPub.Models.Post.Post;
using PrivaPub.Domain.Content; using PrivaPub.Domain.Content;
using PrivaPub.Domain.Timelines;
using PrivaPub.Federation.Actors; using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Rendering; using PrivaPub.Federation.Rendering;
using PrivaPub.Federation.Outbox; using PrivaPub.Federation.Outbox;
@@ -41,6 +42,7 @@ namespace PrivaPub.Services
readonly IDeliveryService _delivery; readonly IDeliveryService _delivery;
readonly IContentRenderer _content; readonly IContentRenderer _content;
readonly IOutboxPublisher _outbox; readonly IOutboxPublisher _outbox;
readonly IFanout _fanout;
readonly IStringLocalizer<GenericRes> _localizer; readonly IStringLocalizer<GenericRes> _localizer;
readonly ILogger<PostsService> _logger; readonly ILogger<PostsService> _logger;
@@ -50,6 +52,7 @@ namespace PrivaPub.Services
IDeliveryService delivery, IDeliveryService delivery,
IContentRenderer content, IContentRenderer content,
IOutboxPublisher outbox, IOutboxPublisher outbox,
IFanout fanout,
IStringLocalizer<GenericRes> localizer, IStringLocalizer<GenericRes> localizer,
ILogger<PostsService> logger) ILogger<PostsService> logger)
{ {
@@ -59,6 +62,7 @@ namespace PrivaPub.Services
_delivery = delivery; _delivery = delivery;
_content = content; _content = content;
_outbox = outbox; _outbox = outbox;
_fanout = fanout;
_localizer = localizer; _localizer = localizer;
_logger = logger; _logger = logger;
} }
@@ -113,6 +117,7 @@ namespace PrivaPub.Services
if (isLocalOnly) if (isLocalOnly)
{ {
await DB.Default.SaveAsync(post, token); await DB.Default.SaveAsync(post, token);
await _fanout.Distribute(post, token);
result.Data = ToView(post); result.Data = ToView(post);
return result; return result;
} }
@@ -126,6 +131,7 @@ namespace PrivaPub.Services
if (parent != default) if (parent != default)
await DB.Default.Update<PostEntity>().MatchID(parent.ID).Modify(b => b.Inc(p => p.RepliesCount, 1)).ExecuteAsync(token); await DB.Default.Update<PostEntity>().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); await _outbox.Publish(author, post, create, token);
if (group is { IsFederated: true } && post.Visibility is PostVisibility.Public or PostVisibility.Unlisted) 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); 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(); 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.SaveAsync(dm, token);
await DB.Default.Update<DmGroup>().MatchID(dmGroup.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token); await DB.Default.Update<DmGroup>().MatchID(dmGroup.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token);
await _fanout.Distribute(dm, token);
if (inboxes.Count > 0) if (inboxes.Count > 0)
await _delivery.Enqueue(author, inboxes, directCreate, token); await _delivery.Enqueue(author, inboxes, directCreate, token);