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
113 lines
4.6 KiB
C#
113 lines
4.6 KiB
C#
using MongoDB.Entities;
|
|
|
|
using PrivaPub.Domain.Statuses;
|
|
using PrivaPub.Federation.Actors;
|
|
using PrivaPub.Federation.Objects;
|
|
using PrivaPub.Models.Federation;
|
|
using PrivaPub.Models.Post;
|
|
using PrivaPub.Models.Social;
|
|
|
|
using PostEntity = PrivaPub.Models.Post.Post;
|
|
|
|
namespace PrivaPub.Federation.Inbox
|
|
{
|
|
public static class RemoteEdits
|
|
{
|
|
public const int MaxRevisions = 20;
|
|
|
|
//newer than what we hold; or, for a first edit, no older and actually different, because GoToSocial's timestamps are
|
|
//whole seconds and an edit made in the second the post was published carries updated == published
|
|
public static bool IsEdit(NoteDocument note, PostEntity post) =>
|
|
note.Updated is { } updated
|
|
&& (updated > (post.EditedAt ?? post.CreationDate)
|
|
|| post.EditedAt == default && updated >= post.CreationDate && Changed(note, post));
|
|
|
|
static bool Changed(NoteDocument note, PostEntity post) =>
|
|
note.ContentHtml != post.ContentHtml || note.SpoilerText != post.SpoilerText || note.Title != post.Title || note.Sensitive != post.HasContentWarning;
|
|
|
|
public static async Task<bool> Apply(PostEntity post, NoteDocument note, string activityId, ILocalActorService localActors, IObjectRecords records,
|
|
IQuoteService quotes, CancellationToken token)
|
|
{
|
|
post.Poll = note.Poll ?? post.Poll;
|
|
post.QuotePolicy = note.QuotePolicy;
|
|
post.Video = note.Video ?? post.Video;
|
|
post.Audio = note.Audio ?? post.Audio;
|
|
post.Event = note.Event ?? post.Event;
|
|
post.Place = note.Place;
|
|
if (!IsEdit(note, post))
|
|
{
|
|
await DB.Default.SaveAsync(post, token);
|
|
await records.Revise(note, activityId, token);
|
|
await Requote(post, note, quotes, token);
|
|
return false;
|
|
}
|
|
|
|
post.Revisions.Add(new PostRevision
|
|
{
|
|
Title = post.Title,
|
|
SpoilerText = post.SpoilerText,
|
|
ContentHtml = post.ContentHtml,
|
|
HasContentWarning = post.HasContentWarning,
|
|
EditedAt = post.EditedAt ?? post.CreationDate
|
|
});
|
|
if (post.Revisions.Count > MaxRevisions)
|
|
post.Revisions.RemoveRange(0, post.Revisions.Count - MaxRevisions);
|
|
|
|
//an edit never narrows who a post was for: what it was addressed to before still holds
|
|
var mentions = await RemotePosts.Mentions(note, post.To.Concat(post.Cc).Concat(note.To).Concat(note.Cc), localActors, token);
|
|
|
|
post.Title = note.Title;
|
|
post.SpoilerText = note.SpoilerText;
|
|
post.Excerpt = note.Excerpt;
|
|
post.Source = note.Source;
|
|
post.Emojis = note.Emojis.ToList();
|
|
post.CoverURL = note.CoverURL;
|
|
post.Link = note.Link;
|
|
post.HasContentWarning = note.Sensitive;
|
|
post.Text = note.ContentHtml;
|
|
post.ContentHtml = note.ContentHtml;
|
|
post.Language = note.Language;
|
|
post.Tags = note.Tags.ToList();
|
|
post.Mentions = mentions;
|
|
post.Media = note.Attachments.ToList();
|
|
post.EditedAt = note.Updated ?? DateTime.UtcNow;
|
|
post.UpdateDate = DateTime.UtcNow;
|
|
await DB.Default.SaveAsync(post, token);
|
|
await Domain.Timelines.Fanout.Edited(post, token);
|
|
await records.Revise(note, activityId, token);
|
|
await Requote(post, note, quotes, token);
|
|
return true;
|
|
}
|
|
|
|
static async Task Requote(PostEntity post, NoteDocument note, IQuoteService quotes, CancellationToken token)
|
|
{
|
|
if (note.QuoteUri != post.QuoteURI || note.QuoteAuthorization != post.QuoteAuthorizationURI || post.QuoteState != QuoteState.Accepted)
|
|
await quotes.Resolve(post, note, token);
|
|
}
|
|
}
|
|
|
|
public static class RemoteDeletes
|
|
{
|
|
public static async Task Tombstone(string objectUri, CancellationToken token) =>
|
|
await DB.Default.Update<DeletedObject>()
|
|
.Match(d => d.ObjectURI == objectUri)
|
|
.Modify(d => d.DeletedAt, DateTime.UtcNow)
|
|
.Option(o => o.IsUpsert = true)
|
|
.ExecuteAsync(token);
|
|
|
|
public static async Task Remove(PostEntity post, string objectUri, CancellationToken token)
|
|
{
|
|
await Tombstone(objectUri, token);
|
|
await Domain.Timelines.Fanout.Deleting(post, token);
|
|
await DB.Default.DeleteAsync<PostEntity>(post.ID);
|
|
await DB.Default.DeleteAsync<ObjectRecord>(r => r.PostId == post.ID || r.ObjectURI == objectUri);
|
|
await DB.Default.DeleteAsync<TimelineEntry>(e => e.PostId == post.ID || e.ReblogOfPostId == post.ID);
|
|
await DB.Default.DeleteAsync<PostEntity>(p => p.ReblogOfPostId == post.ID);
|
|
if (!string.IsNullOrEmpty(post.AnsweringToPostId) && Domain.Privacy.Counted.Reply(post))
|
|
await DB.Default.Update<PostEntity>().MatchID(post.AnsweringToPostId).Modify(b => b.Inc(p => p.RepliesCount, -1)).ExecuteAsync(token);
|
|
if (post.QuoteState == QuoteState.Accepted && !string.IsNullOrEmpty(post.QuotedPostId))
|
|
await DB.Default.Update<PostEntity>().MatchID(post.QuotedPostId).Modify(b => b.Inc(p => p.QuotesCount, -1)).ExecuteAsync(token);
|
|
}
|
|
}
|
|
}
|