using MongoDB.Driver; using MongoDB.Entities; using PrivaPub.Infrastructure.Http; using PrivaPub.Domain.Content; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Moderation; using PrivaPub.Federation.Objects; using PrivaPub.Infrastructure.Ids; using PrivaPub.Infrastructure.Jobs; using PrivaPub.Models.Federation; 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; readonly IObjectRecords _records; readonly ILinkPreviews _previews; public RemotePosts(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDomainBlocks domainBlocks, IJobQueue queue, IObjectRecords records, ILinkPreviews previews) { _previews = previews; _records = records; _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.Arrived(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, ObjectType = note.Type, Title = note.Title, SpoilerText = note.SpoilerText, Excerpt = note.Excerpt, Source = note.Source, Poll = note.Poll, QuotePolicy = note.QuotePolicy, Emojis = note.Emojis.ToList(), CoverURL = note.CoverURL, Link = note.Link, Video = note.Video, Audio = note.Audio, Event = note.Event, 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 || await DB.Default.Find().Match(d => d.ObjectURI == objectUri).ExecuteAnyAsync(token)) return existing; using var scope = HttpScope.For("context"); 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); } await _records.Record(note, post, ObjectPath.Fetched, refetched: false, token); await _previews.Wanted(post, 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; } } }