A thread owner's Add of a reply to its context is the reply passed on

FEP-171b, received (Hubzilla, Streams, Forte): an Add whose target is a collection of the actor's own server, and whose
object is someone's Create, Update or Delete of their own object, is queued as that activity forwarded by the owner, so
it is believed on its FEP-8b32 proof or as its origin has it, and shares its job with the plain forward Hubzilla also
sends. Checked live against Hubzilla (unchanged, 20 checks).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
This commit is contained in:
thepraandClaude Opus 5.5 committed 2026-10-06 10:03:44 +02:00
1 parent b39eadffeb
commit d06f684f0d
7 files changed
+171 -7

No files matched your search

+3
View File
@@ -85,6 +85,9 @@ covered by unit tests written in their documents' shape.
- [FEP-521a: Representing actor's public keys](https://codeberg.org/fediverse/fep/src/branch/main/fep/521a/fep-521a.md) and
[FEP-8b32: Object Integrity Proofs](https://codeberg.org/fediverse/fep/src/branch/main/fep/8b32/fep-8b32.md) (see
"Integrity proofs")
- [FEP-171b: Conversation Containers](https://codeberg.org/fediverse/fep/src/branch/main/fep/171b/fep-171b.md), received:
a thread owner's `Add` of someone's `Create`, `Update` or `Delete` to the thread's context is taken as that activity
passed on by the owner (on its FEP-8b32 proof, or as its origin has it)
- [FEP-8fcf: Followers collection synchronization across servers](https://codeberg.org/fediverse/fep/src/branch/main/fep/8fcf/fep-8fcf.md)
(sent and honoured; see "Followers synchronisation")
@@ -0,0 +1,109 @@
using MongoDB.Entities;
using PrivaPub.Infrastructure.Jobs;
using PrivaPub.Models.Jobs;
using PrivaPub.Models.Social;
using PrivaPub.Tests.Support;
using System.Text.Json.Nodes;
using PostEntity = PrivaPub.Models.Post.Post;
namespace PrivaPub.Tests.Federation
{
// FEP-171b: a thread's owner adds a reply someone else wrote to the thread's context (Hubzilla's Add{Create}), and the
// reply is taken as passed on by the owner, on its own proof; an Add aimed at a collection elsewhere is not a thread's
[Trait("Category", "Integration")]
public sealed class ConversationContainerTests : IAsyncLifetime
{
const string Public = "https://www.w3.org/ns/activitystreams#Public";
Harness _harness;
public async ValueTask InitializeAsync()
{
Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip);
_harness = await Harness.Start();
}
public async ValueTask DisposeAsync()
{
if (_harness != default)
await _harness.DisposeAsync();
}
static CancellationToken Token => TestContext.Current.CancellationToken;
static string Origin(RemoteActor actor) => new Uri(actor.Id).GetLeftPart(UriPartial.Authority);
static JsonObject Create(RemoteActor author, string text, string context, string inReplyTo = default)
{
var noteId = $"{Origin(author)}/item/{Guid.NewGuid():N}";
var note = new JsonObject
{
["id"] = noteId, ["type"] = "Note", ["attributedTo"] = author.Id, ["content"] = $"<p>{text}</p>", ["context"] = context,
["to"] = new JsonArray(Public), ["published"] = DateTime.UtcNow.ToString("O")
};
if (inReplyTo != default)
note["inReplyTo"] = inReplyTo;
return new JsonObject
{
["id"] = $"{Origin(author)}/activity/{Guid.NewGuid():N}", ["type"] = "Create", ["actor"] = author.Id, ["to"] = new JsonArray(Public),
["object"] = note
};
}
async Task<bool> Held(string objectUri) => await DB.Default.Find<PostEntity>().Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(Token);
async Task RunQueued(string innerId)
{
var job = await DB.Default.Find<Job>().Match(j => j.DedupeKey == "inbox|forwarded|" + innerId && j.State == JobState.Pending).ExecuteFirstAsync(Token);
Assert.NotNull(job);
Assert.Equal(JobResult.Done, (await _harness.Processor.Handle(job, Token)).Result);
}
[Fact]
public async Task A_reply_the_threads_owner_adds_to_its_context_is_taken_on_its_proof()
{
var (_, alice) = await _harness.Persona("alice");
var owner = new RemoteActor(_harness.Peer, "owner", ed25519: true);
var commenter = new RemoteActor(_harness.Peer, "commenter", ed25519: true);
await DB.Default.SaveAsync(new Following { AvatarId = alice.Id, TargetActorURI = owner.Id, TargetInboxURL = owner.Id + "/inbox", State = FollowState.Accepted }, Token);
var context = $"{Origin(owner)}/conversation/{Guid.NewGuid():N}";
var root = Create(owner, "a thread", context);
await _harness.Deliver(owner, "/human-centipede", root);
var rootId = root["object"]!["id"]!.GetValue<string>();
Assert.True(await Held(rootId));
var reply = commenter.Prove(Create(commenter, "a reply in it", context, rootId));
var replyId = reply["object"]!["id"]!.GetValue<string>();
await _harness.Deliver(owner, "/human-centipede", owner.Prove(new JsonObject
{
["id"] = $"{Origin(owner)}/activity/{Guid.NewGuid():N}", ["type"] = "Add", ["actor"] = owner.Id, ["object"] = reply,
["target"] = new JsonObject { ["id"] = context, ["type"] = "Collection", ["attributedTo"] = owner.Id }, ["to"] = new JsonArray(Public)
}));
await RunQueued(reply["id"]!.GetValue<string>());
var stored = await DB.Default.Find<PostEntity>().Match(p => p.ObjectURI == replyId).ExecuteFirstAsync(Token);
Assert.NotNull(stored);
Assert.Equal(commenter.Id, stored.ActorURI);
Assert.Equal((await DB.Default.Find<PostEntity>().Match(p => p.ObjectURI == rootId).ExecuteFirstAsync(Token)).ID, stored.AnsweringToPostId);
}
[Fact]
public async Task An_add_to_a_collection_on_another_server_is_no_threads()
{
var owner = new RemoteActor(_harness.Peer, "owner", ed25519: true);
var commenter = new RemoteActor(_harness.Peer, "commenter", _harness.Peer.B, ed25519: true);
var reply = commenter.Prove(Create(commenter, "elsewhere", $"{Origin(commenter)}/conversation/1"));
await _harness.Deliver(owner, "/human-centipede", new JsonObject
{
["id"] = $"{Origin(owner)}/activity/{Guid.NewGuid():N}", ["type"] = "Add", ["actor"] = owner.Id, ["object"] = reply,
["target"] = $"{Origin(commenter)}/conversation/1"
});
Assert.False(await DB.Default.Find<Job>().Match(j => j.DedupeKey == "inbox|forwarded|" + reply["id"]!.GetValue<string>()).ExecuteAnyAsync(Token));
}
}
}
+2 -2
View File
@@ -79,8 +79,8 @@ namespace PrivaPub.Tests.Support
new BlockHandler(Db, Local),
new LockHandler(Db),
new MoveHandler(Remote, Db, Local, Follows, Relationships),
new FeaturedHandler(Featured, add: true, Walls, Remote),
new FeaturedHandler(Featured, add: false, Walls, Remote)
new FeaturedHandler(Featured, add: true, Walls, Remote, Queue),
new FeaturedHandler(Featured, add: false, Walls, Remote, Queue)
};
((AnnounceHandler)Handlers.First(h => h is AnnounceHandler)).Relays = Handlers;
Processor = new InboxProcessor(Remote, Handlers, NullLogger<InboxProcessor>.Instance, Ledger, Local);
@@ -0,0 +1,45 @@
using PrivaPub.Federation.Objects;
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
{
// FEP-171b conversation containers (Hubzilla, Streams, Forte): a thread's owner adds each activity in it to the
// thread's context, a collection of its own, and tells the thread's audience with Add{activity}. The activity it adds
// is another's, so it is taken as forwarded by the owner: a Create, Update or Delete of the actor's own object,
// believed on its FEP-8b32 proof or as its origin has it (Forwarded). Hubzilla also passes the activity on by itself;
// the two are one job.
public static class ConversationContainers
{
public static async Task<bool> Unwrap(JsonNode add, ForeignAvatar owner, IJobQueue queue, CancellationToken token)
{
if (add["object"] is not JsonObject inner || Value(inner, "type") is not ("Create" or "Update" or "Delete") || Id(inner) is not { } innerId)
return false;
var target = add["target"];
var context = Id(target);
if (context == default || !Origin.Same(context, owner.ActorURI)
|| target is JsonObject collection && Id(collection["attributedTo"]) is { } attributedTo && attributedTo != owner.ActorURI)
return false;
var actorUri = Id(inner["actor"]);
if (actorUri == default || !Forwarded.Takeable(Value(inner, "type"), inner, actorUri))
{
Arrival.Drop("container-unusable");
return true;
}
var payload = new InboxPayload(actorUri, inner.ToJsonString(), "shared", ReceivedAt: DateTime.UtcNow, ForwardedBy: owner.ActorURI);
var queued = await queue.Enqueue(JobKind.ProcessInbox, JsonSerializer.Serialize(payload), new Uri(actorUri).Host.ToLowerInvariant(),
"inbox|forwarded|" + innerId, token);
if (queued)
Arrival.Accept("unwrapped");
else
Arrival.Drop("duplicate");
return true;
}
}
}
@@ -14,13 +14,16 @@ namespace PrivaPub.Federation.Inbox.Handlers
readonly IFeaturedPosts _featured;
readonly IWalls _walls;
readonly IRemoteActorService _remoteActors;
readonly Infrastructure.Jobs.IJobQueue _queue;
readonly bool _add;
public FeaturedHandler(IFeaturedPosts featured, bool add, IWalls walls = default, IRemoteActorService remoteActors = default)
public FeaturedHandler(IFeaturedPosts featured, bool add, IWalls walls = default, IRemoteActorService remoteActors = default,
Infrastructure.Jobs.IJobQueue queue = default)
{
_featured = featured;
_walls = walls;
_remoteActors = remoteActors;
_queue = queue;
_add = add;
}
@@ -28,6 +31,8 @@ namespace PrivaPub.Federation.Inbox.Handlers
public async Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token)
{
if (_add && _queue != default && await ConversationContainers.Unwrap(activity, actor, _queue, token))
return;
var target = Objects.ActivityJson.Id(activity["target"]);
// a collection of its own we do not know of: its document read again, at most hourly, in case it was kept before
// we read that collection (a wall)
@@ -97,9 +97,11 @@ namespace PrivaPub.Middleware
.AddSingleton<Federation.Actors.IFeaturedPosts, Federation.Actors.FeaturedPosts>()
.AddSingleton<Federation.Actors.IWalls, Federation.Actors.Walls>()
.AddSingleton<IActivityHandler>(services => new FeaturedHandler(services.GetRequiredService<Federation.Actors.IFeaturedPosts>(), add: true,
services.GetRequiredService<Federation.Actors.IWalls>(), services.GetRequiredService<Federation.Actors.IRemoteActorService>()))
services.GetRequiredService<Federation.Actors.IWalls>(), services.GetRequiredService<Federation.Actors.IRemoteActorService>(),
services.GetRequiredService<IJobQueue>()))
.AddSingleton<IActivityHandler>(services => new FeaturedHandler(services.GetRequiredService<Federation.Actors.IFeaturedPosts>(), add: false,
services.GetRequiredService<Federation.Actors.IWalls>(), services.GetRequiredService<Federation.Actors.IRemoteActorService>()))
services.GetRequiredService<Federation.Actors.IWalls>(), services.GetRequiredService<Federation.Actors.IRemoteActorService>(),
services.GetRequiredService<IJobQueue>()))
.AddSingleton<IActivityHandler, MoveHandler>()
.AddSingleton<IActivityHandler, CreateHandler>()
.AddSingleton<IActivityHandler, DeleteHandler>()
+2 -2
View File
@@ -865,8 +865,8 @@ without a port, before it gives out its OAuth client.
permission") while answering 200.
- **Threads are its owner's:** a comment by another channel in a thread goes to the owner, who passes it on to the
thread's audience as the commenter's own Create, with the commenter's FEP-8b32 proof, and adds each activity to the
thread's `context` (FEP-171b, `Add{Create}`). PrivaPub takes the forwarded Create on its proof; the `Add` is dropped,
the Create having said it already.
thread's `context` (FEP-171b, `Add{Create}`). PrivaPub takes the forwarded Create on its proof, and unwraps the `Add`
as the same activity passed on (one job for both, whichever comes first); an added `Like` is not taken.
- **Signs everything with a proof** (FEP-8b32) and verifies HTTP signatures; NodeInfo comes from its `statistics`
addon.
- **An unfollow changes nothing there** (G-0010): `Activity::unfollow` deletes the setting `system.their_perms`, which