diff --git a/internal/ingest/backfill.go b/internal/ingest/backfill.go index f5e11cb..c1286e5 100644 --- a/internal/ingest/backfill.go +++ b/internal/ingest/backfill.go @@ -284,11 +284,11 @@ func (b *Backfill) materializeOutboxItem(ctx context.Context, item *ap.Object, c // Our own content, walked past. It runs on the UNRESOLVED node, before // resolveEmbedded: our object id is cross-authority with the outbox host, // so resolving would dereference our own origin to fetch back a record we - // already hold — and MaterializePost then calls EnsureActor on its - // attributedTo BEFORE reading any mapping, minting a bridged actor for our - // own persona. A skip, never an error: our post in a community's history is - // an expected item, and failing it would leave every run reporting failures - // and the community never cleanly backfilled. + // already hold — and MaterializePost mints the actor its attributedTo names + // without checking whether the object is bridge-origin, minting a bridged + // actor for our own persona. A skip, never an error: our post in a + // community's history is an expected item, and failing it would leave every + // run reporting failures and the community never cleanly backfilled. // // The question is asked of the unwrapped OBJECT, not the outbox envelope: // the community mints its own Announce/Create ids around our content, so @@ -319,18 +319,16 @@ func (b *Backfill) materializeOutboxItem(ctx context.Context, item *ap.Object, c switch obj.Type { case ap.TypePage, ap.TypeArticle: - res, err := b.mat.MaterializePost(ctx, obj) + res, err := b.mat.MaterializePost(ctx, obj, communityIRI) if err != nil { return false, err } - // A community's host is trusted only for its own posts' counts, and - // the STORED community binding is the authority on which community a - // post lives in. An outbox can list another community's post: on - // first materialization it is bound to its own declared community. - // An outbox can list a copy whose audience was retargeted at the - // walked community: the post stays where it was first stored. Either - // way the walked host has no say over its counts. An unknown binding - // (empty) seeds nothing. + // A community's host is trusted only for its own posts' counts. + // MaterializePost is bound to the walked community, so a post it + // returns without a skip already belongs to it: another community's + // post, or one stored in another community and retargeted at this + // one, is refused above. The comparison is defense in depth on that + // binding; an unknown binding (empty) seeds nothing. if res.CommunityDID != "" && res.CommunityDID == communityDID { b.seedCounts(ctx, obj.ID, communityIRI) } else { @@ -346,7 +344,7 @@ func (b *Backfill) materializeOutboxItem(ctx context.Context, item *ap.Object, c } return true, nil case ap.TypeNote: - if _, err := b.mat.MaterializeComment(ctx, obj); err != nil { + if _, err := b.mat.MaterializeComment(ctx, obj, communityIRI); err != nil { return false, err } return true, nil @@ -458,7 +456,7 @@ func (b *Backfill) backfillReplies(ctx context.Context, post *ap.Object, communi b.logger.Info("backfill reply skipped", "post", post.ID, "reason", "reply was deleted upstream") return nil } - if _, err := b.mat.MaterializeComment(ctx, resolved); err != nil { + if _, err := b.mat.MaterializeComment(ctx, resolved, communityIRI); err != nil { if materialize.IsSkip(err) { b.logger.Info("backfill reply skipped", "post", post.ID, "reason", err.Error()) return nil diff --git a/internal/ingest/backfill_test.go b/internal/ingest/backfill_test.go index b92cb08..794262c 100644 --- a/internal/ingest/backfill_test.go +++ b/internal/ingest/backfill_test.go @@ -14,6 +14,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "tidepool/internal/ap" "tidepool/internal/errors" "tidepool/internal/materialize" ) @@ -161,11 +162,12 @@ func TestBackfillSeedsVoteCounts(t *testing.T) { } // TestBackfillDoesNotSeedPostsOfAnotherCommunity: a community's outbox can -// list a post that declares a DIFFERENT community. The post still lands (in -// its own declared community), but its counts must not be asked of -// the walked community's host — that host has no authority over another -// community's post. The walked community's own post is still seeded, so the -// test cannot pass by seeding nothing. +// list a post that declares a DIFFERENT community. A backfill walks one +// community's history, so content it reaches is bound to that community: the +// foreign post is not materialized at all (its own community's Announce or +// backfill is where it lands), and so it is never seeded from the walked +// community's host either. The walked community's own post is still +// materialized and seeded, so the test cannot pass by doing nothing. func TestBackfillDoesNotSeedPostsOfAnotherCommunity(t *testing.T) { h := newHarness(t) h.subscribeTechnology() @@ -193,18 +195,155 @@ func TestBackfillDoesNotSeedPostsOfAnotherCommunity(t *testing.T) { require.NoError(t, err) require.NoError(t, b.Run(ctx, community, true)) - // Both posts still materialize. - for _, id := range []string{"https://lemmy.world/post/49131386", "https://lemmy.world/post/49122698"} { - mapping, err := h.objects.GetByAPID(ctx, id) - require.NoError(t, err, "outbox post %s must be materialized", id) - assert.Equal(t, materialize.CollectionPostV2, mapping.Collection) - } + mapping, err := h.objects.GetByAPID(ctx, "https://lemmy.world/post/49131386") + require.NoError(t, err, "the walked community's own post must be materialized") + assert.Equal(t, materialize.CollectionPostV2, mapping.Collection) + _, err = h.objects.GetByAPID(ctx, "https://lemmy.world/post/49122698") + assert.True(t, errors.IsNotFound(err), + "a post declaring another community must not be materialized by the walk (err=%v)", err) // Only the walked community's own post is seeded. assert.Equal(t, []string{"https://lemmy.world/post/49131386"}, seeder.seeded, "a post declaring another community must not be seeded from the walked community's host") assert.Equal(t, []string{"https://lemmy.world/c/technology"}, seeder.communities) } +// TestBackfillBindsContentToTheWalkedCommunity: a backfill of technology +// materializes only content that belongs to technology. A post whose origin +// names another community, and a comment — listed in the outbox or in a +// technology post's replies collection — that replies into a co-hosted +// community's thread, are not materialized. A post on another instance that +// names technology still lands, in technology. +func TestBackfillBindsContentToTheWalkedCommunity(t *testing.T) { + const ( + xAuthorID = "https://sopuli.example/u/xPoster" + techPost = "https://lemmy.world/post/tech-replies-1" + ) + page := func(id, author, audience string) map[string]any { + return map[string]any{ + "type": "Page", + "id": id, + "attributedTo": author, + "to": []any{audience, ap.PublicAudience}, + "audience": audience, + "name": "a post at " + id, + "source": map[string]any{"content": "body of " + id, "mediaType": "text/markdown"}, + "published": "2026-07-09T14:00:00.000000Z", + } + } + create := func(obj map[string]any) map[string]any { + return map[string]any{ + "type": "Create", + "id": obj["id"].(string) + "/create", + "actor": obj["attributedTo"], + "object": obj, + } + } + techPostWithReplies := page(techPost, personID, groupID) + techPostWithReplies["replies"] = techPost + "/replies" + + rows := []struct { + name string + // items is technology's outbox. + items []any + // served are the origin documents the walk dereferences. + served map[string]map[string]any + // fetchedPath must be dereferenced by the walk, so a refusal is the + // binding and not a fixture miss. + fetchedPath string + targetID string + // wantCommunityDID is where the target lands; empty means it must not + // be materialized. + wantCommunityDID string + }{ + { + name: "outbox post whose origin names linux", + items: []any{create(page("https://sopuli.example/post/x-into-linux", xAuthorID, groupID))}, + served: map[string]map[string]any{ + "/post/x-into-linux": page("https://sopuli.example/post/x-into-linux", xAuthorID, linuxCommunityID), + }, + fetchedPath: "/post/x-into-linux", + targetID: "https://sopuli.example/post/x-into-linux", + }, + { + name: "outbox comment replying to linux's post", + items: []any{create(note("https://lemmy.world/comment/outbox-into-linux", personID, linuxPostID, + "an outbox reply into linux", "2026-07-09T14:30:00.000000Z"))}, + fetchedPath: "/c/technology/outbox", + targetID: "https://lemmy.world/comment/outbox-into-linux", + }, + { + // The refusal row above, aimed at technology's own post: the same + // embedded outbox Note shape does materialize, so that refusal is the + // binding and not an outbox shape the walk never processes. + name: "outbox comment replying to technology's post", + items: []any{create(note("https://lemmy.world/comment/outbox-into-tech", personID, pageID, + "an outbox reply into technology", "2026-07-09T14:30:00.000000Z"))}, + fetchedPath: "/c/technology/outbox", + targetID: "https://lemmy.world/comment/outbox-into-tech", + wantCommunityDID: testDIDFor("technology", "lemmy.world"), + }, + { + name: "replies-collection comment replying to linux's post", + items: []any{create(techPostWithReplies)}, + served: map[string]map[string]any{ + "/post/tech-replies-1/replies": { + "type": "OrderedCollection", + "id": techPost + "/replies", + "totalItems": 1, + "orderedItems": []any{note("https://lemmy.world/comment/replies-into-linux", personID, + linuxPostID, "a listed reply into linux", "2026-07-09T14:45:00.000000Z")}, + }, + }, + fetchedPath: "/post/tech-replies-1/replies", + targetID: "https://lemmy.world/comment/replies-into-linux", + }, + { + name: "cross-instance post naming technology", + items: []any{create(page("https://sopuli.example/post/x-into-tech", xAuthorID, groupID))}, + served: map[string]map[string]any{ + "/post/x-into-tech": page("https://sopuli.example/post/x-into-tech", xAuthorID, groupID), + }, + fetchedPath: "/post/x-into-tech", + targetID: "https://sopuli.example/post/x-into-tech", + wantCommunityDID: testDIDFor("technology", "lemmy.world"), + }, + } + + for _, row := range rows { + t.Run(row.name, func(t *testing.T) { + h := newHarness(t) + h.subscribeTechnology() + h.serveLemmyWorldContent() + h.coHostedCommunityPost() + h.serveObject("/u/xPoster", person(xAuthorID, "xPoster", nil)) + for path, doc := range row.served { + h.serveObject(path, doc) + } + h.serveObject("/c/technology/outbox", map[string]any{ + "type": "OrderedCollection", + "id": groupID + "/outbox", + "totalItems": len(row.items), + "orderedItems": row.items, + }) + ctx := context.Background() + community, err := h.communities.GetByAPGroupID(ctx, groupID) + require.NoError(t, err) + + require.NoError(t, newBackfill(t, h, 10).Run(ctx, community, true)) + + require.Equal(t, 1, h.hitCount(row.fetchedPath), "the walk dereferences %s", row.fetchedPath) + mapping, err := h.objects.GetByAPID(ctx, row.targetID) + if row.wantCommunityDID == "" { + assert.True(t, errors.IsNotFound(err), + "%s must not be materialized by technology's backfill (err=%v)", row.targetID, err) + return + } + require.NoError(t, err, "%s must be materialized", row.targetID) + assert.Equal(t, row.wantCommunityDID, mapping.CommunityDID) + }) + } +} + // newSeededBackfill is newBackfill with a CountSeeder wired in. func newSeededBackfill(t *testing.T, h *harness, seeder CountSeeder) *Backfill { t.Helper() diff --git a/internal/ingest/consent.go b/internal/ingest/consent.go index ad8639d..b598731 100644 --- a/internal/ingest/consent.go +++ b/internal/ingest/consent.go @@ -420,6 +420,15 @@ func (h *Handler) handleUndoDelete(ctx context.Context, undo, del *ap.Object, si } } + // Content is re-materialized only for the community that announced the + // restore: that community is the binding HandleUpdate checks it against. + // The bare-content drop above leaves only non-content mappings here, and + // those have nothing to re-materialize. + boundCommunityIRI := announcerGroupID(announcer) + if boundCommunityIRI == "" { + return skip(undo.ID, "bare restore of "+targetID+" is not applied: only an announced restore re-materializes") + } + // Pinned to the target's own authority: this fetch's answer is what // authorizes the restore AND what gets written into the repo, so an open // redirect on the origin must fail it rather than both license the restore @@ -440,21 +449,16 @@ func (h *Handler) handleUndoDelete(ctx context.Context, undo, del *ap.Object, si return skip(targetID, fmt.Sprintf( "restored object is a %q but %s is mapped as %s", restored.Type, targetID, mapping.Collection)) } - // An announced restore is bound to the announcing community exactly like - // announced content (materializeContent's guard) and one notch tighter: the - // restored body must NAME a community, and it must be the announcer's. The - // materializer derives the target community from the object's own audience - // and EnsureCommunity()s it, so letting an EMPTY audience pass would let a - // community vouch for a restore into whatever the object turns out to name - // — that vacuous pass is what let a sibling community revive another - // community's soft-deleted comment. Real Lemmy bodies always carry audience, - // so requiring it costs nothing. - if announcer != nil { - if objCommunity := communityIRIFrom(restored); objCommunity != announcer.APGroupID { - return skip(targetID, fmt.Sprintf( - "restored object names community %q but was announced by %s", - objCommunity, announcer.APGroupID)) - } + // The authoritative binding is the materializer's: HandleUpdate refuses a + // restored body whose thread or stored community is not the announcer's. + // This guard is an early filter in front of it, refusing a body that names + // another community (or none) before the tombstone and the mapping's soft + // delete are cleared below, so a refused restore has nothing to undo. Real + // Lemmy bodies always carry audience, so requiring it costs nothing. + if objCommunity := communityIRIFrom(restored); objCommunity != boundCommunityIRI { + return skip(targetID, fmt.Sprintf( + "restored object names community %q but was announced by %s", + objCommunity, boundCommunityIRI)) } if err := h.tombstones.Remove(ctx, targetID, scope); err != nil { @@ -464,7 +468,7 @@ func (h *Handler) handleUndoDelete(ctx context.Context, undo, del *ap.Object, si return fmt.Errorf("ingest: restore mapping for %s: %w", targetID, err) } h.logger.Info("object restored upstream; re-materializing", "ap_id", targetID) - if _, err = h.mat.HandleUpdate(ctx, restored); err != nil { + if _, err = h.mat.HandleUpdate(ctx, restored, boundCommunityIRI); err != nil { // Compensation, for EVERY error class. The mapping's soft delete is // already cleared and its record was deleted from the repo, so leaving // the mapping live strands it WITHOUT a record (downstream diff --git a/internal/ingest/echo_sandwich_test.go b/internal/ingest/echo_sandwich_test.go index 2a7eb5e..412c2b3 100644 --- a/internal/ingest/echo_sandwich_test.go +++ b/internal/ingest/echo_sandwich_test.go @@ -121,6 +121,7 @@ func setupSandwich(t *testing.T, h *harness) sandwich { ATURI: csPostATURI, ID: consume.ActivityID(csUserOrigin, csPostATURI, "create", 0), CommunityAPID: groupID, + CommunityDID: testDIDFor("technology", "lemmy.world"), Snapshot: csPostSnapshot(t), }) mapping, err := h.objects.GetByAPID(ctx, csPostAPID) diff --git a/internal/ingest/handler.go b/internal/ingest/handler.go index 9ae0914..3db89e8 100644 --- a/internal/ingest/handler.go +++ b/internal/ingest/handler.go @@ -25,9 +25,9 @@ const echoDropLogInterval = time.Second // Materializer is the slice of *materialize.Materializer the dispatcher // drives (task 05's entry points). type Materializer interface { - MaterializePost(ctx context.Context, page *ap.Object) (*materialize.Result, error) - MaterializeComment(ctx context.Context, note *ap.Object) (*materialize.Result, error) - HandleUpdate(ctx context.Context, obj *ap.Object) (*materialize.Result, error) + MaterializePost(ctx context.Context, page *ap.Object, communityIRI string) (*materialize.Result, error) + MaterializeComment(ctx context.Context, note *ap.Object, communityIRI string) (*materialize.Result, error) + HandleUpdate(ctx context.Context, obj *ap.Object, communityIRI string) (*materialize.Result, error) // HandleDelete branches actor vs content off a FRESH bridged_actors read; // HandleDeleteRecord is the content-only entry that never can. Callers // that already classified the target (handleDelete, SweepDeleted) use the @@ -550,14 +550,13 @@ func (h *Handler) materializeContent(ctx context.Context, obj *ap.Object, signer // Announced content must belong to the announcing community itself: a // followed community may fan out only its own content, never claim - // another community's (even one co-hosted on the same instance). Since - // the flip the consequence is not a foreign write into a community repo - // — a postv2 goes to its author's repo — but a false BINDING: the - // materializer derives the target community from the object's own - // audience, EnsureCommunity()s it, records it as the mapping's - // community_did and writes that community's acceptance. Without this - // guard an announcer could name any community it likes and hand it both - // visibility over the post and moderation authority over it. + // another community's (even one co-hosted on the same instance). The + // authoritative check is the materializer's, which binds every post and + // comment to the announcer passed below — a post must name it, a comment's + // thread must belong to it, and stored content keeps the community it was + // first bound to. This guard is an early filter in front of that: an + // object naming another community is refused before the materializer + // fetches any ancestor or bridges any author. if objCommunity := communityIRIFrom(obj); objCommunity != "" && objCommunity != announcer { return skip(obj.ID, fmt.Sprintf( "announced object names community %s but was announced by %s", objCommunity, announcer)) @@ -566,15 +565,15 @@ func (h *Handler) materializeContent(ctx context.Context, obj *ap.Object, signer switch obj.Type { case ap.TypePage, ap.TypeArticle: if isUpdate { - _, err = h.mat.HandleUpdate(ctx, obj) + _, err = h.mat.HandleUpdate(ctx, obj, announcer) } else { - _, err = h.mat.MaterializePost(ctx, obj) + _, err = h.mat.MaterializePost(ctx, obj, announcer) } case ap.TypeNote: if isUpdate { - _, err = h.mat.HandleUpdate(ctx, obj) + _, err = h.mat.HandleUpdate(ctx, obj, announcer) } else { - _, err = h.mat.MaterializeComment(ctx, obj) + _, err = h.mat.MaterializeComment(ctx, obj, announcer) } default: return skip(obj.ID, "unsupported content type "+obj.Type) diff --git a/internal/ingest/handler_test.go b/internal/ingest/handler_test.go index e4a23a1..08212ae 100644 --- a/internal/ingest/handler_test.go +++ b/internal/ingest/handler_test.go @@ -91,6 +91,38 @@ func (h *harness) followedCommunity(id, username, instance string) *remoteActor return actor } +// The co-hosted second community of the community-binding tests: followed, +// on the same instance as technology, with one post of its own materialized. +const ( + linuxCommunityID = "https://lemmy.world/c/linux" + linuxPostID = "https://lemmy.world/post/linux-1" +) + +// coHostedCommunityPost subscribes linux — a second community on technology's +// own instance — and materializes linuxPostID through linux's own Announce, so +// a reply naming that post as its parent anchors on a thread that lives in +// linux. +func (h *harness) coHostedCommunityPost() { + h.t.Helper() + linux := h.subscribeCommunityURL(linuxCommunityID, "linux") + const authorID = "https://lemmy.world/u/linuxAuthor" + h.serveObject("/u/linuxAuthor", person(authorID, "linuxAuthor", nil)) + h.announceCreate(linux, "https://lemmy.world/activities/announce/create/linux-1", map[string]any{ + "type": "Page", + "id": linuxPostID, + "attributedTo": authorID, + "to": []any{linuxCommunityID, ap.PublicAudience}, + "audience": linuxCommunityID, + "name": "a post in linux", + "source": map[string]any{"content": "linux body", "mediaType": "text/markdown"}, + "published": "2026-07-09T10:00:00.000000Z", + }) + mapping, err := h.objects.GetByAPID(context.Background(), linuxPostID) + require.NoError(h.t, err, "linux's own post is the premise") + require.Equal(h.t, testDIDFor("linux", "lemmy.world"), mapping.CommunityDID, + "linux's post lives in linux") +} + // tombstoneAnnouncers lists the raw marker rows for an ap id. ExistsFor // cannot answer "whose row is it": a global marker is visible in every scope, // so it masks exactly the per-announcer removals the scoping rules are about. @@ -1048,6 +1080,40 @@ func TestAnnouncedObjectForDifferentCommunityDropped(t *testing.T) { assert.True(t, errors.IsNotFound(err), "the other community must not be bridged") } +// TestAnnouncedCommentReplyingIntoAnotherCommunityDropped: the audience guard +// in materializeContent reads only the comment's own claim, but a comment's +// community is its thread's. A followed community announcing a comment whose +// body names that community yet replies to a post living in a co-hosted +// community must not place the comment in the other community's thread. The +// comment's id is on the author's own host, so its body is re-fetched from +// there: the refusal must hold for the origin's answer, not just for an +// embedded copy. +func TestAnnouncedCommentReplyingIntoAnotherCommunityDropped(t *testing.T) { + h := newHarness(t) + group := h.subscribeTechnology() + h.coHostedCommunityPost() + ctx := context.Background() + + const ( + authorID = "https://sopuli.example/u/threadHopper" + commentID = "https://sopuli.example/comment/thread-hop-1" + ) + h.serveObject("/u/threadHopper", person(authorID, "threadHopper", nil)) + // Names technology, so the leaf audience guard lets it through; replies to + // linux's post. + comment := note(commentID, authorID, linuxPostID, "a reply into linux", "2026-07-09T11:00:00.000000Z") + h.serveObject("/comment/thread-hop-1", comment) + + h.announceCreate(group, "https://lemmy.world/activities/announce/create/thread-hop-1", comment) + + require.Equal(t, 1, h.hitCount("/comment/thread-hop-1"), + "the cross-authority comment is re-fetched from its origin") + _, err := h.objects.GetByAPID(ctx, commentID) + assert.True(t, errors.IsNotFound(err), + "a comment announced by technology must not be materialized into linux's thread (err=%v)", err) + assert.Equal(t, 0, h.firehoseOpCount(materialize.CollectionComment), "no comment record in any repo") +} + // TestUndoDeleteRollsBackWhenRematerializeSkips (Finding 4): if the restore's // re-materialization is declined (a skip), the compensation must re-soft-delete // the mapping and re-record the tombstone — never leave a live mapping without @@ -1439,6 +1505,83 @@ func TestAnnouncedUndoDeleteIntoAnotherCommunityDropped(t *testing.T) { assert.True(t, errors.IsNotFound(err), "the other community must not be bridged") } +// TestAnnouncedRestoreReplyingIntoAnotherCommunityDropped: the restore guard +// checks only the restored body's own audience, but a comment's community is +// its thread's. A comment first materialized in technology's thread is +// deleted, and technology then announces its restore — but the origin now +// serves the body still naming technology while replying to a post that lives +// in co-hosted linux. The restore must not revive the comment into linux's +// thread: the mapping stays soft-deleted and technology's tombstone stays. +func TestAnnouncedRestoreReplyingIntoAnotherCommunityDropped(t *testing.T) { + h := newHarness(t) + group := h.subscribeTechnology() + h.serveLemmyWorldContent() + h.coHostedCommunityPost() + ctx := context.Background() + + const ( + techPostID = "https://lemmy.world/post/restore-hop-post" + commentID = "https://lemmy.world/comment/restore-hop-1" + ) + h.announceCreate(group, "https://lemmy.world/activities/announce/create/restore-hop-post", map[string]any{ + "type": "Page", + "id": techPostID, + "attributedTo": personID, + "to": []any{groupID, ap.PublicAudience}, + "audience": groupID, + "name": "a post in technology", + "source": map[string]any{"content": "technology body", "mediaType": "text/markdown"}, + "published": "2026-07-09T12:00:00.000000Z", + }) + h.announceCreate(group, "https://lemmy.world/activities/announce/create/restore-hop-1", + note(commentID, personID, techPostID, "a reply in technology", "2026-07-09T12:30:00.000000Z")) + mapping, err := h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + require.Equal(t, testDIDFor("technology", "lemmy.world"), mapping.CommunityDID, + "the comment first lands in technology's thread") + + h.announceDelete(group, "https://lemmy.world/activities/announce/delete/restore-hop-1", personID, commentID) + mapping, err = h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + require.True(t, mapping.IsDeleted()) + require.Equal(t, []string{groupID}, h.tombstoneAnnouncers(commentID)) + + // The origin now serves the comment still naming technology, but replying + // to linux's post. + h.serveObject("/comment/restore-hop-1", + note(commentID, personID, linuxPostID, "a reply in technology", "2026-07-09T12:30:00.000000Z")) + hitsBefore := h.hitCount("/comment/restore-hop-1") + require.Equal(t, http.StatusAccepted, h.deliver(group, map[string]any{ + "id": "https://lemmy.world/activities/announce/undo/restore-hop-1", + "type": "Announce", + "actor": groupID, + "audience": groupID, + "object": map[string]any{ + "id": "https://lemmy.world/activities/undo/restore-hop-1", + "type": "Undo", + "actor": personID, + "audience": groupID, + "object": map[string]any{ + "id": "https://lemmy.world/activities/delete/restore-hop-1", + "type": "Delete", + "actor": personID, + "object": commentID, + }, + }, + })) + h.drain() + + require.Equal(t, hitsBefore+1, h.hitCount("/comment/restore-hop-1"), + "the restore re-fetches the comment from its origin") + mapping, err = h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + assert.True(t, mapping.IsDeleted(), "a restore into linux's thread must not revive the mapping") + assert.Equal(t, []string{groupID}, h.tombstoneAnnouncers(commentID), + "technology's tombstone survives the refused restore") + _, _, err = h.manager.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + assert.True(t, errors.IsNotFound(err), "the refused restore must not rewrite the record (err=%v)", err) +} + // TestBareReferenceCreateCannotDodgeScopedTombstone: a bare Create carrying // nothing but {"id": X} names no community at all, so it offers no scope for // the community-scoped marker a delete-before-create left for exactly that id diff --git a/internal/materialize/bridgedstats_test.go b/internal/materialize/bridgedstats_test.go index 1424c9f..3863448 100644 --- a/internal/materialize/bridgedstats_test.go +++ b/internal/materialize/bridgedstats_test.go @@ -27,7 +27,7 @@ func TestSetBridgedStatsStampsRecord(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - created, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + created, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) require.False(t, created.NoOp) @@ -61,7 +61,7 @@ func TestSetBridgedStatsUnchangedCountsNoOp(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, pageID) require.NoError(t, err) @@ -92,7 +92,7 @@ func TestSetBridgedStatsChangedCountsCommits(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, pageID) require.NoError(t, err) @@ -123,7 +123,7 @@ func TestSetBridgedStatsRecordDeleted(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, pageID) require.NoError(t, err) @@ -144,7 +144,7 @@ func TestEditCarriesBridgedStatsForward(t *testing.T) { ctx := context.Background() page := loadFixtureObject(t, "page_lemmy_world.json") - _, err := h.m.MaterializePost(ctx, page) + _, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, pageID) require.NoError(t, err) @@ -156,7 +156,7 @@ func TestEditCarriesBridgedStatsForward(t *testing.T) { // A later edit rebuilds from AP data with no stats field. edited := loadFixtureObject(t, "page_lemmy_world.json") edited.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} - res, err := h.m.HandleUpdate(ctx, edited) + res, err := h.m.HandleUpdate(ctx, edited, groupID) require.NoError(t, err) require.False(t, res.NoOp, "an edited body is a real commit") @@ -180,7 +180,7 @@ func TestUnchangedReingestAfterStatsIsNoOp(t *testing.T) { ctx := context.Background() page := loadFixtureObject(t, "page_lemmy_world.json") - _, err := h.m.MaterializePost(ctx, page) + _, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, pageID) require.NoError(t, err) @@ -190,7 +190,7 @@ func TestUnchangedReingestAfterStatsIsNoOp(t *testing.T) { eventsAfterStamp := len(h.firehoseEvents()) // Re-ingest the identical post (a re-delivery). - again, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + again, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) assert.True(t, again.NoOp, "an unchanged re-ingest after stamping must be a no-op") assert.Equal(t, stamped.CID, again.CID, "carry-forward keeps the CID identical") @@ -283,13 +283,13 @@ func TestCommentEditCarriesReplyRefsForward(t *testing.T) { h.serveObject("/u/alice", person("https://lemmy.world/u/alice", "alice", nil)) ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) c1 := note("https://lemmy.world/comment/1001", "https://lemmy.world/u/alice", pageID, "original comment", "2026-07-07T04:00:00.000000Z") h.serveObject("/comment/1001", c1) - _, err = h.m.MaterializeComment(ctx, objectFromMap(t, c1)) + _, err = h.m.MaterializeComment(ctx, objectFromMap(t, c1), groupID) require.NoError(t, err) c1Mapping, err := h.objects.GetByAPID(ctx, "https://lemmy.world/comment/1001") @@ -313,7 +313,7 @@ func TestCommentEditCarriesReplyRefsForward(t *testing.T) { // directly (its parent is already mapped), so no re-fetch/re-serve is needed. edited := note("https://lemmy.world/comment/1001", "https://lemmy.world/u/alice", pageID, "edited comment body", "2026-07-07T04:00:00.000000Z") - res, err := h.m.HandleUpdate(ctx, objectFromMap(t, edited)) + res, err := h.m.HandleUpdate(ctx, objectFromMap(t, edited), groupID) require.NoError(t, err) require.False(t, res.NoOp, "an edited body is a real commit") diff --git a/internal/materialize/comments.go b/internal/materialize/comments.go index 1ac3c35..615584b 100644 --- a/internal/materialize/comments.go +++ b/internal/materialize/comments.go @@ -24,10 +24,17 @@ const maxAncestorDepth = 50 // have already materialized or the root Page, then the collected ancestors // are materialized oldest-first before the comment itself. The walk carries // a depth cap and a cycle guard, and — because it completes before any -// write — a chain that dead-ends (tombstoned/nobridge/unfetchable ancestor) -// drops the whole subtree without leaving partial ancestors behind from -// this call. -func (m *Materializer) MaterializeComment(ctx context.Context, note *ap.Object) (*Result, error) { +// write — a chain that dead-ends (tombstoned/unfetchable ancestor, one refused +// on its own terms, or one outside the delivering community) drops the whole +// subtree without leaving partial ancestors behind from this call. An +// author's consent (nobridge, deleted) — the leaf author's included — is +// checked only when that record is committed, so a record refused for its +// author still drops the subtree from there down, but the older ancestors +// already committed stay. +func (m *Materializer) MaterializeComment(ctx context.Context, note *ap.Object, communityIRI string) (*Result, error) { + if err := requireBoundCommunityIRI(communityIRI); err != nil { + return nil, err + } if note == nil || note.ID == "" { return nil, errors.NewValidationError("note", "must carry an AP object id") } @@ -38,19 +45,36 @@ func (m *Materializer) MaterializeComment(ctx context.Context, note *ap.Object) if note.InReplyTo == nil || note.InReplyTo.ID == "" { return nil, skip(note.ID, "comment has no inReplyTo") } + // The leaf's own refusals need no ancestor, so they come before the walk: + // a comment dropped for its own reasons must not leave its ancestors, or + // their authors, behind. + draft, err := draftComment(note) + if err != nil { + return nil, err + } - ancestors, err := m.collectUnmappedAncestors(ctx, note) + // A comment already stored belongs to the thread it was first posted in. + // An Update re-parenting it into the delivering community's thread does + // not make it that community's, so it is refused before the new ancestry + // is even walked. commitCommentLeaf repeats the check for every + // record it commits, ancestors included. + if _, err := m.storedCommentMapping(ctx, note.ID, communityIRI); err != nil { + return nil, err + } + + ancestors, err := m.collectUnmappedAncestors(ctx, note, communityIRI) if err != nil { return nil, err } for _, ancestor := range ancestors { - if err := m.materializeAncestor(ctx, ancestor); err != nil { - // A skipped ancestor (nobridge/deleted author, tombstone) takes - // the whole subtree with it — placeholder-free by design. + if err := m.materializeAncestor(ctx, ancestor, communityIRI); err != nil { + // A skipped ancestor (nobridge/deleted author) takes the rest of + // the subtree with it — placeholder-free by design — but the + // older ancestors committed before it stay. return nil, err } } - return m.materializeCommentLeaf(ctx, note) + return m.commitCommentLeaf(ctx, note, draft, communityIRI) } // countAncestorShortCircuit records the anchor as an echo suppression when the @@ -79,18 +103,37 @@ func (m *Materializer) countAncestorShortCircuit(ctx context.Context, parentID s } } +// unmappedAncestor is one object the ancestor walk fetched, with the draft +// its own checks produced when it is a Note (nil for the root Page). +type unmappedAncestor struct { + object *ap.Object + draft *commentDraft +} + // collectUnmappedAncestors walks note's inReplyTo chain upward until it // hits an already-mapped object or the thread's root Page, returning the // unmapped ancestors oldest-first. Nothing is written during the walk. -func (m *Materializer) collectUnmappedAncestors(ctx context.Context, note *ap.Object) ([]*ap.Object, error) { - var chain []*ap.Object +// +// Every fetched Note is drafted as it is reached, so an ancestor that would +// be refused on its own terms refuses the whole chain here, before an older +// ancestor — or its author — is committed on the way to it. +// +// The walk is also where the thread is bound to communityIRI: wherever the +// chain ends — on a mapped parent or on the root Page — that end must belong +// to the delivering community. Deciding it here, before any write, is what +// keeps a refused thread from leaving ancestors, authors or a community +// behind. +func (m *Materializer) collectUnmappedAncestors(ctx context.Context, note *ap.Object, communityIRI string) ([]unmappedAncestor, error) { + var chain []unmappedAncestor seen := map[string]bool{note.ID: true} current := note for { if current.InReplyTo == nil || current.InReplyTo.ID == "" { - // current is the top of the thread (normally the Page). It was - // fetched (unmapped), so it is already in chain; the walk ends. + // current is the top of the thread but not a Page (those end the + // walk where they are fetched). Not reached for a Note: the leaf + // and every fetched ancestor were drafted, and draftComment + // refuses a parentless one. return chain, nil } parentID := current.InReplyTo.ID @@ -102,7 +145,11 @@ func (m *Materializer) collectUnmappedAncestors(ctx context.Context, note *ap.Ob _, _, err := m.objects.ResolveStrongRef(ctx, parentID) switch { case err == nil: - // Anchored: the parent is already materialized. + // Anchored: the parent is already materialized, and its stored + // community decides the whole thread's. + if err := m.requireAnchorInBoundCommunity(ctx, parentID, note.ID, communityIRI); err != nil { + return nil, err + } m.countAncestorShortCircuit(ctx, parentID) return chain, nil case errors.IsTombstoned(err): @@ -129,44 +176,122 @@ func (m *Materializer) collectUnmappedAncestors(ctx context.Context, note *ap.Ob default: return nil, fmt.Errorf("materialize: fetch ancestor %s of %s: %w", parentID, note.ID, err) } - // Bind the self-asserted id to the fetch authority: commitRecord keys - // the ap_objects mapping on parent.ID, so a host serving a body that - // claims another instance's id would forge content under the victim's - // canonical id. Empty id inherits the requested IRI. - if parent.ID == "" { + // Bind the self-asserted id to the requested IRI: commitRecord keys the + // ap_objects mapping on parent.ID, so a body claiming another id — + // another instance's, or a stored object's on the same host — would + // commit content fetched from one IRI under someone else's. Nothing is + // lost by refusing: the child names parentID, so a parent mapped under + // any other id could never resolve it. Empty id inherits the requested + // IRI. + switch { + case parent.ID == "": parent.ID = parentID - } else if !ap.SameAuthority(parent.ID, parentID) { + case !ap.SameAuthority(parent.ID, parentID): return nil, skip(note.ID, fmt.Sprintf("ancestor %s served a cross-authority id %s", parentID, parent.ID)) + case parent.ID != parentID: + return nil, skip(note.ID, + fmt.Sprintf("ancestor %s served a body claiming another id %s", parentID, parent.ID)) + } + if parent.Type == ap.TypePage || parent.Type == ap.TypeArticle { + // A Page is the thread root whatever it replies to, so the walk + // ends here and nothing older is fetched or materialized. Its + // community is compared by IRI, not DID: the root may name a + // community the bridge has not bridged yet, and that is fine when + // it is the delivering one. + chain = append([]unmappedAncestor{{object: parent}}, chain...) + if ref := communityRef(parent); ref == nil || ref.ID != communityIRI { + return nil, skip(note.ID, + fmt.Sprintf("thread root %s is not in the delivering community %s", parent.ID, communityIRI)) + } + return chain, nil + } + if parent.Type != ap.TypeNote { + return nil, skip(note.ID, fmt.Sprintf("ancestor %s has unsupported type %s", parent.ID, parent.Type)) } - chain = append([]*ap.Object{parent}, chain...) + draft, err := draftComment(parent) + if err != nil { + return nil, err + } + chain = append([]unmappedAncestor{{object: parent, draft: draft}}, chain...) current = parent } } -// materializeAncestor writes one fetched ancestor: Pages through the post -// path, Notes as comment leaves (their own parents are guaranteed mapped — -// the chain is processed oldest-first). -func (m *Materializer) materializeAncestor(ctx context.Context, ancestor *ap.Object) error { - switch ancestor.Type { - case ap.TypePage, ap.TypeArticle: - _, err := m.MaterializePost(ctx, ancestor) - return err - case ap.TypeNote: - _, err := m.materializeCommentLeaf(ctx, ancestor) +// requireAnchorInBoundCommunity refuses note when the already-materialized +// object anchorID, on which its ancestry ends, belongs to a community other +// than communityIRI. +func (m *Materializer) requireAnchorInBoundCommunity(ctx context.Context, anchorID, noteID, communityIRI string) error { + anchor, err := m.objects.GetByAPID(ctx, anchorID) + if err != nil { + return fmt.Errorf("materialize: load anchor mapping %s of %s: %w", anchorID, noteID, err) + } + anchorCommunity, err := CommunityDIDOf(ctx, m.repos, anchor) + if err != nil { return err + } + return m.requireBoundCommunity(ctx, anchorCommunity, communityIRI, noteID) +} + +// storedCommentMapping returns commentID's existing mapping, or nil when it +// has none, refusing the comment when the stored one was deleted upstream or +// belongs to a community other than communityIRI. Which community owns a +// comment is decided at first materialization; a later delivery is an edit +// and cannot move it. +func (m *Materializer) storedCommentMapping(ctx context.Context, commentID, communityIRI string) (*store.APObjectMapping, error) { + existing, err := m.objects.GetByAPID(ctx, commentID) + switch { + case err == nil: + case errors.IsNotFound(err): + return nil, nil default: - return skip(ancestor.ID, "ancestor has unsupported type "+ancestor.Type) + return nil, fmt.Errorf("materialize: check mapping for %s: %w", commentID, err) + } + // The refusal commitRecord would make, made before the ancestor walk: a + // deleted comment re-delivered under an unmapped chain must not fetch, + // commit or bridge the authors of ancestors it will never be written + // under. An announced restore clears the soft delete before it gets here. + if existing.IsDeleted() { + return nil, skip(commentID, "object was deleted upstream; not resurrecting") + } + stored, err := CommunityDIDOf(ctx, m.repos, existing) + if err != nil { + return nil, err + } + if err := m.requireBoundCommunity(ctx, stored, communityIRI, commentID); err != nil { + return nil, err + } + return existing, nil +} + +// materializeAncestor writes one fetched ancestor: the root Page through the +// post path, Notes as comment leaves from the draft the walk made (their own +// parents are guaranteed mapped — the chain is processed oldest-first). +func (m *Materializer) materializeAncestor(ctx context.Context, ancestor unmappedAncestor, communityIRI string) error { + if ancestor.draft == nil { + _, err := m.MaterializePost(ctx, ancestor.object, communityIRI) + return err } + _, err := m.commitCommentLeaf(ctx, ancestor.object, ancestor.draft, communityIRI) + return err } -// materializeCommentLeaf writes a single comment whose parent is already -// materialized. -func (m *Materializer) materializeCommentLeaf(ctx context.Context, note *ap.Object) (*Result, error) { +// commentDraft is what a Note yields on its own, before anything is read or +// minted: its record key, its author reference and its content. +type commentDraft struct { + rkey string + authorRef *ap.Object + body string + facets []any +} + +// draftComment runs the checks a comment fails on its own terms, with no +// ancestor, mapping or actor involved. +func draftComment(note *ap.Object) (*commentDraft, error) { // A comment must reply to something. A parentless Note reaching this path // is a thread rooted at a non-Page object (e.g. a Mastodon status that // federated in as a Lemmy comment); drop the subtree rather than deref a - // nil inReplyTo below. + // nil inReplyTo later. if note.InReplyTo == nil || note.InReplyTo.ID == "" { return nil, skip(note.ID, "comment thread roots at a non-Page object") } @@ -183,13 +308,49 @@ func (m *Materializer) materializeCommentLeaf(ctx context.Context, note *ap.Obje if err := requireSameAuthorityAuthor(note, authorRef); err != nil { return nil, err } - author, err := m.EnsureActor(ctx, authorRef) + content := markdownFromObject(note) + if content == "" { + return nil, skip(note.ID, "comment has no content") + } + body, facets := bridgedRichText(content, 3000, 30000) + if body == "" { + // An HTML-only body can reduce to nothing once tags are stripped; a + // comment is nothing but its content, so drop it like a bodiless one. + return nil, skip(note.ID, "comment has no content") + } + return &commentDraft{rkey: rkey, authorRef: authorRef, body: body, facets: facets}, nil +} + +// commitCommentLeaf writes a single drafted comment whose parent is already +// materialized: it binds the comment to its parent's thread and community and +// commits it. The parent must belong to communityIRI; the ancestor walk has +// already established that for the thread, and the check here holds it for +// each record actually committed. +func (m *Materializer) commitCommentLeaf(ctx context.Context, note *ap.Object, draft *commentDraft, communityIRI string) (*Result, error) { + // Resolved before the author is bridged, so a parent outside the bound + // community refuses the comment without minting anyone. + reply, communityDID, err := m.resolveReplyRefs(ctx, note) + if err != nil { + return nil, err + } + if err := m.requireBoundCommunity(ctx, communityDID, communityIRI, note.ID); err != nil { + return nil, err + } + // The stored comment's own community is checked too: an ancestor the walk + // found unmapped can have been stored since by a concurrent delivery, and + // its parent being in the bound community does not make the stored + // comment that community's. + existing, err := m.storedCommentMapping(ctx, note.ID, communityIRI) + if err != nil { + return nil, err + } + author, err := m.EnsureActor(ctx, draft.authorRef) if err != nil { return nil, err } - did, authorDID := author.DID, author.DID - if existing, err := m.objects.GetByAPID(ctx, note.ID); err == nil { + did, rkey, authorDID := author.DID, draft.rkey, author.DID + if existing != nil { // The repo a comment lives in IS its authorship claim, so authorship — // and with it the record's coordinates — is fixed at first // materialization, exactly as MaterializePost pins a post's. attributedTo @@ -206,34 +367,16 @@ func (m *Materializer) materializeCommentLeaf(ctx context.Context, note *ap.Obje if existing.AuthorDID != "" { authorDID = existing.AuthorDID } - } else if !errors.IsNotFound(err) { - return nil, fmt.Errorf("materialize: check mapping for %s: %w", note.ID, err) - } - - reply, communityDID, err := m.resolveReplyRefs(ctx, note) - if err != nil { - return nil, err } - content := markdownFromObject(note) - if content == "" { - return nil, skip(note.ID, "comment has no content") - } - - body, facets := bridgedRichText(content, 3000, 30000) - if body == "" { - // An HTML-only body can reduce to nothing once tags are stripped; a - // comment is nothing but its content, so drop it like a bodiless one. - return nil, skip(note.ID, "comment has no content") - } record := map[string]any{ "$type": CollectionComment, "reply": reply, - "content": body, + "content": draft.body, "createdAt": recordDatetime(note.Published.Time), } - if len(facets) > 0 { - record["facets"] = facets + if len(draft.facets) > 0 { + record["facets"] = draft.facets } if langs := recordLangs(note.Language); len(langs) > 0 { record["langs"] = langs diff --git a/internal/materialize/comments_test.go b/internal/materialize/comments_test.go index 35d58cd..bf55b14 100644 --- a/internal/materialize/comments_test.go +++ b/internal/materialize/comments_test.go @@ -63,7 +63,7 @@ func TestMissingParentChainThreeDeep(t *testing.T) { leaf := serveThread(t, h) ctx := context.Background() - res, err := h.m.MaterializeComment(ctx, leaf) + res, err := h.m.MaterializeComment(ctx, leaf, groupID) require.NoError(t, err) pageMapping, err := h.objects.GetByAPID(ctx, pageID) @@ -137,7 +137,7 @@ func TestCommentCycleGuard(t *testing.T) { h.serveObject("/comment/9001", a) h.serveObject("/comment/9002", b) - _, err := h.m.MaterializeComment(context.Background(), objectFromMap(t, a)) + _, err := h.m.MaterializeComment(context.Background(), objectFromMap(t, a), groupID) require.Error(t, err) assert.True(t, IsSkip(err), "a cycle must be a skip: %v", err) assert.Empty(t, h.firehoseEvents(), "cycle detection happens before any write") @@ -165,7 +165,7 @@ func TestCommentDepthCap(t *testing.T) { "https://lemmy.world/u/alice", fmt.Sprintf("https://lemmy.world/comment/d%d", depth-1), "too deep", "2026-07-07T05:00:00.000000Z") - _, err := h.m.MaterializeComment(context.Background(), objectFromMap(t, leafBody)) + _, err := h.m.MaterializeComment(context.Background(), objectFromMap(t, leafBody), groupID) require.Error(t, err) assert.True(t, IsSkip(err), "over-deep chains must be skipped: %v", err) assert.Empty(t, h.firehoseEvents(), "the cap fires before any write") @@ -179,7 +179,7 @@ func TestTombstonedParentDropsSubtree(t *testing.T) { ctx := context.Background() // Materialize the whole thread, then tombstone the middle comment. - _, err := h.m.MaterializeComment(ctx, leaf) + _, err := h.m.MaterializeComment(ctx, leaf, groupID) require.NoError(t, err) require.NoError(t, h.objects.SoftDelete(ctx, "https://sh.itjust.works/comment/2002")) @@ -191,7 +191,7 @@ func TestTombstonedParentDropsSubtree(t *testing.T) { "https://lemmy.zip/u/carol", "https://sh.itjust.works/comment/2002", "reply to deleted", "2026-07-07T06:00:00.000000Z")) - _, err = h.m.MaterializeComment(ctx, reply) + _, err = h.m.MaterializeComment(ctx, reply, groupID) require.Error(t, err) assert.True(t, IsSkip(err), "tombstoned parent must drop the subtree: %v", err) assert.True(t, errsIsNotFoundFalse(err), "a tombstone is not a missing parent") @@ -208,7 +208,7 @@ func TestCommentWithUnfetchableParent(t *testing.T) { leaf := objectFromMap(t, note("https://lemmy.zip/comment/5005", "https://lemmy.zip/u/carol", "https://lemmy.world/comment/nowhere", "orphan", "2026-07-07T06:00:00.000000Z")) - _, err := h.m.MaterializeComment(context.Background(), leaf) + _, err := h.m.MaterializeComment(context.Background(), leaf, groupID) require.Error(t, err) assert.True(t, IsSkip(err), "unfetchable parent must skip the subtree: %v", err) } diff --git a/internal/materialize/community.go b/internal/materialize/community.go index db46b94..8e78e8d 100644 --- a/internal/materialize/community.go +++ b/internal/materialize/community.go @@ -79,6 +79,58 @@ func CommunityDIDOf(ctx context.Context, records RecordGetter, mapping *store.AP } } +// requireBoundCommunity refuses content that does not belong to the community +// the current delivery is bound to. communityDID is what CommunityDIDOf (or +// resolveReplyRefs, which asks it) answered for the content or its parent; +// communityIRI is the AP Group that delivered it. +// +// The comparison is by DID, so the bound community has to be looked up, and a +// missing communities row fails closed: recreating it here would let a +// delivery mint the very community it is being checked against. An empty +// communityDID is "cannot be determined" and is refused, never matched — see +// CommunityDIDOf. +func (m *Materializer) requireBoundCommunity(ctx context.Context, communityDID, communityIRI, contentID string) error { + // A missing communities row, an empty bound DID and an empty content + // binding are inconsistent state rather than hostile input, so each of + // the three is logged where an operator will see it. + bound, err := m.communities.GetByAPGroupID(ctx, communityIRI) + switch { + case err == nil: + case errors.IsNotFound(err): + m.logger.Warn("bound community has no communities row; content refused", + "ap_id", contentID, "community", communityIRI) + return skip(contentID, "bound community "+communityIRI+" is not bridged; its content cannot be verified") + default: + return fmt.Errorf("materialize: look up bound community %s for %s: %w", communityIRI, contentID, err) + } + if bound.DID == "" { + m.logger.Warn("bound community has an empty DID; content refused", + "ap_id", contentID, "community", communityIRI) + return skip(contentID, "bound community "+communityIRI+" has no DID; its content cannot be verified") + } + if communityDID == "" { + m.logger.Warn("stored content's community binding is empty and cannot be derived; content refused", + "ap_id", contentID, "community", communityIRI) + return skip(contentID, "stored community binding is empty and cannot be derived; it cannot be bound to the delivering community "+communityIRI) + } + if communityDID != bound.DID { + return skip(contentID, fmt.Sprintf("content belongs to community %q, not to the delivering community %s", + communityDID, communityIRI)) + } + return nil +} + +// requireBoundCommunityIRI is the empty-binding guard every content entry +// point opens with. A caller that cannot say which community delivered the +// content has a bug; treating "" as "unbound" would reopen every check the +// binding exists for. +func requireBoundCommunityIRI(communityIRI string) error { + if communityIRI == "" { + return errors.NewValidationError("community_iri", "must name the community that delivered the content") + } + return nil +} + // mappingCommunityDID is the WRITE side of CommunityDIDOf: which community a // record being committed belongs to, for its mapping's community_did column. // The two must agree, so each post era answers from the same thing the read diff --git a/internal/materialize/community_binding_test.go b/internal/materialize/community_binding_test.go new file mode 100644 index 0000000..42b0b33 --- /dev/null +++ b/internal/materialize/community_binding_test.go @@ -0,0 +1,913 @@ +package materialize + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/ap" + "tidepool/internal/errors" + "tidepool/internal/store" + "tidepool/internal/testutil" +) + +// Community binding. Content reaches the bridge because a community delivered +// it, and that community is the only one the delivery can speak for. Every +// entry point therefore takes the bound community's AP Group IRI, and content +// that would land anywhere else — a Page naming another Group, a comment whose +// thread lives in another community, an edit retargeting stored content — is +// refused before anything is minted. +// +// The world these tests share: C and D are CO-HOSTED communities on +// lemmy.world, so nothing about authority tells them apart; the content and +// its authors live on a separate host, lemmy.zip. Every Group and Person a +// call could touch is served as a valid document, so a refusal can only come +// from the binding — never from a missing fixture or a consent marker. +const ( + bindingCommunityC = groupID // https://lemmy.world/c/technology + bindingCommunityD = otherGroupID // https://lemmy.world/c/elsewhere + bindingUnknownGroup = "https://lemmy.world/c/strangers" + bindingAuthor = "https://lemmy.zip/u/xavier" + bindingReplier = "https://lemmy.zip/u/yvonne" + bindingChainAuthor = "https://lemmy.zip/u/zora" +) + +// serveBindingWorld serves the Groups and Persons described above. +func serveBindingWorld(h *harness) { + h.serveFixture("/c/technology", "group_lemmy_world.json") + h.serveObject("/c/elsewhere", group(bindingCommunityD, "elsewhere", nil)) + h.serveObject("/c/strangers", group(bindingUnknownGroup, "strangers", nil)) + h.serveObject("/u/xavier", person(bindingAuthor, "xavier", nil)) + h.serveObject("/u/yvonne", person(bindingReplier, "yvonne", nil)) + h.serveObject("/u/zora", person(bindingChainAuthor, "zora", nil)) +} + +// ensureCommunities bridges the given Groups, the state a followed community +// is in before it delivers anything. +func ensureCommunities(t *testing.T, h *harness, groupIRIs ...string) { + t.Helper() + for _, iri := range groupIRIs { + _, err := h.m.EnsureCommunity(context.Background(), &ap.Object{ID: iri}) + require.NoError(t, err) + } +} + +// noteIn is note() addressed to a given community rather than groupID. +func noteIn(id, author, inReplyTo, community, markdown, published string) map[string]any { + doc := note(id, author, inReplyTo, markdown, published) + doc["audience"] = community + return doc +} + +// contentState is everything a materialization can create or change: one +// entry per mapping (with the CID and community binding it currently holds), +// every bridged actor, every community row, and the firehose length. +type contentState struct { + Mappings []string + Actors []string + Communities []string + Events int +} + +func captureContentState(t *testing.T, h *harness) contentState { + t.Helper() + db := testutil.DB(t) + column := func(query string) []string { + rows, err := db.Query(query) + require.NoError(t, err) + defer func() { _ = rows.Close() }() + var out []string + for rows.Next() { + var value string + require.NoError(t, rows.Scan(&value)) + out = append(out, value) + } + require.NoError(t, rows.Err()) + return out + } + return contentState{ + Mappings: column(`SELECT ap_id || ' ' || cid || ' ' || COALESCE(community_did, '-') || ' ' || + (deleted_at IS NULL)::text FROM ap_objects ORDER BY ap_id`), + Actors: column(`SELECT ap_actor_id FROM bridged_actors ORDER BY ap_actor_id`), + Communities: column(`SELECT ap_group_id FROM communities ORDER BY ap_group_id`), + Events: len(h.firehoseEvents()), + } +} + +func assertActorAbsent(t *testing.T, h *harness, apActorID string) { + t.Helper() + _, err := h.actors.GetByAPActorID(context.Background(), apActorID) + assert.True(t, errors.IsNotFound(err), "actor %s must not be bridged (err=%v)", apActorID, err) +} + +func assertCommunityAbsent(t *testing.T, h *harness, apGroupID string) { + t.Helper() + _, err := h.communities.GetByAPGroupID(context.Background(), apGroupID) + assert.True(t, errors.IsNotFound(err), "community %s must not have a row (err=%v)", apGroupID, err) + assertActorAbsent(t, h, apGroupID) +} + +// clearCommunityColumn turns a mapping into a pre-016 row, whose community +// binding is derived from its record rather than read off the column. +func clearCommunityColumn(t *testing.T, apID string) { + t.Helper() + res, err := testutil.DB(t).Exec(`UPDATE ap_objects SET community_did = NULL WHERE ap_id = $1`, apID) + require.NoError(t, err) + n, err := res.RowsAffected() + require.NoError(t, err) + require.EqualValues(t, 1, n, "no mapping to clear for %s", apID) +} + +// TestPostBoundToAnotherCommunityIsRefused (B1): a Page delivered for C that +// names another Group is not C's to deliver. It is refused before the Group +// or the author is bridged — whether the Group it names is a co-hosted +// community the bridge already follows or one it has never seen. +func TestPostBoundToAnotherCommunityIsRefused(t *testing.T) { + cases := []struct { + name string + audience string + }{ + {name: "page names co-hosted community D", audience: bindingCommunityD}, + {name: "page names an unknown Group", audience: bindingUnknownGroup}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + + const postID = "https://lemmy.zip/post/91001" + post := page(postID, bindingAuthor, tc.audience, "addressed elsewhere", "2026-10-01T10:00:00.000000Z") + h.serveObject("/post/91001", post) + before := captureContentState(t, h) + + res, err := h.m.MaterializePost(ctx, mustObject(t, post), bindingCommunityC) + require.Error(t, err, "C may not deliver a post that names %s", tc.audience) + assert.Nil(t, res) + assert.True(t, IsSkip(err), "a post outside the bound community is a skip, got %v", err) + + assert.Equal(t, before, captureContentState(t, h), "a refused post writes nothing") + assert.Equal(t, 0, countMappings(t, h, postID)) + assertActorAbsent(t, h, bindingAuthor) + assertCommunityAbsent(t, h, bindingUnknownGroup) + }) + } +} + +// TestCommentBoundToAnotherCommunityIsRefused (B3, B4, B6): a comment +// delivered for C must hang in a thread that belongs to C. The thread is +// decided by the comment's ancestry, so whichever way that ancestry leaves C — +// an unmapped root Page naming another Group, an already-mapped parent that +// belongs to D, a fetched ancestor whose body claims a comment stored in D, or +// a Page naming D that itself replies into C — the whole delivery is refused +// before any ancestor is materialized or rewritten, or any author is bridged. A parent whose bound community C has +// no communities row cannot be checked at all, and fails closed. +// +// Every leaf claims C as its own audience, the way a hostile or confused +// delivery would. +func TestCommentBoundToAnotherCommunityIsRefused(t *testing.T) { + const ( + postInD = "https://lemmy.zip/post/94001" + commentInD = "https://lemmy.zip/comment/94002" + leafID = "https://lemmy.zip/comment/94099" + ) + // threadInD materializes a post and a comment in D, delivered by D. + threadInD := func(t *testing.T, h *harness) { + t.Helper() + ctx := context.Background() + post := page(postInD, bindingAuthor, bindingCommunityD, "a thread in D", "2026-10-01T11:00:00.000000Z") + h.serveObject("/post/94001", post) + _, err := h.m.MaterializePost(ctx, mustObject(t, post), bindingCommunityD) + require.NoError(t, err) + comment := noteIn(commentInD, bindingAuthor, postInD, bindingCommunityD, + "a comment in D", "2026-10-01T11:01:00.000000Z") + h.serveObject("/comment/94002", comment) + _, err = h.m.MaterializeComment(ctx, mustObject(t, comment), bindingCommunityD) + require.NoError(t, err) + } + // postInC materializes a post in C, delivered by C, and returns its id. + postInC := func(t *testing.T, h *harness) string { + t.Helper() + const anchorInC = "https://lemmy.zip/post/95001" + post := page(anchorInC, bindingAuthor, bindingCommunityC, "a thread in C", "2026-10-01T10:30:00.000000Z") + h.serveObject("/post/95001", post) + _, err := h.m.MaterializePost(context.Background(), mustObject(t, post), bindingCommunityC) + require.NoError(t, err) + return anchorInC + } + // unmappedRootIn serves (without materializing) a root Page naming the + // given Group and a comment under it, and returns that comment's id. + unmappedRootIn := func(h *harness, community string) string { + const rootID = "https://lemmy.zip/post/93001" + const middleID = "https://lemmy.zip/comment/93002" + h.serveObject("/post/93001", page(rootID, bindingAuthor, community, "a root elsewhere", + "2026-10-01T12:00:00.000000Z")) + h.serveObject("/comment/93002", noteIn(middleID, bindingChainAuthor, rootID, community, + "a reply elsewhere", "2026-10-01T12:01:00.000000Z")) + return middleID + } + + cases := []struct { + name string + // prepare builds the world and returns the leaf's inReplyTo. + prepare func(t *testing.T, h *harness) string + // absent are mappings the refused call would otherwise have created. + absent []string + // verify, when set, checks stored content the refusal must leave alone. + verify func(t *testing.T, h *harness) + // leaf, when set, replaces the default leaf delivered under parentID. + leaf func(parentID string) map[string]any + }{ + { + name: "unmapped root page names co-hosted community D", + prepare: func(t *testing.T, h *harness) string { + return unmappedRootIn(h, bindingCommunityD) + }, + absent: []string{"https://lemmy.zip/post/93001", "https://lemmy.zip/comment/93002"}, + }, + { + name: "unmapped root page names an unknown Group", + prepare: func(t *testing.T, h *harness) string { + return unmappedRootIn(h, bindingUnknownGroup) + }, + absent: []string{"https://lemmy.zip/post/93001", "https://lemmy.zip/comment/93002"}, + }, + { + name: "parent is a mapped post in D", + prepare: func(t *testing.T, h *harness) string { + threadInD(t, h) + return postInD + }, + }, + { + name: "parent is a mapped comment in D", + prepare: func(t *testing.T, h *harness) string { + threadInD(t, h) + return commentInD + }, + }, + { + name: "unmapped comment chain anchors on a comment in D", + prepare: func(t *testing.T, h *harness) string { + threadInD(t, h) + const middleID = "https://lemmy.zip/comment/94003" + h.serveObject("/comment/94003", noteIn(middleID, bindingChainAuthor, commentInD, + bindingCommunityD, "a reply in D", "2026-10-01T11:02:00.000000Z")) + return middleID + }, + absent: []string{"https://lemmy.zip/comment/94003"}, + }, + { + name: "parent is a legacy comment in D with no community_did column", + prepare: func(t *testing.T, h *harness) string { + threadInD(t, h) + clearCommunityColumn(t, postInD) + clearCommunityColumn(t, commentInD) + return commentInD + }, + }, + { + // The leaf's parent is an unmapped IRI on D's comment host whose + // body claims the id of D's stored comment — same authority, so the + // id is accepted — and re-parents it under a post in C. The walk + // would anchor in C and then commit that body over D's comment. + name: "unmapped parent IRI serves a body aliasing a stored comment in D", + prepare: func(t *testing.T, h *harness) string { + threadInD(t, h) + anchorInC := postInC(t, h) + const aliasIRI = "https://lemmy.zip/comment/95010" + h.serveObject("/comment/95010", noteIn(commentInD, bindingAuthor, anchorInC, bindingCommunityC, + "a comment rewritten into C", "2026-10-01T11:01:00.000000Z")) + return aliasIRI + }, + absent: []string{"https://lemmy.zip/comment/95010"}, + verify: func(t *testing.T, h *harness) { + ctx := context.Background() + mapping, err := h.objects.GetByAPID(ctx, commentInD) + require.NoError(t, err) + record, _, err := h.manager.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + require.NoError(t, err) + assert.Equal(t, "a comment in D", record["content"], "D's comment keeps its own content") + resolved, err := CommunityDIDOf(ctx, h.manager, mapping) + require.NoError(t, err) + assert.Equal(t, testDIDFor("elsewhere", "lemmy.world"), resolved, "D's comment stays in D") + }, + }, + { + // A fetched Page ends the walk even when it carries an inReplyTo. + // This one names D and points at an unmapped comment under a post + // in C: the walk must refuse it on fetch, before that older comment + // is fetched or written. + name: "unmapped page naming D replies to an unmapped comment in C", + prepare: func(t *testing.T, h *harness) string { + anchorInC := postInC(t, h) + const ( + pageInD = "https://lemmy.zip/post/95101" + commentInC = "https://lemmy.zip/comment/95102" + ) + h.serveObject("/comment/95102", noteIn(commentInC, bindingChainAuthor, anchorInC, bindingCommunityC, + "a comment in C", "2026-10-01T10:31:00.000000Z")) + pageDoc := page(pageInD, bindingAuthor, bindingCommunityD, "a page in D with a parent", + "2026-10-01T10:32:00.000000Z") + pageDoc["inReplyTo"] = commentInC + h.serveObject("/post/95101", pageDoc) + return pageInD + }, + absent: []string{"https://lemmy.zip/post/95101", "https://lemmy.zip/comment/95102"}, + verify: func(t *testing.T, h *harness) { + assert.NotZero(t, h.hitCount("/post/95101"), "control: the page itself was fetched") + assert.Equal(t, 0, h.hitCount("/comment/95102"), "nothing older than the Page is fetched") + }, + }, + { + // The leaf's parent X is unmapped and serves a Note claiming another + // unmapped id X' on the same host, under a Page in C. Neither id is + // stored, so nothing anchors the alias: committing it would bind + // content fetched from X under X'. Refused before anything is + // minted — even though X' itself serves a valid Note in C. + name: "unmapped parent IRI serves a body aliasing another unmapped id", + prepare: func(t *testing.T, h *harness) string { + const ( + pageInC = "https://lemmy.zip/post/95201" + aliasIRI = "https://lemmy.zip/comment/95202" + aliasedID = "https://lemmy.zip/comment/95203" + ) + h.serveObject("/post/95201", page(pageInC, bindingAuthor, bindingCommunityC, "a page in C", + "2026-10-01T10:50:00.000000Z")) + aliased := noteIn(aliasedID, bindingChainAuthor, pageInC, bindingCommunityC, + "a comment under another id", "2026-10-01T10:51:00.000000Z") + h.serveObject("/comment/95202", aliased) + h.serveObject("/comment/95203", aliased) + return aliasIRI + }, + absent: []string{"https://lemmy.zip/post/95201", "https://lemmy.zip/comment/95202", + "https://lemmy.zip/comment/95203"}, + }, + { + // The leaf is refused on its own terms (attributedTo on another + // host than its id), so the unmapped chain above it — all in C — + // must not be materialized on its way to that refusal. + name: "leaf with a cross-authority author over an unmapped chain in C", + prepare: func(t *testing.T, h *harness) string { + h.serveObject("/u/mallory", person("https://lemmy.world/u/mallory", "mallory", nil)) + return unmappedRootIn(h, bindingCommunityC) + }, + leaf: func(parentID string) map[string]any { + return note(leafID, "https://lemmy.world/u/mallory", parentID, "a reply signed by another host", + "2026-10-01T14:00:00.000000Z") + }, + absent: []string{"https://lemmy.zip/post/93001", "https://lemmy.zip/comment/93002"}, + }, + { + // The refusal sits in the MIDDLE of the unmapped chain: M, between + // A and the leaf, is attributed to a person on another host than + // its own id. Every ancestor is in C, and the chain is refused all + // the same — so neither the Page nor A above M may be committed on + // the way to M's refusal, nor their authors bridged. + name: "middle ancestor with a cross-authority author in an unmapped chain in C", + prepare: func(t *testing.T, h *harness) string { + const middleID = "https://lemmy.zip/comment/93003" + h.serveObject("/u/mallory", person("https://lemmy.world/u/mallory", "mallory", nil)) + h.serveObject("/comment/93003", noteIn(middleID, "https://lemmy.world/u/mallory", + unmappedRootIn(h, bindingCommunityC), bindingCommunityC, "a reply signed by another host", + "2026-10-01T12:02:00.000000Z")) + return middleID + }, + absent: []string{"https://lemmy.zip/post/93001", "https://lemmy.zip/comment/93002", + "https://lemmy.zip/comment/93003"}, + verify: func(t *testing.T, h *harness) { + assertActorAbsent(t, h, bindingAuthor) + assertActorAbsent(t, h, "https://lemmy.world/u/mallory") + }, + }, + { + // Same, for a leaf with no content: a comment is nothing but its + // content, so it will be dropped, and its ancestors with it. + name: "leaf with no content over an unmapped chain in C", + prepare: func(t *testing.T, h *harness) string { + return unmappedRootIn(h, bindingCommunityC) + }, + leaf: func(parentID string) map[string]any { + doc := note(leafID, bindingReplier, parentID, "", "2026-10-01T14:00:00.000000Z") + doc["content"] = "" + delete(doc, "source") + return doc + }, + absent: []string{"https://lemmy.zip/post/93001", "https://lemmy.zip/comment/93002"}, + }, + { + // The alias claims a comment already stored in C and keeps it under + // C's post. Same community, but still an unmapped IRI speaking for a + // stored id: the stored comment is not rewritten and the alias is + // not mapped. + name: "unmapped parent IRI serves a body aliasing a stored comment in C", + prepare: func(t *testing.T, h *harness) string { + anchorInC := postInC(t, h) + const ( + commentInC = "https://lemmy.zip/comment/95301" + aliasIRI = "https://lemmy.zip/comment/95302" + ) + comment := noteIn(commentInC, bindingAuthor, anchorInC, bindingCommunityC, + "a comment in C", "2026-10-01T10:40:00.000000Z") + h.serveObject("/comment/95301", comment) + _, err := h.m.MaterializeComment(context.Background(), mustObject(t, comment), bindingCommunityC) + require.NoError(t, err) + h.serveObject("/comment/95302", noteIn(commentInC, bindingAuthor, anchorInC, bindingCommunityC, + "a comment rewritten through an alias", "2026-10-01T10:40:00.000000Z")) + return aliasIRI + }, + absent: []string{"https://lemmy.zip/comment/95302"}, + verify: func(t *testing.T, h *harness) { + ctx := context.Background() + mapping, err := h.objects.GetByAPID(ctx, "https://lemmy.zip/comment/95301") + require.NoError(t, err) + record, recordCID, err := h.manager.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + require.NoError(t, err) + assert.Equal(t, "a comment in C", record["content"], "C's comment keeps its own content") + assert.Equal(t, mapping.CID, recordCID, "C's comment record is the version its mapping names") + }, + }, + { + // The alias claims a comment in C that has since been deleted. + // A deleted ancestor drops the subtree; it is never re-fetched or + // resurrected through an alias. + name: "unmapped parent IRI serves a body aliasing a deleted comment", + prepare: func(t *testing.T, h *harness) string { + anchorInC := postInC(t, h) + const ( + deletedInC = "https://lemmy.zip/comment/95401" + aliasIRI = "https://lemmy.zip/comment/95402" + ) + comment := noteIn(deletedInC, bindingAuthor, anchorInC, bindingCommunityC, + "a comment since deleted", "2026-10-01T10:45:00.000000Z") + _, err := h.m.MaterializeComment(context.Background(), mustObject(t, comment), bindingCommunityC) + require.NoError(t, err) + require.NoError(t, h.m.HandleDeleteRecord(context.Background(), deletedInC)) + h.serveObject("/comment/95402", noteIn(deletedInC, bindingAuthor, anchorInC, bindingCommunityC, + "a deleted comment brought back", "2026-10-01T10:45:00.000000Z")) + return aliasIRI + }, + absent: []string{"https://lemmy.zip/comment/95402"}, + verify: func(t *testing.T, h *harness) { + mapping, err := h.objects.GetByAPID(context.Background(), "https://lemmy.zip/comment/95401") + require.NoError(t, err) + assert.True(t, mapping.IsDeleted(), "the deleted comment stays deleted") + }, + }, + { + name: "parent is in C but C has no communities row", + prepare: func(t *testing.T, h *harness) string { + const postInC = "https://lemmy.zip/post/96001" + post := page(postInC, bindingAuthor, bindingCommunityC, "a thread in C", "2026-10-01T13:00:00.000000Z") + h.serveObject("/post/96001", post) + _, err := h.m.MaterializePost(context.Background(), mustObject(t, post), bindingCommunityC) + require.NoError(t, err) + res, err := testutil.DB(t).Exec(`DELETE FROM communities WHERE ap_group_id = $1`, bindingCommunityC) + require.NoError(t, err) + n, err := res.RowsAffected() + require.NoError(t, err) + require.EqualValues(t, 1, n, "precondition: C had a communities row to remove") + return postInC + }, + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + + parentID := tc.prepare(t, h) + leaf := note(leafID, bindingReplier, parentID, "a reply delivered by C", "2026-10-01T14:00:00.000000Z") + if tc.leaf != nil { + leaf = tc.leaf(parentID) + } + before := captureContentState(t, h) + + res, err := h.m.MaterializeComment(ctx, mustObject(t, leaf), bindingCommunityC) + require.Error(t, err, "C may not deliver a comment whose thread is not C's") + assert.Nil(t, res) + assert.True(t, IsSkip(err), "a comment outside the bound community is a skip, got %v", err) + + assert.Equal(t, before, captureContentState(t, h), + "a refused comment writes nothing: no ancestor, no community, no author") + assert.Equal(t, 0, countMappings(t, h, leafID)) + for _, apID := range tc.absent { + assert.Equal(t, 0, countMappings(t, h, apID), "ancestor %s must not be materialized", apID) + } + assertActorAbsent(t, h, bindingReplier) + assertActorAbsent(t, h, bindingChainAuthor) + assertCommunityAbsent(t, h, bindingUnknownGroup) + if tc.verify != nil { + tc.verify(t, h) + } + }) + } +} + +// TestDeletedCommentRedeliveredOverUnmappedChainIsRefused: a comment whose +// mapping is soft-deleted stays deleted. Re-delivered by the community it was +// bound to — now replying to an unmapped Note chain in that same community — +// it is a skip refused before the walk: none of the ancestors is fetched or +// materialized, none of their authors is bridged, and the leaf is not +// resurrected. +func TestDeletedCommentRedeliveredOverUnmappedChainIsRefused(t *testing.T) { + const ( + anchorInC = "https://lemmy.zip/post/99001" + leafID = "https://lemmy.zip/comment/99002" + rootInC = "https://lemmy.zip/post/99003" + middleID = "https://lemmy.zip/comment/99004" + ) + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + + // The leaf is first posted in C, under a post in C by its own author, and + // then deleted. + anchor := page(anchorInC, bindingReplier, bindingCommunityC, "a thread in C", "2026-10-01T16:00:00.000000Z") + h.serveObject("/post/99001", anchor) + _, err := h.m.MaterializePost(ctx, mustObject(t, anchor), bindingCommunityC) + require.NoError(t, err) + original := note(leafID, bindingReplier, anchorInC, "a reply delivered by C", "2026-10-01T16:01:00.000000Z") + _, err = h.m.MaterializeComment(ctx, mustObject(t, original), bindingCommunityC) + require.NoError(t, err) + require.NoError(t, h.objects.SoftDelete(ctx, leafID)) + + // The re-delivery hangs it under an unmapped chain in C whose authors the + // bridge has never seen. + h.serveObject("/post/99003", page(rootInC, bindingAuthor, bindingCommunityC, "a root in C", + "2026-10-01T15:00:00.000000Z")) + h.serveObject("/comment/99004", noteIn(middleID, bindingChainAuthor, rootInC, bindingCommunityC, + "a reply in C", "2026-10-01T15:01:00.000000Z")) + redelivered := note(leafID, bindingReplier, middleID, "a reply delivered by C again", + "2026-10-01T16:01:00.000000Z") + before := captureContentState(t, h) + + res, err := h.m.MaterializeComment(ctx, mustObject(t, redelivered), bindingCommunityC) + require.Error(t, err, "a deleted comment is not re-materialized") + assert.Nil(t, res) + assert.True(t, IsSkip(err), "a deleted comment re-delivered is a skip, got %v", err) + + assert.Equal(t, before, captureContentState(t, h), + "a refused comment writes nothing: no ancestor, no author, no event") + assert.Equal(t, 0, countMappings(t, h, rootInC), "the root Page must not be materialized") + assert.Equal(t, 0, countMappings(t, h, middleID), "the middle Note must not be materialized") + assert.Equal(t, 0, h.hitCount("/post/99003"), "the root Page is never fetched") + assert.Equal(t, 0, h.hitCount("/comment/99004"), "the middle Note is never fetched") + assertActorAbsent(t, h, bindingAuthor) + assertActorAbsent(t, h, bindingChainAuthor) + + mapping, err := h.objects.GetByAPID(ctx, leafID) + require.NoError(t, err) + assert.True(t, mapping.IsDeleted(), "the leaf mapping stays deleted") +} + +// TestEmptyBoundCommunityFailsClosed (B6): a caller that cannot say which +// community delivered the content gets a validation error from every content +// entry point, and nothing is written. An empty binding is a programming +// error, never "unbound". +func TestEmptyBoundCommunityFailsClosed(t *testing.T) { + const ( + postID = "https://lemmy.zip/post/97001" + commentID = "https://lemmy.zip/comment/97002" + ) + post := page(postID, bindingAuthor, bindingCommunityC, "a post in C", "2026-10-01T15:00:00.000000Z") + comment := note(commentID, bindingReplier, postID, "a comment in C", "2026-10-01T15:01:00.000000Z") + + cases := []struct { + name string + call func(ctx context.Context, t *testing.T, h *harness) (*Result, error) + }{ + {name: "MaterializePost", call: func(ctx context.Context, t *testing.T, h *harness) (*Result, error) { + return h.m.MaterializePost(ctx, mustObject(t, post), "") + }}, + {name: "MaterializeComment", call: func(ctx context.Context, t *testing.T, h *harness) (*Result, error) { + return h.m.MaterializeComment(ctx, mustObject(t, comment), "") + }}, + {name: "HandleUpdate with a Page", call: func(ctx context.Context, t *testing.T, h *harness) (*Result, error) { + return h.m.HandleUpdate(ctx, mustObject(t, post), "") + }}, + {name: "HandleUpdate with a Note", call: func(ctx context.Context, t *testing.T, h *harness) (*Result, error) { + return h.m.HandleUpdate(ctx, mustObject(t, comment), "") + }}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + h.serveObject("/post/97001", post) + ctx := context.Background() + before := captureContentState(t, h) + + res, err := tc.call(ctx, t, h) + require.Error(t, err, "an empty bound community must be refused") + assert.Nil(t, res) + assert.True(t, errors.IsValidation(err), "an empty bound community is a validation error, got %v", err) + + assert.Equal(t, before, captureContentState(t, h), "a refused call writes nothing") + assert.Equal(t, 0, countMappings(t, h, postID)) + assert.Equal(t, 0, countMappings(t, h, commentID)) + assertCommunityAbsent(t, h, bindingCommunityC) + assertActorAbsent(t, h, bindingAuthor) + assertActorAbsent(t, h, bindingReplier) + }) + } +} + +// TestCommentChainRootedInBoundCommunityMaterializes (B7, guard): the binding +// refuses foreign threads, not unfamiliar ones. A comment delivered for C +// whose whole three-level ancestry is unmapped, rooted at a Page naming C, +// still pulls in every ancestor — and bridges C itself when C has no row yet, +// because the root names the bound community. +func TestCommentChainRootedInBoundCommunityMaterializes(t *testing.T) { + h := newHarness(t) + leaf := serveThread(t, h) + ctx := context.Background() + + _, err := h.communities.GetByAPGroupID(ctx, groupID) + require.True(t, errors.IsNotFound(err), "precondition: C has no communities row yet (err=%v)", err) + + _, err = h.m.MaterializeComment(ctx, leaf, groupID) + require.NoError(t, err) + + communityC := testDIDFor("technology", "lemmy.world") + community, err := h.communities.GetByAPGroupID(ctx, groupID) + require.NoError(t, err, "the root names the bound community, so C is bridged") + assert.Equal(t, communityC, community.DID) + + for _, apID := range []string{ + pageID, + "https://lemmy.world/comment/1001", + "https://sh.itjust.works/comment/2002", + "https://lemmy.zip/comment/3003", + } { + mapping, err := h.objects.GetByAPID(ctx, apID) + require.NoError(t, err, "%s must be materialized", apID) + resolved, err := CommunityDIDOf(ctx, h.manager, mapping) + require.NoError(t, err) + assert.Equal(t, communityC, resolved, "%s belongs to C", apID) + } +} + +// staleMappingObjects answers GetByAPID for one id as though that id had no +// mapping yet, and delegates everything else to the real store: the reads a +// delivery makes just before a concurrent delivery of the same id commits. +// +// staleReads bounds the stale window. Zero answers every read of staleAPID +// stale, so the concurrent commit is never seen. N answers the first N reads +// stale and every later one from the real store: the concurrent commit landed +// after the delivery's own binding checks and before its commit read the +// mapping. Reads are counted, not timed, so the interleaving is deterministic. +type staleMappingObjects struct { + store.APObjects + staleAPID string + staleReads int + reads int +} + +func (s *staleMappingObjects) GetByAPID(ctx context.Context, apID string) (*store.APObjectMapping, error) { + if apID == s.staleAPID { + s.reads++ + if s.staleReads == 0 || s.reads <= s.staleReads { + return nil, errors.NewNotFoundError("ap_object", apID) + } + } + return s.APObjects.GetByAPID(ctx, apID) +} + +// TestFirstMaterializationRaceKeepsRecordAndMapping: two deliveries of one +// unmapped comment, bound to D and to C, can both read "no mapping". When D's +// commits first, C's must not overwrite the record under the mapping D's +// commit bound: the comment stays D's — record, reply refs and mapping alike. +// The stale read is simulated deterministically by a Materializer whose store +// reports the comment unmapped. +func TestFirstMaterializationRaceKeepsRecordAndMapping(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + + const ( + postInD = "https://lemmy.zip/post/98001" + postInC = "https://lemmy.zip/post/98002" + commentID = "https://lemmy.zip/comment/98003" + published = "2026-10-01T16:01:00.000000Z" + ) + for _, p := range []struct{ path, id, community string }{ + {"/post/98001", postInD, bindingCommunityD}, + {"/post/98002", postInC, bindingCommunityC}, + } { + doc := page(p.id, bindingAuthor, p.community, "a thread", "2026-10-01T16:00:00.000000Z") + h.serveObject(p.path, doc) + _, err := h.m.MaterializePost(ctx, mustObject(t, doc), p.community) + require.NoError(t, err) + } + postInDMapping, err := h.objects.GetByAPID(ctx, postInD) + require.NoError(t, err) + + inD := noteIn(commentID, bindingReplier, postInD, bindingCommunityD, "a comment in D", published) + first, err := h.m.MaterializeComment(ctx, mustObject(t, inD), bindingCommunityD) + require.NoError(t, err) + + inC := noteIn(commentID, bindingReplier, postInC, bindingCommunityC, "a comment moved into C", published) + h.serveObject("/comment/98003", inC) + before := captureContentState(t, h) + + racing := *h.m + racing.objects = &staleMappingObjects{APObjects: h.objects, staleAPID: commentID} + res, err := racing.MaterializeComment(ctx, mustObject(t, inC), bindingCommunityC) + require.Error(t, err, "the losing delivery must not commit over the winner's record") + assert.Nil(t, res) + assert.True(t, IsSkip(err), "losing a first-materialization race is a skip, got %v", err) + + assert.Equal(t, before, captureContentState(t, h), "the losing delivery writes nothing") + mapping, err := h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + assert.Equal(t, testDIDFor("elsewhere", "lemmy.world"), mapping.CommunityDID) + assert.Equal(t, first.ATURI, mapping.ATURI) + assert.Equal(t, first.CID, mapping.CID) + record, recordCID, err := h.manager.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + require.NoError(t, err) + assert.Equal(t, first.CID, recordCID, "the record is still the version the mapping names") + assert.Equal(t, "a comment in D", record["content"]) + reply, _ := record["reply"].(map[string]any) + parent, _ := reply["parent"].(map[string]any) + assert.Equal(t, postInDMapping.ATURI, parent["uri"], "the comment still replies to D's post") +} + +// TestPartiallyStaleRaceKeepsRecordAndMapping: the same race with a narrower +// window. The losing delivery, bound to C, runs every binding check while the +// content is still unmapped; the winner, bound to D, commits; only then does +// the loser's commit read the mapping. Finding it there is not licence to +// commit as an edit of D's content: the delivery was checked as a first +// materialization for C, and nothing checked it against D's binding. It is +// refused, and D's record and mapping are left exactly as the winner wrote +// them. +// +// staleReads is how many reads of the content's mapping the delivery makes +// BEFORE its commit reads it, counted on the sequential path: a post reads it +// once (MaterializePost's existing-mapping check); a comment twice +// (MaterializeComment's stored-mapping check and commitCommentLeaf's). The +// read after those is the commit's own, and it sees the winner. A fix that +// re-reads the mapping before the commit only moves that read into the +// winner's view earlier, which must refuse just the same. +func TestPartiallyStaleRaceKeepsRecordAndMapping(t *testing.T) { + const ( + postInD = "https://lemmy.zip/post/98201" + postInC = "https://lemmy.zip/post/98202" + racedPost = "https://lemmy.zip/post/98203" + racedNote = "https://lemmy.zip/comment/98204" + ) + communityD := testDIDFor("elsewhere", "lemmy.world") + + cases := []struct { + name string + apID string + staleReads int + // win is the winner's delivery, bound to D. + win func(ctx context.Context, t *testing.T, h *harness) *Result + // lose is the losing delivery, bound to C, through the racing store. + lose func(ctx context.Context, t *testing.T, h *harness, m *Materializer) (*Result, error) + // verify checks the winner's record survived. + verify func(t *testing.T, h *harness, record map[string]any) + }{ + { + name: "postv2 post", + apID: racedPost, + staleReads: 1, + win: func(ctx context.Context, t *testing.T, h *harness) *Result { + doc := page(racedPost, bindingAuthor, bindingCommunityD, "a post in D", "2026-10-01T16:10:00.000000Z") + res, err := h.m.MaterializePost(ctx, mustObject(t, doc), bindingCommunityD) + require.NoError(t, err) + return res + }, + lose: func(ctx context.Context, t *testing.T, h *harness, m *Materializer) (*Result, error) { + doc := page(racedPost, bindingAuthor, bindingCommunityC, "a post moved into C", + "2026-10-01T16:10:00.000000Z") + h.serveObject("/post/98203", doc) + return m.MaterializePost(ctx, mustObject(t, doc), bindingCommunityC) + }, + verify: func(t *testing.T, h *harness, record map[string]any) { + assert.Equal(t, "a post in D", record["title"]) + assert.Equal(t, communityD, record["community"], "the post still names D") + }, + }, + { + name: "comment", + apID: racedNote, + staleReads: 2, + win: func(ctx context.Context, t *testing.T, h *harness) *Result { + for _, p := range []struct{ path, id, community string }{ + {"/post/98201", postInD, bindingCommunityD}, + {"/post/98202", postInC, bindingCommunityC}, + } { + doc := page(p.id, bindingAuthor, p.community, "a thread", "2026-10-01T16:00:00.000000Z") + h.serveObject(p.path, doc) + _, err := h.m.MaterializePost(ctx, mustObject(t, doc), p.community) + require.NoError(t, err) + } + doc := noteIn(racedNote, bindingReplier, postInD, bindingCommunityD, "a comment in D", + "2026-10-01T16:11:00.000000Z") + res, err := h.m.MaterializeComment(ctx, mustObject(t, doc), bindingCommunityD) + require.NoError(t, err) + return res + }, + lose: func(ctx context.Context, t *testing.T, h *harness, m *Materializer) (*Result, error) { + doc := noteIn(racedNote, bindingReplier, postInC, bindingCommunityC, "a comment moved into C", + "2026-10-01T16:11:00.000000Z") + h.serveObject("/comment/98204", doc) + return m.MaterializeComment(ctx, mustObject(t, doc), bindingCommunityC) + }, + verify: func(t *testing.T, h *harness, record map[string]any) { + assert.Equal(t, "a comment in D", record["content"]) + postInDMapping, err := h.objects.GetByAPID(context.Background(), postInD) + require.NoError(t, err) + reply, _ := record["reply"].(map[string]any) + parent, _ := reply["parent"].(map[string]any) + assert.Equal(t, postInDMapping.ATURI, parent["uri"], "the comment still replies to D's post") + }, + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + + first := tc.win(ctx, t, h) + before := captureContentState(t, h) + + racing := *h.m + racing.objects = &staleMappingObjects{APObjects: h.objects, staleAPID: tc.apID, staleReads: tc.staleReads} + res, err := tc.lose(ctx, t, h, &racing) + require.Error(t, err, "the losing delivery must not commit over the winner's record") + assert.Nil(t, res) + assert.True(t, IsSkip(err), "losing a first-materialization race is a skip, got %v", err) + + assert.Equal(t, before, captureContentState(t, h), "the losing delivery writes nothing") + mapping, err := h.objects.GetByAPID(ctx, tc.apID) + require.NoError(t, err) + assert.Equal(t, communityD, mapping.CommunityDID) + assert.Equal(t, first.ATURI, mapping.ATURI) + assert.Equal(t, first.CID, mapping.CID) + record, recordCID, err := h.manager.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + require.NoError(t, err) + assert.Equal(t, first.CID, recordCID, "the record is still the version the mapping names") + tc.verify(t, h, record) + }) + } +} + +// TestDeletedLegacyCommentKeepsItsCommunityForRestore: a pre-016 comment +// mapping has no community_did, and its community is derived from the records. +// Deleting the comment deletes the very record that derivation reads, so the +// delete persists the derived community on the mapping first — otherwise a +// later restore, bound to the comment's own community, could no longer prove +// the comment is that community's and would be refused. +func TestDeletedLegacyCommentKeepsItsCommunityForRestore(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC) + ctx := context.Background() + + const ( + postID = "https://lemmy.zip/post/98101" + commentID = "https://lemmy.zip/comment/98102" + ) + post := page(postID, bindingAuthor, bindingCommunityC, "a thread in C", "2026-10-01T17:00:00.000000Z") + h.serveObject("/post/98101", post) + _, err := h.m.MaterializePost(ctx, mustObject(t, post), bindingCommunityC) + require.NoError(t, err) + comment := noteIn(commentID, bindingReplier, postID, bindingCommunityC, "a comment in C", + "2026-10-01T17:01:00.000000Z") + h.serveObject("/comment/98102", comment) + _, err = h.m.MaterializeComment(ctx, mustObject(t, comment), bindingCommunityC) + require.NoError(t, err) + clearCommunityColumn(t, commentID) + + require.NoError(t, h.m.HandleDeleteRecord(ctx, commentID)) + deleted, err := h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + require.True(t, deleted.IsDeleted(), "precondition: the delete soft-deleted the mapping") + assert.Equal(t, testDIDFor("technology", "lemmy.world"), deleted.CommunityDID, + "the delete persists the community it derived before the record is gone") + + // The restore path: clear the soft delete, then re-materialize the body. + require.NoError(t, h.objects.Restore(ctx, commentID)) + _, err = h.m.HandleUpdate(ctx, mustObject(t, comment), bindingCommunityC) + require.NoError(t, err, "a restore bound to the comment's own community re-materializes it") + + restored, err := h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + assert.False(t, restored.IsDeleted()) + assert.Equal(t, testDIDFor("technology", "lemmy.world"), restored.CommunityDID) + record, _, err := h.manager.GetRecord(ctx, restored.DID, restored.Collection, restored.RKey) + require.NoError(t, err, "the comment's record is back") + assert.Equal(t, "a comment in C", record["content"]) +} diff --git a/internal/materialize/forged_attribution_test.go b/internal/materialize/forged_attribution_test.go index 0e227bd..7627c47 100644 --- a/internal/materialize/forged_attribution_test.go +++ b/internal/materialize/forged_attribution_test.go @@ -29,13 +29,13 @@ func TestCommentEditCannotReattributeAuthor(t *testing.T) { h.serveObject("/u/victim", person("https://lemmy.world/u/victim", "victim", nil)) ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) const commentID = "https://lemmy.world/comment/80001" original := note(commentID, "https://lemmy.world/u/alice", pageID, "alice wrote this", "2026-07-08T16:00:00.000000Z") - _, err = h.m.MaterializeComment(ctx, objectFromMap(t, original)) + _, err = h.m.MaterializeComment(ctx, objectFromMap(t, original), groupID) require.NoError(t, err) aliceDID := testDIDFor("alice", "lemmy.world") @@ -51,7 +51,7 @@ func TestCommentEditCannotReattributeAuthor(t *testing.T) { // so the authority check below cannot be what refuses it. forged := note(commentID, "https://lemmy.world/u/victim", pageID, "alice wrote this (edited)", "2026-07-08T16:00:00.000000Z") - _, err = h.m.HandleUpdate(ctx, objectFromMap(t, forged)) + _, err = h.m.HandleUpdate(ctx, objectFromMap(t, forged), groupID) require.NoError(t, err, "a re-attributed edit must not error — it must simply not re-attribute") after, err := h.objects.GetByAPID(ctx, commentID) @@ -86,7 +86,7 @@ func TestCommentCrossAuthorityAttributionRefused(t *testing.T) { // The thread root is real and materialized, so a refusal below cannot be // the missing-parent protocol talking. - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) const victimIRI = "https://lemmy.world/u/victim" @@ -94,7 +94,7 @@ func TestCommentCrossAuthorityAttributionRefused(t *testing.T) { h.serveObject("/u/victim", person(victimIRI, "victim", nil)) res, err := h.m.MaterializeComment(ctx, objectFromMap(t, - note(forgedID, victimIRI, pageID, "words the victim never wrote", "2026-07-08T16:10:00.000000Z"))) + note(forgedID, victimIRI, pageID, "words the victim never wrote", "2026-07-08T16:10:00.000000Z")), groupID) require.Nil(t, res) require.Error(t, err) assert.True(t, IsSkip(err), "cross-authority attribution must be a skip, got %v", err) @@ -115,7 +115,7 @@ func TestPostCrossAuthorityAttributionRefused(t *testing.T) { const forgedID = "https://evil.example/post/1" res, err := h.m.MaterializePost(ctx, mustObject(t, - page(forgedID, personID, groupID, "not their post", "2026-07-08T16:20:00.000000Z"))) + page(forgedID, personID, groupID, "not their post", "2026-07-08T16:20:00.000000Z")), groupID) require.Nil(t, res) require.Error(t, err) assert.True(t, IsSkip(err), "cross-authority attribution must be a skip, got %v", err) diff --git a/internal/materialize/golden_test.go b/internal/materialize/golden_test.go index e287bc9..f7d4214 100644 --- a/internal/materialize/golden_test.go +++ b/internal/materialize/golden_test.go @@ -58,7 +58,7 @@ func TestGoldenPostAndProfiles(t *testing.T) { h := newHarness(t) h.serveLemmyWorldFixtures() - _, err := h.m.MaterializePost(context.Background(), loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(context.Background(), loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) assertGolden(t, "community_profile", h.recordFor(t, groupID)) @@ -87,7 +87,7 @@ func TestGoldenPostRichText(t *testing.T) { "```go\nfmt.Println(\"press\")\n```", } - _, err := h.m.MaterializePost(context.Background(), page) + _, err := h.m.MaterializePost(context.Background(), page, groupID) require.NoError(t, err) assertGolden(t, "post_richtext", h.recordFor(t, pageID)) @@ -110,7 +110,7 @@ func TestGoldenComment(t *testing.T) { "It has always been this way.", "2026-07-07T05:00:00.000000Z")) - _, err := h.m.MaterializeComment(context.Background(), loadFixtureObject(t, "note_lemmy_zip.json")) + _, err := h.m.MaterializeComment(context.Background(), loadFixtureObject(t, "note_lemmy_zip.json"), groupID) require.NoError(t, err) assertGolden(t, "comment", h.recordFor(t, noteID)) diff --git a/internal/materialize/hardening_test.go b/internal/materialize/hardening_test.go index c5749cc..b979a37 100644 --- a/internal/materialize/hardening_test.go +++ b/internal/materialize/hardening_test.go @@ -36,7 +36,7 @@ func TestDeleteActor_AccountEventVoteAndBlobScrub(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) communityDID := testDIDFor("technology", "lemmy.world") authorDID := testDIDFor("LeftLeaningFreedomFighters", "lemmy.world") @@ -112,7 +112,7 @@ func TestSuppressActor_ScrubsVotesToo(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) require.NoError(t, h.m.SuppressActor(ctx, personID)) @@ -236,7 +236,7 @@ func TestDeleteActor_BlobScrubFailureIsRetryable(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) authorDID := testDIDFor("LeftLeaningFreedomFighters", "lemmy.world") mapping, err := h.objects.GetByAPID(ctx, pageID) diff --git a/internal/materialize/materialize_test.go b/internal/materialize/materialize_test.go index 19be65e..1b9a4b5 100644 --- a/internal/materialize/materialize_test.go +++ b/internal/materialize/materialize_test.go @@ -117,6 +117,9 @@ type harness struct { mux *http.ServeMux fixtures *httptest.Server scrubbed *recordingScrubber + // hits counts GETs per served path (see hitCount). + hitsMu sync.Mutex + hits map[string]int } // recordingScrubber records ScrubVoter calls (the task-11 vote-scrub hook). @@ -185,6 +188,7 @@ func newHarness(t *testing.T) *harness { t: t, m: m, manager: manager, objects: objects, actors: actors, communities: communities, mux: mux, fixtures: fixtures, scrubbed: scrubbed, + hits: map[string]int{}, } // Every pictrs-style image path serves fixed bytes by extension. mux.HandleFunc("/pictrs/image/", func(w http.ResponseWriter, r *http.Request) { @@ -217,11 +221,22 @@ func (h *harness) serveObject(path string, obj map[string]any) { func (h *harness) serveJSON(path string, body []byte) { h.mux.HandleFunc("GET "+path, func(w http.ResponseWriter, r *http.Request) { + h.hitsMu.Lock() + h.hits[path]++ + h.hitsMu.Unlock() w.Header().Set("Content-Type", ap.ContentTypeActivityJSON) _, _ = w.Write(body) }) } +// hitCount is how many times a path registered through serveJSON (and so +// serveObject/serveFixture) has been fetched. +func (h *harness) hitCount(path string) int { + h.hitsMu.Lock() + defer h.hitsMu.Unlock() + return h.hits[path] +} + // serveLemmyWorldFixtures registers the standard page/person/group trio. func (h *harness) serveLemmyWorldFixtures() { h.serveFixture("/post/49131386", "page_lemmy_world.json") @@ -298,7 +313,7 @@ func TestMaterializePostEndToEnd(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - res, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + res, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) require.False(t, res.NoOp) @@ -363,7 +378,7 @@ func TestEmissionOrdering(t *testing.T) { h := newHarness(t) h.serveLemmyWorldFixtures() - _, err := h.m.MaterializePost(context.Background(), loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(context.Background(), loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) seqOf := func(collection string) int64 { @@ -396,12 +411,12 @@ func TestIdempotentRematerialize(t *testing.T) { ctx := context.Background() page := loadFixtureObject(t, "page_lemmy_world.json") - first, err := h.m.MaterializePost(ctx, page) + first, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) require.False(t, first.NoOp) eventsBefore := len(h.firehoseEvents()) - second, err := h.m.MaterializePost(ctx, page) + second, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) assert.True(t, second.NoOp, "identical re-materialization must be a repo no-op") assert.Equal(t, first.ATURI, second.ATURI) @@ -427,7 +442,7 @@ func TestNobridgeActorNeverMaterialized(t *testing.T) { pageObj := loadFixtureObject(t, "page_lemmy_world.json") pageObj.AttributedTo = ap.Refs{ap.Object{ID: "https://lemmy.world/u/optout"}} - _, err := h.m.MaterializePost(context.Background(), pageObj) + _, err := h.m.MaterializePost(context.Background(), pageObj, groupID) require.Error(t, err) assert.True(t, IsSkip(err), "nobridge must be a skip, not a failure: %v", err) @@ -443,7 +458,7 @@ func TestPostWithoutPublishedSkipped(t *testing.T) { pageObj := loadFixtureObject(t, "page_lemmy_world.json") pageObj.Published = nil - _, err := h.m.MaterializePost(context.Background(), pageObj) + _, err := h.m.MaterializePost(context.Background(), pageObj, groupID) require.Error(t, err) assert.True(t, IsSkip(err)) } @@ -532,7 +547,7 @@ func TestPublicOnlyAudienceFallsBackToAddressing(t *testing.T) { pageObj.Audience = ap.Audience{spelling} pageObj.To = ap.Audience{groupID, ap.PublicAudience} - res, err := h.m.MaterializePost(ctx, pageObj) + res, err := h.m.MaterializePost(ctx, pageObj, groupID) require.NoError(t, err, "audience %q must not be mistaken for a community", spelling) assert.Equal(t, testDIDFor("LeftLeaningFreedomFighters", "lemmy.world"), res.DID, "the post lands in the author's repo") diff --git a/internal/materialize/materializer.go b/internal/materialize/materializer.go index 24fe3ac..c62688d 100644 --- a/internal/materialize/materializer.go +++ b/internal/materialize/materializer.go @@ -116,6 +116,20 @@ func IsSkip(err error) bool { return stderrors.Is(err, ErrSkipped) } func skip(apID, reason string) error { return &SkipError{APID: apID, Reason: reason} } +// errCommunityBindingRaced marks a content commit whose mapping upsert came +// back bound to a different community than the one the delivery was bound to: +// a concurrent first materialization of the same object bound it first, and the +// COALESCE in the upsert kept that binding. Returned from inside the commit +// transaction so the record write rolls back with it; commitRecord reports it +// as a skip. +var errCommunityBindingRaced = stderrors.New("materialize: content was bound to another community concurrently") + +// isContentCollection reports whether collection holds community content — +// posts of either era and comments — as opposed to profiles. +func isContentCollection(collection string) bool { + return collection == CollectionPost || collection == CollectionPostV2 || collection == CollectionComment +} + // requireSameAuthorityAuthor refuses content that attributes itself to an actor // on a DIFFERENT authority than the object's own id. // @@ -356,6 +370,8 @@ type Result struct { // for comments the author's repo IS did); communityDID records which // community's content it is, for the // membership binding announced deletes and announced votes authorize against. +// For content it is the community THIS delivery was checked against, and a +// mapping bound to any other community refuses the commit. func (m *Materializer) commitRecord(ctx context.Context, did, collection, rkey string, record map[string]any, obj *ap.Object, authorDID, communityDID string) (*Result, error) { // Don't resurrect deleted content. AP delivery is unordered, so a Create // or Update can arrive (or be re-delivered) after a Delete already @@ -390,9 +406,16 @@ func (m *Materializer) commitRecord(ctx context.Context, did, collection, rkey s if existing.IsDeleted() { return nil, skip(obj.ID, "object was deleted upstream; not resurrecting") } - carryForward = collection == CollectionPost || - collection == CollectionPostV2 || - collection == CollectionComment + carryForward = isContentCollection(collection) + // The caller checked this delivery against communityDID, possibly + // while the content was still unmapped. A mapping that has appeared + // since, bound elsewhere, is a concurrent first materialization that + // won: committing now would carry ITS binding forward and pass as an + // edit of content nothing checked this delivery against. A NULL + // binding (pre-016 rows) is filled below instead. + if carryForward && existing.CommunityDID != "" && existing.CommunityDID != communityDID { + return nil, skip(obj.ID, "a concurrent materialization bound it to another community first") + } storedCommunityDID = existing.CommunityDID storedThreadRoot = existing.ThreadRootATURI } else if !errors.IsNotFound(err) { @@ -442,6 +465,16 @@ func (m *Materializer) commitRecord(ctx context.Context, did, collection, rkey s } return fmt.Errorf("materialize: map %s: %w", obj.ID, mapErr) } + // The mapping read above can be stale: two first materializations of + // one object, bound to different communities, both find no mapping. + // The upsert keeps whichever binding landed first, so a stored + // binding other than the delivery's means this record would sit under + // another community's mapping — refuse it, record write included. The + // delivery's community, not mapping.CommunityDID: that one may have + // been carried forward from the very binding being checked. + if isContentCollection(collection) && stored.CommunityDID != communityDID { + return errCommunityBindingRaced + } return nil } @@ -452,6 +485,9 @@ func (m *Materializer) commitRecord(ctx context.Context, did, collection, rkey s return nil, err } res, err := m.repos.PutRecordTx(ctx, did, collection, rkey, record, putMapping) + if stderrors.Is(err, errCommunityBindingRaced) { + return nil, skip(obj.ID, "a concurrent materialization bound it to another community first") + } if err != nil { return nil, fmt.Errorf("materialize: put %s/%s/%s for %s: %w", did, collection, rkey, obj.ID, err) } @@ -479,6 +515,9 @@ func (m *Materializer) commitRecord(ctx context.Context, did, collection, rkey s } return nil, fmt.Errorf("materialize: put %s: record kept changing across %d attempts: %w", obj.ID, maxStatsCommitAttempts, err) } + if stderrors.Is(err, errCommunityBindingRaced) { + return nil, skip(obj.ID, "a concurrent materialization bound it to another community first") + } if err != nil { return nil, fmt.Errorf("materialize: put %s/%s/%s for %s: %w", did, collection, rkey, obj.ID, err) } diff --git a/internal/materialize/posts.go b/internal/materialize/posts.go index 2b095e5..ef507d2 100644 --- a/internal/materialize/posts.go +++ b/internal/materialize/posts.go @@ -23,7 +23,10 @@ import ( // how the two eras coexist: a post first materialized under the deprecated // collection keeps being updated there (no migration is planned; Coves // indexes both), while everything new is a postv2. -func (m *Materializer) MaterializePost(ctx context.Context, page *ap.Object) (*Result, error) { +func (m *Materializer) MaterializePost(ctx context.Context, page *ap.Object, communityIRI string) (*Result, error) { + if err := requireBoundCommunityIRI(communityIRI); err != nil { + return nil, err + } if page == nil || page.ID == "" { return nil, errors.NewValidationError("page", "must carry an AP object id") } @@ -50,6 +53,28 @@ func (m *Materializer) MaterializePost(ctx context.Context, page *ap.Object) (*R if groupRef == nil { return nil, skip(page.ID, "post names no community (no audience/to group IRI)") } + // The delivering community can only speak for itself. Checked before the + // Group is bridged, so a Page naming some other Group cannot get it a repo. + if groupRef.ID != communityIRI { + return nil, skip(page.ID, "post names community "+groupRef.ID+", not the delivering community "+communityIRI) + } + existing, err := m.objects.GetByAPID(ctx, page.ID) + switch { + case err == nil: + // A post's community is fixed at first materialization, so the + // audience above proves nothing about a post already stored: an edit + // retargeted at C is still D's post, and C may not re-commit it. + stored, storedErr := CommunityDIDOf(ctx, m.repos, existing) + if storedErr != nil { + return nil, storedErr + } + if err := m.requireBoundCommunity(ctx, stored, communityIRI, page.ID); err != nil { + return nil, err + } + case errors.IsNotFound(err): + default: + return nil, fmt.Errorf("materialize: check mapping for %s: %w", page.ID, err) + } community, err := m.EnsureCommunity(ctx, groupRef) if err != nil { return nil, err @@ -60,7 +85,7 @@ func (m *Materializer) MaterializePost(ctx context.Context, page *ap.Object) (*R } did, collection, authorDID := author.DID, CollectionPostV2, author.DID - if existing, err := m.objects.GetByAPID(ctx, page.ID); err == nil { + if existing != nil { did, collection, rkey = existing.DID, existing.Collection, existing.RKey // The repo a postv2 lives in IS its authorship claim, so authorship is // fixed at first materialization. attributedTo on an updated Page is @@ -72,8 +97,6 @@ func (m *Materializer) MaterializePost(ctx context.Context, page *ap.Object) (*R if existing.AuthorDID != "" { authorDID = existing.AuthorDID } - } else if !errors.IsNotFound(err) { - return nil, fmt.Errorf("materialize: check mapping for %s: %w", page.ID, err) } // The blob DID is the repo the record lands in, not the community: a blob diff --git a/internal/materialize/postv2_acceptance_test.go b/internal/materialize/postv2_acceptance_test.go index 0c1bddf..1f18836 100644 --- a/internal/materialize/postv2_acceptance_test.go +++ b/internal/materialize/postv2_acceptance_test.go @@ -42,7 +42,7 @@ func TestAcceptanceHealsOnRedelivery(t *testing.T) { rkey, err := recordRKey(page) require.NoError(t, err) - first, err := h.m.MaterializePost(ctx, page) + first, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) require.False(t, first.NoOp) @@ -64,7 +64,7 @@ func TestAcceptanceHealsOnRedelivery(t *testing.T) { } // The queue redelivers the same Create. - second, err := h.m.MaterializePost(ctx, page) + second, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) assert.True(t, second.NoOp, "an unchanged redelivery is an idempotent no-op at the postv2 commit") @@ -95,7 +95,7 @@ func TestAcceptanceIsNotOnTheMappingSpine(t *testing.T) { rkey, err := recordRKey(page) require.NoError(t, err) - _, err = h.m.MaterializePost(ctx, page) + _, err = h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) communityDID := testDIDFor("technology", "lemmy.world") @@ -154,7 +154,7 @@ func TestAcceptanceRedeliveryDoesNotChurnCommunityRepo(t *testing.T) { rkey, err := recordRKey(page) require.NoError(t, err) - first, err := h.m.MaterializePost(ctx, page) + first, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) require.False(t, first.NoOp) @@ -168,7 +168,7 @@ func TestAcceptanceRedeliveryDoesNotChurnCommunityRepo(t *testing.T) { eventsBefore := eventsForDID(t, h, communityDID) // Redeliver the identical post. - second, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + second, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) assert.True(t, second.NoOp, "an unchanged redelivery is a no-op at the postv2 commit") @@ -224,7 +224,7 @@ func TestAcceptanceRefusedWhenRemovalLandsAfterTheGuard(t *testing.T) { page := loadFixtureObject(t, "page_lemmy_world.json") rkey, err := recordRKey(page) require.NoError(t, err) - _, err = h.m.MaterializePost(ctx, page) + _, err = h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) communityDID := testDIDFor("technology", "lemmy.world") @@ -252,7 +252,7 @@ func TestAcceptanceRefusedWhenRemovalLandsAfterTheGuard(t *testing.T) { } // Redelivery of the same post drives acceptPost with the stale answer. - _, err = h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err = h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err, "a refused acceptance is not an error: the community's removal simply stands") @@ -281,7 +281,7 @@ func TestAcceptanceRefusedWhenRemovalLandsAfterTheGuard(t *testing.T) { _, err = h.manager.DeleteRecord(ctx, communityDID, CollectionRemoval, acceptanceRKey) require.NoError(t, err) - _, err = h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err = h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) _, _, err = h.manager.GetRecord(ctx, communityDID, CollectionAcceptance, acceptanceRKey) diff --git a/internal/materialize/postv2_create_test.go b/internal/materialize/postv2_create_test.go index dc91d9b..3b78b12 100644 --- a/internal/materialize/postv2_create_test.go +++ b/internal/materialize/postv2_create_test.go @@ -4,12 +4,14 @@ import ( "context" "net/http" "testing" + "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "tidepool/internal/ap" "tidepool/internal/errors" + "tidepool/internal/store" "tidepool/internal/testutil" ) @@ -17,12 +19,9 @@ import ( // shared pageID so each test owns its object graph. const ( imagePageID = "https://lemmy.world/post/70001" - hijackPageID = "https://lemmy.world/post/70002" statsPageID = "https://lemmy.world/post/70003" namePageID = "https://lemmy.world/post/70004" hijackAuthorPageID = "https://lemmy.world/post/70005" - threadPageID = "https://lemmy.world/post/70010" - otherThreadPageID = "https://lemmy.world/post/70011" otherGroupID = "https://lemmy.world/c/elsewhere" embedImageURL = "https://lemmy.world/media/postv2-embed.png" ) @@ -77,7 +76,7 @@ func TestPostV2EmbedBlobsLandInAuthorRepo(t *testing.T) { "name": "a picture", }} - _, err := h.m.MaterializePost(ctx, mustObject(t, imagePost)) + _, err := h.m.MaterializePost(ctx, mustObject(t, imagePost), groupID) require.NoError(t, err) communityDID := testDIDFor("technology", "lemmy.world") @@ -136,7 +135,7 @@ func TestPostV2DisplayNameComesFromFetchedActorNotInlineRef(t *testing.T) { post := pageAttributedToInline(namePageID, personID, inlineLie, groupID, "a post with a lying inline author", "2026-07-08T13:00:00.000000Z") - _, err := h.m.MaterializePost(ctx, mustObject(t, post)) + _, err := h.m.MaterializePost(ctx, mustObject(t, post), groupID) require.NoError(t, err) record := h.recordFor(t, namePageID) @@ -156,7 +155,7 @@ func TestPostV2DisplayNameComesFromFetchedActorNotInlineRef(t *testing.T) { post := pageAttributedToInline(namePageID, "https://lemmy.world/u/nameless", inlineLie, groupID, "a post by a nameless author", "2026-07-08T13:00:00.000000Z") - _, err := h.m.MaterializePost(ctx, mustObject(t, post)) + _, err := h.m.MaterializePost(ctx, mustObject(t, post), groupID) require.NoError(t, err) record := h.recordFor(t, namePageID) @@ -168,52 +167,118 @@ func TestPostV2DisplayNameComesFromFetchedActorNotInlineRef(t *testing.T) { }) } -// TestPostV2CommunityIsImmutableAcrossUpdates (B2): postv2.community is -// immutable. Coves' consumers DISCARD any update event that changes it, so a -// bridge that re-derived the community from a hostile or merely edited -// `audience` would emit an event Coves throws away — silently freezing the -// post at its pre-edit version — and, worse, would relocate the record. -// The stored community must be carried forward from the stored record. +// TestPostV2CommunityIsImmutableAcrossUpdates (B2): a post's community is +// decided once, when it is first materialized, and a later delivery cannot +// move it. postv2.community is immutable to Coves (its consumers DISCARD any +// update event that changes it), and a post's community is what authorizes +// announced deletes and votes against it. +// +// The community a delivery speaks for is fixed by the delivery itself — the +// community that announced it — never read off the edited audience. So when C +// re-delivers a post stored in D, even one whose audience now names C, it is +// not C's post: the delivery is refused, and the stored record, its CID and +// its mapping are left exactly as they were. That holds for a postv2, for a +// postv2 whose mapping predates the community_did column, and for a post +// written in the deprecated era into D's own repo. func TestPostV2CommunityIsImmutableAcrossUpdates(t *testing.T) { - h := newHarness(t) - h.serveLemmyWorldFixtures() - h.serveObject("/c/elsewhere", group(otherGroupID, "elsewhere", nil)) - ctx := context.Background() - - original := page(hijackPageID, personID, groupID, "a post in technology", "2026-07-08T11:00:00.000000Z") - _, err := h.m.MaterializePost(ctx, mustObject(t, original)) - require.NoError(t, err) - - before, err := h.objects.GetByAPID(ctx, hijackPageID) - require.NoError(t, err) - - communityDID := testDIDFor("technology", "lemmy.world") - otherCommunityDID := testDIDFor("elsewhere", "lemmy.world") - require.NotEqual(t, communityDID, otherCommunityDID) - - // The upstream edit now names a DIFFERENT group. - hijacked := page(hijackPageID, personID, otherGroupID, "a post in technology", "2026-07-08T11:00:00.000000Z") - _, err = h.m.HandleUpdate(ctx, mustObject(t, hijacked)) - require.NoError(t, err, "a retargeted audience must not error — it must simply not retarget the record") - - after, err := h.objects.GetByAPID(ctx, hijackPageID) - require.NoError(t, err) - assert.Equal(t, before.DID, after.DID, "the record must not be MOVED by an audience change") - assert.Equal(t, before.RKey, after.RKey) - - record, _, err := h.manager.GetRecord(ctx, after.DID, after.Collection, after.RKey) - require.NoError(t, err) - assert.Equal(t, communityDID, record["community"], - "community is immutable: it must be carried forward from the stored record, "+ - "NOT re-derived from the changed audience (Coves discards any event that changes it)") - - // And nothing may have been planted in the hijacked community's repo. - _, _, err = h.manager.GetRecord(ctx, otherCommunityDID, testPostV2Collection, after.RKey) - assert.True(t, errors.IsNotFound(err), - "no postv2 record may appear in the retargeted community's repo (err=%v)", err) - _, _, err = h.manager.GetRecord(ctx, otherCommunityDID, testLegacyPostCollection, after.RKey) - assert.True(t, errors.IsNotFound(err), - "no post record of any era may appear in the retargeted community's repo (err=%v)", err) + const postID = "https://lemmy.zip/post/92001" + const title = "a post in D" + publishedAt := time.Date(2026, 10, 1, 9, 0, 0, 0, time.UTC) + const published = "2026-10-01T09:00:00.000000Z" + communityD := testDIDFor("elsewhere", "lemmy.world") + + storedInD := func(t *testing.T, h *harness) { + t.Helper() + _, err := h.m.MaterializePost(context.Background(), + mustObject(t, page(postID, bindingAuthor, bindingCommunityD, title, published)), bindingCommunityD) + require.NoError(t, err) + } + cases := []struct { + name string + prepare func(t *testing.T, h *harness) + }{ + {name: "postv2 stored in D", prepare: storedInD}, + {name: "postv2 stored in D with no community_did column", prepare: func(t *testing.T, h *harness) { + storedInD(t, h) + clearCommunityColumn(t, postID) + }}, + {name: "legacy community.post in D's repo", prepare: func(t *testing.T, h *harness) { + ctx := context.Background() + community, err := h.m.EnsureCommunity(ctx, &ap.Object{ID: bindingCommunityD}) + require.NoError(t, err) + author, err := h.m.EnsureActor(ctx, &ap.Object{ID: bindingAuthor}) + require.NoError(t, err) + rkey, err := recordRKey(mustObject(t, page(postID, bindingAuthor, bindingCommunityD, title, published))) + require.NoError(t, err) + commit, err := h.manager.PutRecord(ctx, community.DID, CollectionPost, rkey, map[string]any{ + "$type": CollectionPost, + "community": community.DID, + "author": author.DID, + "createdAt": recordDatetime(publishedAt), + "title": title, + }) + require.NoError(t, err) + // Pre-016: the legacy era never filled community_did. + _, err = h.objects.PutMapping(ctx, store.APObjectMapping{ + APID: postID, + APType: "Page", + OriginInstance: "lemmy.zip", + Origin: store.OriginFediverse, + DID: community.DID, + AuthorDID: author.DID, + Collection: CollectionPost, + RKey: rkey, + CID: commit.RecordCID, + PublishedAt: &publishedAt, + }) + require.NoError(t, err) + }}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + tc.prepare(t, h) + + before, err := h.objects.GetByAPID(ctx, postID) + require.NoError(t, err) + recordBefore, cidBefore, err := h.manager.GetRecord(ctx, before.DID, before.Collection, before.RKey) + require.NoError(t, err) + require.Equal(t, communityD, recordBefore["community"], "precondition: the post is stored in D") + stateBefore := captureContentState(t, h) + + // C re-delivers it, edited, with an audience retargeted at C. + retargeted := mustObject(t, page(postID, bindingAuthor, bindingCommunityC, title, published)) + retargeted.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} + res, err := h.m.HandleUpdate(ctx, retargeted, bindingCommunityC) + require.Error(t, err, "C may not re-deliver a post that is stored in D") + assert.Nil(t, res) + assert.True(t, IsSkip(err), "a re-delivery bound to another community is a skip, got %v", err) + + after, err := h.objects.GetByAPID(ctx, postID) + require.NoError(t, err) + assert.Equal(t, before.DID, after.DID, "the record must not be moved") + assert.Equal(t, before.Collection, after.Collection) + assert.Equal(t, before.RKey, after.RKey) + assert.Equal(t, before.CID, after.CID, "the mapping must still name the stored version") + assert.Equal(t, before.CommunityDID, after.CommunityDID, "the community binding must not change") + + recordAfter, cidAfter, err := h.manager.GetRecord(ctx, after.DID, after.Collection, after.RKey) + require.NoError(t, err) + assert.Equal(t, cidBefore, cidAfter, "the stored record must not be re-committed") + assert.Equal(t, recordBefore, recordAfter) + assert.Equal(t, communityD, recordAfter["community"]) + + resolved, err := CommunityDIDOf(ctx, h.manager, after) + require.NoError(t, err) + assert.Equal(t, communityD, resolved, "the post still belongs to D") + + assert.Equal(t, stateBefore, captureContentState(t, h), + "a refused re-delivery writes nothing: no record, no acceptance, no mapping change") + }) + } } // TestPostV2EditCarriesBridgedStatsForward (B3): the vote refresher stamps @@ -235,7 +300,7 @@ func TestPostV2EditCarriesBridgedStatsForward(t *testing.T) { ctx := context.Background() original := page(statsPageID, personID, groupID, "a post with votes", "2026-07-08T12:00:00.000000Z") - _, err := h.m.MaterializePost(ctx, mustObject(t, original)) + _, err := h.m.MaterializePost(ctx, mustObject(t, original), groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, statsPageID) @@ -246,7 +311,7 @@ func TestPostV2EditCarriesBridgedStatsForward(t *testing.T) { // The upstream edit changes the body only. edited := mustObject(t, page(statsPageID, personID, groupID, "a post with votes", "2026-07-08T12:00:00.000000Z")) edited.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} - res, err := h.m.HandleUpdate(ctx, edited) + res, err := h.m.HandleUpdate(ctx, edited, groupID) require.NoError(t, err) require.False(t, res.NoOp, "an edited body is a real commit") @@ -287,7 +352,7 @@ func TestPostV2AuthorIsImmutableAcrossUpdates(t *testing.T) { original := page(hijackAuthorPageID, personID, groupID, "a post by its real author", "2026-07-08T14:00:00.000000Z") - _, err := h.m.MaterializePost(ctx, mustObject(t, original)) + _, err := h.m.MaterializePost(ctx, mustObject(t, original), groupID) require.NoError(t, err) before, err := h.objects.GetByAPID(ctx, hijackAuthorPageID) @@ -301,7 +366,7 @@ func TestPostV2AuthorIsImmutableAcrossUpdates(t *testing.T) { // The edit now claims a different author. hijacked := page(hijackAuthorPageID, "https://lemmy.world/u/impostor", groupID, "a post by its real author", "2026-07-08T14:00:00.000000Z") - _, err = h.m.HandleUpdate(ctx, mustObject(t, hijacked)) + _, err = h.m.HandleUpdate(ctx, mustObject(t, hijacked), groupID) require.NoError(t, err, "a retargeted attributedTo must not error — it must simply not retarget the record") @@ -325,68 +390,93 @@ func TestPostV2AuthorIsImmutableAcrossUpdates(t *testing.T) { "no postv2 may appear in the claimed author's repo (err=%v)", err) } -// TestCommentCommunityIsImmutableAcrossUpdates (F2): a comment's community is -// fixed by the thread it was posted in, and an edit may not move it. +// TestCommentCommunityIsImmutableAcrossUpdates (B5): a comment's community is +// fixed by the thread it was posted in, and a later delivery may not move it. // -// This is the comment counterpart of B2, and it is a SECURITY property rather -// than a tidiness one. A comment's community_did is what authorizes announced -// deletes and binds announced votes: whoever the mapping says owns the comment -// may moderate it. `inReplyTo` on an updated Note is attacker-influenced, so a -// rebuild that re-derived the community from it would let a delivery hand -// community B moderation authority over a comment posted in community A — -// community A's members' content, moderated by a community they never posted -// to, with no moderator action on A's side at all. +// This is a SECURITY property rather than a tidiness one. A comment's +// community is what authorizes announced deletes and binds announced votes: +// whoever owns the comment may moderate it. So when C delivers an Update for a +// comment stored in D's thread — re-parented into a post in C, addressed to C — +// the comment is still D's, C cannot speak for it, and the delivery is refused +// with the record, its CID and its mapping unchanged. A comment whose mapping +// predates the community_did column is just as much D's. func TestCommentCommunityIsImmutableAcrossUpdates(t *testing.T) { - h := newHarness(t) - h.serveLemmyWorldFixtures() - h.serveObject("/c/elsewhere", group(otherGroupID, "elsewhere", nil)) - ctx := context.Background() - - communityA := testDIDFor("technology", "lemmy.world") - communityB := testDIDFor("elsewhere", "lemmy.world") - require.NotEqual(t, communityA, communityB) - - // A thread in community A, and an unrelated post in community B. - postA := page(threadPageID, personID, groupID, "thread root in A", "2026-07-08T15:00:00.000000Z") - h.serveObject("/post/70010", postA) - _, err := h.m.MaterializePost(ctx, mustObject(t, postA)) - require.NoError(t, err) - - postB := page(otherThreadPageID, personID, otherGroupID, "thread root in B", "2026-07-08T15:01:00.000000Z") - h.serveObject("/post/70011", postB) - _, err = h.m.MaterializePost(ctx, mustObject(t, postB)) - require.NoError(t, err) - - // A comment in A's thread. - const commentID = "https://lemmy.world/comment/70012" - comment := note(commentID, personID, threadPageID, "a comment in A", "2026-07-08T15:02:00.000000Z") - h.serveObject("/comment/70012", comment) - _, err = h.m.MaterializeComment(ctx, mustObject(t, comment)) - require.NoError(t, err) - - before, err := h.objects.GetByAPID(ctx, commentID) - require.NoError(t, err) - require.Equal(t, communityA, before.CommunityDID, "precondition: the comment belongs to community A") - - // The edit re-parents the comment into community B's thread. - hijacked := note(commentID, personID, otherThreadPageID, "a comment in A (edited)", - "2026-07-08T15:02:00.000000Z") - _, err = h.m.HandleUpdate(ctx, mustObject(t, hijacked)) - require.NoError(t, err, "a re-parented comment must not error — it must simply not be re-parented") - - after, err := h.objects.GetByAPID(ctx, commentID) - require.NoError(t, err) - assert.Equal(t, communityA, after.CommunityDID, - "the comment's community must stay A: community_did is what authorizes announced deletes "+ - "and binds announced votes, so moving it hands B moderation authority over A's content") - - // CommunityDIDOf is the exact function ingest's announced-delete - // authorization and votes' announced-vote binding both consult, so - // asserting it here is asserting who may moderate this comment. - resolved, err := CommunityDIDOf(ctx, h.manager, after) - require.NoError(t, err) - assert.Equal(t, communityA, resolved, - "the community that may moderate this comment must still be A, not the one the edit named") - assert.NotEqual(t, communityB, resolved, - "community B must NOT have acquired authority over a comment posted in A") + const ( + postInC = "https://lemmy.zip/post/95001" + postInD = "https://lemmy.zip/post/95002" + commentID = "https://lemmy.zip/comment/95003" + published = "2026-10-01T16:02:00.000000Z" + ) + communityD := testDIDFor("elsewhere", "lemmy.world") + + cases := []struct { + name string + legacyColumn bool + }{ + {name: "comment in D"}, + {name: "comment in D with no community_did column", legacyColumn: true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + h := newHarness(t) + serveBindingWorld(h) + ensureCommunities(t, h, bindingCommunityC, bindingCommunityD) + ctx := context.Background() + + // A thread in C, and a thread in D holding the comment. + threadC := page(postInC, bindingAuthor, bindingCommunityC, "thread root in C", "2026-10-01T16:00:00.000000Z") + h.serveObject("/post/95001", threadC) + _, err := h.m.MaterializePost(ctx, mustObject(t, threadC), bindingCommunityC) + require.NoError(t, err) + threadD := page(postInD, bindingAuthor, bindingCommunityD, "thread root in D", "2026-10-01T16:01:00.000000Z") + h.serveObject("/post/95002", threadD) + _, err = h.m.MaterializePost(ctx, mustObject(t, threadD), bindingCommunityD) + require.NoError(t, err) + comment := noteIn(commentID, bindingReplier, postInD, bindingCommunityD, "a comment in D", published) + h.serveObject("/comment/95003", comment) + _, err = h.m.MaterializeComment(ctx, mustObject(t, comment), bindingCommunityD) + require.NoError(t, err) + if tc.legacyColumn { + clearCommunityColumn(t, commentID) + } + + before, err := h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + resolvedBefore, err := CommunityDIDOf(ctx, h.manager, before) + require.NoError(t, err) + require.Equal(t, communityD, resolvedBefore, "precondition: the comment belongs to D") + recordBefore, cidBefore, err := h.manager.GetRecord(ctx, before.DID, before.Collection, before.RKey) + require.NoError(t, err) + stateBefore := captureContentState(t, h) + + // C delivers an edit re-parenting the comment into C's thread. + hijacked := noteIn(commentID, bindingReplier, postInC, bindingCommunityC, + "a comment in D (edited)", published) + res, err := h.m.HandleUpdate(ctx, mustObject(t, hijacked), bindingCommunityC) + require.Error(t, err, "C may not re-deliver a comment that is stored in D") + assert.Nil(t, res) + assert.True(t, IsSkip(err), "a re-delivery bound to another community is a skip, got %v", err) + + after, err := h.objects.GetByAPID(ctx, commentID) + require.NoError(t, err) + assert.Equal(t, before.DID, after.DID) + assert.Equal(t, before.RKey, after.RKey) + assert.Equal(t, before.CID, after.CID, "the mapping must still name the stored version") + assert.Equal(t, before.CommunityDID, after.CommunityDID, "the community binding must not change") + + recordAfter, cidAfter, err := h.manager.GetRecord(ctx, after.DID, after.Collection, after.RKey) + require.NoError(t, err) + assert.Equal(t, cidBefore, cidAfter, "the stored record must not be re-committed") + assert.Equal(t, recordBefore, recordAfter, "reply refs and content must be untouched") + + // CommunityDIDOf is what ingest's announced-delete authorization and + // votes' announced-vote binding both consult. + resolved, err := CommunityDIDOf(ctx, h.manager, after) + require.NoError(t, err) + assert.Equal(t, communityD, resolved, + "the community that may moderate this comment must still be D") + + assert.Equal(t, stateBefore, captureContentState(t, h), "a refused re-delivery writes nothing") + }) + } } diff --git a/internal/materialize/postv2_flip_test.go b/internal/materialize/postv2_flip_test.go index 19fa2fb..bd8526c 100644 --- a/internal/materialize/postv2_flip_test.go +++ b/internal/materialize/postv2_flip_test.go @@ -60,7 +60,7 @@ func TestOuterAcceptance_PostV2Flip(t *testing.T) { rkey, err := recordRKey(page) require.NoError(t, err) - res, err := h.m.MaterializePost(ctx, page) + res, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) require.False(t, res.NoOp) diff --git a/internal/materialize/postv2_repin_test.go b/internal/materialize/postv2_repin_test.go index a4cbd05..2db3ea9 100644 --- a/internal/materialize/postv2_repin_test.go +++ b/internal/materialize/postv2_repin_test.go @@ -33,7 +33,7 @@ func TestUpstreamEditRepinsAcceptance(t *testing.T) { rkey, err := recordRKey(page) require.NoError(t, err) - first, err := h.m.MaterializePost(ctx, page) + first, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) require.False(t, first.NoOp) @@ -53,7 +53,7 @@ func TestUpstreamEditRepinsAcceptance(t *testing.T) { // The author edits the post upstream. edited := loadFixtureObject(t, "page_lemmy_world.json") edited.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} - res, err := h.m.HandleUpdate(ctx, edited) + res, err := h.m.HandleUpdate(ctx, edited, groupID) require.NoError(t, err) require.False(t, res.NoOp, "an edited body is a real commit") require.NotEqual(t, first.CID, res.CID, "the edit must move the post's CID") @@ -82,7 +82,7 @@ func TestRedeliveredEditDoesNotChurnEitherRepo(t *testing.T) { ctx := context.Background() page := loadFixtureObject(t, "page_lemmy_world.json") - _, err := h.m.MaterializePost(ctx, page) + _, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) communityDID := testDIDFor("technology", "lemmy.world") @@ -90,7 +90,7 @@ func TestRedeliveredEditDoesNotChurnEitherRepo(t *testing.T) { editedOnce := loadFixtureObject(t, "page_lemmy_world.json") editedOnce.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} - _, err = h.m.HandleUpdate(ctx, editedOnce) + _, err = h.m.HandleUpdate(ctx, editedOnce, groupID) require.NoError(t, err) authorEvents := eventsForDID(t, h, authorDID) @@ -103,7 +103,7 @@ func TestRedeliveredEditDoesNotChurnEitherRepo(t *testing.T) { // The queue redelivers the identical Update. editedAgain := loadFixtureObject(t, "page_lemmy_world.json") editedAgain.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} - again, err := h.m.HandleUpdate(ctx, editedAgain) + again, err := h.m.HandleUpdate(ctx, editedAgain, groupID) require.NoError(t, err) assert.True(t, again.NoOp, "an unchanged redelivered edit is a no-op at the postv2 commit") @@ -131,7 +131,7 @@ func TestStaleAcceptanceRepinnedDespiteNoOpStamp(t *testing.T) { rkey, err := recordRKey(page) require.NoError(t, err) - first, err := h.m.MaterializePost(ctx, page) + first, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) require.False(t, first.NoOp) diff --git a/internal/materialize/richtext_test.go b/internal/materialize/richtext_test.go index 0826292..f93d845 100644 --- a/internal/materialize/richtext_test.go +++ b/internal/materialize/richtext_test.go @@ -565,7 +565,7 @@ func TestMaterializePostStoresRichText(t *testing.T) { MediaType: "text/markdown", } - _, err := h.m.MaterializePost(context.Background(), page) + _, err := h.m.MaterializePost(context.Background(), page, groupID) require.NoError(t, err) record := h.recordFor(t, pageID) @@ -622,7 +622,7 @@ func TestMaterializeCommentStoresRichText(t *testing.T) { pageID, "> quoted\n\n**bold** reply", "2026-07-07T04:00:00.000000Z") h.serveObject("/comment/9001", c1) - _, err := h.m.MaterializeComment(context.Background(), objectFromMap(t, c1)) + _, err := h.m.MaterializeComment(context.Background(), objectFromMap(t, c1), groupID) require.NoError(t, err) record := h.recordFor(t, commentID) @@ -834,7 +834,7 @@ func TestMaterializePostStoresSpoilerFacet(t *testing.T) { MediaType: "text/markdown", } - _, err := h.m.MaterializePost(context.Background(), page) + _, err := h.m.MaterializePost(context.Background(), page, groupID) require.NoError(t, err) record := h.recordFor(t, pageID) diff --git a/internal/materialize/security_test.go b/internal/materialize/security_test.go index ff6ab01..de01c48 100644 --- a/internal/materialize/security_test.go +++ b/internal/materialize/security_test.go @@ -80,7 +80,7 @@ func TestCommentThreadRootedAtNote_SkipsWithoutPanic(t *testing.T) { child := note("https://lemmy.zip/comment/child", "https://lemmy.zip/u/carol", "https://lemmy.zip/comment/root", "child", "2024-01-02T01:00:00.000000Z") - res, err := h.m.MaterializeComment(ctx, mustObject(t, child)) + res, err := h.m.MaterializeComment(ctx, mustObject(t, child), groupID) require.Nil(t, res) require.Error(t, err) assert.True(t, IsSkip(err), "parentless-Note root must be a skip, got %v", err) @@ -110,7 +110,7 @@ func TestAncestorCrossAuthorityID_Skips(t *testing.T) { child := note("https://lemmy.world/comment/child", personID, "https://lemmy.world/comment/parent", "child", "2024-01-02T01:00:00.000000Z") - res, err := h.m.MaterializeComment(ctx, mustObject(t, child)) + res, err := h.m.MaterializeComment(ctx, mustObject(t, child), groupID) require.Nil(t, res) require.Error(t, err) assert.True(t, IsSkip(err), "cross-authority ancestor id must be a skip, got %v", err) @@ -126,13 +126,13 @@ func TestCreateAfterDelete_DoesNotResurrect(t *testing.T) { ctx := context.Background() pageObj := loadFixtureObject(t, "page_lemmy_world.json") - _, err := h.m.MaterializePost(ctx, pageObj) + _, err := h.m.MaterializePost(ctx, pageObj, groupID) require.NoError(t, err) require.NoError(t, h.m.HandleDelete(ctx, pageObj.ID)) before := len(h.firehoseEvents()) - res, err := h.m.MaterializePost(ctx, pageObj) + res, err := h.m.MaterializePost(ctx, pageObj, groupID) require.Nil(t, res) require.Error(t, err) assert.True(t, IsSkip(err), "re-create after delete must be a skip, got %v", err) @@ -164,7 +164,7 @@ func TestNobridgeOnRefresh_ScrubsExistingContent(t *testing.T) { h.serveObject("/c/general", group(groupIRI, "general", nil)) // First pass: bridge + materialize the post (author = scrubme). - _, err := h.m.MaterializePost(ctx, mustObject(t, page(pageIRI, personIRI, groupIRI, "hello", "2024-02-01T00:00:00.000000Z"))) + _, err := h.m.MaterializePost(ctx, mustObject(t, page(pageIRI, personIRI, groupIRI, "hello", "2024-02-01T00:00:00.000000Z")), groupIRI) require.NoError(t, err) authorDID := testDIDFor("scrubme", "lemmy.world") @@ -203,7 +203,7 @@ func TestKnownPersonAsCommunity_Skips(t *testing.T) { // A post claiming that person's IRI as its community. pg := page("https://lemmy.world/post/2002", personIRI, personIRI, "x", "2024-02-02T00:00:00.000000Z") - res, err := h.m.MaterializePost(ctx, mustObject(t, pg)) + res, err := h.m.MaterializePost(ctx, mustObject(t, pg), personIRI) require.Nil(t, res) require.Error(t, err) assert.True(t, IsSkip(err), "person-as-community must be a skip, got %v", err) diff --git a/internal/materialize/updates.go b/internal/materialize/updates.go index e9807e9..df2e5dc 100644 --- a/internal/materialize/updates.go +++ b/internal/materialize/updates.go @@ -18,7 +18,12 @@ import ( // and comments flow through the same paths as creation (an Update for a // never-seen object simply materializes it), and actor/community updates // force a profile refresh regardless of the TTL. -func (m *Materializer) HandleUpdate(ctx context.Context, obj *ap.Object) (*Result, error) { +// +// communityIRI is the AP Group that delivered the update. It binds a Page, +// Article or Note exactly as MaterializePost and MaterializeComment do (empty +// is a validation error there), and is ignored for a Person or Group profile +// refresh, which is not community content. +func (m *Materializer) HandleUpdate(ctx context.Context, obj *ap.Object, communityIRI string) (*Result, error) { if obj == nil || obj.ID == "" { return nil, errors.NewValidationError("object", "must carry an AP object id") } @@ -34,9 +39,9 @@ func (m *Materializer) HandleUpdate(ctx context.Context, obj *ap.Object) (*Resul } return m.resultForMapping(ctx, obj.ID) case ap.TypePage, ap.TypeArticle: - return m.MaterializePost(ctx, obj) + return m.MaterializePost(ctx, obj, communityIRI) case ap.TypeNote: - return m.MaterializeComment(ctx, obj) + return m.MaterializeComment(ctx, obj, communityIRI) default: return nil, skip(obj.ID, "unsupported Update object type "+obj.Type) } @@ -293,6 +298,26 @@ func (m *Materializer) deleteMapping(ctx context.Context, mapping *store.APObjec return err } } + // A mapping written before migration 016 has no community_did, and its + // community is derived from the very record about to be deleted. Persist + // the derived binding first, or a later restore bound to the content's + // own community could no longer prove it is that community's. An + // underivable one is deleted as before and stays unbound — and so is one + // whose records cannot be read, or whose binding cannot be written, right + // now: the binding only keeps a later restore possible, and must not block + // the delete, nor with it a scrub. + if mapping.CommunityDID == "" { + communityDID, err := CommunityDIDOf(ctx, m.repos, mapping) + if err != nil { + m.logger.Warn("community of deleted content could not be derived; deleting without keeping a binding", + "ap_id", mapping.APID, "error", err) + } else if communityDID != "" { + if err := m.objects.SetCommunityDIDIfUnset(ctx, mapping.APID, communityDID); err != nil { + m.logger.Warn("community binding could not be kept; deleting without keeping a binding", + "ap_id", mapping.APID, "error", err) + } + } + } if _, err := m.repos.DeleteRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey); err != nil && !errors.IsNotFound(err) { return fmt.Errorf("materialize: delete record %s: %w", mapping.ATURI, err) } diff --git a/internal/materialize/updates_test.go b/internal/materialize/updates_test.go index 5b66ccb..00832d8 100644 --- a/internal/materialize/updates_test.go +++ b/internal/materialize/updates_test.go @@ -22,12 +22,12 @@ func TestUpdateRePutsSameRkey(t *testing.T) { ctx := context.Background() page := loadFixtureObject(t, "page_lemmy_world.json") - first, err := h.m.MaterializePost(ctx, page) + first, err := h.m.MaterializePost(ctx, page, groupID) require.NoError(t, err) edited := loadFixtureObject(t, "page_lemmy_world.json") edited.Source = &ap.Source{Content: "edited body text", MediaType: "text/markdown"} - second, err := h.m.HandleUpdate(ctx, edited) + second, err := h.m.HandleUpdate(ctx, edited, groupID) require.NoError(t, err) assert.Equal(t, first.ATURI, second.ATURI, "updates re-put under the same at-uri") @@ -55,7 +55,7 @@ func TestUpdatePersonRefreshesProfile(t *testing.T) { updated := loadFixtureObject(t, "person_lemmy_world.json") updated.Name = "Renamed Neelix" - res, err := h.m.HandleUpdate(ctx, updated) + res, err := h.m.HandleUpdate(ctx, updated, groupID) require.NoError(t, err) record, _, err := h.manager.GetRecord(ctx, res.DID, CollectionActorProfile, ProfileRKey) @@ -70,7 +70,7 @@ func TestHandleDeleteObject(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) mapping, err := h.objects.GetByAPID(ctx, pageID) require.NoError(t, err) @@ -97,7 +97,7 @@ func TestDeleteActorScrubsEverything(t *testing.T) { ctx := context.Background() // carol authors the leaf comment; the fixture person authors the post. - _, err := h.m.MaterializeComment(ctx, leaf) + _, err := h.m.MaterializeComment(ctx, leaf, groupID) require.NoError(t, err) authorAP := personID // fixture person: authored the post in the community repo @@ -191,7 +191,7 @@ func TestListByActorDIDSpine(t *testing.T) { h.serveLemmyWorldFixtures() ctx := context.Background() - _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json")) + _, err := h.m.MaterializePost(ctx, loadFixtureObject(t, "page_lemmy_world.json"), groupID) require.NoError(t, err) authorDID := testDIDFor("LeftLeaningFreedomFighters", "lemmy.world") diff --git a/internal/store/ap_objects.go b/internal/store/ap_objects.go index c2db4bd..d8a7024 100644 --- a/internal/store/ap_objects.go +++ b/internal/store/ap_objects.go @@ -264,6 +264,19 @@ func (r *postgresAPObjects) MarkCommunityAnnounced(ctx context.Context, apID str return nil } +func (r *postgresAPObjects) SetCommunityDIDIfUnset(ctx context.Context, apID, communityDID string) error { + if _, err := syntax.ParseDID(communityDID); err != nil { + return errors.NewValidationError("community_did", err.Error()) + } + _, err := r.db.ExecContext(ctx, + `UPDATE ap_objects SET community_did = $2 WHERE ap_id = $1 AND community_did IS NULL`, + apID, communityDID) + if err != nil { + return fmt.Errorf("set community_did for ap_object %q: %w", apID, err) + } + return nil +} + // validateMapping checks the atproto identifiers with indigo's syntax // package, defaults Origin to fediverse, and derives ATURI from // (DID, Collection, RKey). diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index f1ef3a8..0326f7b 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -146,6 +146,13 @@ type APObjects interface { // PutMapping upserts leave it alone. A missing mapping is an error // satisfying errors.IsNotFound. MarkCommunityAnnounced(ctx context.Context, apID string) error + + // SetCommunityDIDIfUnset records communityDID as the mapping's community + // binding only when it has none: a binding already made is never + // overwritten. It is how a binding derived from a record is kept before + // that record is deleted. Setting an already-bound or missing mapping is + // a no-op success. + SetCommunityDIDIfUnset(ctx context.Context, apID, communityDID string) error } // BridgedActors registers fediverse actors bridged into atproto and their diff --git a/internal/votes/e2e_test.go b/internal/votes/e2e_test.go index e49eab2..f0fab4b 100644 --- a/internal/votes/e2e_test.go +++ b/internal/votes/e2e_test.go @@ -51,17 +51,17 @@ const ( // materializer — votes must stay bridge-side. type stubMaterializer struct{ t *testing.T } -func (s *stubMaterializer) MaterializePost(context.Context, *ap.Object) (*materialize.Result, error) { +func (s *stubMaterializer) MaterializePost(context.Context, *ap.Object, string) (*materialize.Result, error) { s.t.Fatal("votes must never reach MaterializePost") return nil, nil } -func (s *stubMaterializer) MaterializeComment(context.Context, *ap.Object) (*materialize.Result, error) { +func (s *stubMaterializer) MaterializeComment(context.Context, *ap.Object, string) (*materialize.Result, error) { s.t.Fatal("votes must never reach MaterializeComment") return nil, nil } -func (s *stubMaterializer) HandleUpdate(context.Context, *ap.Object) (*materialize.Result, error) { +func (s *stubMaterializer) HandleUpdate(context.Context, *ap.Object, string) (*materialize.Result, error) { s.t.Fatal("votes must never reach HandleUpdate") return nil, nil }