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
172 lines
6.6 KiB
C#
172 lines
6.6 KiB
C#
using MongoDB.Driver;
|
|
using MongoDB.Entities;
|
|
|
|
using PrivaPub.Domain.Timelines;
|
|
using PrivaPub.Federation.Actors;
|
|
using PrivaPub.Federation.Moderation;
|
|
using PrivaPub.Federation.Objects;
|
|
using PrivaPub.Federation.Outbox;
|
|
using PrivaPub.Infrastructure.Ids;
|
|
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;
|
|
|
|
using System.Text.Json.Nodes;
|
|
|
|
using static PrivaPub.Federation.Inbox.ForeignMembers;
|
|
using static PrivaPub.Federation.Objects.ActivityJson;
|
|
|
|
using PostEntity = PrivaPub.Models.Post.Post;
|
|
|
|
namespace PrivaPub.Federation.Inbox.Handlers
|
|
{
|
|
public class CreateHandler : IActivityHandler
|
|
{
|
|
readonly DbEntities _dbEntities;
|
|
readonly ILocalActorService _localActors;
|
|
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, IFanout fanout, IRemotePosts remotePosts)
|
|
{
|
|
_fanout = fanout;
|
|
_remotePosts = remotePosts;
|
|
_dbEntities = dbEntities;
|
|
_localActors = localActors;
|
|
_remoteActors = remoteActors;
|
|
_delivery = delivery;
|
|
_domainBlocks = domainBlocks;
|
|
}
|
|
|
|
public string Type => "Create";
|
|
|
|
public async Task Handle(JsonNode activity, ForeignAvatar author, CancellationToken token)
|
|
{
|
|
var node = activity["object"];
|
|
if (node is not JsonObject || !Origin.Same(Id(node), author.ActorURI))
|
|
{
|
|
using var fetched = await _remoteActors.FetchObject(Id(node), token);
|
|
node = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText());
|
|
}
|
|
var note = NoteParser.Parse(node);
|
|
if (note == default || note.AttributedTo != author.ActorURI)
|
|
return;
|
|
if (await _dbEntities.Posts.Match(p => p.ObjectURI == note.Id).ExecuteAnyAsync(token))
|
|
return;
|
|
|
|
var to = note.To.Concat(Strings(activity["to"])).Distinct(StringComparer.Ordinal).ToList();
|
|
var cc = note.Cc.Concat(Strings(activity["cc"])).Distinct(StringComparer.Ordinal).ToList();
|
|
var visibility = Addressing.Classify(to, cc, author.FollowersURL);
|
|
|
|
var addressed = to.Concat(cc)
|
|
.Concat(Addresses(activity))
|
|
.Concat(note.Mentions.Select(m => m.ActorURI))
|
|
.Append(note.Audience)
|
|
.Where(a => a != default && !Addressing.IsPublic(a) && a != author.FollowersURL)
|
|
.Distinct(StringComparer.Ordinal)
|
|
.ToList();
|
|
var localTargets = new List<LocalActor>();
|
|
foreach (var uri in addressed)
|
|
{
|
|
var local = await _localActors.FindByUri(uri, token);
|
|
if (local is { IsFederated: true } && localTargets.All(l => l.Id != local.Id))
|
|
localTargets.Add(local);
|
|
}
|
|
var persons = localTargets.Where(t => t.Kind == LocalActorKind.Person).ToList();
|
|
|
|
var parent = string.IsNullOrEmpty(note.InReplyTo)
|
|
? default
|
|
: await _dbEntities.Posts.Match(p => p.ObjectURI == note.InReplyTo && !p.DeletedAt.HasValue).ExecuteFirstAsync(token);
|
|
var group = visibility is PostVisibility.Public or PostVisibility.Unlisted
|
|
? localTargets.FirstOrDefault(t => t.Kind == LocalActorKind.Group)
|
|
: default;
|
|
if (group != default && !await IsAcceptedFollower(group, author.ActorURI, token))
|
|
group = default;
|
|
|
|
var repliesToLocal = parent is { IsFederatedCopy: false } && visibility != PostVisibility.Direct;
|
|
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;
|
|
if (parent == default && visibility != PostVisibility.Direct)
|
|
parent = await _remotePosts.Parent(note.InReplyTo, 1, token);
|
|
|
|
DmGroup conversation = default;
|
|
if (visibility == PostVisibility.Direct)
|
|
{
|
|
var participants = persons.Select(p => p.Uri)
|
|
.Concat(addressed.Where(a => !IsCollection(a)))
|
|
.Append(author.ActorURI)
|
|
.ToList();
|
|
conversation = await FindOrCreateConversation(participants, Origin.Same(note.Context, author.ActorURI) ? note.Context : default, token);
|
|
}
|
|
|
|
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);
|
|
}
|
|
catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey)
|
|
{
|
|
return;
|
|
}
|
|
|
|
if (parent != default)
|
|
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)
|
|
await DB.Default.Update<DmGroup>().MatchID(conversation.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token);
|
|
if (group != default)
|
|
{
|
|
var announce = ActivityPubRenderer.Announce(group, note.Id, $"announce-{post.ID}");
|
|
await _delivery.EnqueueToFollowers(group, announce, token);
|
|
}
|
|
}
|
|
|
|
async Task<DmGroup> FindOrCreateConversation(List<string> participantUris, string context, CancellationToken token)
|
|
{
|
|
var members = new List<GroupMember>();
|
|
foreach (var uri in participantUris.Distinct(StringComparer.Ordinal))
|
|
{
|
|
var local = await _localActors.FindByUri(uri, token);
|
|
var member = local != default && local.Kind == LocalActorKind.Person
|
|
? new GroupMember { AvatarId = local.Id }
|
|
: new GroupMember { AvatarId = uri, IsForeign = true };
|
|
if (members.All(m => m.AvatarId != member.AvatarId || m.IsForeign != member.IsForeign))
|
|
members.Add(member);
|
|
}
|
|
|
|
var key = DmGroup.KeyOf(members);
|
|
var match = await _dbEntities.DmGroups.Match(g => g.ParticipantsKey == key && !g.DeletionAt.HasValue).ExecuteFirstAsync(token);
|
|
if (match != default)
|
|
return match;
|
|
|
|
var conversation = new DmGroup { Members = members, ConversationURI = context, ParticipantsKey = key };
|
|
await DB.Default.SaveAsync(conversation, token);
|
|
return conversation;
|
|
}
|
|
|
|
static bool IsCollection(string uri) =>
|
|
uri.EndsWith("/followers", StringComparison.Ordinal) || uri.EndsWith("/following", StringComparison.Ordinal);
|
|
|
|
static IEnumerable<string> Strings(JsonNode node) => node switch
|
|
{
|
|
JsonArray array => array.Select(Id).Where(id => id != default),
|
|
JsonNode single when Id(single) is { } id => new[] { id },
|
|
_ => Enumerable.Empty<string>()
|
|
};
|
|
}
|
|
}
|