Mastodon's streaming API: /api/v1/streaming as a WebSocket (streams
subscribed in the URL or by message) and /api/v1/streaming/{stream} as
server-sent events, with health and the URL advertised. The user stream
tells posts reaching the persona's home (not those an exclusive list keeps
apart, which its list stream tells), notifications, edits and deletions;
public, hashtag and list streams tell what belongs in them. An in-process
hub carries ids only; each connection maps a post or a notification for its
own persona as it sends it, so nothing it may not see, or whose author it
blocked or muted, goes out. A deletion reaches only the streams that showed
the post. The token comes as access_token, header or WebSocket protocol.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
190 lines
8.5 KiB
C#
190 lines
8.5 KiB
C#
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);
|
|
if (!string.IsNullOrEmpty(post.AudienceURI))
|
|
foreach (var member in await _dbEntities.Followings.Match(f => f.TargetActorURI == post.AudienceURI && f.State == FollowState.Accepted).ExecuteAsync(token))
|
|
recipients.Add(member.AvatarId);
|
|
// a public post (not a boost, not a reply) comes to those following one of its hashtags, as on Mastodon
|
|
if (post.Visibility == PostVisibility.Public && post.Tags.Count > 0 && string.IsNullOrEmpty(post.ReblogOfPostId)
|
|
&& string.IsNullOrEmpty(post.InReplyToURI) && string.IsNullOrEmpty(post.AnsweringToPostId))
|
|
{
|
|
var tags = post.Tags.ToList();
|
|
foreach (var followed in await DB.Default.Find<FollowedTag>().Match(t => tags.Contains(t.Name)).ExecuteAsync(token))
|
|
recipients.Add(followed.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;
|
|
}
|
|
|
|
recipients.ExceptWith(await Relationships.Hidden.RecipientsHiding(recipients, post.ActorURI, token));
|
|
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);
|
|
if (Streams.Anyone)
|
|
await Stream(post, recipients, token);
|
|
}
|
|
|
|
// the post to the streams it belongs in: the homes it reached (but not where an exclusive list keeps its author
|
|
// apart), the lists of those homes that hold its author, and the public feeds and hashtags when it is public
|
|
static async Task Stream(Post post, IReadOnlyCollection<string> recipients, CancellationToken token)
|
|
{
|
|
var update = new Streams.Event("update", post.ID);
|
|
var author = post.AuthorAccountId ?? post.GroupUserId;
|
|
var listening = recipients.Where(Streams.Listening).ToList();
|
|
if (listening.Count > 0)
|
|
{
|
|
var memberships = await DB.Default.Find<PersonaListMember>().Match(m => listening.Contains(m.AvatarId) && m.AccountId == author).ExecuteAsync(token);
|
|
var listIds = memberships.Select(m => m.ListId).Distinct().ToList();
|
|
var exclusive = listIds.Count == 0
|
|
? new HashSet<string>()
|
|
: (await DB.Default.Find<PersonaList>().Match(l => listIds.Contains(l.ID) && l.Exclusive).ExecuteAsync(token)).Select(l => l.ID).ToHashSet();
|
|
foreach (var avatarId in listening)
|
|
{
|
|
var lists = memberships.Where(m => m.AvatarId == avatarId).Select(m => m.ListId).ToList();
|
|
if (!lists.Any(exclusive.Contains))
|
|
Streams.ToPersona(avatarId, new Streams.Key(Streams.User), update);
|
|
foreach (var listId in lists)
|
|
Streams.ToPersona(avatarId, new Streams.Key(Streams.List, listId), update);
|
|
}
|
|
}
|
|
if (post.Visibility != PostVisibility.Public || !string.IsNullOrEmpty(post.ReblogOfPostId) || !Privacy.VisibilityPolicy.Shown(post))
|
|
return;
|
|
var local = !post.IsFederatedCopy;
|
|
if (!post.IsLocalOnly)
|
|
Streams.ToAll(new Streams.Key(Streams.Public), update);
|
|
Streams.ToAll(new Streams.Key(local ? Streams.PublicLocal : Streams.PublicRemote), update);
|
|
foreach (var tag in post.Tags.Distinct())
|
|
{
|
|
if (!post.IsLocalOnly)
|
|
Streams.ToAll(new Streams.Key(Streams.Hashtag, tag), update);
|
|
if (local)
|
|
Streams.ToAll(new Streams.Key(Streams.HashtagLocal, tag), update);
|
|
}
|
|
}
|
|
|
|
// a post's deletion, told before its home entries go: to the homes holding it or a boost of it (the boost's id
|
|
// there), to its author's own connections, and to the public feeds and hashtags when it was public. Nobody else
|
|
// hears of it, not even its id.
|
|
public static async Task Deleting(Post post, CancellationToken token)
|
|
{
|
|
if (!Streams.Anyone || post == default)
|
|
return;
|
|
var entries = await DB.Default.Find<TimelineEntry>().Match(e => e.PostId == post.ID || e.ReblogOfPostId == post.ID).ExecuteAsync(token);
|
|
foreach (var entry in entries.Where(e => Streams.Listening(e.AvatarId)))
|
|
Streams.ToPersonaEverywhere(entry.AvatarId, new Streams.Event("delete", entry.PostId));
|
|
var author = post.AuthorAccountId ?? post.GroupUserId;
|
|
if (!post.IsFederatedCopy && !string.IsNullOrEmpty(author))
|
|
Streams.ToPersonaEverywhere(author, new Streams.Event("delete", post.ID));
|
|
if (post.Visibility != PostVisibility.Public || !string.IsNullOrEmpty(post.ReblogOfPostId))
|
|
return;
|
|
var gone = new Streams.Event("delete", post.ID);
|
|
var local = !post.IsFederatedCopy;
|
|
Streams.ToAll(new Streams.Key(Streams.Public), gone);
|
|
Streams.ToAll(new Streams.Key(local ? Streams.PublicLocal : Streams.PublicRemote), gone);
|
|
foreach (var tag in post.Tags.Distinct())
|
|
{
|
|
Streams.ToAll(new Streams.Key(Streams.Hashtag, tag), gone);
|
|
if (local)
|
|
Streams.ToAll(new Streams.Key(Streams.HashtagLocal, tag), gone);
|
|
}
|
|
}
|
|
|
|
// an edit to the homes holding the post and, when it is public, to the public feeds (status.update)
|
|
public static async Task Edited(Post post, CancellationToken token)
|
|
{
|
|
if (!Streams.Anyone)
|
|
return;
|
|
var update = new Streams.Event("status.update", post.ID);
|
|
var homes = (await DB.Default.Find<TimelineEntry>().Match(e => e.PostId == post.ID).ExecuteAsync(token)).Select(e => e.AvatarId).Distinct();
|
|
foreach (var avatarId in homes.Where(Streams.Listening))
|
|
Streams.ToPersona(avatarId, new Streams.Key(Streams.User), update);
|
|
if (post.Visibility == PostVisibility.Public && Privacy.VisibilityPolicy.Shown(post))
|
|
{
|
|
if (!post.IsLocalOnly)
|
|
Streams.ToAll(new Streams.Key(Streams.Public), update);
|
|
Streams.ToAll(new Streams.Key(post.IsFederatedCopy ? Streams.PublicRemote : Streams.PublicLocal), update);
|
|
}
|
|
}
|
|
|
|
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);
|
|
}
|
|
}
|
|
}
|