diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index 85fd8bb..8863603 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -63,6 +63,32 @@ task documents and git history rather than this list. subtract Tidepool's written-back tally during subsequent Lemmy re-seeds. - Moderation federation and DMs. +## Postv2 flip (task 19) + +- **Legacy-set migration is deferred deliberately** (product decision + 2026-08-11): pre-flip bridged posts stay on the deprecated + `social.coves.community.post` collection in community repos forever; + every era-sensitive path dispatches on the mapping's collection. A + future migration needs Coves' PRD §11 remap decision (comment/vote + refs point at the old URIs) plus an old→new ledger. Recorded + consequence: Coves' legacy-drain gate never fires while bridged + legacy posts exist. +- **postv2 provenance costs one extra actor fetch per post + create/edit**: `originalAuthor.displayName` is read from the fetched, + authority-bound actor document on every materialization (deterministic + records; inline `attributedTo` is never trusted). The zero-extra-fetch + alternative is persisting the display name on `bridged_actors`, + refreshed with the profile — a migration plus plumbing. Revisit if + actor-fetch egress on post paths matters at scale. +- **`SetBridgedStats`' NoOp result carries `mapping.CID`**, which can be + stale relative to the record the unchanged-counts read just observed. + Callers currently ignore it; worth aligning with the read's CID. +- **No standalone acceptance-reconcile sweep**: every acceptance heal is + driven by an event on the post (redelivery, edit, stats sweep). A + quiet community's crash-window acceptance gap persists until any such + event. The decision-19 reconciliation job (task 18) is the natural + home for a periodic acceptance-vs-record pin audit. + ## Production rollout - Bluesky's public relay accepts new PDS hosts, but its documented default diff --git a/lexicons/README.md b/lexicons/README.md new file mode 100644 index 0000000..f0b4a94 --- /dev/null +++ b/lexicons/README.md @@ -0,0 +1,66 @@ +# Vendored Coves lexicons + +These JSON files are re-synced from the Coves repository via +`scripts/sync-lexicons.sh` (see that script's header for the ownership +rules — `social/coves/bridge/` is Tidepool-owned and survives re-syncs). +Records the bridge writes are validated against this exact set. + +## Bridge conventions for lexicon-unconstrained fields + +`social.coves.community.postv2` deliberately leaves `originalAuthor`, +`federatedFrom`, and `location` as `type: unknown` — consumers MUST +tolerate any object there and MUST NOT treat them as authorship claims +(authorship is the repository the record lives in). Tidepool is the +first writer of these fields; the shapes below are the bridge's +published convention (additive forever — fields may be added, never +renamed or repurposed): + +### `originalAuthor` — who wrote this on the origin platform + +```json +{ + "apId": "https://lemmy.world/u/ExampleUser", + "handle": "ExampleUser", + "instance": "lemmy.world", + "displayName": "Example Display Name" +} +``` + +- `apId` — the actor's canonical ActivityPub IRI (authority-bound at + materialization; never taken from inline objects). +- `handle` — the BARE origin-platform username (the local part, original + casing, derived from the canonical actor IRI). NOT the bridged + atproto handle: `instance` sits beside it, and the bridged handle is + derivable from the repo DID. +- `instance` — the origin host. +- `displayName` — optional; present only when the FETCHED, + authority-bound actor document asserts a `name`. Never sourced from + inline `attributedTo` objects on content (attacker-influenced), and + omitted rather than backfilled when the actor document has none. + +### `federatedFrom` — where this record was bridged from + +```json +{ + "platform": "lemmy", + "instance": "lemmy.world", + "apId": "https://lemmy.world/post/12345" +} +``` + +- `platform` — the origin software family (`"lemmy"` for anything + speaking Lemmy's dialect; PieFed/Mbin get their own values if and + when quirk handling diverges). +- `instance` — the origin host of the object. +- `apId` — the object's canonical AP id (the same id `ap_objects` maps). + +## Acceptance / removal record keys + +`social.coves.community.acceptance` and `social.coves.community.removal` +use the digest record key from Coves' PRD_AUTHOR_OWNED_POSTS §3.2: the +unpadded lowercase base32 (RFC 4648 standard alphabet) of the SHA-256 +digest of the canonical subject at-uri — fixed 52 chars, raw bytes +hashed with no normalization. The Go implementation is +`internal/materialize.SubjectRKey`, golden-pinned byte-for-byte against +Coves' test vectors: a divergence forks acceptance identity between the +two engines writing into community repos. diff --git a/tests/e2e/bridge_test.go b/tests/e2e/bridge_test.go index d98fe26..0c3159a 100644 --- a/tests/e2e/bridge_test.go +++ b/tests/e2e/bridge_test.go @@ -4,6 +4,7 @@ package e2e import ( "fmt" + "strings" "testing" "time" ) @@ -72,7 +73,7 @@ func TestPost_ActorProfileThenPost(t *testing.T) { user := h.registerUser(t, username) cursor := cursorNow() - l := h.newListener(t, cursor, colActorProfile, colPost) + l := h.newListener(t, cursor, colActorProfile, colPostV2, colAcceptance) title := "Hello from " + username // A link post. The url stays on the compose network (LOCAL-ONLY): Lemmy @@ -92,21 +93,67 @@ func TestPost_ActorProfileThenPost(t *testing.T) { t.Errorf("actor.profile rkey = %q, want %q", profileEv.Commit.RKey, rkeySelf) } - postEv := l.await("community.post create", func(e *jsEvent) bool { + postEv := l.await("community.postv2 create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && got == title }) - if got := recordField(t, postEv.Commit.Record, "author"); got != profileEv.Did { - t.Errorf("post author = %q, want the author's DID %q", got, profileEv.Did) + // Authorship IS the repo (PLAN.md decision 20): the post is committed + // into the AUTHOR's repo and carries no in-record author field at all. + if postEv.Did != profileEv.Did { + t.Errorf("postv2 landed in repo %q, want the AUTHOR's repo %q", postEv.Did, profileEv.Did) + } + if _, ok := fieldOf(postEv.Commit.Record, "author"); ok { + t.Errorf("postv2 carries an author field — authorship is the repo it lives in: %s", + truncate(postEv.Commit.Record, 300)) } if got := recordField(t, postEv.Commit.Record, "community"); got != sub.DID { - t.Errorf("post community = %q, want repo DID %q (posts live in the community's repo)", got, sub.DID) + t.Errorf("post community = %q, want the community's DID %q", got, sub.DID) } if got := recordField(t, postEv.Commit.Record, "embed", "external", "uri"); got != linkURL { t.Errorf("post embed.external.uri = %q, want the shared link %q", got, linkURL) } + + // Bridged provenance: who wrote it upstream, and where it came from. + // The lexicon leaves both shapes open, so THIS is the bridge's published + // convention and the suite is where it is pinned. + if got := recordField(t, postEv.Commit.Record, "originalAuthor", "apId"); !strings.Contains(got, username) { + t.Errorf("originalAuthor.apId = %q, want the origin actor IRI for %s", got, username) + } + if got := recordField(t, postEv.Commit.Record, "originalAuthor", "instance"); got != "lemmy" { + t.Errorf("originalAuthor.instance = %q, want %q", got, "lemmy") + } + if got := recordField(t, postEv.Commit.Record, "originalAuthor", "handle"); got != username { + t.Errorf("originalAuthor.handle = %q, want the ORIGIN username %q (not the bridged handle)", got, username) + } + if got := recordField(t, postEv.Commit.Record, "federatedFrom", "platform"); got != "lemmy" { + t.Errorf("federatedFrom.platform = %q, want %q", got, "lemmy") + } + if got := recordField(t, postEv.Commit.Record, "federatedFrom", "instance"); got != "lemmy" { + t.Errorf("federatedFrom.instance = %q, want %q", got, "lemmy") + } + + // The community's acceptance is what makes the post VISIBLE in it: + // Coves renders community surfaces from acceptance records alone, so a + // postv2 without one is a post nobody can see. It rides the firehose as + // its own commit in the COMMUNITY's repo, at the subject digest rkey. + acceptEv := l.await("community.acceptance create", func(e *jsEvent) bool { + uri, _ := fieldOf(e.Commit.Record, "subject", "uri") + return e.Commit.Collection == colAcceptance && uri == postEv.atURI() + }) + if acceptEv.Did != sub.DID { + t.Errorf("acceptance landed in repo %q, want the COMMUNITY's repo %q", acceptEv.Did, sub.DID) + } + if want := subjectRKey(postEv.atURI()); acceptEv.Commit.RKey != want { + t.Errorf("acceptance rkey = %q, want the subject digest %q — Coves derives the same key, "+ + "and a mismatch silently forks acceptance identity between the two engines", + acceptEv.Commit.RKey, want) + } + if got := recordField(t, acceptEv.Commit.Record, "subject", "cid"); got != postEv.Commit.CID { + t.Errorf("acceptance subject.cid = %q, want the post's cid %q", got, postEv.Commit.CID) + } + h.awaitAcceptanceConverged(t, sub.DID, postEv.Did, postEv.Commit.RKey) } // Scenario 3: comment + nested reply → comment records in the AUTHOR's repo @@ -122,7 +169,7 @@ func TestComments_StrongRefsResolve(t *testing.T) { user := h.registerUser(t, username) cursor := cursorNow() - l := h.newListener(t, cursor, colActorProfile, colPost, colComment) + l := h.newListener(t, cursor, colActorProfile, colPostV2, colComment) title := "Comment thread " + h.suffix post := user.createPost(t, community.ID, title, "root post") @@ -138,8 +185,13 @@ func TestComments_StrongRefsResolve(t *testing.T) { postEv := l.await("post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && got == title + return e.Commit.Collection == colPostV2 && e.Did == authorDID && got == title }) + // Post and comments now share one repo (the author's), so the thread's + // community linkage lives in the record rather than in the repo DID. + if got := recordField(t, postEv.Commit.Record, "community"); got != sub.DID { + t.Errorf("post community = %q, want %q", got, sub.DID) + } postURI, postCID := postEv.atURI(), postEv.Commit.CID comment := user.createComment(t, post.ID, 0, "top-level comment") @@ -199,14 +251,18 @@ func TestUpdateAndDelete(t *testing.T) { user := h.registerUser(t, username) cursor := cursorNow() - l := h.newListener(t, cursor, colPost, colComment) + l := h.newListener(t, cursor, colPostV2, colAcceptance, colComment) title := "Editable " + h.suffix post := user.createPost(t, community.ID, title, "original body") postEv := l.await("post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Commit.Operation == opCreate && - e.Did == sub.DID && got == title + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && + got == title + }) + l.await("acceptance for the new post", func(e *jsEvent) bool { + uri, _ := fieldOf(e.Commit.Record, "subject", "uri") + return e.Commit.Collection == colAcceptance && uri == postEv.atURI() }) comment := user.createComment(t, post.ID, 0, "doomed comment") @@ -218,8 +274,8 @@ func TestUpdateAndDelete(t *testing.T) { user.editPost(t, post.ID, "edited body") updateEv := l.await("post update", func(e *jsEvent) bool { - return e.Commit.Collection == colPost && e.Commit.Operation == opUpdate && - e.Did == sub.DID && e.Commit.RKey == postEv.Commit.RKey + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opUpdate && + e.Did == postEv.Did && e.Commit.RKey == postEv.Commit.RKey }) if got := recordField(t, updateEv.Commit.Record, "content"); got != "edited body" { t.Errorf("updated post content = %q, want %q", got, "edited body") @@ -228,6 +284,23 @@ func TestUpdateAndDelete(t *testing.T) { t.Error("update event carries the same cid as the create — no new commit?") } + // An edit moves the post's CID, so the community's acceptance must be + // RE-PINNED to the new version — at the same rkey, and without + // restamping createdAt (the community accepted this post once; re-pinning + // the version it accepts is not a new acceptance). An acceptance left on + // the old CID means Coves stops rendering the post the moment its author + // edits it. + repinEv := l.await("acceptance repin after the edit", func(e *jsEvent) bool { + cid, _ := fieldOf(e.Commit.Record, "subject", "cid") + return e.Commit.Collection == colAcceptance && e.Did == sub.DID && + e.Commit.RKey == subjectRKey(postEv.atURI()) && cid == updateEv.Commit.CID + }) + if repinEv.Commit.Operation != opUpdate { + t.Errorf("acceptance repin arrived as %q, want %q — the record is updated in place", + repinEv.Commit.Operation, opUpdate) + } + h.awaitAcceptanceConverged(t, sub.DID, postEv.Did, postEv.Commit.RKey) + user.deleteComment(t, comment.ID) deleteEv := l.await("comment delete", func(e *jsEvent) bool { return e.Commit.Collection == colComment && e.Commit.Operation == opDelete && @@ -259,8 +332,11 @@ func TestVotes_SideChannelOnly(t *testing.T) { post := author.createPost(t, community.ID, title, "vote on me") postEv := l.await("post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && got == title + return e.Commit.Collection == colPostV2 && got == title }) + if got := recordField(t, postEv.Commit.Record, "community"); got != sub.DID { + t.Errorf("post community = %q, want %q", got, sub.DID) + } postURI := postEv.atURI() // Upvote → aggregates show it. @@ -362,7 +438,7 @@ func TestBackfill_PreexistingPosts(t *testing.T) { voter.likePost(t, votedPost.ID, 1) cursor := cursorNow() - l := h.newListener(t, cursor, colCommunityProfile, colActorProfile, colPost) + l := h.newListener(t, cursor, colCommunityProfile, colActorProfile, colPostV2) sub := h.subscribeCommunity(t, "!"+name+"@lemmy") @@ -386,9 +462,9 @@ func TestBackfill_PreexistingPosts(t *testing.T) { switch e.Commit.Collection { case colActorProfile: return true // consume every profile to track the author set - case colPost: + case colPostV2: title, _ := fieldOf(e.Commit.Record, "title") - return e.Did == sub.DID && e.Commit.Operation == opCreate && titles[title] + return e.Commit.Operation == opCreate && titles[title] } return false }) @@ -397,7 +473,9 @@ func TestBackfill_PreexistingPosts(t *testing.T) { continue } title := recordField(t, ev.Commit.Record, "title") - postAuthors[recordField(t, ev.Commit.Record, "author")] = title + // The repo IS the authorship claim now — there is no author field to + // read, so the post's own DID is the author whose profile must exist. + postAuthors[ev.Did] = title if title == votedTitle { votedURI = ev.atURI() } @@ -446,14 +524,14 @@ func TestRestart_ReplayIsIdempotent(t *testing.T) { author.createPost(t, community.ID, title, "pre-restart post") cursor := cursorNow() - l := h.newListener(t, cursor, colCommunityProfile, colActorProfile, colPost) + l := h.newListener(t, cursor, colCommunityProfile, colActorProfile, colPostV2) handle := "!" + name + "@lemmy" sub := h.subscribeCommunity(t, handle) firstEv := l.await("pre-restart post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Did == sub.DID && e.Commit.Collection == colPost && got == title + return e.Commit.Collection == colPostV2 && got == title }) if firstEv.Commit.Operation != opCreate { t.Fatalf("pre-restart post arrived as %q, want %q", firstEv.Commit.Operation, opCreate) @@ -506,7 +584,7 @@ func TestRestart_ReplayIsIdempotent(t *testing.T) { // non-idempotent rebuild that duplicated a record's CONTENT would also carry // it and slip past — the modulo comparison still catches that (its content // differs), so this assertion keeps its teeth. - l2 := h.newListener(t, cursor, colCommunityProfile, colActorProfile, colPost) + l2 := h.newListener(t, cursor, colCommunityProfile, colActorProfile, colPostV2) title2 := "Post-restart " + h.suffix author.createPost(t, community.ID, title2, "after the bounce") @@ -538,7 +616,7 @@ func TestRestart_ReplayIsIdempotent(t *testing.T) { if !pureStats { seen[keyOf(ev)]++ } - if ev.Did == sub.DID && ev.Commit.Collection == colPost && !pureStats { + if ev.Commit.Collection == colPostV2 && !pureStats { switch got, _ := fieldOf(ev.Commit.Record, "title"); got { case gapTitle: sawGap, gapKey = true, keyOf(ev) @@ -563,7 +641,7 @@ func TestRestart_ReplayIsIdempotent(t *testing.T) { // The replayed pre-restart post must appear EXACTLY once: zero would // mean Jetstream lost history, twice would mean the backfill redo // re-committed it (deterministic-rkey idempotency broken). - firstKey := fmt.Sprintf("%s %s/%s", sub.DID, colPost, firstEv.Commit.RKey) + firstKey := fmt.Sprintf("%s %s/%s", firstEv.Did, colPostV2, firstEv.Commit.RKey) if n := seen[firstKey]; n != 1 { t.Errorf("pre-restart post replayed %d times (want exactly 1): %s", n, firstKey) } @@ -585,8 +663,8 @@ func TestRestart_ReplayIsIdempotent(t *testing.T) { func TestBurst_ConcurrentIngestionExactlyOnce(t *testing.T) { h := newHarness(t) - commA, subA := setupSubscribedCommunity(t, h, "ba") - commB, subB := setupSubscribedCommunity(t, h, "bb") + commA, _ := setupSubscribedCommunity(t, h, "ba") + commB, _ := setupSubscribedCommunity(t, h, "bb") users := []*lemmyClient{ h.registerUser(t, h.uniqueName(t, "ivy")), @@ -595,7 +673,7 @@ func TestBurst_ConcurrentIngestionExactlyOnce(t *testing.T) { } cursor := cursorNow() - l := h.newListener(t, cursor, colPost) + l := h.newListener(t, cursor, colPostV2) const perCommunity = 6 titles := make(map[string]bool, 2*perCommunity) @@ -615,20 +693,26 @@ func TestBurst_ConcurrentIngestionExactlyOnce(t *testing.T) { // stay pure now that non-matches are buffered and rescanned), and it is // inherently order-agnostic — which the relay's cross-repo reordering // demands anyway. + // Scoped by TITLE, not by repo DID: postv2 records land in each author's + // own repo (of which there are three here), so the two community DIDs no + // longer identify this scenario's posts. The titles are unique per run + // and are already the identity the exactly-once accounting turns on. seen := map[string]int{} keyTitle := map[string]string{} count := func(e *jsEvent) { if e.Kind != kindCommit || e.Commit == nil { return } - if e.Did != subA.DID && e.Did != subB.DID { + if e.Commit.Collection != colPostV2 { + return + } + title, ok := fieldOf(e.Commit.Record, "title") + if !ok || !titles[title] { return } key := fmt.Sprintf("%s %s/%s", e.Did, e.Commit.Collection, e.Commit.RKey) seen[key]++ - if title, ok := fieldOf(e.Commit.Record, "title"); ok && titles[title] { - keyTitle[key] = title - } + keyTitle[key] = title } // Every post arrives… @@ -643,9 +727,6 @@ func TestBurst_ConcurrentIngestionExactlyOnce(t *testing.T) { if ev.Kind != kindCommit || ev.Commit == nil || ev.Commit.Operation != opCreate { continue } - if ev.Did != subA.DID && ev.Did != subB.DID { - continue - } if title, ok := fieldOf(ev.Commit.Record, "title"); ok && titles[title] { matched[title] = true } diff --git a/tests/e2e/helpers.go b/tests/e2e/helpers.go index 099ac9f..8511d9d 100644 --- a/tests/e2e/helpers.go +++ b/tests/e2e/helpers.go @@ -19,6 +19,8 @@ package e2e import ( "bytes" "context" + "crypto/sha256" + "encoding/base32" "encoding/json" "fmt" "io" @@ -117,14 +119,43 @@ const sweepTimeout = 3 * time.Minute // ── Firehose vocabulary ──────────────────────────────────────────────────── -// Collections the bridge emits (task 05's materializer). +// Collections the bridge emits (task 05's materializer, task 19's flip). const ( colCommunityProfile = "social.coves.community.profile" colActorProfile = "social.coves.actor.profile" - colPost = "social.coves.community.post" - colComment = "social.coves.community.comment" + // colPost is the DEPRECATED post collection: a post in the COMMUNITY's + // repo carrying an in-record author. Nothing is created under it since + // the author-owned flip (PLAN.md decision 20), but records written + // before it are never migrated — Coves indexes both collections + // indefinitely — so the suite must still tolerate it on the firehose. + colPost = "social.coves.community.post" + // colPostV2 is the post collection after the flip: the record lives in + // the AUTHOR's repo and only NAMES its community. + colPostV2 = "social.coves.community.postv2" + // colAcceptance is the community's attestation that makes a postv2 + // visible in it; colRemoval is the moderation record that replaces it. + // Both live in the COMMUNITY's repo at the subject digest rkey. + colAcceptance = "social.coves.community.acceptance" + colRemoval = "social.coves.community.removal" + colComment = "social.coves.community.comment" ) +// subjectRKey derives an acceptance/removal record key from its subject's +// at-uri: unpadded lowercase base32 of the SHA-256 digest, a fixed 52 +// characters. +// +// Re-derived HERE rather than called from internal/materialize on purpose. +// This suite's job is to catch a change in the production derivation, and a +// test that computed the key with the same function under test would follow +// it silently wherever it went — including into a fork with Coves, whose +// engine writes acceptance records into these same community repos. The +// golden vectors that pin the encoding itself live in +// internal/materialize/subject_rkey_test.go. +func subjectRKey(subjectURI string) string { + digest := sha256.Sum256([]byte(subjectURI)) + return strings.ToLower(base32.StdEncoding.WithPadding(base32.NoPadding).EncodeToString(digest[:])) +} + // Event kinds, operations, and rkeys as they appear on the Jetstream wire. // Consts, not inline strings: a typo'd operation in an await predicate would // not fail compilation — it would burn a 90s timeout instead. @@ -152,6 +183,9 @@ var expectedCollections = map[string]bool{ colCommunityProfile: true, colActorProfile: true, colPost: true, + colPostV2: true, + colAcceptance: true, + colRemoval: true, colComment: true, } @@ -497,6 +531,32 @@ func (c *lemmyClient) deleteComment(t *testing.T, commentID int) { } } +// removePost is a MODERATOR action: Lemmy federates it as +// Announce{Delete{post}} whose inner Delete carries a `summary` — present +// (and EMPTY when no reason is given) is exactly what distinguishes it from +// the author's own delete, which carries no summary key at all. Verified +// against Lemmy 0.19.20 on the wire. +func (c *lemmyClient) removePost(t *testing.T, postID int, removed bool, reason string) { + t.Helper() + body := map[string]any{"post_id": postID, "removed": removed} + if reason != "" { + body["reason"] = reason + } + if err := c.do(http.MethodPost, "/api/v3/post/remove", body, nil); err != nil { + t.Fatalf("remove post %d (removed=%v): %v", postID, removed, err) + } +} + +// deletePost is the AUTHOR's own delete — no summary on the wire, so the +// bridge must read it as a self-delete and write no moderation record. +func (c *lemmyClient) deletePost(t *testing.T, postID int) { + t.Helper() + if err := c.do(http.MethodPost, "/api/v3/post/delete", + map[string]any{"post_id": postID, "deleted": true}, nil); err != nil { + t.Fatalf("delete post %d: %v", postID, err) + } +} + // likePost casts a vote: score 1 (up), -1 (down), 0 (retract). func (c *lemmyClient) likePost(t *testing.T, postID, score int) { t.Helper() @@ -1036,6 +1096,135 @@ func bridgedHandle(username string) string { return b.String() + ".lemmy." + bridgeHostname() } +// bridgeRecord is one record read back from the bridge's repo surface. +type bridgeRecord struct { + URI string `json:"uri"` + CID string `json:"cid"` + Value map[string]any `json:"value"` +} + +// bridgeGetRecord reads one record through com.atproto.repo.getRecord — the +// END-STATE surface, as distinct from the firehose. Both matter to this +// suite: the firehose proves the transition was PUBLISHED (a consumer that +// only ever tails the stream must be able to follow the moderation state), +// and the record read proves the repo actually converged. found=false is a +// clean not-found, not a transport failure (which fails the test). +func (h *harness) bridgeGetRecord(t *testing.T, did, collection, rkey string) (rec bridgeRecord, found bool) { + t.Helper() + res, err := h.bridgeXRPC("/xrpc/com.atproto.repo.getRecord", url.Values{ + "repo": {did}, "collection": {collection}, "rkey": {rkey}, + }) + if err != nil { + t.Fatalf("getRecord(%s/%s/%s): %v", did, collection, rkey, err) + } + switch res.status { + case http.StatusOK: + if err := json.Unmarshal(res.body, &rec); err != nil { + t.Fatalf("getRecord(%s/%s/%s): decode: %v", did, collection, rkey, err) + } + return rec, true + case http.StatusBadRequest, http.StatusNotFound: + return bridgeRecord{}, false + default: + t.Fatalf("getRecord(%s/%s/%s): unexpected status %d: %s", + did, collection, rkey, res.status, truncate(res.body, 200)) + return bridgeRecord{}, false + } +} + +// awaitRecordGone polls the repo surface until a record disappears. The +// firehose delete op and the repo read are separate observations of one +// commit, and the suite asserts both — but Jetstream delivery and the +// bridge's own read path settle independently, so the end-state check gets a +// bounded poll rather than a single racy read. +func (h *harness) awaitRecordGone(t *testing.T, did, collection, rkey, desc string) { + t.Helper() + deadline := time.Now().Add(eventTimeout) + for time.Now().Before(deadline) { + if _, found := h.bridgeGetRecord(t, did, collection, rkey); !found { + return + } + time.Sleep(250 * time.Millisecond) + } + t.Fatalf("%s: %s/%s/%s still present after %s", desc, did, collection, rkey, eventTimeout) +} + +// awaitRecordPresent is awaitRecordGone's positive twin. +func (h *harness) awaitRecordPresent(t *testing.T, did, collection, rkey, desc string) bridgeRecord { + t.Helper() + deadline := time.Now().Add(eventTimeout) + for time.Now().Before(deadline) { + if rec, found := h.bridgeGetRecord(t, did, collection, rkey); found { + return rec + } + time.Sleep(250 * time.Millisecond) + } + t.Fatalf("%s: %s/%s/%s never appeared within %s", desc, did, collection, rkey, eventTimeout) + return bridgeRecord{} +} + +// awaitAcceptanceConverged asserts the END STATE of the acceptance flow: the +// community's acceptance exists at the subject digest rkey and pins the +// post's CURRENT version. +// +// Convergence, not a fixed CID, is the assertion. The pinned CID is a moving +// target by design — every edit and every vote-stats stamp rewrites the post +// and drives a repin — so comparing against a CID captured earlier races any +// stamp that lands in between and fails on correct behavior. What must always +// hold is that the two agree once the system settles: an acceptance naming a +// version the post no longer has is exactly the state the lexicon tells Coves +// not to render, i.e. the post silently gone from its community. +// +// The per-transition assertions (a repin followed THIS edit, THIS stamp) live +// on the firehose awaits, where the specific commit is identified. +func (h *harness) awaitAcceptanceConverged(t *testing.T, communityDID, postRepoDID, postRKey string) { + t.Helper() + deadline := time.Now().Add(eventTimeout) + lastPost, lastPin := "", "" + for time.Now().Before(deadline) { + post, found := h.bridgeGetRecord(t, postRepoDID, colPostV2, postRKey) + if !found { + t.Fatalf("post %s/%s/%s is gone — cannot check its acceptance", + postRepoDID, colPostV2, postRKey) + } + acceptance, found := h.bridgeGetRecord(t, communityDID, colAcceptance, subjectRKey(post.URI)) + if found { + pinnedURI, _ := softString(acceptance.Value, "subject", "uri") + pinnedCID, _ := softString(acceptance.Value, "subject", "cid") + if pinnedURI != post.URI { + t.Fatalf("acceptance subject.uri = %q, want %q", pinnedURI, post.URI) + } + if _, ok := softString(acceptance.Value, "createdAt"); !ok { + t.Fatalf("acceptance carries no createdAt: %+v", acceptance.Value) + } + lastPost, lastPin = post.CID, pinnedCID + if pinnedCID == post.CID { + return + } + } + time.Sleep(500 * time.Millisecond) + } + t.Fatalf("acceptance never converged on the post's current version within %s: pins %q, post is at %q — "+ + "a stale pin means Coves stops rendering the post", eventTimeout, lastPin, lastPost) +} + +// softString reads a nested string without failing the test when it is +// absent (the polling loops above need to retry, not abort). +func softString(root map[string]any, path ...string) (string, bool) { + var cur any = root + for _, key := range path { + m, ok := cur.(map[string]any) + if !ok { + return "", false + } + if cur, ok = m[key]; !ok { + return "", false + } + } + s, ok := cur.(string) + return s, ok +} + // bridgeGetBlob fetches a stored blob from the bridge's // com.atproto.sync.getBlob (the AppView's media path for bridged images), // returning the bytes and the served Content-Type. @@ -1243,6 +1432,10 @@ type jsListener struct { // listeners each see a per-repo-ordered stream of their own. Only // touched from the test goroutine. lastRev map[string]string + // revPaths tracks which record paths have already been seen at the + // CURRENT rev per repo, so a multi-op commit's several events (which + // share one rev) are told apart from the same record emitted twice. + revPaths map[string]map[string]bool mu sync.Mutex readErr error // readLoop's terminal error, nil on deliberate close @@ -1311,12 +1504,13 @@ func (h *harness) newListener(t *testing.T, cursorMicros int64, collections ...s time.Sleep(2 * time.Second) } l := &jsListener{ - t: t, - conn: conn, - events: make(chan *jsEvent, 1024), - closed: make(chan struct{}), - done: make(chan struct{}), - lastRev: map[string]string{}, + t: t, + conn: conn, + events: make(chan *jsEvent, 1024), + closed: make(chan struct{}), + done: make(chan struct{}), + lastRev: map[string]string{}, + revPaths: map[string]map[string]bool{}, } go l.readLoop() t.Cleanup(l.close) @@ -1395,7 +1589,9 @@ func (l *jsListener) close() { // lands in; // - every create/update record must validate against the vendored Coves // lexicons (deletes carry no record); -// - commit revs must be strictly increasing per repo DID (see lastRev): +// - commit revs must be non-decreasing per repo DID (see lastRev), with +// equal revs admitted only across DISTINCT records — a multi-op commit +// reaches Jetstream as one event per op, all carrying its single rev: // the relay guarantees per-repo order even though cross-repo order is // lost, and this is the one place every consumed event passes through. func (l *jsListener) vetEvent(ev *jsEvent) { @@ -1404,11 +1600,36 @@ func (l *jsListener) vetEvent(ev *jsEvent) { return } if !expectedCollections[ev.Commit.Collection] { - l.t.Fatalf("unexpected collection on firehose: %s — only community/actor profiles, posts, and comments may ever appear (votes never become records)", ev) - } - if prev, ok := l.lastRev[ev.Did]; ok && ev.Commit.Rev <= prev { - l.t.Fatalf("per-repo rev order violated on firehose: %s has rev %q after rev %q — bigsky preserves per-repo commit order, so this is a relay/bridge ordering bug", - ev, ev.Commit.Rev, prev) + l.t.Fatalf("unexpected collection on firehose: %s — only community/actor profiles, posts (both eras), acceptances, removals, and comments may ever appear (votes never become records)", ev) + } + // Per-repo rev order. NON-decreasing, not strictly increasing: one + // commit may carry SEVERAL ops (the moderation transitions commit + // delete-acceptance + write-removal together so the firehose never shows + // a half-completed action), and Jetstream flattens a commit into one + // event per op — so those events legitimately share a rev. Equal revs are + // therefore admitted, but only for DISTINCT records: the same record + // twice at one rev, or any rev going backwards, is still the ordering bug + // this check exists to catch. + if prev, ok := l.lastRev[ev.Did]; ok { + path := ev.Commit.Collection + "/" + ev.Commit.RKey + switch { + case ev.Commit.Rev < prev: + l.t.Fatalf("per-repo rev order violated on firehose: %s has rev %q after rev %q — bigsky preserves per-repo commit order, so this is a relay/bridge ordering bug", + ev, ev.Commit.Rev, prev) + case ev.Commit.Rev == prev: + if l.revPaths[ev.Did][path] { + l.t.Fatalf("record %s emitted twice at the same rev %q on repo %s: %s — one commit may carry several ops, but never two ops on one record", + path, ev.Commit.Rev, ev.Did, ev) + } + default: + l.revPaths[ev.Did] = map[string]bool{} + } + if l.revPaths[ev.Did] == nil { + l.revPaths[ev.Did] = map[string]bool{} + } + l.revPaths[ev.Did][path] = true + } else { + l.revPaths[ev.Did] = map[string]bool{ev.Commit.Collection + "/" + ev.Commit.RKey: true} } l.lastRev[ev.Did] = ev.Commit.Rev if op := ev.Commit.Operation; op == opCreate || op == opUpdate { @@ -1431,11 +1652,22 @@ func (l *jsListener) vetEvent(ev *jsEvent) { // perform no content edits on the affected records, so an update-with-stats // there is unambiguously a stats emission — the helper is not a general // "is this only a stats change" classifier. +// isContentCollection reports whether a collection carries bridged CONTENT — +// the records the vote-stats refresher stamps bridgedStats onto. Posts of +// BOTH eras qualify: the deprecated collection because its records still +// exist and still get swept, postv2 because it is what every new post is. +// Acceptance and removal are deliberately excluded: they are the community's +// attestations ABOUT content, they carry no bridgedStats, and an acceptance +// repin riding a stats sweep is a different event from the stamp itself. +func isContentCollection(collection string) bool { + return collection == colPost || collection == colPostV2 || collection == colComment +} + func isBridgedStatsUpdate(ev *jsEvent) bool { if ev.Kind != kindCommit || ev.Commit == nil || ev.Commit.Operation != opUpdate { return false } - if ev.Commit.Collection != colPost && ev.Commit.Collection != colComment { + if !isContentCollection(ev.Commit.Collection) { return false } if len(ev.Commit.Record) == 0 { @@ -1477,7 +1709,7 @@ func (s *statsDedup) isPureStatsEmission(ev *jsEvent, key string) bool { if !seen || ev.Commit.Operation != opUpdate { return false } - if ev.Commit.Collection != colPost && ev.Commit.Collection != colComment { + if !isContentCollection(ev.Commit.Collection) { return false } return recordEqualsModuloBridgedStats(prev, ev.Commit.Record) @@ -1732,3 +1964,25 @@ func (h *harness) restartTidepool(t *testing.T) { time.Sleep(time.Second) } } + +// isPostCollection reports whether a collection holds POST records of either +// era — the deprecated community-repo collection or the author-repo postv2. +func isPostCollection(collection string) bool { + return collection == colPost || collection == colPostV2 +} + +// isAcceptanceRepinOf reports whether a commit is an update to the acceptance +// of one specific post — the community-repo event a stats stamp or an edit on +// that post produces. Scoped to the subject on purpose: it is used to excuse a +// legitimate settle inside a negative window, and an excuse that matched ANY +// acceptance would excuse exactly the leak the window exists to catch. +func isAcceptanceRepinOf(ev *jsEvent, postURI string) bool { + if ev.Kind != kindCommit || ev.Commit == nil || ev.Commit.Operation != opUpdate { + return false + } + if ev.Commit.Collection != colAcceptance || ev.Commit.RKey != subjectRKey(postURI) { + return false + } + uri, ok := fieldOf(ev.Commit.Record, "subject", "uri") + return ok && uri == postURI +} diff --git a/tests/e2e/lifecycle_test.go b/tests/e2e/lifecycle_test.go index f35ba4d..c59608c 100644 --- a/tests/e2e/lifecycle_test.go +++ b/tests/e2e/lifecycle_test.go @@ -66,7 +66,7 @@ func TestConsent_NobridgeLifecycle(t *testing.T) { l.await("phase-1 control post", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && got == control1Title + return e.Commit.Collection == colPostV2 && got == control1Title }) assertNoMarkedOutput := func(phase, forbiddenTitle string) { t.Helper() @@ -105,9 +105,16 @@ func TestConsent_NobridgeLifecycle(t *testing.T) { markedDID := profileEv.Did resumeEv := l.await("resumed post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && got == resumeTitle }) + if resumeEv.Did != markedDID { + t.Errorf("resumed post landed in repo %q, want the re-opted-in author's repo %q", + resumeEv.Did, markedDID) + } + if got := recordField(t, resumeEv.Commit.Record, "community"); got != sub.DID { + t.Errorf("resumed post community = %q, want %q", got, sub.DID) + } if did, res := h.bridgeResolveHandle(t, handle); res.status != 200 || did != markedDID { t.Fatalf("phase 2: handle %s should resolve to %s after opt-in (status %d, did %q)", handle, markedDID, res.status, did) @@ -127,8 +134,8 @@ func TestConsent_NobridgeLifecycle(t *testing.T) { // the actor profile (author repo), each on the exact rkey observed at // create time. l.await("scrub delete of the resumed post", func(e *jsEvent) bool { - return e.Commit.Collection == colPost && e.Commit.Operation == opDelete && - e.Did == sub.DID && e.Commit.RKey == resumeEv.Commit.RKey + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opDelete && + e.Did == resumeEv.Did && e.Commit.RKey == resumeEv.Commit.RKey }) l.await("scrub delete of the actor profile", func(e *jsEvent) bool { return e.Commit.Collection == colActorProfile && e.Commit.Operation == opDelete && @@ -141,7 +148,7 @@ func TestConsent_NobridgeLifecycle(t *testing.T) { control.createPost(t, community.ID, control2Title, "positive control after scrub") l.await("phase-3 control post", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && got == control2Title + return e.Commit.Collection == colPostV2 && got == control2Title }) assertNoMarkedOutput("phase 3 (bridged actor opted out)", scrubTriggerTitle) @@ -191,7 +198,7 @@ func TestDeleteActor_ScrubsAndTombstones(t *testing.T) { // the community's own profile (rkey `self`) must SURVIVE its member's // deletion, and a delete of it could not be seen through a narrower // filter. - l := h.newListener(t, cursor, colActorProfile, colPost, colComment, colCommunityProfile) + l := h.newListener(t, cursor, colActorProfile, colPostV2, colAcceptance, colComment, colCommunityProfile) title := "Doomed post " + h.suffix post := user.createPost(t, community.ID, title, "author will self-delete") @@ -203,7 +210,7 @@ func TestDeleteActor_ScrubsAndTombstones(t *testing.T) { authorDID := profileEv.Did postEv := l.await("post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && + return e.Commit.Collection == colPostV2 && e.Did == authorDID && e.Commit.Operation == opCreate && got == title }) user.createComment(t, post.ID, 0, "doomed comment") @@ -229,9 +236,16 @@ func TestDeleteActor_ScrubsAndTombstones(t *testing.T) { user.deleteAccount(t) // Scrub delete-commits for all three records, on their observed rkeys. - l.await("delete of the post (community repo)", func(e *jsEvent) bool { - return e.Commit.Collection == colPost && e.Commit.Operation == opDelete && - e.Did == sub.DID && e.Commit.RKey == postEv.Commit.RKey + l.await("delete of the post (author repo)", func(e *jsEvent) bool { + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opDelete && + e.Did == authorDID && e.Commit.RKey == postEv.Commit.RKey + }) + // The community's attestation must go with it: an acceptance whose + // subject no longer exists is a community vouching for a record that + // isn't there. + l.await("delete of the post's acceptance (community repo)", func(e *jsEvent) bool { + return e.Commit.Collection == colAcceptance && e.Commit.Operation == opDelete && + e.Did == sub.DID && e.Commit.RKey == subjectRKey(postEv.atURI()) }) l.await("delete of the comment (author repo)", func(e *jsEvent) bool { return e.Commit.Collection == colComment && e.Commit.Operation == opDelete && @@ -363,13 +377,26 @@ func TestUnsubscribe_StopsBridging(t *testing.T) { // Prove the doomed community flows BEFORE unsubscribing (otherwise the // negative below would pass vacuously on a broken subscription). + // + // A post is now TWO commits in TWO repos — the postv2 in the author's + // and the acceptance in the community's — so awaiting only the postv2 + // leaves the acceptance in flight, and the Undo{Follow} then races it + // into the negative window below. That is correct production behaviour + // (an in-flight materialization finishes; the unfollow fences FUTURE + // deliveries, exactly as v1 behaved for an in-flight post), so the + // scenario has to settle the pair before it starts asserting absence. preCursor := cursorNow() - pre := h.newListener(t, preCursor, colPost) + pre := h.newListener(t, preCursor, colPostV2, colAcceptance) preTitle := "Pre-unsubscribe " + h.suffix user.createPost(t, unsCommunity.ID, preTitle, "flows while subscribed") preEv := pre.await("pre-unsubscribe post", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == unsSub.DID && got == preTitle + return e.Commit.Collection == colPostV2 && got == preTitle + }) + preAcceptEv := pre.await("pre-unsubscribe post's acceptance", func(e *jsEvent) bool { + uri, _ := fieldOf(e.Commit.Record, "subject", "uri") + return e.Commit.Collection == colAcceptance && e.Did == unsSub.DID && + e.Commit.RKey == subjectRKey(preEv.atURI()) && uri == preEv.atURI() }) pre.close() @@ -392,47 +419,49 @@ func TestUnsubscribe_StopsBridging(t *testing.T) { // therefore the only cursor that reliably bounds this window; the one // known replayed event (the pre-post itself) is excluded by timestamp // in the sweep below. - l := h.newListener(t, preEv.TimeUs) + l := h.newListener(t, preAcceptEv.TimeUs) deadTitle := "Dead post " + h.suffix user.createPost(t, unsCommunity.ID, deadTitle, "must not bridge") liveTitle := "Live post " + h.suffix user.createPost(t, ctlCommunity.ID, liveTitle, "control keeps flowing") - l.await("control community post", func(e *jsEvent) bool { + ctlEv := l.await("control community post", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == ctlSub.DID && + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && got == liveTitle }) + if got := recordField(t, ctlEv.Commit.Record, "community"); got != ctlSub.DID { + t.Errorf("control post community = %q, want the still-subscribed community %q", got, ctlSub.DID) + } // Bounded negative: nothing on the unsubscribed community's repo - // STRICTLY NEWER than the pre-post (the newest thing it legitimately - // emitted), and the dead post's title nowhere (belt against it landing - // in a wrong repo). + // STRICTLY NEWER than that repo's own newest legitimate event (the + // pre-post's acceptance), and the dead post's title nowhere (belt + // against it landing in a wrong repo). for _, ev := range l.drain(negativeWindow) { if ev.Kind != kindCommit || ev.Commit == nil { continue } - if ev.Did == unsSub.DID && ev.TimeUs > preEv.TimeUs { - // A trailing bridgedStats UPDATE on an ALREADY-bridged post is the - // vote-stats refresher settling seeded counts (SEED_COUNTS_FROM_API - // is on), not new content bridged after the Undo{Follow} — - // unsubscribe stops NEW content, it does not roll back stats the - // aggregates already hold. Tolerated, but ONLY as a pure stats - // settle: carry-forward ships the whole record, so a genuine content - // change could hide inside a stats-shaped update. Pin the pre-post's - // title AND content unchanged so a real edit cannot pass as one. - if isBridgedStatsUpdate(ev) { - if got, _ := fieldOf(ev.Commit.Record, "title"); got != preTitle { - t.Errorf("stats-shaped update on unsubscribed repo %s changed the title to %q (want the pre-post %q): %s", - unsSub.DID, got, preTitle, ev) - } - if got, _ := fieldOf(ev.Commit.Record, "content"); got != "flows while subscribed" { - t.Errorf("stats-shaped update on unsubscribed repo %s changed the content to %q (want the pre-post body): %s", - unsSub.DID, got, ev) - } - } else { + if ev.Did == unsSub.DID && ev.TimeUs > preAcceptEv.TimeUs { + // ONE late event is still legitimate here: the vote-stats + // refresher settling the pre-post's seeded counts + // (SEED_COUNTS_FROM_API is on) rewrites that post — in the + // AUTHOR's repo — which moves its CID, and its acceptance is + // re-pinned to follow, in THIS repo. Unsubscribe stops NEW + // content; it does not roll back a settle already in flight. + // Tolerated ONLY for the pre-post's own subject, so an + // acceptance for anything else (the dead post above all) still + // fails. The stats stamp itself can no longer appear on this + // repo at all — content lives in the author's repo now — so + // there is nothing else to excuse. + if !isAcceptanceRepinOf(ev, preEv.atURI()) { t.Errorf("unsubscribed community repo %s emitted after Undo{Follow}: %s", unsSub.DID, ev) + } else { + // Logged, not silent: if this stops firing across runs the + // tolerance is dead weight, and a tolerance nobody can see + // exercised is indistinguishable from a hole. + t.Logf("tolerated in-flight stats repin of the pre-post's acceptance: %s", ev) } } if got, _ := fieldOf(ev.Commit.Record, "title"); got == deadTitle { diff --git a/tests/e2e/media_test.go b/tests/e2e/media_test.go index 06e0a20..dc278ac 100644 --- a/tests/e2e/media_test.go +++ b/tests/e2e/media_test.go @@ -109,7 +109,7 @@ func TestImagePost_EmbedImagesAndNSFWLabel(t *testing.T) { t.Logf("uploaded image: %s", imageURL) cursor := cursorNow() - l := h.newListener(t, cursor, colPost) + l := h.newListener(t, cursor, colPostV2) title := "Image post " + h.suffix const altText = "tiny e2e gradient" @@ -117,7 +117,7 @@ func TestImagePost_EmbedImagesAndNSFWLabel(t *testing.T) { ev := l.await("image post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && got == title }) rec := decodeRecord(t, ev.Commit.Record) @@ -162,12 +162,19 @@ func TestImagePost_EmbedImagesAndNSFWLabel(t *testing.T) { } // The blob is stored and served through the bridge's getBlob (the - // AppView's media path): posts live in the community repo, so the blob - // does too. Byte length must match the record's size claim; the bytes - // must decode as the image we uploaded (pict-rs 0.5 serves the original - // alias unmodified, but dimensions — not byte identity — are the - // contract worth pinning against a future pict-rs that re-encodes). - data, contentType := h.bridgeGetBlob(t, sub.DID, blobCID) + // AppView's media path). It lives in the repo holding the RECORD — the + // author's, since the flip — because a blob ref resolves against that + // repo and nowhere else; a blob left in the community's repo would be + // unresolvable for every consumer. Byte length must match the record's + // size claim; the bytes must decode as the image we uploaded (pict-rs + // 0.5 serves the original alias unmodified, but dimensions — not byte + // identity — are the contract worth pinning against a future pict-rs + // that re-encodes). + data, contentType := h.bridgeGetBlob(t, ev.Did, blobCID) + if _, found := h.bridgeGetRecord(t, sub.DID, colPostV2, ev.Commit.RKey); found { + t.Errorf("the image post also exists in the community repo %s — postv2 records live in "+ + "the author's repo only", sub.DID) + } if int64(len(data)) != int64(size) { t.Errorf("getBlob returned %d bytes, record claims size %d", len(data), int64(size)) } diff --git a/tests/e2e/moderation_test.go b/tests/e2e/moderation_test.go new file mode 100644 index 0000000..dc44018 --- /dev/null +++ b/tests/e2e/moderation_test.go @@ -0,0 +1,191 @@ +//go:build e2e + +package e2e + +import ( + "testing" +) + +// Scenario 18 (task 19): the moderation loop against a REAL Lemmy. +// +// This is the flow the postv2 flip exists for. A community's visibility +// decisions are records in the COMMUNITY's repo — acceptance while a post is +// visible, removal once a moderator takes it down — and they replace each +// other in ONE commit so the firehose never carries a half-completed +// moderation action. The author's post is never touched by a removal: it is +// the author's record, in the author's repo, and a community removing it from +// its own surface says nothing about the author's ownership of it. +// +// The three transitions asserted here are the three Lemmy actually sends +// (verified on the wire against 0.19.20): +// +// - mod remove WITH a reason → Announce{Delete} with summary "spam" +// - mod restore → Announce{Undo{Delete}} +// - author deletes their own post → Announce{Delete} with NO summary key +// +// Both the FIREHOSE (a consumer tailing the stream must be able to follow the +// state) and the END-STATE repo reads are asserted: they are two independent +// observations of the same commits, and a bridge that emitted the right +// events while converging to the wrong repo state would pass only one of them. +func TestModeration_RemoveRestoreAndSelfDelete(t *testing.T) { + h := newHarness(t) + community, sub := setupSubscribedCommunity(t, h, "mod") + + username := h.uniqueName(t, "molly") + user := h.registerUser(t, username) + + cursor := cursorNow() + l := h.newListener(t, cursor, colPostV2, colAcceptance, colRemoval) + + title := "Moderated post " + h.suffix + post := user.createPost(t, community.ID, title, "will be removed and restored") + + postEv := l.await("postv2 create", func(e *jsEvent) bool { + got, _ := fieldOf(e.Commit.Record, "title") + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && got == title + }) + postURI := postEv.atURI() + digest := subjectRKey(postURI) + + acceptEv := l.await("acceptance create", func(e *jsEvent) bool { + uri, _ := fieldOf(e.Commit.Record, "subject", "uri") + return e.Commit.Collection == colAcceptance && e.Did == sub.DID && uri == postURI + }) + if acceptEv.Commit.RKey != digest { + t.Fatalf("acceptance rkey = %q, want the subject digest %q", acceptEv.Commit.RKey, digest) + } + + // ── 1. Moderator removes the post, with a reason ─────────────────────── + const reason = "spam wave" + h.admin.removePost(t, post.ID, true, reason) + + removalEv := l.await("removal record written", func(e *jsEvent) bool { + uri, _ := fieldOf(e.Commit.Record, "subject", "uri") + return e.Commit.Collection == colRemoval && e.Did == sub.DID && uri == postURI + }) + if removalEv.Commit.RKey != digest { + t.Errorf("removal rkey = %q, want the SAME subject digest the acceptance used (%q) — "+ + "one derivation per subject is what lets a restore find and delete it", + removalEv.Commit.RKey, digest) + } + if got := recordField(t, removalEv.Commit.Record, "code"); got != "moderator-discretion" { + t.Errorf("removal code = %q, want %q (Lemmy sends no machine-readable code, so the "+ + "open knownValues set's default applies)", got, "moderator-discretion") + } + if got := recordField(t, removalEv.Commit.Record, "reason"); got != reason { + t.Errorf("removal reason = %q, want the moderator's text %q", got, reason) + } + if got := recordField(t, removalEv.Commit.Record, "subject", "cid"); got != postEv.Commit.CID { + t.Errorf("removal subject.cid = %q, want the version present at removal time %q", + got, postEv.Commit.CID) + } + + // The acceptance is withdrawn in the same breath. + l.await("acceptance deleted by the removal", func(e *jsEvent) bool { + return e.Commit.Collection == colAcceptance && e.Commit.Operation == opDelete && + e.Did == sub.DID && e.Commit.RKey == digest + }) + + // End state: removed, not accepted — and the AUTHOR's post untouched. + h.awaitRecordGone(t, sub.DID, colAcceptance, digest, "acceptance after removal") + h.awaitRecordPresent(t, sub.DID, colRemoval, digest, "removal after removal") + if _, found := h.bridgeGetRecord(t, postEv.Did, colPostV2, postEv.Commit.RKey); !found { + t.Error("the author's postv2 was deleted by a community removal — removal is " + + "community-scoped; the record belongs to the author") + } + + // ── 2. Moderator restores it ─────────────────────────────────────────── + h.admin.removePost(t, post.ID, false, "") + + l.await("removal deleted by the restore", func(e *jsEvent) bool { + return e.Commit.Collection == colRemoval && e.Commit.Operation == opDelete && + e.Did == sub.DID && e.Commit.RKey == digest + }) + reacceptEv := l.await("fresh acceptance after the restore", func(e *jsEvent) bool { + uri, _ := fieldOf(e.Commit.Record, "subject", "uri") + return e.Commit.Collection == colAcceptance && e.Did == sub.DID && + e.Commit.RKey == digest && uri == postURI + }) + if reacceptEv.Commit.Operation == opDelete { + t.Fatalf("expected the restore to WRITE an acceptance, got %s", reacceptEv) + } + + h.awaitRecordGone(t, sub.DID, colRemoval, digest, "removal after restore") + restored := h.awaitRecordPresent(t, sub.DID, colAcceptance, digest, "acceptance after restore") + // The fresh acceptance pins whatever version the post is on NOW. + current, found := h.bridgeGetRecord(t, postEv.Did, colPostV2, postEv.Commit.RKey) + if !found { + t.Fatal("the post vanished across the restore") + } + if got := stringAt(t, restored.Value, "subject", "cid"); got != current.CID { + t.Errorf("restored acceptance pins cid %q, want the post's current cid %q", got, current.CID) + } + + // ── 3. The AUTHOR deletes their own post ─────────────────────────────── + // No summary on the wire, so this is not moderation: the post goes, its + // acceptance goes with it, and NO removal record may be fabricated + // against an author nobody moderated. + user.deletePost(t, post.ID) + + l.await("postv2 delete after the author's own delete", func(e *jsEvent) bool { + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opDelete && + e.Did == postEv.Did && e.Commit.RKey == postEv.Commit.RKey + }) + l.await("acceptance delete after the author's own delete", func(e *jsEvent) bool { + return e.Commit.Collection == colAcceptance && e.Commit.Operation == opDelete && + e.Did == sub.DID && e.Commit.RKey == digest + }) + + h.awaitRecordGone(t, postEv.Did, colPostV2, postEv.Commit.RKey, "postv2 after self-delete") + h.awaitRecordGone(t, sub.DID, colAcceptance, digest, "acceptance after self-delete") + if rec, found := h.bridgeGetRecord(t, sub.DID, colRemoval, digest); found { + t.Errorf("a self-delete wrote a removal record — author deletion is not moderation, and "+ + "this puts a moderation action in the community's log against someone who was never "+ + "moderated: %+v", rec.Value) + } + + // Nothing else may arrive for this subject in the trailing window — in + // particular no re-acceptance resurrecting a post the author deleted. + for _, ev := range l.drain(negativeWindow) { + if ev.Kind != kindCommit || ev.Commit == nil { + continue + } + if ev.Commit.RKey == digest && ev.Commit.Operation != opDelete { + t.Errorf("post-deletion write for the deleted subject: %s", ev) + } + } +} + +// TestMixedEra_LegacyPostDispatch is DELIBERATELY SKIPPED at this tier, and +// the skip is the honest answer rather than a gap nobody noticed. +// +// The mixed-era contract — a pre-flip post keeps v1 semantics (record in the +// COMMUNITY's repo under the deprecated collection, no acceptance, no removal +// on moderation) while new posts flow postv2 — needs a LEGACY post to exist. +// Nothing in this stack can produce one: the flip is in the binary under +// test, so every post it materializes is a postv2, and legacy records exist +// only in repos written by a pre-flip build. The bridge exposes no +// record-seeding seam either (the admin API is communities, backfill, reemit, +// sweep-deleted, metrics) — by design. +// +// The two ways to manufacture one here would both be worse than skipping: +// writing rows straight into the bridge's postgres, or adding a test-only +// write endpoint. Each fabricates the state this scenario is supposed to +// OBSERVE, so what it would prove is that the fabrication matches the +// assertion — while adding a production seam that exists solely to be +// bypassed by tests. +// +// Era dispatch is covered where the seam is legitimate — the unit tier, which +// constructs a legacy record and mapping directly against the store: +// +// internal/materialize/postv2_repin_test.go TestLegacyPostStatsWritesNoAcceptance +// internal/ingest/moderation_test.go TestLegacyPostModRemovalKeepsV1Semantics +// internal/votes/postv2_binding_test.go (legacy binding by repo DID) +// +// If a future task ever needs true end-to-end mixed-era coverage, the honest +// route is a stack that boots a PRE-FLIP image, materializes a post, then +// upgrades the binary in place — real legacy records, real pipeline. +func TestMixedEra_LegacyPostDispatch(t *testing.T) { + t.Skip("no sanctioned seam to seed a pre-flip legacy post in the e2e stack; " + + "era dispatch is covered at the unit tier (see this test's comment)") +} diff --git a/tests/e2e/votes_hammer_test.go b/tests/e2e/votes_hammer_test.go index 7936428..a40da57 100644 --- a/tests/e2e/votes_hammer_test.go +++ b/tests/e2e/votes_hammer_test.go @@ -41,7 +41,7 @@ func TestVoteHammer_ConcurrentVotersExactAggregates(t *testing.T) { post := author.createPost(t, community.ID, title, "vote target") postEv := l.await("hammer post create", func(e *jsEvent) bool { got, _ := fieldOf(e.Commit.Record, "title") - return e.Commit.Collection == colPost && e.Did == sub.DID && + return e.Commit.Collection == colPostV2 && e.Commit.Operation == opCreate && got == title }) postURI := postEv.atURI() @@ -96,6 +96,25 @@ func TestVoteHammer_ConcurrentVotersExactAggregates(t *testing.T) { }) awaitAggregates(t, h, postURI, 0, upVoters) + // The stats sweep folds the settled counts onto the POST record (author + // repo), which moves its CID — so the community's acceptance must be + // re-pinned to the version that now exists. This is the bridgedStats + // exception to the acceptance flow (PRD §5.5: we are the engine for our + // own communities, so a stats stamp repins synchronously rather than + // re-running admission), and it is the one cross-repo write a vote sweep + // performs. Left unpinned, every sweep would drop the post out of its + // community until somebody edited it. + statsEv := l.await("bridgedStats stamp on the hammered post", func(e *jsEvent) bool { + return e.Commit != nil && e.Did == postEv.Did && e.Commit.Collection == colPostV2 && + e.Commit.RKey == postEv.Commit.RKey && isBridgedStatsUpdate(e) + }) + l.await("acceptance repin following the stats stamp", func(e *jsEvent) bool { + cid, _ := fieldOf(e.Commit.Record, "subject", "cid") + return e.Commit.Collection == colAcceptance && e.Did == sub.DID && + e.Commit.RKey == subjectRKey(postURI) && cid == statsEv.Commit.CID + }) + h.awaitAcceptanceConverged(t, sub.DID, postEv.Did, postEv.Commit.RKey) + // Belt on locked decision 7: everything consumed during the hammer was // already vetted (vetEvent Fatalfs on any unexpected collection inside // await AND drain), so the drain itself IS the assertion — it forces diff --git a/tests/e2e/zz_sweep_test.go b/tests/e2e/zz_sweep_test.go index dedda29..6f7d9fe 100644 --- a/tests/e2e/zz_sweep_test.go +++ b/tests/e2e/zz_sweep_test.go @@ -84,7 +84,11 @@ func TestZZ_SuiteEndSweep(t *testing.T) { total++ repos[ev.Did] = true counts[ev.Commit.Collection+" "+ev.Commit.Operation]++ - if ev.Commit.Collection == colPost && ev.Commit.Operation == opCreate { + // Both post eras count: the suite creates postv2 records now, + // but a replay reaching far enough back can still legally + // surface pre-flip legacy posts, and either one proves the + // replayed history is real. + if isPostCollection(ev.Commit.Collection) && ev.Commit.Operation == opCreate { postCreates++ } if ev.Commit.Operation == opDelete { @@ -128,7 +132,7 @@ func TestZZ_SuiteEndSweep(t *testing.T) { } if postCreates == 0 { t.Errorf("sweep saw no %s create — a full suite run always emits posts, so the replayed history is incomplete: %v", - colPost, counts) + colPostV2, counts) } if deleteOps == 0 { t.Errorf("sweep saw no delete op in any collection — a full suite run always emits deletes (edits/deletes, scrubs), so the replayed history is incomplete: %v",