From 796f06acf400c20ae52a03cd678a3a5d60c2855c Mon Sep 17 00:00:00 2001 From: thepra Date: Thu, 1 Oct 2026 11:12:34 +0200 Subject: [PATCH] The inbox answers once it has verified, and processes from the job queue InboxService is split the way the roadmap lays out Federation/Inbox: - InboxReceiver reads and verifies the request exactly as before, runs the checks that need no fetch (the activity id's origin, a Follow of a missing or local-only actor, an Undo of someone else's activity, an embedded object attributed to someone else), queues a ProcessInbox job and answers 202. The job's dedupe key is the activity id, so a peer that delivers the same activity twice is processed once. - InboxProcessor (two at a time, eight attempts) loads the verified actor and hands the activity to the handler for its type. - Handlers/{Follow,Undo,Create,Delete,Update}Handler are the old methods, unchanged except that they no longer produce status codes; the JSON helpers live in Objects/ActivityJson and the group membership helpers in Inbox/ForeignMembers. A slow fetch of an object or a remote actor now delays the job, not the sender's HTTP request. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB --- .../Federation/InboxScenarioTests.cs | 51 ++- .../Controllers/PeasantsController.cs | 4 +- PrivaPub/Federation/Inbox/ForeignMembers.cs | 37 ++ .../Inbox/Handlers/CreateHandler.cs | 169 +++++++ .../Inbox/Handlers/DeleteHandler.cs | 68 +++ .../Inbox/Handlers/FollowHandler.cs | 75 ++++ .../Federation/Inbox/Handlers/UndoHandler.cs | 70 +++ .../Inbox/Handlers/UpdateHandler.cs | 70 +++ PrivaPub/Federation/Inbox/InboxProcessor.cs | 54 +++ PrivaPub/Federation/Inbox/InboxReceiver.cs | 126 ++++++ PrivaPub/Federation/Inbox/InboxService.cs | 413 ------------------ PrivaPub/Federation/Objects/ActivityJson.cs | 39 ++ .../Middleware/SocialPubConfigurations.cs | 9 +- 13 files changed, 755 insertions(+), 430 deletions(-) create mode 100644 PrivaPub/Federation/Inbox/ForeignMembers.cs create mode 100644 PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs create mode 100644 PrivaPub/Federation/Inbox/Handlers/DeleteHandler.cs create mode 100644 PrivaPub/Federation/Inbox/Handlers/FollowHandler.cs create mode 100644 PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs create mode 100644 PrivaPub/Federation/Inbox/Handlers/UpdateHandler.cs create mode 100644 PrivaPub/Federation/Inbox/InboxProcessor.cs create mode 100644 PrivaPub/Federation/Inbox/InboxReceiver.cs delete mode 100644 PrivaPub/Federation/Inbox/InboxService.cs create mode 100644 PrivaPub/Federation/Objects/ActivityJson.cs diff --git a/PrivaPub.Tests/Federation/InboxScenarioTests.cs b/PrivaPub.Tests/Federation/InboxScenarioTests.cs index 40779ed..1ef2ead 100644 --- a/PrivaPub.Tests/Federation/InboxScenarioTests.cs +++ b/PrivaPub.Tests/Federation/InboxScenarioTests.cs @@ -5,6 +5,8 @@ using MongoDB.Entities; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Inbox; +using PrivaPub.Federation.Inbox.Handlers; +using PrivaPub.Models.Jobs; using PrivaPub.Federation.Outbox; using PrivaPub.Infrastructure.Jobs; using PrivaPub.Models; @@ -28,7 +30,8 @@ namespace PrivaPub.Tests.Federation Peer _peer; LocalActorService _local; - InboxService _inbox; + InboxReceiver _receiver; + InboxProcessor _processor; public async ValueTask InitializeAsync() { @@ -37,7 +40,18 @@ namespace PrivaPub.Tests.Federation var cache = new MemoryCache(new MemoryCacheOptions()); _local = new LocalActorService(new DbEntities(), new StaticOptions(new AppConfiguration { BackendBaseAddress = Base })); var remote = new RemoteActorService(Peer.Http(cache), _local, cache, new DbEntities()); - _inbox = new InboxService(new DbEntities(), _local, remote, new DeliveryService(new DbEntities(), new JobQueue()), NullLogger.Instance); + var queue = new JobQueue(); + var delivery = new DeliveryService(new DbEntities(), queue); + var db = new DbEntities(); + _receiver = new InboxReceiver(_local, remote, queue, NullLogger.Instance); + _processor = new InboxProcessor(remote, new IActivityHandler[] + { + new FollowHandler(db, _local, remote, delivery), + new UndoHandler(db, _local, remote, delivery), + new CreateHandler(db, _local, remote, delivery), + new DeleteHandler(db, _local, remote, delivery), + new UpdateHandler(db, _local, remote, delivery) + }, NullLogger.Instance); } public async ValueTask DisposeAsync() @@ -78,6 +92,17 @@ namespace PrivaPub.Tests.Federation }; } + async Task Deliver(RemoteActor sender, string path, JsonNode activity) + { + var token = TestContext.Current.CancellationToken; + var result = await _receiver.Receive(sender.Post(Host, path, activity), default, token); + var dedupe = "inbox|" + (activity is JsonObject ? activity["id"]?.GetValue() : default); + var job = await DB.Default.Find().Match(j => j.DedupeKey == dedupe).ExecuteFirstAsync(token); + if (job != default) + Assert.Equal(JobResult.Done, (await _processor.Handle(job, token)).Result); + return result; + } + static string Origin(string uri) => new Uri(uri).GetLeftPart(UriPartial.Authority); [Fact] @@ -89,8 +114,8 @@ namespace PrivaPub.Tests.Federation var mallory = new RemoteActor(_peer, "mallory"); var context = $"{_peer.A}/contexts/{Guid.NewGuid():N}"; - var first = await _inbox.Receive(bob.Post(Host, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri, context)), alice, token); - var injected = await _inbox.Receive(mallory.Post(Host, $"/peasants/{alice.UserName}/mouth", DirectCreate(mallory, alice.Uri, context)), alice, token); + var first = await Deliver(bob, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri, context)); + var injected = await Deliver(mallory, $"/peasants/{alice.UserName}/mouth", DirectCreate(mallory, alice.Uri, context)); Assert.Equal(202, first.StatusCode); Assert.Equal(202, injected.StatusCode); @@ -108,8 +133,8 @@ namespace PrivaPub.Tests.Federation var alice = await LocalAvatar("alice"); var bob = new RemoteActor(_peer, "bob"); - await _inbox.Receive(bob.Post(Host, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri)), alice, token); - await _inbox.Receive(bob.Post(Host, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri)), alice, token); + await Deliver(bob, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri)); + await Deliver(bob, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri)); var dms = await DB.Default.Find().Match(p => p.ActorURI == bob.Id).ExecuteAsync(token); Assert.Equal(2, dms.Count); @@ -125,7 +150,7 @@ namespace PrivaPub.Tests.Federation var create = DirectCreate(mallory, alice.Uri); create["id"] = $"{_peer.B}/activities/{Guid.NewGuid():N}"; - var result = await _inbox.Receive(mallory.Post(Host, $"/peasants/{alice.UserName}/mouth", create), alice, token); + var result = await Deliver(mallory, $"/peasants/{alice.UserName}/mouth", create); Assert.Equal(400, result.StatusCode); } @@ -138,8 +163,7 @@ namespace PrivaPub.Tests.Federation var mallory = new RemoteActor(_peer, "mallory"); var victim = new RemoteActor(_peer, "victim"); - var result = await _inbox.Receive(mallory.Post(Host, $"/peasants/{alice.UserName}/mouth", - DirectCreate(mallory, alice.Uri, attributedTo: victim.Id)), alice, token); + var result = await Deliver(mallory, $"/peasants/{alice.UserName}/mouth", DirectCreate(mallory, alice.Uri, attributedTo: victim.Id)); Assert.Equal(400, result.StatusCode); Assert.False(await DB.Default.Find().Match(p => p.ActorURI == victim.Id).ExecuteAnyAsync(token)); @@ -152,8 +176,7 @@ namespace PrivaPub.Tests.Federation var alice = await LocalAvatar("alice"); var mallory = new RemoteActor(_peer, "mallory"); - var result = await _inbox.Receive(mallory.Post(Host, $"/peasants/{alice.UserName}/mouth", - DirectCreate(mallory, alice.Uri, objectOrigin: _peer.B)), alice, token); + var result = await Deliver(mallory, $"/peasants/{alice.UserName}/mouth", DirectCreate(mallory, alice.Uri, objectOrigin: _peer.B)); Assert.Equal(202, result.StatusCode); Assert.False(await DB.Default.Find().Match(p => p.ActorURI == mallory.Id).ExecuteAnyAsync(token)); @@ -168,7 +191,7 @@ namespace PrivaPub.Tests.Federation var request = mallory.Post(Host, $"/peasants/{alice.UserName}/mouth", DirectCreate(mallory, alice.Uri)); request.Headers["Signature"] = request.Headers["Signature"].ToString().Replace("signature=\"", "signature=\"AAAA"); - Assert.Equal(401, (await _inbox.Receive(request, alice, token)).StatusCode); + Assert.Equal(401, (await _receiver.Receive(request, alice, token)).StatusCode); } [Fact] @@ -179,7 +202,7 @@ namespace PrivaPub.Tests.Federation var mallory = new RemoteActor(_peer, "mallory"); foreach (var junk in new JsonNode[] { new JsonArray(1, 2), JsonValue.Create("x"), new JsonObject { ["type"] = "Create" } }) - Assert.Equal(400, (await _inbox.Receive(mallory.Post(Host, $"/peasants/{alice.UserName}/mouth", junk), alice, token)).StatusCode); + Assert.Equal(400, (await Deliver(mallory, $"/peasants/{alice.UserName}/mouth", junk)).StatusCode); } [Fact] @@ -200,7 +223,7 @@ namespace PrivaPub.Tests.Federation }; Assert.False(actor.IsFederated); - Assert.Equal(404, (await _inbox.Receive(bob.Post(Host, "/human-centipede", follow), default, token)).StatusCode); + Assert.Equal(404, (await Deliver(bob, "/human-centipede", follow)).StatusCode); } } } diff --git a/PrivaPub/Federation/Controllers/PeasantsController.cs b/PrivaPub/Federation/Controllers/PeasantsController.cs index a46708c..b0521d0 100644 --- a/PrivaPub/Federation/Controllers/PeasantsController.cs +++ b/PrivaPub/Federation/Controllers/PeasantsController.cs @@ -23,11 +23,11 @@ namespace PrivaPub.Federation.Controllers const int OutboxSize = 20; readonly ILocalActorService _localActors; - readonly IInboxService _inbox; + readonly IInboxReceiver _inbox; readonly DbEntities _dbEntities; readonly ILogger _logger; - public PeasantsController(ILocalActorService localActors, IInboxService inbox, DbEntities dbEntities, + public PeasantsController(ILocalActorService localActors, IInboxReceiver inbox, DbEntities dbEntities, ILogger logger) { _localActors = localActors; diff --git a/PrivaPub/Federation/Inbox/ForeignMembers.cs b/PrivaPub/Federation/Inbox/ForeignMembers.cs new file mode 100644 index 0000000..e7b9f96 --- /dev/null +++ b/PrivaPub/Federation/Inbox/ForeignMembers.cs @@ -0,0 +1,37 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Group; + +using GroupEntity = PrivaPub.Models.Group.Group; + +namespace PrivaPub.Federation.Inbox +{ + public static class ForeignMembers + { + public static async Task IsAcceptedFollower(LocalActor group, string actorUri, CancellationToken token) => + await DB.Default.Find() + .Match(f => f.LocalActorId == group.Id && f.LocalActorKind == LocalActorKind.Group && f.ActorURI == actorUri && f.IsAccepted) + .ExecuteAnyAsync(token); + + public static async Task AddForeignMember(string groupId, string actorUri, CancellationToken token) + { + var group = await DB.Default.Find().MatchID(groupId).ExecuteFirstAsync(token); + if (group == default || group.Members.Any(m => m.IsForeign && m.AvatarId == actorUri)) + return; + group.Members.Add(new GroupMember { AvatarId = actorUri, IsForeign = true }); + group.UpdatedAt = DateTime.UtcNow; + await DB.Default.SaveAsync(group, token); + } + + public static async Task RemoveForeignMember(string groupId, string actorUri, CancellationToken token) + { + var group = await DB.Default.Find().MatchID(groupId).ExecuteFirstAsync(token); + if (group == default || group.Members.RemoveAll(m => m.IsForeign && m.AvatarId == actorUri) == 0) + return; + group.UpdatedAt = DateTime.UtcNow; + await DB.Default.SaveAsync(group, token); + } + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs b/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs new file mode 100644 index 0000000..9b1df55 --- /dev/null +++ b/PrivaPub/Federation/Inbox/Handlers/CreateHandler.cs @@ -0,0 +1,169 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Outbox; +using PrivaPub.Federation.Rendering; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Group; +using PrivaPub.Models.Post; +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 DmPostEntity = PrivaPub.Models.Post.DmPost; +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; + + public CreateHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery) + { + _dbEntities = dbEntities; + _localActors = localActors; + _remoteActors = remoteActors; + _delivery = delivery; + } + + public string Type => "Create"; + + public async Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token) + { + var create = activity; + var author = actor; + + var note = create["object"]; + if (note is not JsonObject || !Origin.Same(Id(note), author.ActorURI)) + { + using var fetched = await _remoteActors.FetchObject(Id(note), token); + note = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()); + } + if (note is not JsonObject || Value(note, "type") is not ("Note" or "Article" or "Page" or "Question")) + return; + + var objectUri = Id(note); + if (string.IsNullOrEmpty(objectUri) || Id(note["attributedTo"]) != author.ActorURI) + return; + + var addressed = Addresses(note).Concat(Addresses(create)).Distinct(StringComparer.Ordinal).ToList(); + var isPublic = addressed.Contains(ActivityPubRenderer.Public) || addressed.Contains("as:Public") || addressed.Contains("Public"); + var inReplyTo = Id(note["inReplyTo"]); + + var localTargets = new List(); + foreach (var uri in addressed.Concat(new[] { Id(note["audience"]) }).Where(u => u != default).Distinct()) + { + var local = await _localActors.FindByUri(uri, token); + if (local is { IsFederated: true } && localTargets.All(l => l.Id != local.Id)) + localTargets.Add(local); + } + + if (isPublic || addressed.Any(a => a == author.ActorURI + "/followers" || a.EndsWith("/followers"))) + { + var group = localTargets.FirstOrDefault(t => t.Kind == LocalActorKind.Group); + if (group != default && !await IsAcceptedFollower(group, author.ActorURI, token)) + group = default; + if (group == default && localTargets.All(t => t.Kind != LocalActorKind.Person)) + return; + + if (await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(token)) + return; + + var html = ContentSanitizer.Html(Value(note, "content")); + var post = new PostEntity + { + ObjectURI = objectUri, + ActorURI = author.ActorURI, + GroupId = group?.Id, + Title = Value(note, "summary") ?? Value(note, "name"), + Text = html, + ContentHtml = html, + ContentFormat = ContentFormat.Html, + HasContentWarning = note["sensitive"] is JsonValue sensitive && sensitive.TryGetValue(out var s) && s, + AnsweringToPostId = await LocalPostId(inReplyTo, token) ?? inReplyTo, + IsFederatedCopy = true, + CreationDate = DateTime.TryParse(Value(note, "published"), out var published) ? published.ToUniversalTime() : DateTime.UtcNow + }; + await DB.Default.SaveAsync(post, token); + + if (group != default) + { + var announce = ActivityPubRenderer.Announce(group, objectUri, $"announce-{post.ID}"); + await _delivery.EnqueueToFollowers(group, announce, token); + } + return; + } + + var recipients = localTargets.Where(t => t.Kind == LocalActorKind.Person).ToList(); + if (recipients.Count == 0) + return; + if (await _dbEntities.DmPosts.Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(token)) + return; + + var participants = recipients.Select(r => r.Uri) + .Concat(addressed.Where(a => a != ActivityPubRenderer.Public && !a.EndsWith("/followers") && !a.EndsWith("/following"))) + .Append(author.ActorURI) + .ToList(); + var context = Value(note, "context") ?? Value(note, "conversation"); + var dmGroup = await FindOrCreateDmGroup(participants, Origin.Same(context, author.ActorURI) ? context : default, token); + var dmHtml = ContentSanitizer.Html(Value(note, "content")); + var dm = new DmPostEntity + { + ObjectURI = objectUri, + ActorURI = author.ActorURI, + GroupId = dmGroup.ID, + Title = Value(note, "summary"), + Text = dmHtml, + ContentHtml = dmHtml, + ContentFormat = ContentFormat.Html, + HasContentWarning = note["sensitive"] is JsonValue dmSensitive && dmSensitive.TryGetValue(out var ds) && ds, + AnsweringToPostId = inReplyTo, + IsFederatedCopy = true, + CreationDate = DateTime.TryParse(Value(note, "published"), out var dmPublished) ? dmPublished.ToUniversalTime() : DateTime.UtcNow + }; + await DB.Default.SaveAsync(dm, token); + await DB.Default.Update().MatchID(dmGroup.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token); + return; + } + +async Task FindOrCreateDmGroup(List participantUris, string context, CancellationToken token) + { + var members = new List(); + 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 dmGroup = new DmGroup { Members = members, ConversationURI = context, ParticipantsKey = key }; + await DB.Default.SaveAsync(dmGroup, token); + return dmGroup; + } + +async Task LocalPostId(string uri, CancellationToken token) + { + if (string.IsNullOrEmpty(uri) || !uri.StartsWith(_localActors.BaseAddress + "/peasants/", StringComparison.OrdinalIgnoreCase)) + return default; + var id = uri[(uri.LastIndexOf('/') + 1)..]; + return await _dbEntities.Posts.MatchID(id).ExecuteAnyAsync(token) ? id : default; + } + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/DeleteHandler.cs b/PrivaPub/Federation/Inbox/Handlers/DeleteHandler.cs new file mode 100644 index 0000000..2a6e6f9 --- /dev/null +++ b/PrivaPub/Federation/Inbox/Handlers/DeleteHandler.cs @@ -0,0 +1,68 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Outbox; +using PrivaPub.Federation.Rendering; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Group; +using PrivaPub.Models.Post; +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 DmPostEntity = PrivaPub.Models.Post.DmPost; +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Federation.Inbox.Handlers +{ + public class DeleteHandler : IActivityHandler + { + readonly DbEntities _dbEntities; + readonly ILocalActorService _localActors; + readonly IRemoteActorService _remoteActors; + readonly IDeliveryService _delivery; + + public DeleteHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery) + { + _dbEntities = dbEntities; + _localActors = localActors; + _remoteActors = remoteActors; + _delivery = delivery; + } + + public string Type => "Delete"; + + public async Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token) + { + var delete = activity; + + var objectUri = Id(delete["object"]); + if (!Origin.Same(objectUri, actor.ActorURI)) + return; + + if (objectUri == actor.ActorURI) + { + actor.DeletionAt = DateTime.UtcNow; + actor.AccountState = AvatarAccountState.Deleted; + await DB.Default.SaveAsync(actor, token); + var followers = await _dbEntities.Followers.Match(f => f.ActorURI == actor.ActorURI).ExecuteAsync(token); + foreach (var follower in followers) + { + await DB.Default.DeleteAsync(follower.ID); + if (follower.LocalActorKind == LocalActorKind.Group) + await RemoveForeignMember(follower.LocalActorId, actor.ActorURI, token); + } + return; + } + + await DB.Default.DeleteAsync(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI); + await DB.Default.DeleteAsync(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI); + return; + } + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/FollowHandler.cs b/PrivaPub/Federation/Inbox/Handlers/FollowHandler.cs new file mode 100644 index 0000000..f015fe0 --- /dev/null +++ b/PrivaPub/Federation/Inbox/Handlers/FollowHandler.cs @@ -0,0 +1,75 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Outbox; +using PrivaPub.Federation.Rendering; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Group; +using PrivaPub.Models.Post; +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 DmPostEntity = PrivaPub.Models.Post.DmPost; +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Federation.Inbox.Handlers +{ + public class FollowHandler : IActivityHandler + { + readonly DbEntities _dbEntities; + readonly ILocalActorService _localActors; + readonly IRemoteActorService _remoteActors; + readonly IDeliveryService _delivery; + + public FollowHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery) + { + _dbEntities = dbEntities; + _localActors = localActors; + _remoteActors = remoteActors; + _delivery = delivery; + } + + public string Type => "Follow"; + + public async Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token) + { + var follow = activity; + var follower = actor; + + var target = await _localActors.FindByUri(Id(follow["object"]), token); + if (target is not { IsFederated: true } || target.Kind == LocalActorKind.Application) + return; + + var existing = await _dbEntities.Followers + .Match(f => f.LocalActorId == target.Id && f.LocalActorKind == target.Kind && f.ActorURI == follower.ActorURI) + .ExecuteFirstAsync(token); + var record = existing ?? new Follower + { + LocalActorId = target.Id, + LocalActorKind = target.Kind, + ActorURI = follower.ActorURI + }; + record.InboxURL = follower.InboxURL; + record.SharedInboxURL = follower.SharedInboxURL; + record.FollowActivityURI = Id(follow); + record.IsAccepted = existing?.IsAccepted == true || !target.ManuallyApprovesFollowers; + await DB.Default.SaveAsync(record, token); + + if (!record.IsAccepted) + return; + + if (target.Kind == LocalActorKind.Group) + await AddForeignMember(target.Id, follower.ActorURI, token); + + var accept = ActivityPubRenderer.Accept(target, follow, $"accept-{record.ID}-{DateTime.UtcNow.Ticks}"); + await _delivery.Enqueue(target, new[] { follower.InboxURL }, accept, token); + return; + } + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs b/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs new file mode 100644 index 0000000..d955000 --- /dev/null +++ b/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs @@ -0,0 +1,70 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Outbox; +using PrivaPub.Federation.Rendering; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Group; +using PrivaPub.Models.Post; +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 DmPostEntity = PrivaPub.Models.Post.DmPost; +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Federation.Inbox.Handlers +{ + public class UndoHandler : IActivityHandler + { + readonly DbEntities _dbEntities; + readonly ILocalActorService _localActors; + readonly IRemoteActorService _remoteActors; + readonly IDeliveryService _delivery; + + public UndoHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery) + { + _dbEntities = dbEntities; + _localActors = localActors; + _remoteActors = remoteActors; + _delivery = delivery; + } + + public string Type => "Undo"; + + public async Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token) + { + var undo = activity; + + var inner = undo["object"]; + var innerType = inner is JsonObject ? Value(inner, "type") : default; + var innerId = Id(inner); + + if (inner is JsonObject && Id(inner["actor"]) != actor.ActorURI) + return; + if (innerType is not (null or "Follow")) + return; + + var followers = await _dbEntities.Followers.Match(f => f.ActorURI == actor.ActorURI).ExecuteAsync(token); + var targetUri = inner is JsonObject ? Id(inner["object"]) : default; + foreach (var follower in followers) + { + var matchesActivity = innerId != default && follower.FollowActivityURI == innerId; + var target = targetUri == default ? default : await _localActors.FindByUri(targetUri, token); + var matchesTarget = target != default && target.Id == follower.LocalActorId && target.Kind == follower.LocalActorKind; + if (!matchesActivity && !matchesTarget) + continue; + + await DB.Default.DeleteAsync(follower.ID); + if (follower.LocalActorKind == LocalActorKind.Group) + await RemoveForeignMember(follower.LocalActorId, actor.ActorURI, token); + } + return; + } + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/UpdateHandler.cs b/PrivaPub/Federation/Inbox/Handlers/UpdateHandler.cs new file mode 100644 index 0000000..92daf6d --- /dev/null +++ b/PrivaPub/Federation/Inbox/Handlers/UpdateHandler.cs @@ -0,0 +1,70 @@ +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Outbox; +using PrivaPub.Federation.Rendering; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Group; +using PrivaPub.Models.Post; +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 DmPostEntity = PrivaPub.Models.Post.DmPost; +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Federation.Inbox.Handlers +{ + public class UpdateHandler : IActivityHandler + { + readonly DbEntities _dbEntities; + readonly ILocalActorService _localActors; + readonly IRemoteActorService _remoteActors; + readonly IDeliveryService _delivery; + + public UpdateHandler(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, IDeliveryService delivery) + { + _dbEntities = dbEntities; + _localActors = localActors; + _remoteActors = remoteActors; + _delivery = delivery; + } + + public string Type => "Update"; + + public async Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token) + { + var update = activity; + + var inner = update["object"]; + if (Id(inner) == actor.ActorURI) + { + await _remoteActors.GetActor(actor.ActorURI, refresh: true, token); + return; + } + if (inner is not JsonObject || Id(inner["attributedTo"]) != actor.ActorURI || !Origin.Same(Id(inner), actor.ActorURI)) + return; + + var objectUri = Id(inner); + var text = ContentSanitizer.Html(Value(inner, "content")); + await DB.Default.Update() + .Match(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI) + .Modify(p => p.Text, text) + .Modify(p => p.ContentHtml, text) + .Modify(p => p.UpdateDate, DateTime.UtcNow) + .ExecuteAsync(token); + await DB.Default.Update() + .Match(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI) + .Modify(p => p.Text, text) + .Modify(p => p.ContentHtml, text) + .Modify(p => p.UpdateDate, DateTime.UtcNow) + .ExecuteAsync(token); + return; + } + } +} diff --git a/PrivaPub/Federation/Inbox/InboxProcessor.cs b/PrivaPub/Federation/Inbox/InboxProcessor.cs new file mode 100644 index 0000000..1738950 --- /dev/null +++ b/PrivaPub/Federation/Inbox/InboxProcessor.cs @@ -0,0 +1,54 @@ +using PrivaPub.Federation.Actors; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.User; + +using System.Text.Json; +using System.Text.Json.Nodes; + +using static PrivaPub.Federation.Objects.ActivityJson; + +namespace PrivaPub.Federation.Inbox +{ + public interface IActivityHandler + { + string Type { get; } + Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token); + } + + public class InboxProcessor : IJobHandler + { + readonly IRemoteActorService _remoteActors; + readonly IReadOnlyDictionary _handlers; + readonly ILogger _logger; + + public InboxProcessor(IRemoteActorService remoteActors, IEnumerable handlers, ILogger logger) + { + _remoteActors = remoteActors; + _handlers = handlers.ToDictionary(h => h.Type, StringComparer.Ordinal); + _logger = logger; + } + + public JobKind Kind => JobKind.ProcessInbox; + public int Concurrency => 2; + public int MaxAttempts => 8; + public int PerHostLimit => 2; + + public async Task Handle(Job job, CancellationToken token) + { + var payload = JsonSerializer.Deserialize(job.Payload); + var activity = payload == default ? default : JsonNode.Parse(payload.Activity); + var type = Value(activity, "type"); + if (type == default || !_handlers.TryGetValue(type, out var handler)) + return JobOutcome.Done; + + var actor = await _remoteActors.GetActor(payload.ActorURI, refresh: false, token); + if (actor == default) + return JobOutcome.Retry("the actor could not be loaded"); + + await handler.Handle(activity, actor, token); + _logger.LogInformation("Processed {Type} {Id} from {Actor}", type, Id(activity), actor.ActorURI); + return JobOutcome.Done; + } + } +} diff --git a/PrivaPub/Federation/Inbox/InboxReceiver.cs b/PrivaPub/Federation/Inbox/InboxReceiver.cs new file mode 100644 index 0000000..404fb43 --- /dev/null +++ b/PrivaPub/Federation/Inbox/InboxReceiver.cs @@ -0,0 +1,126 @@ +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Signing; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Jobs; + +using System.Text.Json; +using System.Text.Json.Nodes; + +using static PrivaPub.Federation.Objects.ActivityJson; + +namespace PrivaPub.Federation.Inbox +{ + public sealed record InboxResult(int StatusCode, string Error = default); + + public sealed record InboxPayload(string ActorURI, string Activity); + + public interface IInboxReceiver + { + Task Receive(HttpRequest request, LocalActor recipient, CancellationToken token); + } + + public class InboxReceiver : IInboxReceiver + { + const int MaxBodyBytes = 1024 * 1024; + + readonly ILocalActorService _localActors; + readonly IRemoteActorService _remoteActors; + readonly IJobQueue _queue; + readonly ILogger _logger; + + public InboxReceiver(ILocalActorService localActors, IRemoteActorService remoteActors, IJobQueue queue, ILogger logger) + { + _localActors = localActors; + _remoteActors = remoteActors; + _queue = queue; + _logger = logger; + } + + public async Task Receive(HttpRequest request, LocalActor recipient, CancellationToken token) + { + if (request.ContentLength > MaxBodyBytes) + return new(StatusCodes.Status413PayloadTooLarge); + + using var buffer = new MemoryStream(); + await request.Body.CopyToAsync(buffer, token); + if (buffer.Length > MaxBodyBytes) + return new(StatusCodes.Status413PayloadTooLarge); + var body = buffer.ToArray(); + + JsonNode activity; + try + { + activity = JsonNode.Parse(body); + } + catch (JsonException) + { + return new(StatusCodes.Status400BadRequest, "the body is not JSON"); + } + if (activity is not JsonObject) + return new(StatusCodes.Status400BadRequest, "the body is not an activity"); + + var type = Value(activity, "type"); + var actorUri = Id(activity["actor"]); + if (string.IsNullOrEmpty(type) || string.IsNullOrEmpty(actorUri)) + return new(StatusCodes.Status400BadRequest, "type and actor are required"); + + var parameters = HttpSignatures.Parse(request.Headers["Signature"].ToString()); + if (parameters == default) + return new(StatusCodes.Status401Unauthorized, "missing or unreadable Signature header"); + + var requestProblem = HttpSignatures.CheckRequest(request, parameters, body); + if (requestProblem != default) + return new(StatusCodes.Status401Unauthorized, requestProblem); + + var signingString = HttpSignatures.SigningString(request, parameters); + var keyOwner = await _remoteActors.GetActorByKeyId(parameters.KeyId, refresh: false, token); + if (keyOwner == default && type == "Delete" && Id(activity["object"]) == actorUri) + return new(StatusCodes.Status202Accepted); + if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature)) + { + keyOwner = await _remoteActors.GetActorByKeyId(parameters.KeyId, refresh: true, token); + if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature)) + return new(StatusCodes.Status401Unauthorized, "the signature does not verify"); + } + + if (!string.Equals(keyOwner.ActorURI, actorUri, StringComparison.Ordinal)) + return new(StatusCodes.Status401Unauthorized, "the activity's actor is not the key's owner"); + + var shapeProblem = await ShapeProblem(type, activity, actorUri, token); + if (shapeProblem != default) + return shapeProblem; + + var activityId = Id(activity); + await _queue.Enqueue(JobKind.ProcessInbox, JsonSerializer.Serialize(new InboxPayload(actorUri, activity.ToJsonString())), + new Uri(actorUri).Host.ToLowerInvariant(), activityId == default ? default : "inbox|" + activityId, token); + _logger.LogInformation("Inbox {Recipient}: {Type} from {Actor} queued", recipient?.Handle ?? "shared", type, actorUri); + return new(StatusCodes.Status202Accepted); + } + + async Task ShapeProblem(string type, JsonNode activity, string actorUri, CancellationToken token) + { + var activityId = Id(activity); + if (activityId != default && !Origin.Same(activityId, actorUri)) + return new(StatusCodes.Status400BadRequest, "the activity's id is not on its actor's origin"); + + var inner = activity["object"]; + switch (type) + { + case "Follow": + var target = await _localActors.FindByUri(Id(inner), token); + if (target is not { IsFederated: true } || target.Kind == LocalActorKind.Application) + return new(StatusCodes.Status404NotFound, "no such local actor"); + break; + case "Undo" when inner is JsonObject && Id(inner["actor"]) != actorUri: + return new(StatusCodes.Status400BadRequest, "an actor can only undo its own activities"); + case "Create" or "Update" when inner is JsonObject && Origin.Same(Id(inner), actorUri) + && inner["attributedTo"] != default && Id(inner["attributedTo"]) != actorUri + && Value(inner, "type") is not ("Person" or "Service" or "Application" or "Group" or "Organization"): + return new(StatusCodes.Status400BadRequest, "the object is not attributed to the actor"); + } + return default; + } + } +} diff --git a/PrivaPub/Federation/Inbox/InboxService.cs b/PrivaPub/Federation/Inbox/InboxService.cs deleted file mode 100644 index 8748d32..0000000 --- a/PrivaPub/Federation/Inbox/InboxService.cs +++ /dev/null @@ -1,413 +0,0 @@ -using MongoDB.Entities; - -using PrivaPub.Models.Federation; -using PrivaPub.Models.Group; -using PrivaPub.Models.Post; -using PrivaPub.Models.User; -using PrivaPub.StaticServices; - -using System.Text.Json; -using System.Text.Json.Nodes; - -using DmPostEntity = PrivaPub.Models.Post.DmPost; -using GroupEntity = PrivaPub.Models.Group.Group; -using PostEntity = PrivaPub.Models.Post.Post; -using PrivaPub.Federation.Actors; -using PrivaPub.Federation.Signing; -using PrivaPub.Federation.Rendering; -using PrivaPub.Federation.Objects; -using PrivaPub.Federation.Outbox; - -namespace PrivaPub.Federation.Inbox -{ - public sealed record InboxResult(int StatusCode, string Error = default); - - public interface IInboxService - { - Task Receive(HttpRequest request, LocalActor recipient, CancellationToken token); - } - - public class InboxService : IInboxService - { - const int MaxBodyBytes = 1024 * 1024; - - readonly DbEntities _dbEntities; - readonly ILocalActorService _localActors; - readonly IRemoteActorService _remoteActors; - readonly IDeliveryService _delivery; - readonly ILogger _logger; - - public InboxService(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors, - IDeliveryService delivery, ILogger logger) - { - _dbEntities = dbEntities; - _localActors = localActors; - _remoteActors = remoteActors; - _delivery = delivery; - _logger = logger; - } - - public async Task Receive(HttpRequest request, LocalActor recipient, CancellationToken token) - { - if (request.ContentLength > MaxBodyBytes) - return new(StatusCodes.Status413PayloadTooLarge); - - using var buffer = new MemoryStream(); - await request.Body.CopyToAsync(buffer, token); - if (buffer.Length > MaxBodyBytes) - return new(StatusCodes.Status413PayloadTooLarge); - var body = buffer.ToArray(); - - JsonNode activity; - try - { - activity = JsonNode.Parse(body); - } - catch (JsonException) - { - return new(StatusCodes.Status400BadRequest, "the body is not JSON"); - } - if (activity is not JsonObject) - return new(StatusCodes.Status400BadRequest, "the body is not an activity"); - - var type = Value(activity, "type"); - var actorUri = Id(activity["actor"]); - if (string.IsNullOrEmpty(type) || string.IsNullOrEmpty(actorUri)) - return new(StatusCodes.Status400BadRequest, "type and actor are required"); - - var parameters = HttpSignatures.Parse(request.Headers["Signature"].ToString()); - if (parameters == default) - return new(StatusCodes.Status401Unauthorized, "missing or unreadable Signature header"); - - var requestProblem = HttpSignatures.CheckRequest(request, parameters, body); - if (requestProblem != default) - return new(StatusCodes.Status401Unauthorized, requestProblem); - - var signingString = HttpSignatures.SigningString(request, parameters); - var keyOwner = await _remoteActors.GetActorByKeyId(parameters.KeyId, refresh: false, token); - if (keyOwner == default && type == "Delete" && Id(activity["object"]) == actorUri) - return new(StatusCodes.Status202Accepted); - if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature)) - { - keyOwner = await _remoteActors.GetActorByKeyId(parameters.KeyId, refresh: true, token); - if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature)) - return new(StatusCodes.Status401Unauthorized, "the signature does not verify"); - } - - if (!string.Equals(keyOwner.ActorURI, actorUri, StringComparison.Ordinal)) - return new(StatusCodes.Status401Unauthorized, "the activity's actor is not the key's owner"); - - var activityId = Id(activity); - if (activityId != default && !Origin.Same(activityId, actorUri)) - return new(StatusCodes.Status400BadRequest, "the activity's id is not on its actor's origin"); - - _logger.LogInformation("Inbox {Recipient}: {Type} from {Actor}", recipient?.Handle ?? "shared", type, actorUri); - - return type switch - { - "Follow" => await Follow(activity, keyOwner, token), - "Undo" => await Undo(activity, keyOwner, token), - "Create" => await Create(activity, keyOwner, token), - "Delete" => await Delete(activity, keyOwner, token), - "Update" => await Update(activity, keyOwner, token), - _ => new(StatusCodes.Status202Accepted) - }; - } - - async Task Follow(JsonNode follow, ForeignAvatar follower, CancellationToken token) - { - var target = await _localActors.FindByUri(Id(follow["object"]), token); - if (target is not { IsFederated: true } || target.Kind == LocalActorKind.Application) - return new(StatusCodes.Status404NotFound, "no such local actor"); - - var existing = await _dbEntities.Followers - .Match(f => f.LocalActorId == target.Id && f.LocalActorKind == target.Kind && f.ActorURI == follower.ActorURI) - .ExecuteFirstAsync(token); - var record = existing ?? new Follower - { - LocalActorId = target.Id, - LocalActorKind = target.Kind, - ActorURI = follower.ActorURI - }; - record.InboxURL = follower.InboxURL; - record.SharedInboxURL = follower.SharedInboxURL; - record.FollowActivityURI = Id(follow); - record.IsAccepted = existing?.IsAccepted == true || !target.ManuallyApprovesFollowers; - await DB.Default.SaveAsync(record, token); - - if (!record.IsAccepted) - return new(StatusCodes.Status202Accepted); - - if (target.Kind == LocalActorKind.Group) - await AddForeignMember(target.Id, follower.ActorURI, token); - - var accept = ActivityPubRenderer.Accept(target, follow, $"accept-{record.ID}-{DateTime.UtcNow.Ticks}"); - await _delivery.Enqueue(target, new[] { follower.InboxURL }, accept, token); - return new(StatusCodes.Status202Accepted); - } - - async Task Undo(JsonNode undo, ForeignAvatar actor, CancellationToken token) - { - var inner = undo["object"]; - var innerType = inner is JsonObject ? Value(inner, "type") : default; - var innerId = Id(inner); - - if (inner is JsonObject && Id(inner["actor"]) != actor.ActorURI) - return new(StatusCodes.Status400BadRequest, "an actor can only undo its own activities"); - if (innerType is not (null or "Follow")) - return new(StatusCodes.Status202Accepted); - - var followers = await _dbEntities.Followers.Match(f => f.ActorURI == actor.ActorURI).ExecuteAsync(token); - var targetUri = inner is JsonObject ? Id(inner["object"]) : default; - foreach (var follower in followers) - { - var matchesActivity = innerId != default && follower.FollowActivityURI == innerId; - var target = targetUri == default ? default : await _localActors.FindByUri(targetUri, token); - var matchesTarget = target != default && target.Id == follower.LocalActorId && target.Kind == follower.LocalActorKind; - if (!matchesActivity && !matchesTarget) - continue; - - await DB.Default.DeleteAsync(follower.ID); - if (follower.LocalActorKind == LocalActorKind.Group) - await RemoveForeignMember(follower.LocalActorId, actor.ActorURI, token); - } - return new(StatusCodes.Status202Accepted); - } - - async Task Create(JsonNode create, ForeignAvatar author, CancellationToken token) - { - var note = create["object"]; - if (note is not JsonObject || !Origin.Same(Id(note), author.ActorURI)) - { - using var fetched = await _remoteActors.FetchObject(Id(note), token); - note = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()); - } - if (note is not JsonObject || Value(note, "type") is not ("Note" or "Article" or "Page" or "Question")) - return new(StatusCodes.Status202Accepted); - - var objectUri = Id(note); - if (string.IsNullOrEmpty(objectUri) || Id(note["attributedTo"]) != author.ActorURI) - return new(StatusCodes.Status400BadRequest, "the object is not attributed to the actor"); - - var addressed = Addresses(note).Concat(Addresses(create)).Distinct(StringComparer.Ordinal).ToList(); - var isPublic = addressed.Contains(ActivityPubRenderer.Public) || addressed.Contains("as:Public") || addressed.Contains("Public"); - var inReplyTo = Id(note["inReplyTo"]); - - var localTargets = new List(); - foreach (var uri in addressed.Concat(new[] { Id(note["audience"]) }).Where(u => u != default).Distinct()) - { - var local = await _localActors.FindByUri(uri, token); - if (local is { IsFederated: true } && localTargets.All(l => l.Id != local.Id)) - localTargets.Add(local); - } - - if (isPublic || addressed.Any(a => a == author.ActorURI + "/followers" || a.EndsWith("/followers"))) - { - var group = localTargets.FirstOrDefault(t => t.Kind == LocalActorKind.Group); - if (group != default && !await IsAcceptedFollower(group, author.ActorURI, token)) - group = default; - if (group == default && localTargets.All(t => t.Kind != LocalActorKind.Person)) - return new(StatusCodes.Status202Accepted); - - if (await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(token)) - return new(StatusCodes.Status202Accepted); - - var html = ContentSanitizer.Html(Value(note, "content")); - var post = new PostEntity - { - ObjectURI = objectUri, - ActorURI = author.ActorURI, - GroupId = group?.Id, - Title = Value(note, "summary") ?? Value(note, "name"), - Text = html, - ContentHtml = html, - ContentFormat = ContentFormat.Html, - HasContentWarning = note["sensitive"] is JsonValue sensitive && sensitive.TryGetValue(out var s) && s, - AnsweringToPostId = await LocalPostId(inReplyTo, token) ?? inReplyTo, - IsFederatedCopy = true, - CreationDate = DateTime.TryParse(Value(note, "published"), out var published) ? published.ToUniversalTime() : DateTime.UtcNow - }; - await DB.Default.SaveAsync(post, token); - - if (group != default) - { - var announce = ActivityPubRenderer.Announce(group, objectUri, $"announce-{post.ID}"); - await _delivery.EnqueueToFollowers(group, announce, token); - } - return new(StatusCodes.Status202Accepted); - } - - var recipients = localTargets.Where(t => t.Kind == LocalActorKind.Person).ToList(); - if (recipients.Count == 0) - return new(StatusCodes.Status202Accepted); - if (await _dbEntities.DmPosts.Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(token)) - return new(StatusCodes.Status202Accepted); - - var participants = recipients.Select(r => r.Uri) - .Concat(addressed.Where(a => a != ActivityPubRenderer.Public && !a.EndsWith("/followers") && !a.EndsWith("/following"))) - .Append(author.ActorURI) - .ToList(); - var context = Value(note, "context") ?? Value(note, "conversation"); - var dmGroup = await FindOrCreateDmGroup(participants, Origin.Same(context, author.ActorURI) ? context : default, token); - var dmHtml = ContentSanitizer.Html(Value(note, "content")); - var dm = new DmPostEntity - { - ObjectURI = objectUri, - ActorURI = author.ActorURI, - GroupId = dmGroup.ID, - Title = Value(note, "summary"), - Text = dmHtml, - ContentHtml = dmHtml, - ContentFormat = ContentFormat.Html, - HasContentWarning = note["sensitive"] is JsonValue dmSensitive && dmSensitive.TryGetValue(out var ds) && ds, - AnsweringToPostId = inReplyTo, - IsFederatedCopy = true, - CreationDate = DateTime.TryParse(Value(note, "published"), out var dmPublished) ? dmPublished.ToUniversalTime() : DateTime.UtcNow - }; - await DB.Default.SaveAsync(dm, token); - await DB.Default.Update().MatchID(dmGroup.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token); - return new(StatusCodes.Status202Accepted); - } - - async Task Delete(JsonNode delete, ForeignAvatar actor, CancellationToken token) - { - var objectUri = Id(delete["object"]); - if (!Origin.Same(objectUri, actor.ActorURI)) - return new(StatusCodes.Status202Accepted); - - if (objectUri == actor.ActorURI) - { - actor.DeletionAt = DateTime.UtcNow; - actor.AccountState = AvatarAccountState.Deleted; - await DB.Default.SaveAsync(actor, token); - var followers = await _dbEntities.Followers.Match(f => f.ActorURI == actor.ActorURI).ExecuteAsync(token); - foreach (var follower in followers) - { - await DB.Default.DeleteAsync(follower.ID); - if (follower.LocalActorKind == LocalActorKind.Group) - await RemoveForeignMember(follower.LocalActorId, actor.ActorURI, token); - } - return new(StatusCodes.Status202Accepted); - } - - await DB.Default.DeleteAsync(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI); - await DB.Default.DeleteAsync(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI); - return new(StatusCodes.Status202Accepted); - } - - async Task Update(JsonNode update, ForeignAvatar actor, CancellationToken token) - { - var inner = update["object"]; - if (Id(inner) == actor.ActorURI) - { - await _remoteActors.GetActor(actor.ActorURI, refresh: true, token); - return new(StatusCodes.Status202Accepted); - } - if (inner is not JsonObject || Id(inner["attributedTo"]) != actor.ActorURI || !Origin.Same(Id(inner), actor.ActorURI)) - return new(StatusCodes.Status202Accepted); - - var objectUri = Id(inner); - var text = ContentSanitizer.Html(Value(inner, "content")); - await DB.Default.Update() - .Match(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI) - .Modify(p => p.Text, text) - .Modify(p => p.ContentHtml, text) - .Modify(p => p.UpdateDate, DateTime.UtcNow) - .ExecuteAsync(token); - await DB.Default.Update() - .Match(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI) - .Modify(p => p.Text, text) - .Modify(p => p.ContentHtml, text) - .Modify(p => p.UpdateDate, DateTime.UtcNow) - .ExecuteAsync(token); - return new(StatusCodes.Status202Accepted); - } - - async Task IsAcceptedFollower(LocalActor group, string actorUri, CancellationToken token) => - await _dbEntities.Followers - .Match(f => f.LocalActorId == group.Id && f.LocalActorKind == LocalActorKind.Group && f.ActorURI == actorUri && f.IsAccepted) - .ExecuteAnyAsync(token); - - async Task AddForeignMember(string groupId, string actorUri, CancellationToken token) - { - var group = await _dbEntities.Groups.MatchID(groupId).ExecuteFirstAsync(token); - if (group == default || group.Members.Any(m => m.IsForeign && m.AvatarId == actorUri)) - return; - group.Members.Add(new GroupMember { AvatarId = actorUri, IsForeign = true }); - group.UpdatedAt = DateTime.UtcNow; - await DB.Default.SaveAsync(group, token); - } - - async Task RemoveForeignMember(string groupId, string actorUri, CancellationToken token) - { - var group = await _dbEntities.Groups.MatchID(groupId).ExecuteFirstAsync(token); - if (group == default || group.Members.RemoveAll(m => m.IsForeign && m.AvatarId == actorUri) == 0) - return; - group.UpdatedAt = DateTime.UtcNow; - await DB.Default.SaveAsync(group, token); - } - - async Task FindOrCreateDmGroup(List participantUris, string context, CancellationToken token) - { - var members = new List(); - 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 dmGroup = new DmGroup { Members = members, ConversationURI = context, ParticipantsKey = key }; - await DB.Default.SaveAsync(dmGroup, token); - return dmGroup; - } - - async Task LocalPostId(string uri, CancellationToken token) - { - if (string.IsNullOrEmpty(uri) || !uri.StartsWith(_localActors.BaseAddress + "/peasants/", StringComparison.OrdinalIgnoreCase)) - return default; - var id = uri[(uri.LastIndexOf('/') + 1)..]; - return await _dbEntities.Posts.MatchID(id).ExecuteAnyAsync(token) ? id : default; - } - - static IEnumerable Addresses(JsonNode node) - { - foreach (var field in new[] { "to", "cc", "bto", "bcc", "audience" }) - { - switch (node[field]) - { - case JsonArray array: - foreach (var item in array) - { - var id = Id(item); - if (id != default) - yield return id; - } - break; - case JsonNode single when Id(single) is { } id: - yield return id; - break; - } - } - } - - public static string Id(JsonNode node) => node switch - { - JsonValue value when value.TryGetValue(out var text) => text, - JsonObject obj => Value(obj, "id") ?? Value(obj, "href"), - JsonArray array => array.Select(Id).FirstOrDefault(i => i != default), - _ => default - }; - - public static string Value(JsonNode node, string property) => - node is JsonObject obj && obj[property] is JsonValue value && value.TryGetValue(out var text) ? text : default; - } -} diff --git a/PrivaPub/Federation/Objects/ActivityJson.cs b/PrivaPub/Federation/Objects/ActivityJson.cs new file mode 100644 index 0000000..9dda2d9 --- /dev/null +++ b/PrivaPub/Federation/Objects/ActivityJson.cs @@ -0,0 +1,39 @@ +using System.Text.Json.Nodes; + +namespace PrivaPub.Federation.Objects +{ + public static class ActivityJson + { + public static string Id(JsonNode node) => node switch + { + JsonValue value when value.TryGetValue(out var text) => text, + JsonObject obj => Value(obj, "id") ?? Value(obj, "href"), + JsonArray array => array.Select(Id).FirstOrDefault(i => i != default), + _ => default + }; + + public static string Value(JsonNode node, string property) => + node is JsonObject obj && obj[property] is JsonValue value && value.TryGetValue(out var text) ? text : default; + + public static IEnumerable Addresses(JsonNode node) + { + foreach (var field in new[] { "to", "cc", "bto", "bcc", "audience" }) + { + switch (node[field]) + { + case JsonArray array: + foreach (var item in array) + { + var id = Id(item); + if (id != default) + yield return id; + } + break; + case JsonNode single when Id(single) is { } id: + yield return id; + break; + } + } + } + } +} diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index a2b87d2..e3e5be8 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -14,6 +14,7 @@ using PrivaPub.Services.ClientToServer.Public; using PrivaPub.Federation.Actors; using PrivaPub.Federation.Outbox; using PrivaPub.Federation.Inbox; +using PrivaPub.Federation.Inbox.Handlers; using PrivaPub.Infrastructure.Http; using PrivaPub.Infrastructure.Jobs; using Microsoft.Extensions.Options; @@ -51,7 +52,13 @@ namespace PrivaPub.Middleware .AddSingleton() .AddSingleton() .AddSingleton() - .AddTransient() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton() .AddSingleton()